Data Engineering Command Center
Complete methodology for designing, building, operating, and scaling data pipelines and infrastructure. Zero dependencies — pure agent skill.
Phase 1: Data Architecture Assessment
Before building anything, understand the landscape.
Architecture Brief
project_name: ""
business_context: ""
data_consumers:
- team: ""
use_case: "" # analytics | ML | operational | reporting | reverse-ETL
latency_requirement: "" # real-time (<1s) | near-real-time (<5min) | batch (hourly+)
query_pattern: "" # ad-hoc | scheduled | API | dashboard
current_state:
sources: [] # list every system producing data
storage: [] # where data lives today
pain_points: [] # what's broken, slow, unreliable
data_volume:
current_gb_per_day: 0
growth_rate_percent: 0
retention_months: 0
constraints:
budget_monthly_usd: 0
team_size: 0
skill_level: "" # junior | mid | senior | mixed
compliance: [] # GDPR, HIPAA, SOX, PCI, none
cloud_provider: "" # AWS | GCP | Azure | multi | on-prem
Architecture Pattern Decision Matrix
| Signal |
Pattern |
When to Use |
| All consumers need data hourly+ |
Batch ETL |
Reporting, warehousing, most analytics |
| Some need <5 min latency |
Micro-batch |
Dashboard freshness, near-real-time analytics |
| Events need <1s processing |
Streaming |
Fraud detection, real-time pricing, alerts |
| Need both batch + streaming |
Lambda |
When batch accuracy + real-time speed both matter |
| Want to simplify Lambda |
Kappa |
When you can reprocess from stream replay |
| Data lake + warehouse combined |
Lakehouse |
When you need both cheap storage + fast SQL |
| Sources change independently |
Data Mesh |
Large orgs, domain-owned data, >5 teams |
| ML is primary consumer |
Feature Store |
ML-heavy orgs with feature reuse needs |
Technology Selection Guide
Orchestration
| Tool |
Best For |
Avoid When |
| Airflow |
Complex DAGs, Python-native teams, mature ecosystem |
Simple pipelines (<5 tasks) |
| Dagster |
Software-defined assets, strong typing, dev experience |
Legacy team resistant to new paradigms |
| Prefect |
Dynamic workflows, cloud-native, Python-first |
Need on-prem with no cloud dependency |
| dbt |
SQL transformations, ELT, analytics engineering |
Non-SQL transforms, streaming |
| Temporal |
Long-running workflows, retry-heavy, microservices |
Simple ETL, small teams |
| Cron + scripts |
<3 pipelines, solo engineer, simple schedules |
Anything with dependencies or retries |
Processing
| Tool |
Best For |
Avoid When |
| Spark |
>100GB, complex transforms, ML pipelines |
<10GB (overkill), real-time streaming |
| DuckDB |
Local analytics, <100GB, SQL on files |
Distributed processing, production streaming |
| Polars |
Single-node, Rust-speed, <50GB, DataFrames |
Distributed, need Spark ecosystem |
| Pandas |
<1GB, quick analysis, prototyping |
Production pipelines, anything >5GB |
| Flink |
True streaming, event-time processing |
Batch-only, small team (steep learning curve) |
| SQL (warehouse) |
ELT in Snowflake/BigQuery/Redshift |
Complex ML transforms, binary data |
Storage
| Tool |
Best For |
Avoid When |
| Snowflake |
Analytics, separation of compute/storage, multi-cloud |
Tight budget, real-time OLTP |
| BigQuery |
GCP-native, serverless, large-scale analytics |
Multi-cloud, need fine-grained cost control |
| Redshift |
AWS-native, existing AWS ecosystem |
Elastic scaling needs, multi-cloud |
| Databricks |
ML + analytics unified, Spark-native, lakehouse |
Pure SQL analytics, small data |
| PostgreSQL |
OLTP + light analytics, <500GB, budget-conscious |
>1TB analytics, real-time dashboards on large data |
| S3/GCS/ADLS |
Raw data lake, cheap storage, any format |
Direct SQL queries (need compute layer) |
| Delta Lake/Iceberg |
Table format on data lake, ACID on files |
Simple file storage, no lakehouse need |
Phase 2: Data Modeling
Modeling Methodology Decision
| Approach |
Best For |
Key Concept |
| Kimball (Dimensional) |
BI/reporting, star schemas |
Facts + Dimensions, business-process-centric |
| Inmon (3NF) |
Enterprise data warehouse, single source of truth |
Normalized, subject-area-centric |
| Data Vault 2.0 |
Agile warehousing, auditability, multiple sources |
Hubs + Links + Satellites, insert-only |
| One Big Table (OBT) |
Simple analytics, few joins, dashboard performance |
Pre-joined, denormalized, fast queries |
| Activity Schema |
Event analytics, product analytics |
Entity + Activity + Feature columns |
Dimensional Model Template
fact_table:
name: "fact_[business_process]"
grain: "" # one row = one [what]?
grain_statement: "One row per [transaction/event/snapshot] at [time grain]"
measures:
- name: ""
type: "" # additive | semi-additive | non-additive
aggregation: "" # SUM | AVG | COUNT | MIN | MAX | COUNT DISTINCT
business_definition: ""
degenerate_dimensions: [] # IDs stored in fact (order_number, invoice_id)
foreign_keys: [] # links to dimension tables
dimensions:
- name: "dim_[entity]"
type: "" # Type 1 (overwrite) | Type 2 (history) | Type 3 (previous value)
natural_key: "" # business key from source
surrogate_key: "" # warehouse-generated key
attributes:
- name: ""
source: ""
scd_type: "" # 1 | 2 | 3
hierarchy: [] # e.g., [country, region, city, store]
SCD Type Decision Guide
| Scenario |
SCD Type |
Implementation |
| Don't care about history |
Type 1 |
UPDATE in place |
| Need full history |
Type 2 |
New row + valid_from/valid_to + is_current flag |
| Only need previous value |
Type 3 |
Add previous_[column] |
| Track changes with timestamps |
Type 4 |
Mini-dimension (history table) |
| Hybrid: some attrs Type 1, some Type 2 |
Type 6 |
Combine 1+2+3 in one table |
Default recommendation: Type 2 for anything business-critical (customer status, product price, employee department). Type 1 for everything else.
Naming Conventions
| Object |
Convention |
Example |
| Raw/staging tables |
raw_[source]_[table] |
raw_stripe_payments |
| Staging models |
stg_[source]__[entity] |
stg_stripe__payments |
| Intermediate models |
int_[entity]_[verb] |
int_orders_pivoted |
| Mart/fact tables |
fct_[business_process] |
fct_orders |
| Dimension tables |
dim_[entity] |
dim_customers |
| Metrics/aggregates |
mrt_[domain]_[metric] |
mrt_sales_daily |
| Snapshots |
snp_[entity]_[grain] |
snp_inventory_daily |
| Columns: boolean |
is_[state] or has_[thing] |
is_active, has_subscription |
| Columns: timestamp |
[event]_at |
created_at, shipped_at |
| Columns: date |
[event]_date |
order_date |
| Columns: ID |
[entity]_id |
customer_id |
| Columns: amount |
[thing]_amount |
order_amount |
| Columns: count |
[thing]_count |
line_item_count |
Phase 3: Pipeline Design Patterns
Universal Pipeline Template
pipeline:
name: ""
owner: ""
schedule: "" # cron expression
sla_minutes: 0 # max acceptable runtime
tier: "" # 1 (critical) | 2 (important) | 3 (nice-to-have)
extract:
source_system: ""
connection: ""
strategy: "" # full | incremental | CDC | log-based
incremental_key: "" # column for incremental (e.g., updated_at)
watermark_storage: "" # where to persist last-extracted position
transform:
engine: "" # SQL | Spark | Python | dbt
stages:
- name: "clean"
operations: [] # dedupe, null handling, type casting, trimming
- name: "conform"
operations: [] # standardize codes, currencies, timezones
- name: "enrich"
operations: [] # lookups, calculations, derived fields
- name: "aggregate"
operations: [] # rollups, pivots, window functions
load:
target_system: ""
strategy: "" # append | upsert | merge | truncate-reload | partition-swap
merge_keys: []
partition_key: ""
clustering_keys: []
quality_gates:
pre_load: [] # checks before writing
post_load: [] # checks after writing
error_handling:
strategy: "" # fail-fast | dead-letter | retry | skip-and-alert
max_retries: 3
retry_delay_seconds: 300
alert_channels: []
Extraction Strategy Decision Tree
Is the source database?
├── Yes → Does it support CDC?
│ ├── Yes → Use CDC (Debezium, AWS DMS, Fivetran)
│ │ Best for: high-volume, low-latency, minimal source impact
│ └── No → Does it have a reliable updated_at column?
│ ├── Yes → Incremental extraction on updated_at
│ │ ⚠️ Won't catch hard deletes — add periodic full reconciliation
│ └── No → Full extraction
│ Only viable for small tables (<1M rows)
├── Is it an API?
│ ├── Supports webhooks? → Event-driven ingestion
│ ├── Has cursor/pagination? → Incremental with cursor bookmark
│ └── No pagination? → Full pull with rate-limit handling
├── Is it files (S3, SFTP, email)?
│ └── Event-triggered (S3 notification, file watcher)
│ Validate: schema, completeness, filename pattern
└── Is it streaming (Kafka, Kinesis, Pub/Sub)?
└── Consumer group with offset management
Key decisions: at-least-once vs exactly-once, consumer lag alerting
Load Strategy Decision
| Strategy |
When |
Trade-off |
| Append |
Event/log data, immutable facts |
Simple but grows forever — partition + retain |
| Upsert/Merge |
Dimension updates, SCD Type 1 |
Handles updates but slower on large tables |
| Truncate-Reload |
Small tables (<1M), reference data |
Simple but window of missing data |
| Partition Swap |
Large fact tables, daily loads |
Atomic, fast, but needs partition alignment |
| Soft Delete |
Need audit trail of deletions |
Adds complexity to every downstream query |
Idempotency Rules (NON-NEGOTIABLE)
Every pipeline MUST be re-runnable without side effects:
- Use MERGE/UPSERT, never blind INSERT for mutable data
- Partition-swap for immutable data — drop partition + reload
- Store watermarks externally — not in the pipeline code
- Dedup at ingestion — use source natural keys
- Test by running twice — output must be identical both times
Phase 4: Data Quality Framework
Quality Dimensions
| Dimension |
Definition |
Example Check |
| Completeness |
No missing values where required |
NOT NULL on required fields, row count within range |
| Uniqueness |
No unexpected duplicates |
Primary key uniqueness, natural key uniqueness |
| Validity |
Values within expected domain |
Enum checks, range checks, regex patterns |
| Accuracy |
Data matches real-world truth |
Cross-system reconciliation, manual spot checks |
| Freshness |
Data arrives on time |
MAX(loaded_at) > NOW() - INTERVAL '2 hours' |
| Consistency |
Same data agrees across systems |
Sum reconciliation between source and target |
Quality Check Templates
-- Completeness: Required fields not null
SELECT COUNT(*) AS null_violations
FROM {table}
WHERE {required_column} IS NULL;
-- Threshold: 0
-- Uniqueness: No duplicate primary keys
SELECT {pk_column}, COUNT(*) AS dupe_count
FROM {table}
GROUP BY {pk_column}
HAVING COUNT(*) > 1;
-- Threshold: 0
-- Freshness: Data arrived within SLA
SELECT CASE
WHEN MAX({timestamp_col}) > CURRENT_TIMESTAMP - INTERVAL '{sla_hours} hours'
THEN 'PASS' ELSE 'FAIL'
END AS freshness_check
FROM {table};
-- Volume: Row count within expected range
SELECT CASE
WHEN COUNT(*) BETWEEN {min_expected} AND {max_expected}
THEN 'PASS' ELSE 'FAIL'
END AS volume_check
FROM {table}
WHERE {partition_col} = '{run_date}';
-- Referential: FK integrity
SELECT COUNT(*) AS orphan_count
FROM {fact_table} f
LEFT JOIN {dim_table} d ON f.{fk} = d.{pk}
WHERE d.{pk} IS NULL;
-- Threshold: 0
-- Distribution: No unexpected skew
SELECT {column}, COUNT(*) AS cnt,
ROUND(100.0 * COUNT(*) / SUM(COUNT(*)) OVER (), 2) AS pct
FROM {table}
GROUP BY {column}
ORDER BY cnt DESC;
-- Alert if any single value > {max_pct}%
-- Cross-system reconciliation
SELECT
(SELECT SUM(amount) FROM source_system.orders WHERE date = '{date}') AS source_total,
(SELECT SUM(amount) FROM warehouse.fct_orders WHERE order_date = '{date}') AS target_total,
ABS(source_total - target_total) AS variance;
-- Threshold: variance < 0.01 * source_total (1%)
Data Contract Template
contract:
name: ""
version: ""
owner: "" # team responsible for producing this data
consumers: [] # teams consuming this data
sla:
freshness_hours: 0
availability_percent: 99.9
support_hours: "" # business-hours | 24x7
schema:
- column: ""
type: ""
nullable: false
description: ""
business_definition: ""
pii: false
checks:
- type: "" # not_null | unique | range | enum | regex | custom
params: {}
breaking_change_policy: "" # notify-30-days | version-bump | never-break
notification_channel: ""
Quality Severity Levels
| Level |
Definition |
Response |
| P0 — Critical |
Data corruption, wrong numbers in production dashboards, compliance data wrong |
Stop pipeline, alert immediately, rollback if possible |
| P1 — High |
Missing data for key reports, SLA breach, >5% of records affected |
Alert team, fix within 4 hours, post-mortem required |
| P2 — Medium |
Non-critical field quality, <1% records affected, no downstream impact |
Fix in next sprint, add monitoring to prevent recurrence |
| P3 — Low |
Cosmetic issues, edge cases, non-critical data |
Backlog, fix when convenient |
Phase 5: Performance Optimization
SQL Optimization Checklist
| Problem |
Fix |
Impact |
| Full table scan |
Add/use partition pruning |
10-100x faster |
| Large joins |
Pre-aggregate before joining |
5-50x faster |
| SELECT * |
Select only needed columns |
2-10x faster (columnar stores) |
| Correlated subquery |
Rewrite as JOIN or window function |
10-100x faster |
| DISTINCT on large result |
Fix upstream duplication instead |
2-5x faster |
| ORDER BY without LIMIT |
Add LIMIT or remove if not needed |
Prevents memory spills |
| String operations in WHERE |
Pre-compute, use lookup table |
Enables index usage |
| Multiple passes over same data |
Combine with CASE WHEN + GROUP BY |
2-5x faster |
| NOT IN with NULLs |
Use NOT EXISTS or LEFT JOIN IS NULL |
Correctness + performance |
Spark Optimization Guide
| Problem |
Solution |
| Shuffle-heavy joins |
Broadcast small table (broadcast(df)) if <100MB |
| Data skew |
Salt the skewed key: add random prefix, join on salted key, aggregate |
| Small files |
Coalesce output: .coalesce(target_files) or use adaptive query execution |
| Too many partitions |
spark.sql.shuffle.partitions = 2-3x cluster cores |
| OOM errors |
Increase spark.executor.memory, reduce partition size, spill to disk |
| Slow writes |
Use Parquet with snappy, partition by date, avoid small writes |
| Repeated computation |
.cache() or .persist() DataFrames used >1 time |
| Complex transformations |
Push down predicates, filter early, select early |
Partitioning Strategy
| Data Type |
Partition Key |
Why |
| Transactional/event |
Date (daily or monthly) |
Most queries filter by time range |
| Multi-tenant |
Tenant ID + date |
Isolate tenant queries, time-range pruning |
| Geospatial |
Region + date |
Regional queries are common |
| Log data |
Date + hour |
High volume needs finer partitions |
| Reference/dimension |
Don't partition |
Too small, full scan is fine |
Rules:
- Target 100MB-1GB per partition (compressed)
- <10,000 total partitions per table
- Never partition on high-cardinality columns (user_id)
- Always include partition key in WHERE clauses
Phase 6: Data Governance & Cataloging
Data Classification
| Level |
Examples |
Controls |
| Public |
Product catalog, published stats |
No restrictions |
| Internal |
Aggregated metrics, non-PII analytics |
Auth required, audit logging |
| Confidential |
Customer PII, financial records, HR data |
Encryption, column-level access, masking |
| Restricted |
SSN, payment cards, health records, passwords |
Encryption at rest + transit, tokenization, audit every access, retention limits |
PII Handling Rules
- Identify: Scan all sources for PII columns (name, email, phone, SSN, DOB, address, IP)
- Classify: Tag each with sensitivity level
- Minimize: Only ingest PII you actually need
- Protect:
- Hash or tokenize in staging (SHA-256 with salt for pseudonymization)
- Dynamic masking for non-privileged users
- Column-level encryption for restricted data
- Retain: Auto-delete after retention period
- Audit: Log every query touching PII columns
- Right to delete: Build a deletion pipeline that propagates across all derived tables
Data Catalog Entry Template
dataset:
name: ""
description: ""
owner_team: ""
steward: "" # person responsible for quality
domain: "" # sales | marketing | finance | product | engineering
tier: "" # gold (trusted) | silver (cleaned) | bronze (raw)
lineage:
sources: [] # upstream datasets/systems
transformations: "" # brief description of key transforms
downstream: [] # who consumes this
refresh:
schedule: ""
sla_hours: 0
last_successful_run: ""
quality:
tests: [] # list of quality checks
last_score: 0 # 0-100
known_issues: []
access:
classification: "" # public | internal | confidential | restricted
pii_columns: []
access_request_process: "" # how to get access
usage:
avg_daily_queries: 0
top_consumers: []
cost_monthly_usd: 0
Phase 7: Pipeline Monitoring & Alerting
Pipeline Health Dashboard
dashboard:
pipeline_metrics:
- metric: "Pipeline Success Rate"
formula: "successful_runs / total_runs * 100"
target: ">99%"
alert_threshold: "<95%"
- metric: "Average Runtime"
formula: "avg(end_time - start_time) over 7 days"
target: "<SLA"
alert_threshold: ">80% of SLA"
- metric: "Data Freshness"
formula: "NOW() - MAX(loaded_at)"
target: "<SLA hours"
alert_threshold: ">SLA"
- metric: "Data Volume Variance"
formula: "abs(today_rows - avg_7d_rows) / avg_7d_rows * 100"
target: "<20%"
alert_threshold: ">50%"
- metric: "Quality Check Pass Rate"
formula: "passed_checks / total_checks * 100"
target: "100%"
alert_threshold: "<95%"
- metric: "Failed Pipeline Count"
formula: "count where status = failed in last 24h"
target: "0"
alert_threshold: ">0"
- metric: "Backfill Queue"
formula: "count of pending backfill requests"
target: "0"
alert_threshold: ">5"
- metric: "Infrastructure Cost"
formula: "compute + storage + egress"
target: "<budget"
alert_threshold: ">110% budget"
Alert Severity
| Severity |
Condition |
Response Time |
Example |
| P0 |
Revenue/compliance impacting |
15 min |
Payment pipeline down, regulatory report delayed |
| P1 |
Business-critical dashboard stale |
1 hour |
Executive dashboard >4h stale |
| P2 |
Non-critical pipeline failed |
4 hours |
Marketing attribution delayed |
| P3 |
Warning/degradation |
Next business day |
Pipeline 80% of SLA, minor quality drift |
Structured Logging Standard
Every pipeline run MUST log:
{
"pipeline_name": "",
"run_id": "",
"started_at": "",
"completed_at": "",
"status": "success|failed|partial",
"stage": "",
"rows_extracted": 0,
"rows_transformed": 0,
"rows_loaded": 0,
"rows_rejected": 0,
"quality_checks_passed": 0,
"quality_checks_failed": 0,
"duration_seconds": 0,
"error_message": "",
"watermark_before": "",
"watermark_after": ""
}
Phase 8: Testing Strategy
Pipeline Test Pyramid
| Layer |
What to Test |
How |
When |
| Unit |
Individual transforms, business logic |
pytest with fixtures, dbt unit tests |
Every PR |
| Integration |
Source connectivity, schema compatibility |
Test against staging/dev environment |
Daily + PR |
| Contract |
Schema hasn't changed, data types stable |
Schema registry, contract tests |
Every pipeline run |
| Data Quality |
Completeness, uniqueness, freshness, validity |
Quality framework checks |
Every pipeline run |
| E2E |
Full pipeline produces correct output |
Golden dataset comparison |
Weekly + release |
| Performance |
Runtime within SLA, no regression |
Benchmark against historical runs |
Weekly |
dbt Testing Checklist
# For every model, define at minimum:
models:
- name: fct_orders
columns:
- name: order_id
tests:
- unique
- not_null
- name: customer_id
tests:
- not_null
- relationships:
to: ref('dim_customers')
field: customer_id
- name: order_amount
tests:
- not_null
- dbt_utils.accepted_range:
min_value: 0
max_value: 1000000
- name: order_status
tests:
- accepted_values:
values: ['pending', 'confirmed', 'shipped', 'delivered', 'cancelled']
- name: ordered_at
tests:
- not_null
- dbt_utils.recency:
datepart: day
field: ordered_at
interval: 2
Backfill Protocol
When you need to reprocess historical data:
- Scope: Define exact date range and affected tables
- Impact assessment: What downstream models/dashboards will be affected?
- Communication: Notify consumers of temporary data inconsistency
- Isolation: Run backfill in separate compute to avoid impacting current pipelines
- Validation: Compare row counts and key metrics pre/post backfill
- Execution: Process in reverse-chronological order (most recent first)
- Monitoring: Watch for resource spikes, duplicate creation
- Verification: Reconcile against source after completion
- Documentation: Log what was backfilled, why, and any anomalies found
Phase 9: Cost Optimization
Cloud Cost Reduction Strategies
| Strategy |
Savings |
Effort |
| Right-size compute (auto-scaling) |
20-40% |
Low |
| Use spot/preemptible instances for batch |
60-80% |
Medium |
| Compress data (Parquet + Snappy/Zstd) |
50-80% storage |
Low |
| Lifecycle policies (hot → warm → cold → archive) |
40-70% storage |
Low |
| Eliminate unused tables/pipelines |
10-30% |
Low |
| Optimize query patterns (partition pruning) |
30-60% compute |
Medium |
| Reserved capacity for steady-state |
30-50% |
Medium |
| Cache expensive queries |
20-50% compute |
Medium |
Cost Allocation Template
cost_tracking:
by_pipeline:
- pipeline: ""
compute_monthly_usd: 0
storage_monthly_usd: 0
egress_monthly_usd: 0
total: 0
cost_per_row: 0 # total / rows_processed
business_value: "" # what revenue/decision does this enable?
roi_justified: true # is the cost worth it?
optimization_opportunities:
- description: ""
estimated_savings_usd: 0
effort: "" # low | medium | high
priority: 0 # 1 = do now
Cost Red Flags
- Single pipeline >30% of total spend
- Cost per row increasing month-over-month
- Tables with 0 queries in 30 days
- Dev/staging environments running 24/7
- Full table scans on >1TB tables
- Uncompressed data in cloud storage
- Cross-region data transfer
Phase 10: Operational Runbooks
Pipeline Failure Triage
Pipeline failed →
1. Check error message in logs
├── Connection timeout → Check source availability, network, credentials
├── Schema mismatch → Source schema changed → update extract + notify
├── Data quality check failed → Investigate source data, check for anomalies
├── Out of memory → Increase resources or optimize query
├── Permission denied → Check IAM roles, token expiry
├── Duplicate key violation → Check idempotency, investigate source dupes
└── Timeout (SLA breach) → Check data volume spike, query plan, cluster health
2. Determine impact
├── What dashboards/reports are affected?
├── What's the data freshness SLA?
└── Who needs to be notified?
3. Fix
├── Transient (network, timeout) → Retry
├── Data issue → Fix source data, re-run with quality gate override if safe
├── Schema change → Update pipeline, backfill if needed
└── Infrastructure → Scale up, file ticket with cloud provider
4. Post-fix
├── Verify data correctness
├── Update runbook with new failure mode
└── Add monitoring/alerting to catch earlier next time
Schema Change Management
When a source system changes schema:
- Detect: Schema comparison check in extraction pipeline (hash schema, compare to registered)
- Classify:
- Additive (new column): Usually safe — add to pipeline, backfill if needed
- Rename: Map old → new in transform, update downstream
- Type change: Assess compatibility, may need cast or historical rebuild
- Column removed: Critical — breaks queries, need immediate attention
- Test: Run pipeline in dry-run mode with new schema
- Deploy: Update transforms, quality checks, documentation
- Communicate: Notify downstream consumers via data contract channel
Disaster Recovery
| Scenario |
RPO |
RTO |
Recovery Steps |
| Pipeline code lost |
0 (git) |
1h |
Redeploy from git, restore orchestrator state |
| Warehouse data corrupted |
Varies |
4h |
Restore from Time Travel/snapshot, re-run affected pipelines |
| Source system down |
N/A |
Wait |
Queue extractions, catch up with incremental once restored |
| Cloud region outage |
24h |
8h |
Failover to DR region if configured, else wait |
| Credential compromise |
0 |
2h |
Rotate all credentials, audit access logs, re-run affected pipelines |
Phase 11: Advanced Patterns
Slowly Changing Dimension Type 2 (SQL Template)
-- Merge pattern for SCD Type 2
MERGE INTO dim_customer AS target
USING (
SELECT * FROM stg_customers
WHERE updated_at > (SELECT MAX(valid_from) FROM dim_customer)
) AS source
ON target.customer_natural_key = source.customer_id
AND target.is_current = TRUE
-- Update: close old record
WHEN MATCHED AND (
target.customer_name != source.name OR
target.customer_status != source.status
-- list all Type 2 tracked columns
) THEN UPDATE SET
is_current = FALSE,
valid_to = CURRENT_TIMESTAMP
-- Insert: new record (both new customers and changed ones)
WHEN NOT MATCHED THEN INSERT (
customer_natural_key, customer_name, customer_status,
valid_from, valid_to, is_current
) VALUES (
source.customer_id, source.name, source.status,
CURRENT_TIMESTAMP, '9999-12-31', TRUE
);
-- Then insert new versions of changed records
INSERT INTO dim_customer (
customer_natural_key, customer_name, customer_status,
valid_from, valid_to, is_current
)
SELECT customer_id, name, status,
CURRENT_TIMESTAMP, '9999-12-31', TRUE
FROM stg_customers s
WHERE EXISTS (
SELECT 1 FROM dim_customer d
WHERE d.customer_natural_key = s.customer_id
AND d.is_current = FALSE
AND d.valid_to = CURRENT_TIMESTAMP
);
CDC with Debezium (Architecture Pattern)
Source DB → Debezium Connector → Kafka Topic →
├── Stream processor (Flink/Spark Streaming) → Target DB
├── S3 sink connector → Data Lake (raw)
└── Elasticsearch sink → Search index
Key decisions:
- Topic per table or single topic: Per table (easier routing, independent scaling)
- Schema registry: Always use (Confluent Schema Registry or AWS Glue)
- Serialization: Avro (compact + schema evolution) or Protobuf (strict + fast)
- Offset management: Connector manages; monitor consumer lag
Feature Store Pattern
feature_store:
entity: "customer"
entity_key: "customer_id"
features:
- name: "total_orders_30d"
description: "Total orders in last 30 days"
type: "INT"
source: "fct_orders"
computation: "batch" # batch | streaming | on-demand
freshness: "daily"
ttl_hours: 48
- name: "avg_order_value_90d"
description: "Average order value last 90 days"
type: "FLOAT"
source: "fct_orders"
computation: "batch"
freshness: "daily"
ttl_hours: 48
- name: "last_login_minutes_ago"
description: "Minutes since last login event"
type: "INT"
source: "events_stream"
computation: "streaming"
freshness: "real-time"
ttl_hours: 1
serving:
online: true # low-latency feature serving (Redis/DynamoDB)
offline: true # batch feature retrieval for training
point_in_time_correct: true # prevent feature leakage in ML training
Data Mesh Principles
If operating at scale (>5 data teams):
- Domain ownership: Each business domain owns its data products (not central data team)
- Data as a product: Treat datasets like products — SLAs, documentation, discoverability
- Self-serve platform: Central team builds the platform, domains build on top
- Federated governance: Standards and interoperability maintained centrally, implementation decentralized
When NOT to use Data Mesh:
- <5 data producers/consumers
- Small team (<20 engineers total)
- Single business domain
- Early-stage company (over-engineering)
Quality Scoring Rubric (0-100)
| Dimension |
Weight |
Scoring |
| Pipeline Reliability |
20 |
0=frequent failures, 10=some failures with manual recovery, 20=99.5%+ success rate with auto-retry |
| Data Quality |
20 |
0=no checks, 10=basic null/unique checks, 20=comprehensive quality framework with contracts |
| Performance |
15 |
0=regularly breaches SLA, 8=meets SLA, 15=well under SLA with optimization |
| Documentation |
10 |
0=none, 5=basic README, 10=full catalog entries with lineage and business definitions |
| Monitoring |
15 |
0=no alerts, 8=failure alerts only, 15=proactive monitoring with dashboards and anomaly detection |
| Testing |
10 |
0=no tests, 5=basic smoke tests, 10=full test pyramid (unit+integration+contract+E2E) |
| Cost Efficiency |
10 |
0=no cost tracking, 5=tracked, 10=optimized with ROI justification per pipeline |
Scoring guide:
- 0-40: Critical gaps — prioritize pipeline reliability and data quality
- 41-60: Functional but fragile — add monitoring, testing, documentation
- 61-80: Solid — optimize performance, cost, governance
- 81-100: Excellent — maintain, innovate, mentor
Edge Cases & Gotchas
Timezone Traps
- Store everything in UTC. Convert only at presentation layer
- Event timestamps: use event time, not processing time
- Daylight saving:
TIMESTAMP WITH TIME ZONE, never WITHOUT
- Late-arriving data: watermark strategy + allowed lateness window
Late-Arriving Data
- Define maximum acceptable lateness per source
- Reprocess affected partitions when late data arrives
- Track late arrival rate as a quality metric
- Consider separate "late data" pipeline that patches in
Exactly-Once Processing
- True exactly-once is expensive. Most systems need at-least-once + idempotent writes
- Use transaction IDs or natural keys for deduplication
- Kafka: use idempotent producer + transactional consumer
- Database: MERGE/UPSERT on natural key
Schema Evolution
- Forward compatible: New code reads old data (safe to deploy new readers first)
- Backward compatible: Old code reads new data (safe to deploy new writers first)
- Full compatible: Both directions (safest, most restrictive)
- Use Avro or Protobuf with schema registry for streaming data
Multi-Tenant Data
- Tenant ID in every table, every query, every log
- Row-level security in warehouse
- Separate compute per tenant (or at least isolation)
- Never join across tenants without explicit business reason
- Tenant-aware backfill (don't rebuild all tenants for one tenant's issue)
Data Lake Anti-Patterns
- "Data Swamp": ingesting everything with no organization or catalog → only ingest what has a known consumer
- Small files: thousands of <1MB files → compact regularly (target 100MB-1GB)
- No table format: raw Parquet/CSV without Delta/Iceberg → loses ACID, schema evolution, time travel
- No access controls: single bucket, everyone admin → implement IAM per domain/team
Natural Language Commands
Say any of these to activate specific workflows:
- "Design a data pipeline for [source] to [target]" → Full pipeline template with extraction strategy, transforms, load pattern, quality checks
- "Model [entity/domain] for analytics" → Dimensional model with fact/dimension tables, grain, measures, SCD types
- "Optimize this query/pipeline" → Performance analysis with specific recommendations
- "Set up data quality for [table/pipeline]" → Quality framework with checks, contracts, monitoring
- "Audit our data infrastructure" → Full assessment using scoring rubric
- "Help with [Spark/Airflow/dbt/Kafka] issue" → Troubleshooting with technology-specific guidance
- "Design a data catalog for our org" → Catalog template with governance, classification, lineage
- "Plan a data migration from [old] to [new]" → Migration plan with validation, rollback, parallel-run
- "Set up monitoring for our pipelines" → Dashboard template with alerts, logging standards, runbooks
- "Review our data costs" → Cost analysis with optimization strategies and ROI framework
- "Handle schema change in [source]" → Change management protocol with impact assessment
- "Backfill [table] for [date range]" → Backfill protocol with validation and communication plan
1---2name: data-engineering-command-center3description: Complete methodology for designing, building, operating, and scaling data pipelines and infrastructure. Zero dependencies — pure agent skill.4---5
6# Data Engineering Command Center
7
8Complete methodology for designing, building, operating, and scaling data pipelines and infrastructure. Zero dependencies — pure agent skill.
9
10---
11
12## Phase 1: Data Architecture Assessment
13
14Before building anything, understand the landscape.
15
16### Architecture Brief
17
18```yaml
19project_name: ""
20business_context: ""
21data_consumers:
22 - team: ""
23 use_case: "" # analytics | ML | operational | reporting | reverse-ETL
24 latency_requirement: "" # real-time (<1s) | near-real-time (<5min) | batch (hourly+)
25 query_pattern: "" # ad-hoc | scheduled | API | dashboard
26
27current_state:
28 sources: [] # list every system producing data
29 storage: [] # where data lives today
30 pain_points: [] # what's broken, slow, unreliable
31 data_volume:
32 current_gb_per_day: 0
33 growth_rate_percent: 0
34 retention_months: 0
35
36constraints:
37 budget_monthly_usd: 0
38 team_size: 0
39 skill_level: "" # junior | mid | senior | mixed
40 compliance: [] # GDPR, HIPAA, SOX, PCI, none
41 cloud_provider: "" # AWS | GCP | Azure | multi | on-prem
42```
43
44### Architecture Pattern Decision Matrix
45
46| Signal | Pattern | When to Use |
47|--------|---------|-------------|
48| All consumers need data hourly+ | **Batch ETL** | Reporting, warehousing, most analytics |
49| Some need <5 min latency | **Micro-batch** | Dashboard freshness, near-real-time analytics |
50| Events need <1s processing | **Streaming** | Fraud detection, real-time pricing, alerts |
51| Need both batch + streaming | **Lambda** | When batch accuracy + real-time speed both matter |
52| Want to simplify Lambda | **Kappa** | When you can reprocess from stream replay |
53| Data lake + warehouse combined | **Lakehouse** | When you need both cheap storage + fast SQL |
54| Sources change independently | **Data Mesh** | Large orgs, domain-owned data, >5 teams |
55| ML is primary consumer | **Feature Store** | ML-heavy orgs with feature reuse needs |
56
57### Technology Selection Guide
58
59#### Orchestration
60
61| Tool | Best For | Avoid When |
62|------|----------|------------|
63| **Airflow** | Complex DAGs, Python-native teams, mature ecosystem | Simple pipelines (<5 tasks) |
64| **Dagster** | Software-defined assets, strong typing, dev experience | Legacy team resistant to new paradigms |
65| **Prefect** | Dynamic workflows, cloud-native, Python-first | Need on-prem with no cloud dependency |
66| **dbt** | SQL transformations, ELT, analytics engineering | Non-SQL transforms, streaming |
67| **Temporal** | Long-running workflows, retry-heavy, microservices | Simple ETL, small teams |
68| **Cron + scripts** | <3 pipelines, solo engineer, simple schedules | Anything with dependencies or retries |
69
70#### Processing
71
72| Tool | Best For | Avoid When |
73|------|----------|------------|
74| **Spark** | >100GB, complex transforms, ML pipelines | <10GB (overkill), real-time streaming |
75| **DuckDB** | Local analytics, <100GB, SQL on files | Distributed processing, production streaming |
76| **Polars** | Single-node, Rust-speed, <50GB, DataFrames | Distributed, need Spark ecosystem |
77| **Pandas** | <1GB, quick analysis, prototyping | Production pipelines, anything >5GB |
78| **Flink** | True streaming, event-time processing | Batch-only, small team (steep learning curve) |
79| **SQL (warehouse)** | ELT in Snowflake/BigQuery/Redshift | Complex ML transforms, binary data |
80
81#### Storage
82
83| Tool | Best For | Avoid When |
84|------|----------|------------|
85| **Snowflake** | Analytics, separation of compute/storage, multi-cloud | Tight budget, real-time OLTP |
86| **BigQuery** | GCP-native, serverless, large-scale analytics | Multi-cloud, need fine-grained cost control |
87| **Redshift** | AWS-native, existing AWS ecosystem | Elastic scaling needs, multi-cloud |
88| **Databricks** | ML + analytics unified, Spark-native, lakehouse | Pure SQL analytics, small data |
89| **PostgreSQL** | OLTP + light analytics, <500GB, budget-conscious | >1TB analytics, real-time dashboards on large data |
90| **S3/GCS/ADLS** | Raw data lake, cheap storage, any format | Direct SQL queries (need compute layer) |
91| **Delta Lake/Iceberg** | Table format on data lake, ACID on files | Simple file storage, no lakehouse need |
92
93---
94
95## Phase 2: Data Modeling
96
97### Modeling Methodology Decision
98
99| Approach | Best For | Key Concept |
100|----------|----------|-------------|
101| **Kimball (Dimensional)** | BI/reporting, star schemas | Facts + Dimensions, business-process-centric |
102| **Inmon (3NF)** | Enterprise data warehouse, single source of truth | Normalized, subject-area-centric |
103| **Data Vault 2.0** | Agile warehousing, auditability, multiple sources | Hubs + Links + Satellites, insert-only |
104| **One Big Table (OBT)** | Simple analytics, few joins, dashboard performance | Pre-joined, denormalized, fast queries |
105| **Activity Schema** | Event analytics, product analytics | Entity + Activity + Feature columns |
106
107### Dimensional Model Template
108
109```yaml
110fact_table:
111 name: "fact_[business_process]"
112 grain: "" # one row = one [what]?
113 grain_statement: "One row per [transaction/event/snapshot] at [time grain]"
114 measures:
115 - name: ""
116 type: "" # additive | semi-additive | non-additive
117 aggregation: "" # SUM | AVG | COUNT | MIN | MAX | COUNT DISTINCT
118 business_definition: ""
119 degenerate_dimensions: [] # IDs stored in fact (order_number, invoice_id)
120 foreign_keys: [] # links to dimension tables
121
122dimensions:
123 - name: "dim_[entity]"
124 type: "" # Type 1 (overwrite) | Type 2 (history) | Type 3 (previous value)
125 natural_key: "" # business key from source
126 surrogate_key: "" # warehouse-generated key
127 attributes:
128 - name: ""
129 source: ""
130 scd_type: "" # 1 | 2 | 3
131 hierarchy: [] # e.g., [country, region, city, store]
132```
133
134### SCD Type Decision Guide
135
136| Scenario | SCD Type | Implementation |
137|----------|----------|----------------|
138| Don't care about history | **Type 1** | UPDATE in place |
139| Need full history | **Type 2** | New row + valid_from/valid_to + is_current flag |
140| Only need previous value | **Type 3** | Add previous_[column] |
141| Track changes with timestamps | **Type 4** | Mini-dimension (history table) |
142| Hybrid: some attrs Type 1, some Type 2 | **Type 6** | Combine 1+2+3 in one table |
143
144**Default recommendation:** Type 2 for anything business-critical (customer status, product price, employee department). Type 1 for everything else.
145
146### Naming Conventions
147
148| Object | Convention | Example |
149|--------|-----------|---------|
150| Raw/staging tables | `raw_[source]_[table]` | `raw_stripe_payments` |
151| Staging models | `stg_[source]__[entity]` | `stg_stripe__payments` |
152| Intermediate models | `int_[entity]_[verb]` | `int_orders_pivoted` |
153| Mart/fact tables | `fct_[business_process]` | `fct_orders` |
154| Dimension tables | `dim_[entity]` | `dim_customers` |
155| Metrics/aggregates | `mrt_[domain]_[metric]` | `mrt_sales_daily` |
156| Snapshots | `snp_[entity]_[grain]` | `snp_inventory_daily` |
157| Columns: boolean | `is_[state]` or `has_[thing]` | `is_active`, `has_subscription` |
158| Columns: timestamp | `[event]_at` | `created_at`, `shipped_at` |
159| Columns: date | `[event]_date` | `order_date` |
160| Columns: ID | `[entity]_id` | `customer_id` |
161| Columns: amount | `[thing]_amount` | `order_amount` |
162| Columns: count | `[thing]_count` | `line_item_count` |
163
164---
165
166## Phase 3: Pipeline Design Patterns
167
168### Universal Pipeline Template
169
170```yaml
171pipeline:
172 name: ""
173 owner: ""
174 schedule: "" # cron expression
175 sla_minutes: 0 # max acceptable runtime
176 tier: "" # 1 (critical) | 2 (important) | 3 (nice-to-have)
177
178 extract:
179 source_system: ""
180 connection: ""
181 strategy: "" # full | incremental | CDC | log-based
182 incremental_key: "" # column for incremental (e.g., updated_at)
183 watermark_storage: "" # where to persist last-extracted position
184
185 transform:
186 engine: "" # SQL | Spark | Python | dbt
187 stages:
188 - name: "clean"
189 operations: [] # dedupe, null handling, type casting, trimming
190 - name: "conform"
191 operations: [] # standardize codes, currencies, timezones
192 - name: "enrich"
193 operations: [] # lookups, calculations, derived fields
194 - name: "aggregate"
195 operations: [] # rollups, pivots, window functions
196
197 load:
198 target_system: ""
199 strategy: "" # append | upsert | merge | truncate-reload | partition-swap
200 merge_keys: []
201 partition_key: ""
202 clustering_keys: []
203
204 quality_gates:
205 pre_load: [] # checks before writing
206 post_load: [] # checks after writing
207
208 error_handling:
209 strategy: "" # fail-fast | dead-letter | retry | skip-and-alert
210 max_retries: 3
211 retry_delay_seconds: 300
212 alert_channels: []
213```
214
215### Extraction Strategy Decision Tree
216
217```
218Is the source database?
219├── Yes → Does it support CDC?
220│ ├── Yes → Use CDC (Debezium, AWS DMS, Fivetran)
221│ │ Best for: high-volume, low-latency, minimal source impact
222│ └── No → Does it have a reliable updated_at column?
223│ ├── Yes → Incremental extraction on updated_at
224│ │ ⚠️ Won't catch hard deletes — add periodic full reconciliation
225│ └── No → Full extraction
226│ Only viable for small tables (<1M rows)
227├── Is it an API?
228│ ├── Supports webhooks? → Event-driven ingestion
229│ ├── Has cursor/pagination? → Incremental with cursor bookmark
230│ └── No pagination? → Full pull with rate-limit handling
231├── Is it files (S3, SFTP, email)?
232│ └── Event-triggered (S3 notification, file watcher)
233│ Validate: schema, completeness, filename pattern
234└── Is it streaming (Kafka, Kinesis, Pub/Sub)?
235 └── Consumer group with offset management
236 Key decisions: at-least-once vs exactly-once, consumer lag alerting
237```
238
239### Load Strategy Decision
240
241| Strategy | When | Trade-off |
242|----------|------|-----------|
243| **Append** | Event/log data, immutable facts | Simple but grows forever — partition + retain |
244| **Upsert/Merge** | Dimension updates, SCD Type 1 | Handles updates but slower on large tables |
245| **Truncate-Reload** | Small tables (<1M), reference data | Simple but window of missing data |
246| **Partition Swap** | Large fact tables, daily loads | Atomic, fast, but needs partition alignment |
247| **Soft Delete** | Need audit trail of deletions | Adds complexity to every downstream query |
248
249### Idempotency Rules (NON-NEGOTIABLE)
250
251Every pipeline MUST be re-runnable without side effects:
252
2531. **Use MERGE/UPSERT, never blind INSERT** for mutable data
2542. **Partition-swap for immutable data** — drop partition + reload
2553. **Store watermarks externally** — not in the pipeline code
2564. **Dedup at ingestion** — use source natural keys
2575. **Test by running twice** — output must be identical both times
258
259---
260
261## Phase 4: Data Quality Framework
262
263### Quality Dimensions
264
265| Dimension | Definition | Example Check |
266|-----------|-----------|---------------|
267| **Completeness** | No missing values where required | `NOT NULL` on required fields, row count within range |
268| **Uniqueness** | No unexpected duplicates | Primary key uniqueness, natural key uniqueness |
269| **Validity** | Values within expected domain | Enum checks, range checks, regex patterns |
270| **Accuracy** | Data matches real-world truth | Cross-system reconciliation, manual spot checks |
271| **Freshness** | Data arrives on time | `MAX(loaded_at) > NOW() - INTERVAL '2 hours'` |
272| **Consistency** | Same data agrees across systems | Sum reconciliation between source and target |
273
274### Quality Check Templates
275
276```sql
277-- Completeness: Required fields not null
278SELECT COUNT(*) AS null_violations
279FROM {table}
280WHERE {required_column} IS NULL;
281-- Threshold: 0
282
283-- Uniqueness: No duplicate primary keys
284SELECT {pk_column}, COUNT(*) AS dupe_count
285FROM {table}
286GROUP BY {pk_column}
287HAVING COUNT(*) > 1;
288-- Threshold: 0
289
290-- Freshness: Data arrived within SLA
291SELECT CASE
292 WHEN MAX({timestamp_col}) > CURRENT_TIMESTAMP - INTERVAL '{sla_hours} hours'
293 THEN 'PASS' ELSE 'FAIL'
294END AS freshness_check
295FROM {table};
296
297-- Volume: Row count within expected range
298SELECT CASE
299 WHEN COUNT(*) BETWEEN {min_expected} AND {max_expected}
300 THEN 'PASS' ELSE 'FAIL'
301END AS volume_check
302FROM {table}
303WHERE {partition_col} = '{run_date}';
304
305-- Referential: FK integrity
306SELECT COUNT(*) AS orphan_count
307FROM {fact_table} f
308LEFT JOIN {dim_table} d ON f.{fk} = d.{pk}
309WHERE d.{pk} IS NULL;
310-- Threshold: 0
311
312-- Distribution: No unexpected skew
313SELECT {column}, COUNT(*) AS cnt,
314 ROUND(100.0 * COUNT(*) / SUM(COUNT(*)) OVER (), 2) AS pct
315FROM {table}
316GROUP BY {column}
317ORDER BY cnt DESC;
318-- Alert if any single value > {max_pct}%
319
320-- Cross-system reconciliation
321SELECT
322 (SELECT SUM(amount) FROM source_system.orders WHERE date = '{date}') AS source_total,
323 (SELECT SUM(amount) FROM warehouse.fct_orders WHERE order_date = '{date}') AS target_total,
324 ABS(source_total - target_total) AS variance;
325-- Threshold: variance < 0.01 * source_total (1%)
326```
327
328### Data Contract Template
329
330```yaml
331contract:
332 name: ""
333 version: ""
334 owner: "" # team responsible for producing this data
335 consumers: [] # teams consuming this data
336 sla:
337 freshness_hours: 0
338 availability_percent: 99.9
339 support_hours: "" # business-hours | 24x7
340
341 schema:
342 - column: ""
343 type: ""
344 nullable: false
345 description: ""
346 business_definition: ""
347 pii: false
348 checks:
349 - type: "" # not_null | unique | range | enum | regex | custom
350 params: {}
351
352 breaking_change_policy: "" # notify-30-days | version-bump | never-break
353 notification_channel: ""
354```
355
356### Quality Severity Levels
357
358| Level | Definition | Response |
359|-------|-----------|----------|
360| **P0 — Critical** | Data corruption, wrong numbers in production dashboards, compliance data wrong | Stop pipeline, alert immediately, rollback if possible |
361| **P1 — High** | Missing data for key reports, SLA breach, >5% of records affected | Alert team, fix within 4 hours, post-mortem required |
362| **P2 — Medium** | Non-critical field quality, <1% records affected, no downstream impact | Fix in next sprint, add monitoring to prevent recurrence |
363| **P3 — Low** | Cosmetic issues, edge cases, non-critical data | Backlog, fix when convenient |
364
365---
366
367## Phase 5: Performance Optimization
368
369### SQL Optimization Checklist
370
371| Problem | Fix | Impact |
372|---------|-----|--------|
373| Full table scan | Add/use partition pruning | 10-100x faster |
374| Large joins | Pre-aggregate before joining | 5-50x faster |
375| SELECT * | Select only needed columns | 2-10x faster (columnar stores) |
376| Correlated subquery | Rewrite as JOIN or window function | 10-100x faster |
377| DISTINCT on large result | Fix upstream duplication instead | 2-5x faster |
378| ORDER BY without LIMIT | Add LIMIT or remove if not needed | Prevents memory spills |
379| String operations in WHERE | Pre-compute, use lookup table | Enables index usage |
380| Multiple passes over same data | Combine with CASE WHEN + GROUP BY | 2-5x faster |
381| NOT IN with NULLs | Use NOT EXISTS or LEFT JOIN IS NULL | Correctness + performance |
382
383### Spark Optimization Guide
384
385| Problem | Solution |
386|---------|----------|
387| Shuffle-heavy joins | Broadcast small table (`broadcast(df)`) if <100MB |
388| Data skew | Salt the skewed key: add random prefix, join on salted key, aggregate |
389| Small files | Coalesce output: `.coalesce(target_files)` or use adaptive query execution |
390| Too many partitions | `spark.sql.shuffle.partitions` = 2-3x cluster cores |
391| OOM errors | Increase `spark.executor.memory`, reduce partition size, spill to disk |
392| Slow writes | Use Parquet with snappy, partition by date, avoid small writes |
393| Repeated computation | `.cache()` or `.persist()` DataFrames used >1 time |
394| Complex transformations | Push down predicates, filter early, select early |
395
396### Partitioning Strategy
397
398| Data Type | Partition Key | Why |
399|-----------|--------------|-----|
400| Transactional/event | Date (daily or monthly) | Most queries filter by time range |
401| Multi-tenant | Tenant ID + date | Isolate tenant queries, time-range pruning |
402| Geospatial | Region + date | Regional queries are common |
403| Log data | Date + hour | High volume needs finer partitions |
404| Reference/dimension | Don't partition | Too small, full scan is fine |
405
406**Rules:**
407- Target 100MB-1GB per partition (compressed)
408- <10,000 total partitions per table
409- Never partition on high-cardinality columns (user_id)
410- Always include partition key in WHERE clauses
411
412---
413
414## Phase 6: Data Governance & Cataloging
415
416### Data Classification
417
418| Level | Examples | Controls |
419|-------|---------|----------|
420| **Public** | Product catalog, published stats | No restrictions |
421| **Internal** | Aggregated metrics, non-PII analytics | Auth required, audit logging |
422| **Confidential** | Customer PII, financial records, HR data | Encryption, column-level access, masking |
423| **Restricted** | SSN, payment cards, health records, passwords | Encryption at rest + transit, tokenization, audit every access, retention limits |
424
425### PII Handling Rules
426
4271. **Identify:** Scan all sources for PII columns (name, email, phone, SSN, DOB, address, IP)
4282. **Classify:** Tag each with sensitivity level
4293. **Minimize:** Only ingest PII you actually need
4304. **Protect:**
431 - Hash or tokenize in staging (SHA-256 with salt for pseudonymization)
432 - Dynamic masking for non-privileged users
433 - Column-level encryption for restricted data
4345. **Retain:** Auto-delete after retention period
4356. **Audit:** Log every query touching PII columns
4367. **Right to delete:** Build a deletion pipeline that propagates across all derived tables
437
438### Data Catalog Entry Template
439
440```yaml
441dataset:
442 name: ""
443 description: ""
444 owner_team: ""
445 steward: "" # person responsible for quality
446 domain: "" # sales | marketing | finance | product | engineering
447 tier: "" # gold (trusted) | silver (cleaned) | bronze (raw)
448
449 lineage:
450 sources: [] # upstream datasets/systems
451 transformations: "" # brief description of key transforms
452 downstream: [] # who consumes this
453
454 refresh:
455 schedule: ""
456 sla_hours: 0
457 last_successful_run: ""
458
459 quality:
460 tests: [] # list of quality checks
461 last_score: 0 # 0-100
462 known_issues: []
463
464 access:
465 classification: "" # public | internal | confidential | restricted
466 pii_columns: []
467 access_request_process: "" # how to get access
468
469 usage:
470 avg_daily_queries: 0
471 top_consumers: []
472 cost_monthly_usd: 0
473```
474
475---
476
477## Phase 7: Pipeline Monitoring & Alerting
478
479### Pipeline Health Dashboard
480
481```yaml
482dashboard:
483 pipeline_metrics:
484 - metric: "Pipeline Success Rate"
485 formula: "successful_runs / total_runs * 100"
486 target: ">99%"
487 alert_threshold: "<95%"
488
489 - metric: "Average Runtime"
490 formula: "avg(end_time - start_time) over 7 days"
491 target: "<SLA"
492 alert_threshold: ">80% of SLA"
493
494 - metric: "Data Freshness"
495 formula: "NOW() - MAX(loaded_at)"
496 target: "<SLA hours"
497 alert_threshold: ">SLA"
498
499 - metric: "Data Volume Variance"
500 formula: "abs(today_rows - avg_7d_rows) / avg_7d_rows * 100"
501 target: "<20%"
502 alert_threshold: ">50%"
503
504 - metric: "Quality Check Pass Rate"
505 formula: "passed_checks / total_checks * 100"
506 target: "100%"
507 alert_threshold: "<95%"
508
509 - metric: "Failed Pipeline Count"
510 formula: "count where status = failed in last 24h"
511 target: "0"
512 alert_threshold: ">0"
513
514 - metric: "Backfill Queue"
515 formula: "count of pending backfill requests"
516 target: "0"
517 alert_threshold: ">5"
518
519 - metric: "Infrastructure Cost"
520 formula: "compute + storage + egress"
521 target: "<budget"
522 alert_threshold: ">110% budget"
523```
524
525### Alert Severity
526
527| Severity | Condition | Response Time | Example |
528|----------|-----------|---------------|---------|
529| **P0** | Revenue/compliance impacting | 15 min | Payment pipeline down, regulatory report delayed |
530| **P1** | Business-critical dashboard stale | 1 hour | Executive dashboard >4h stale |
531| **P2** | Non-critical pipeline failed | 4 hours | Marketing attribution delayed |
532| **P3** | Warning/degradation | Next business day | Pipeline 80% of SLA, minor quality drift |
533
534### Structured Logging Standard
535
536Every pipeline run MUST log:
537
538```json
539{
540 "pipeline_name": "",
541 "run_id": "",
542 "started_at": "",
543 "completed_at": "",
544 "status": "success|failed|partial",
545 "stage": "",
546 "rows_extracted": 0,
547 "rows_transformed": 0,
548 "rows_loaded": 0,
549 "rows_rejected": 0,
550 "quality_checks_passed": 0,
551 "quality_checks_failed": 0,
552 "duration_seconds": 0,
553 "error_message": "",
554 "watermark_before": "",
555 "watermark_after": ""
556}
557```
558
559---
560
561## Phase 8: Testing Strategy
562
563### Pipeline Test Pyramid
564
565| Layer | What to Test | How | When |
566|-------|-------------|-----|------|
567| **Unit** | Individual transforms, business logic | pytest with fixtures, dbt unit tests | Every PR |
568| **Integration** | Source connectivity, schema compatibility | Test against staging/dev environment | Daily + PR |
569| **Contract** | Schema hasn't changed, data types stable | Schema registry, contract tests | Every pipeline run |
570| **Data Quality** | Completeness, uniqueness, freshness, validity | Quality framework checks | Every pipeline run |
571| **E2E** | Full pipeline produces correct output | Golden dataset comparison | Weekly + release |
572| **Performance** | Runtime within SLA, no regression | Benchmark against historical runs | Weekly |
573
574### dbt Testing Checklist
575
576```yaml
577# For every model, define at minimum:
578models:
579 - name: fct_orders
580 columns:
581 - name: order_id
582 tests:
583 - unique
584 - not_null
585 - name: customer_id
586 tests:
587 - not_null
588 - relationships:
589 to: ref('dim_customers')
590 field: customer_id
591 - name: order_amount
592 tests:
593 - not_null
594 - dbt_utils.accepted_range:
595 min_value: 0
596 max_value: 1000000
597 - name: order_status
598 tests:
599 - accepted_values:
600 values: ['pending', 'confirmed', 'shipped', 'delivered', 'cancelled']
601 - name: ordered_at
602 tests:
603 - not_null
604 - dbt_utils.recency:
605 datepart: day
606 field: ordered_at
607 interval: 2
608```
609
610### Backfill Protocol
611
612When you need to reprocess historical data:
613
6141. **Scope:** Define exact date range and affected tables
6152. **Impact assessment:** What downstream models/dashboards will be affected?
6163. **Communication:** Notify consumers of temporary data inconsistency
6174. **Isolation:** Run backfill in separate compute to avoid impacting current pipelines
6185. **Validation:** Compare row counts and key metrics pre/post backfill
6196. **Execution:** Process in reverse-chronological order (most recent first)
6207. **Monitoring:** Watch for resource spikes, duplicate creation
6218. **Verification:** Reconcile against source after completion
6229. **Documentation:** Log what was backfilled, why, and any anomalies found
623
624---
625
626## Phase 9: Cost Optimization
627
628### Cloud Cost Reduction Strategies
629
630| Strategy | Savings | Effort |
631|----------|---------|--------|
632| Right-size compute (auto-scaling) | 20-40% | Low |
633| Use spot/preemptible instances for batch | 60-80% | Medium |
634| Compress data (Parquet + Snappy/Zstd) | 50-80% storage | Low |
635| Lifecycle policies (hot → warm → cold → archive) | 40-70% storage | Low |
636| Eliminate unused tables/pipelines | 10-30% | Low |
637| Optimize query patterns (partition pruning) | 30-60% compute | Medium |
638| Reserved capacity for steady-state | 30-50% | Medium |
639| Cache expensive queries | 20-50% compute | Medium |
640
641### Cost Allocation Template
642
643```yaml
644cost_tracking:
645 by_pipeline:
646 - pipeline: ""
647 compute_monthly_usd: 0
648 storage_monthly_usd: 0
649 egress_monthly_usd: 0
650 total: 0
651 cost_per_row: 0 # total / rows_processed
652 business_value: "" # what revenue/decision does this enable?
653 roi_justified: true # is the cost worth it?
654
655 optimization_opportunities:
656 - description: ""
657 estimated_savings_usd: 0
658 effort: "" # low | medium | high
659 priority: 0 # 1 = do now
660```
661
662### Cost Red Flags
663
664- Single pipeline >30% of total spend
665- Cost per row increasing month-over-month
666- Tables with 0 queries in 30 days
667- Dev/staging environments running 24/7
668- Full table scans on >1TB tables
669- Uncompressed data in cloud storage
670- Cross-region data transfer
671
672---
673
674## Phase 10: Operational Runbooks
675
676### Pipeline Failure Triage
677
678```
679Pipeline failed →
6801. Check error message in logs
681 ├── Connection timeout → Check source availability, network, credentials
682 ├── Schema mismatch → Source schema changed → update extract + notify
683 ├── Data quality check failed → Investigate source data, check for anomalies
684 ├── Out of memory → Increase resources or optimize query
685 ├── Permission denied → Check IAM roles, token expiry
686 ├── Duplicate key violation → Check idempotency, investigate source dupes
687 └── Timeout (SLA breach) → Check data volume spike, query plan, cluster health
688
6892. Determine impact
690 ├── What dashboards/reports are affected?
691 ├── What's the data freshness SLA?
692 └── Who needs to be notified?
693
6943. Fix
695 ├── Transient (network, timeout) → Retry
696 ├── Data issue → Fix source data, re-run with quality gate override if safe
697 ├── Schema change → Update pipeline, backfill if needed
698 └── Infrastructure → Scale up, file ticket with cloud provider
699
7004. Post-fix
701 ├── Verify data correctness
702 ├── Update runbook with new failure mode
703 └── Add monitoring/alerting to catch earlier next time
704```
705
706### Schema Change Management
707
708When a source system changes schema:
709
7101. **Detect:** Schema comparison check in extraction pipeline (hash schema, compare to registered)
7112. **Classify:**
712 - **Additive** (new column): Usually safe — add to pipeline, backfill if needed
713 - **Rename**: Map old → new in transform, update downstream
714 - **Type change**: Assess compatibility, may need cast or historical rebuild
715 - **Column removed**: Critical — breaks queries, need immediate attention
7163. **Test:** Run pipeline in dry-run mode with new schema
7174. **Deploy:** Update transforms, quality checks, documentation
7185. **Communicate:** Notify downstream consumers via data contract channel
719
720### Disaster Recovery
721
722| Scenario | RPO | RTO | Recovery Steps |
723|----------|-----|-----|----------------|
724| Pipeline code lost | 0 (git) | 1h | Redeploy from git, restore orchestrator state |
725| Warehouse data corrupted | Varies | 4h | Restore from Time Travel/snapshot, re-run affected pipelines |
726| Source system down | N/A | Wait | Queue extractions, catch up with incremental once restored |
727| Cloud region outage | 24h | 8h | Failover to DR region if configured, else wait |
728| Credential compromise | 0 | 2h | Rotate all credentials, audit access logs, re-run affected pipelines |
729
730---
731
732## Phase 11: Advanced Patterns
733
734### Slowly Changing Dimension Type 2 (SQL Template)
735
736```sql
737-- Merge pattern for SCD Type 2
738MERGE INTO dim_customer AS target
739USING (
740 SELECT * FROM stg_customers
741 WHERE updated_at > (SELECT MAX(valid_from) FROM dim_customer)
742) AS source
743ON target.customer_natural_key = source.customer_id
744 AND target.is_current = TRUE
745
746-- Update: close old record
747WHEN MATCHED AND (
748 target.customer_name != source.name OR
749 target.customer_status != source.status
750 -- list all Type 2 tracked columns
751) THEN UPDATE SET
752 is_current = FALSE,
753 valid_to = CURRENT_TIMESTAMP
754
755-- Insert: new record (both new customers and changed ones)
756WHEN NOT MATCHED THEN INSERT (
757 customer_natural_key, customer_name, customer_status,
758 valid_from, valid_to, is_current
759) VALUES (
760 source.customer_id, source.name, source.status,
761 CURRENT_TIMESTAMP, '9999-12-31', TRUE
762);
763
764-- Then insert new versions of changed records
765INSERT INTO dim_customer (
766 customer_natural_key, customer_name, customer_status,
767 valid_from, valid_to, is_current
768)
769SELECT customer_id, name, status,
770 CURRENT_TIMESTAMP, '9999-12-31', TRUE
771FROM stg_customers s
772WHERE EXISTS (
773 SELECT 1 FROM dim_customer d
774 WHERE d.customer_natural_key = s.customer_id
775 AND d.is_current = FALSE
776 AND d.valid_to = CURRENT_TIMESTAMP
777);
778```
779
780### CDC with Debezium (Architecture Pattern)
781
782```
783Source DB → Debezium Connector → Kafka Topic →
784 ├── Stream processor (Flink/Spark Streaming) → Target DB
785 ├── S3 sink connector → Data Lake (raw)
786 └── Elasticsearch sink → Search index
787```
788
789Key decisions:
790- **Topic per table** or **single topic**: Per table (easier routing, independent scaling)
791- **Schema registry**: Always use (Confluent Schema Registry or AWS Glue)
792- **Serialization**: Avro (compact + schema evolution) or Protobuf (strict + fast)
793- **Offset management**: Connector manages; monitor consumer lag
794
795### Feature Store Pattern
796
797```yaml
798feature_store:
799 entity: "customer"
800 entity_key: "customer_id"
801
802 features:
803 - name: "total_orders_30d"
804 description: "Total orders in last 30 days"
805 type: "INT"
806 source: "fct_orders"
807 computation: "batch" # batch | streaming | on-demand
808 freshness: "daily"
809 ttl_hours: 48
810
811 - name: "avg_order_value_90d"
812 description: "Average order value last 90 days"
813 type: "FLOAT"
814 source: "fct_orders"
815 computation: "batch"
816 freshness: "daily"
817 ttl_hours: 48
818
819 - name: "last_login_minutes_ago"
820 description: "Minutes since last login event"
821 type: "INT"
822 source: "events_stream"
823 computation: "streaming"
824 freshness: "real-time"
825 ttl_hours: 1
826
827 serving:
828 online: true # low-latency feature serving (Redis/DynamoDB)
829 offline: true # batch feature retrieval for training
830 point_in_time_correct: true # prevent feature leakage in ML training
831```
832
833### Data Mesh Principles
834
835If operating at scale (>5 data teams):
836
8371. **Domain ownership**: Each business domain owns its data products (not central data team)
8382. **Data as a product**: Treat datasets like products — SLAs, documentation, discoverability
8393. **Self-serve platform**: Central team builds the platform, domains build on top
8404. **Federated governance**: Standards and interoperability maintained centrally, implementation decentralized
841
842**When NOT to use Data Mesh:**
843- <5 data producers/consumers
844- Small team (<20 engineers total)
845- Single business domain
846- Early-stage company (over-engineering)
847
848---
849
850## Quality Scoring Rubric (0-100)
851
852| Dimension | Weight | Scoring |
853|-----------|--------|---------|
854| **Pipeline Reliability** | 20 | 0=frequent failures, 10=some failures with manual recovery, 20=99.5%+ success rate with auto-retry |
855| **Data Quality** | 20 | 0=no checks, 10=basic null/unique checks, 20=comprehensive quality framework with contracts |
856| **Performance** | 15 | 0=regularly breaches SLA, 8=meets SLA, 15=well under SLA with optimization |
857| **Documentation** | 10 | 0=none, 5=basic README, 10=full catalog entries with lineage and business definitions |
858| **Monitoring** | 15 | 0=no alerts, 8=failure alerts only, 15=proactive monitoring with dashboards and anomaly detection |
859| **Testing** | 10 | 0=no tests, 5=basic smoke tests, 10=full test pyramid (unit+integration+contract+E2E) |
860| **Cost Efficiency** | 10 | 0=no cost tracking, 5=tracked, 10=optimized with ROI justification per pipeline |
861
862**Scoring guide:**
863- 0-40: Critical gaps — prioritize pipeline reliability and data quality
864- 41-60: Functional but fragile — add monitoring, testing, documentation
865- 61-80: Solid — optimize performance, cost, governance
866- 81-100: Excellent — maintain, innovate, mentor
867
868---
869
870## Edge Cases & Gotchas
871
872### Timezone Traps
873- Store everything in UTC. Convert only at presentation layer
874- Event timestamps: use event time, not processing time
875- Daylight saving: `TIMESTAMP WITH TIME ZONE`, never `WITHOUT`
876- Late-arriving data: watermark strategy + allowed lateness window
877
878### Late-Arriving Data
879- Define maximum acceptable lateness per source
880- Reprocess affected partitions when late data arrives
881- Track late arrival rate as a quality metric
882- Consider separate "late data" pipeline that patches in
883
884### Exactly-Once Processing
885- True exactly-once is expensive. Most systems need at-least-once + idempotent writes
886- Use transaction IDs or natural keys for deduplication
887- Kafka: use idempotent producer + transactional consumer
888- Database: MERGE/UPSERT on natural key
889
890### Schema Evolution
891- **Forward compatible**: New code reads old data (safe to deploy new readers first)
892- **Backward compatible**: Old code reads new data (safe to deploy new writers first)
893- **Full compatible**: Both directions (safest, most restrictive)
894- Use Avro or Protobuf with schema registry for streaming data
895
896### Multi-Tenant Data
897- Tenant ID in every table, every query, every log
898- Row-level security in warehouse
899- Separate compute per tenant (or at least isolation)
900- Never join across tenants without explicit business reason
901- Tenant-aware backfill (don't rebuild all tenants for one tenant's issue)
902
903### Data Lake Anti-Patterns
904- "Data Swamp": ingesting everything with no organization or catalog → only ingest what has a known consumer
905- Small files: thousands of <1MB files → compact regularly (target 100MB-1GB)
906- No table format: raw Parquet/CSV without Delta/Iceberg → loses ACID, schema evolution, time travel
907- No access controls: single bucket, everyone admin → implement IAM per domain/team
908
909---
910
911## Natural Language Commands
912
913Say any of these to activate specific workflows:
914
9151. **"Design a data pipeline for [source] to [target]"** → Full pipeline template with extraction strategy, transforms, load pattern, quality checks
9162. **"Model [entity/domain] for analytics"** → Dimensional model with fact/dimension tables, grain, measures, SCD types
9173. **"Optimize this query/pipeline"** → Performance analysis with specific recommendations
9184. **"Set up data quality for [table/pipeline]"** → Quality framework with checks, contracts, monitoring
9195. **"Audit our data infrastructure"** → Full assessment using scoring rubric
9206. **"Help with [Spark/Airflow/dbt/Kafka] issue"** → Troubleshooting with technology-specific guidance
9217. **"Design a data catalog for our org"** → Catalog template with governance, classification, lineage
9228. **"Plan a data migration from [old] to [new]"** → Migration plan with validation, rollback, parallel-run
9239. **"Set up monitoring for our pipelines"** → Dashboard template with alerts, logging standards, runbooks
92410. **"Review our data costs"** → Cost analysis with optimization strategies and ROI framework
92511. **"Handle schema change in [source]"** → Change management protocol with impact assessment
92612. **"Backfill [table] for [date range]"** → Backfill protocol with validation and communication plan