Data Pipeline Engineer
Expert data engineer specializing in ETL/ELT pipelines, streaming architectures, data warehousing, and modern data stack implementation.
Quick Start
- Identify sources - data formats, volumes, freshness requirements
- Choose architecture - Medallion (Bronze/Silver/Gold), Lambda, or Kappa
- Design layers - staging → intermediate → marts (dbt pattern)
- Add quality gates - Great Expectations or dbt tests at each layer
- Orchestrate - Airflow DAGs with sensors and retries
- Monitor - lineage, freshness, anomaly detection
Core Capabilities
| Capability |
Technologies |
Key Patterns |
| Batch Processing |
Spark, dbt, Databricks |
Incremental, partitioning, Delta/Iceberg |
| Stream Processing |
Kafka, Flink, Spark Streaming |
Watermarks, exactly-once, windowing |
| Orchestration |
Airflow, Dagster, Prefect |
DAG design, sensors, task groups |
| Data Modeling |
dbt, SQL |
Kimball, Data Vault, SCD |
| Data Quality |
Great Expectations, dbt tests |
Validation suites, freshness |
Architecture Patterns
Medallion Architecture (Recommended)
BRONZE (Raw) → Exact source copy, schema-on-read, partitioned by ingestion
↓ Cleaning, Deduplication
SILVER (Cleansed) → Validated, standardized, business logic applied
↓ Aggregation, Enrichment
GOLD (Business) → Dimensional models, aggregates, ready for BI/ML
Lambda vs Kappa
- Lambda: Batch + Stream layers → merged serving layer (complex but complete)
- Kappa: Stream-only with replay → simpler but requires robust streaming
Reference Examples
Full implementation examples in ./references/:
| File |
Description |
dbt-project-structure.md |
Complete dbt layout with staging, intermediate, marts |
airflow-dag.py |
Production DAG with sensors, task groups, quality checks |
spark-streaming.py |
Kafka-to-Delta processor with windowing |
great-expectations-suite.json |
Comprehensive data quality expectation suite |
Anti-Patterns (10 Critical Mistakes)
1. Full Table Refreshes
Symptom: Truncate and rebuild entire tables every run
Fix: Use incremental models with is_incremental(), partition by date
2. Tight Coupling to Source Schemas
Symptom: Pipeline breaks when upstream adds/removes columns
Fix: Explicit source contracts, select only needed columns in staging
3. Monolithic DAGs
Symptom: One 200-task DAG running 8 hours
Fix: Domain-specific DAGs, ExternalTaskSensor for dependencies
4. No Data Quality Gates
Symptom: Bad data reaches production before detection
Fix: Great Expectations or dbt tests at each layer, block on failures
5. Processing Before Archiving
Symptom: Raw data transformed without preserving original
Fix: Always land raw in Bronze first, make transformations reproducible
6. Hardcoded Dates in Queries
Symptom: Manual updates needed for date filters
Fix: Use Airflow templating (e.g., ds variable) or dynamic date functions
7. Missing Watermarks in Streaming
Symptom: Unbounded state growth, OOM in long-running jobs
Fix: Add withWatermark() to handle late-arriving data
8. No Retry/Backoff Strategy
Symptom: Transient failures cause DAG failures
Fix: retries=3, retry_exponential_backoff=True, max_retry_delay
9. Undocumented Data Lineage
Symptom: No one knows where data comes from or who uses it
Fix: dbt docs, data catalog integration, column-level lineage
10. Testing Only in Production
Symptom: Bugs discovered by stakeholders, not engineers
Fix: dbt --target dev, sample datasets, CI/CD for models
Quality Checklist
Pipeline Design:
Data Quality:
Orchestration:
Operations:
Validation Script
Run ./scripts/validate-pipeline.sh to check:
- dbt project structure and conventions
- Airflow DAG best practices
- Spark job configurations
- Data quality setup
External Resources
1---2name: data-pipeline-engineer3description: Expert data engineer for ETL/ELT pipelines, streaming, data warehousing. Activate on: data pipeline, ETL, ELT, data warehouse, Spark, Kafka, Airflow, dbt, data modeling, star schema, streaming data, batch processing, data quality. NOT for: API design (use api-architect), ML training (use ML skills), dashboards (use design skills).4license: Apache-2.05---6
7# Data Pipeline Engineer
8
9Expert data engineer specializing in ETL/ELT pipelines, streaming architectures, data warehousing, and modern data stack implementation.
10
11## Quick Start
12
131. **Identify sources** - data formats, volumes, freshness requirements
142. **Choose architecture** - Medallion (Bronze/Silver/Gold), Lambda, or Kappa
153. **Design layers** - staging → intermediate → marts (dbt pattern)
164. **Add quality gates** - Great Expectations or dbt tests at each layer
175. **Orchestrate** - Airflow DAGs with sensors and retries
186. **Monitor** - lineage, freshness, anomaly detection
19
20## Core Capabilities
21
22| Capability | Technologies | Key Patterns |
23|------------|--------------|--------------|
24| **Batch Processing** | Spark, dbt, Databricks | Incremental, partitioning, Delta/Iceberg |
25| **Stream Processing** | Kafka, Flink, Spark Streaming | Watermarks, exactly-once, windowing |
26| **Orchestration** | Airflow, Dagster, Prefect | DAG design, sensors, task groups |
27| **Data Modeling** | dbt, SQL | Kimball, Data Vault, SCD |
28| **Data Quality** | Great Expectations, dbt tests | Validation suites, freshness |
29
30## Architecture Patterns
31
32### Medallion Architecture (Recommended)
33```
34BRONZE (Raw) → Exact source copy, schema-on-read, partitioned by ingestion
35 ↓ Cleaning, Deduplication
36SILVER (Cleansed) → Validated, standardized, business logic applied
37 ↓ Aggregation, Enrichment
38GOLD (Business) → Dimensional models, aggregates, ready for BI/ML
39```
40
41### Lambda vs Kappa
42- **Lambda**: Batch + Stream layers → merged serving layer (complex but complete)
43- **Kappa**: Stream-only with replay → simpler but requires robust streaming
44
45## Reference Examples
46
47Full implementation examples in `./references/`:
48
49| File | Description |
50|------|-------------|
51| `dbt-project-structure.md` | Complete dbt layout with staging, intermediate, marts |
52| `airflow-dag.py` | Production DAG with sensors, task groups, quality checks |
53| `spark-streaming.py` | Kafka-to-Delta processor with windowing |
54| `great-expectations-suite.json` | Comprehensive data quality expectation suite |
55
56## Anti-Patterns (10 Critical Mistakes)
57
58### 1. Full Table Refreshes
59**Symptom**: Truncate and rebuild entire tables every run
60**Fix**: Use incremental models with `is_incremental()`, partition by date
61
62### 2. Tight Coupling to Source Schemas
63**Symptom**: Pipeline breaks when upstream adds/removes columns
64**Fix**: Explicit source contracts, select only needed columns in staging
65
66### 3. Monolithic DAGs
67**Symptom**: One 200-task DAG running 8 hours
68**Fix**: Domain-specific DAGs, ExternalTaskSensor for dependencies
69
70### 4. No Data Quality Gates
71**Symptom**: Bad data reaches production before detection
72**Fix**: Great Expectations or dbt tests at each layer, block on failures
73
74### 5. Processing Before Archiving
75**Symptom**: Raw data transformed without preserving original
76**Fix**: Always land raw in Bronze first, make transformations reproducible
77
78### 6. Hardcoded Dates in Queries
79**Symptom**: Manual updates needed for date filters
80**Fix**: Use Airflow templating (e.g., `ds` variable) or dynamic date functions
81
82### 7. Missing Watermarks in Streaming
83**Symptom**: Unbounded state growth, OOM in long-running jobs
84**Fix**: Add `withWatermark()` to handle late-arriving data
85
86### 8. No Retry/Backoff Strategy
87**Symptom**: Transient failures cause DAG failures
88**Fix**: `retries=3`, `retry_exponential_backoff=True`, `max_retry_delay`
89
90### 9. Undocumented Data Lineage
91**Symptom**: No one knows where data comes from or who uses it
92**Fix**: dbt docs, data catalog integration, column-level lineage
93
94### 10. Testing Only in Production
95**Symptom**: Bugs discovered by stakeholders, not engineers
96**Fix**: dbt `--target dev`, sample datasets, CI/CD for models
97
98## Quality Checklist
99
100**Pipeline Design:**
101- [ ] Incremental processing where possible
102- [ ] Idempotent transformations (re-runnable safely)
103- [ ] Partitioning strategy defined and documented
104- [ ] Backfill procedures documented
105
106**Data Quality:**
107- [ ] Tests at Bronze layer (schema, nulls, ranges)
108- [ ] Tests at Silver layer (business rules, referential integrity)
109- [ ] Tests at Gold layer (aggregation checks, trend monitoring)
110- [ ] Anomaly detection for volumes and distributions
111
112**Orchestration:**
113- [ ] Retry and alerting configured
114- [ ] SLAs defined and monitored
115- [ ] Cross-DAG dependencies use sensors
116- [ ] max_active_runs prevents parallel conflicts
117
118**Operations:**
119- [ ] Data lineage documented
120- [ ] Runbooks for common failures
121- [ ] Monitoring dashboards for pipeline health
122- [ ] On-call procedures defined
123
124## Validation Script
125
126Run `./scripts/validate-pipeline.sh` to check:
127- dbt project structure and conventions
128- Airflow DAG best practices
129- Spark job configurations
130- Data quality setup
131
132## External Resources
133
134- [dbt Best Practices](https://docs.getdbt.com/guides/best-practices)
135- [Airflow Best Practices](https://airflow.apache.org/docs/apache-airflow/stable/best-practices.html)
136- [Great Expectations Docs](https://docs.greatexpectations.io/)
137- [Delta Lake Guide](https://docs.delta.io/latest/index.html)
138- [Kafka Streams](https://kafka.apache.org/documentation/streams/)