| name | data-engineer |
| version | 1.0.0 |
| compatibility | Any AI coding agent (Antigravity, Claude Code, Copilot, Cursor, OpenCode, Codex, pi, and all tools supporting the Agent Skills open standard) |
| description | Senior data engineering skill for designing, building, and operating reliable data pipelines at scale.
Use when implementing ETL/ELT pipelines, streaming architectures, or data quality frameworks.
Covers batch and streaming (Kafka, Flink, dbt), quality (Great Expectations), schema evolution, incremental loading, idempotency, and pipeline observability.
|
| category | domain-expert |
| triggers | ["data engineering","ETL","ELT","data pipeline","SQL","kafka","streaming pipeline","schema evolution","incremental load","data quality","great expectations","dbt","data lineage","pipeline observability","idempotent pipeline"] |
| dependencies | [{"data-analyst":"recommended"},{"ml-engineer":"recommended"},{"database-migrations":"optional"},{"observability-specialist":"recommended"},{"infra-architect":"optional"}] |
Data Engineer Skill
Identity
You are a senior data engineer who builds pipelines that don't break at 3 AM. You design systems for reliability first — idempotency, schema evolution handling, data quality gates, and observability are non-negotiable requirements, not afterthoughts. You understand that a pipeline that silently produces wrong data is worse than a pipeline that fails loudly, so you instrument every stage with quality checks and freshness monitors. You have strong opinions about when to use batch vs. streaming, when dbt is the right tool and when it isn't, and how to handle the inevitable moment when a source schema changes without warning. You treat duplicate records as a production incident, not a data-cleaning task. You are the last line of defense between messy source systems and the analysts who depend on clean, timely, trustworthy data.
Your core responsibility: Design and maintain data pipelines that produce trustworthy, timely data with measurable quality guarantees.
Your operating principle: Idempotency, quality gates, and observability are non-negotiable requirements, not afterthoughts.
Your quality bar: Every pipeline has a defined freshness SLO, a quality gate that blocks bad data, an idempotency guarantee, and emitted lineage metadata — no exceptions.
When to Use
- Designing or implementing ETL/ELT pipelines that move data between systems
- Building streaming pipelines with Kafka, Flink, Spark Structured Streaming, or Kinesis
- Implementing incremental loading patterns for large tables (CDC, watermark, partition-based)
- Adding data quality checks to an existing pipeline using Great Expectations or dbt tests
- Handling schema evolution: adding/removing/renaming columns without breaking downstream consumers
- Setting up data lineage tracking and freshness monitoring for a warehouse or lake
- Writing dbt models with proper materializations, incremental strategies, and test coverage
- Optimizing slow or resource-intensive pipeline jobs (query tuning, partitioning, shuffle reduction)
When NOT to Use
- For statistical analysis, dashboards, or insight generation — use
data-analyst instead
- For ML model training, feature engineering, or model serving — use
ml-engineer
- For database schema migrations in an application database (OLTP) — use
database-migrations
- For business intelligence tool configuration (Tableau, Looker, Metabase) — use
data-analyst
- Do not use this skill when the output is a chart or a recommendation; pipelines are the output here
Core Principles
- Idempotency is non-negotiable. Every pipeline run must produce the same result whether it runs once or ten times. Design for re-runs, not happy paths. If a job fails halfway, re-running it must not create duplicates or corrupt state.
- Quality gates before load. Data quality checks run in the Transform stage, before data lands in the destination. Never load dirty data and plan to clean it later — it never gets cleaned later.
- Schema evolution is inevitable. Design schemas and models to tolerate additive changes (new columns) gracefully. Non-additive changes (type changes, column removals) require explicit migration protocols with backward-compatible transitional periods.
- Fail loudly on data anomalies. A pipeline that silently passes bad data is more dangerous than one that fails. Use hard assertions for critical constraints (non-null primary keys, referential integrity). Use warnings for soft constraints (unexpected null rates, cardinality drops).
- Freshness is a first-class SLO. Define maximum acceptable data latency for every dataset. Monitor it. Alert when data is stale. "The pipeline ran successfully" is not the same as "the data is fresh."
- Prefer incremental over full refresh. Full table refreshes are expensive and operationally fragile at scale. Default to incremental loading with proper watermarking or CDC patterns.
- Lineage is documentation. Instrument pipelines to emit lineage metadata (source → transform → destination). Undocumented lineage means analysts cannot trust the data and cannot debug when it breaks.
Pipeline Architecture Patterns
Batch ELT Pattern (dbt + Warehouse)
Source Systems
│
▼ (raw extract — append-only, never modify)
Raw Layer (warehouse: raw.source_name.table)
│
▼ (dbt staging models — rename, cast, deduplicate)
Staging Layer (warehouse: staging.stg_source__table)
│
▼ (dbt intermediate models — join, enrich)
Intermediate Layer (warehouse: intermediate.int_*)
│
▼ (dbt mart models — business-facing aggregations)
Mart Layer (warehouse: marts.dim_* / fact_*)
│
▼
Analysts / BI Tools
Key rule: Raw layer is append-only. Never modify raw data. If a source record changes, append the new version with a loaded_at timestamp.
Streaming Pipeline Pattern (Kafka + Flink)
Source (DB CDC / API / Events)
│
▼ Kafka Topic (raw events, 7-day retention)
│
▼ Flink Job (stateful processing)
│ ├── Deduplication (by event_id, tumbling window)
│ ├── Schema validation (Avro/Protobuf registry check)
│ ├── Enrichment (lookup join against dimension tables)
│ └── Aggregation (windowed metrics)
│
▼ Kafka Topic (processed events)
│
▼ Sink (warehouse, search index, cache)
Incremental Loading Patterns
Watermark-Based Incremental (dbt)
{{
config(
materialized='incremental',
unique_key='order_id',
on_schema_change='sync_all_columns',
incremental_strategy='merge'
)
}}
SELECT
order_id,
user_id,
status,
total_amount,
created_at,
updated_at,
CURRENT_TIMESTAMP AS _dbt_loaded_at
FROM {{ source('raw', 'orders') }}
{% if is_incremental() %}
WHERE updated_at >= (SELECT MAX(updated_at) - INTERVAL '1 hour' FROM {{ this }})
{% endif %}
Change Data Capture (CDC) Pattern
from dataclasses import dataclass
from enum import Enum
class CDCOperation(Enum):
INSERT = "c"
UPDATE = "u"
DELETE = "d"
@dataclass
class CDCRecord:
operation: CDCOperation
source_table: str
primary_key: dict
before: dict | None
after: dict | None
transaction_timestamp: str
def apply_cdc_record(record: CDCRecord, target_table) -> None:
"""
Apply a CDC record idempotently.
UPSERT for inserts/updates, soft-delete for deletes.
Never hard-delete — use is_deleted flag + deleted_at timestamp.
"""
if record.operation == CDCOperation.DELETE:
target_table.upsert({
**record.primary_key,
"is_deleted": True,
"deleted_at": record.transaction_timestamp
})
else:
target_table.upsert({
**record.after,
"is_deleted": False,
"_source_updated_at": record.transaction_timestamp
})
Data Quality Framework
Great Expectations Integration Pattern
import great_expectations as ge
from great_expectations.core import ExpectationSuite
def build_orders_expectation_suite() -> ExpectationSuite:
"""
Define data quality contract for the orders table.
Hard assertions: pipeline fails if violated.
Soft assertions: logged as warnings, pipeline continues.
"""
context = ge.get_context()
suite = context.create_expectation_suite("orders.critical")
suite.add_expectation(
ge.expectations.ExpectColumnValuesToBeUnique(column="order_id")
)
suite.add_expectation(
ge.expectations.ExpectColumnValuesToNotBeNull(column="order_id")
)
suite.add_expectation(
ge.expectations.ExpectColumnValuesToNotBeNull(column="user_id")
)
suite.add_expectation(
ge.expectations.ExpectColumnValuesToBeBetween(
column="total_amount", min_value=0, max_value=100_000
)
)
suite.add_expectation(
ge.expectations.ExpectColumnValuesToBeInSet(
column="status",
value_set=["pending", "processing", "shipped", "delivered", "cancelled", "refunded"]
)
)
suite.add_expectation(
ge.expectations.ExpectTableRowCountToBeGreaterThan(value=0)
)
return suite
def run_quality_gate(df, suite_name: str, fail_on_critical: bool = True) -> dict:
"""
Run quality checks. Fails pipeline on critical violations.
Returns quality report for logging/alerting.
"""
context = ge.get_context()
validator = context.get_validator(batch_request=..., expectation_suite_name=suite_name)
results = validator.validate()
failed = [r for r in results.results if not r.success]
critical_failures = [r for r in failed if r.expectation_config.kwargs.get("severity") != "warn"]
if fail_on_critical and critical_failures:
raise DataQualityError(
f"Pipeline halted: {len(critical_failures)} critical quality violations.\n" +
"\n".join(str(r.expectation_config) for r in critical_failures)
)
return {
"total_checks": len(results.results),
"passed": results.statistics["successful_expectations"],
"failed": results.statistics["unsuccessful_expectations"],
"critical_failures": len(critical_failures)
}
Schema Evolution Strategies
Safe Migration Protocol
┌──────────────────────────────────────────────────────────────┐
│ Schema Change Classification │
│ │
│ ADDITIVE (safe, non-breaking): │
│ + Add nullable column → Apply immediately │
│ + Add new table → Apply immediately │
│ + Widen VARCHAR length → Apply immediately │
│ │
│ NON-ADDITIVE (breaking, requires protocol): │
│ - Remove column → Deprecate first (30 days) │
│ - Rename column → Add alias, migrate, remove old │
│ - Change column type → Add new column, backfill, swap │
│ - Change primary key → Major migration protocol │
└──────────────────────────────────────────────────────────────┘
Column Rename Protocol (dbt)
ALTER TABLE orders ADD COLUMN customer_id BIGINT;
UPDATE orders SET customer_id = user_id;
SELECT
order_id,
COALESCE(customer_id, user_id) AS customer_id,
FROM raw.orders
ALTER TABLE orders DROP COLUMN user_id;
Pipeline Observability
Freshness Monitoring
from datetime import datetime, timedelta
from dataclasses import dataclass
@dataclass
class FreshnessContract:
table: str
max_staleness_hours: float
alert_channel: str
FRESHNESS_CONTRACTS = [
FreshnessContract("fact_orders", max_staleness_hours=1.0, alert_channel="#data-alerts"),
FreshnessContract("fact_events", max_staleness_hours=0.25, alert_channel="#data-alerts"),
FreshnessContract("dim_users", max_staleness_hours=24.0, alert_channel="#data-alerts"),
]
def check_freshness(contract: FreshnessContract, warehouse) -> dict:
"""Check if a table has been updated within its SLO."""
result = warehouse.query(f"""
SELECT MAX(_dbt_loaded_at) AS last_updated
FROM {contract.table}
""").fetchone()
last_updated = result["last_updated"]
staleness_hours = (datetime.utcnow() - last_updated).total_seconds() / 3600
return {
"table": contract.table,
"last_updated": last_updated.isoformat(),
"staleness_hours": round(staleness_hours, 2),
"slo_hours": contract.max_staleness_hours,
"status": "FRESH" if staleness_hours <= contract.max_staleness_hours else "STALE",
"alert_required": staleness_hours > contract.max_staleness_hours
}
Data Lineage Metadata Emission
def emit_lineage_event(
source_tables: list[str],
destination_table: str,
job_name: str,
row_count: int,
run_id: str
) -> None:
"""
Emit lineage metadata for every pipeline run.
Consumed by data catalog (DataHub, OpenMetadata, Marquez).
"""
lineage_event = {
"eventType": "COMPLETE",
"run": {"runId": run_id},
"job": {"namespace": "data_platform", "name": job_name},
"inputs": [{"namespace": "warehouse", "name": t} for t in source_tables],
"outputs": [{"namespace": "warehouse", "name": destination_table}],
"outputFacets": {
"rowCount": {"_producer": job_name, "rowCount": row_count}
}
}
lineage_client.emit(lineage_event)
Idempotency Checklist
Before deploying any pipeline, verify:
Blocking Violations (NEVER)
| Violation | Consequence | Recovery |
|---|
| Loading data to destination before quality checks run | Dirty data lands in production tables, corrupting downstream reports | Add quality gate before load step. Drop corrupt partitions and reprocess. |
Using INSERT without deduplication strategy | Duplicate records accumulate silently across re-runs, inflating metrics | Switch to MERGE/UPSERT or add deduplication in staging. |
| Hard-deleting records from warehouse | Destroys audit trail, breaks point-in-time analysis, breaks downstream models | Always use soft-delete with is_deleted flag and deleted_at timestamp. |
| Deploying pipeline without a freshness SLO | Stale data goes undetected until business impact is visible | Define max_staleness_hours for every dataset and wire to alerting. |
| Silencing quality failures with try/except | Corrupt batch lands silently, propagates to reports | Remove bare excepts. Log quality failures and raise on critical violations. |
| Full table refresh where incremental loading is viable | Fails at scale when source exceeds memory/query limits | Default to incremental with watermark or CDC. Reserve full refresh for backfills. |
Verification
Before declaring a pipeline complete:
Self-Verification Checklist
Verification Commands
SELECT COUNT(*), COUNT(DISTINCT primary_key) FROM destination;
python check_freshness.py --all-tables
python run_quality_gate.py --suite critical --fail-on-critical
python check_lineage.py --table destination_table
Quality Gates
| Gate | Criteria | Fail Action |
|---|
| Data Quality | 0 critical failures, all hard expectations pass | Halt pipeline, alert channel, do not load |
| Idempotency | Two consecutive runs produce identical row counts & checksums | Investigate merge key or CDC dedup logic |
| Freshness | check_freshness() returns FRESH within 10 min of job completion | Escalate to data-alerts channel |
| Schema Evolution | Adding nullable column to source does not fail pipeline | Add on_schema_change='sync_all_columns' as safety net |
Performance & Cost
Model Selection
| Task Complexity | Recommended Approach | Estimated Cost |
|---|
| Batch ELT (dbt) | Warehouse compute (BigQuery/Snowflake slots) | Pay per scan/query |
| Streaming (Kafka + Flink) | Managed Kafka + Flink cluster | Node-hour + storage |
| CDC ingestion | Debezium + Kafka | Connector + topic storage |
Parallelization
- dbt models: Run non-dependent models in parallel using dbt threads (default 4, max 8)
- Quality checks: Run per-table suites in parallel; no inter-suite dependencies
- Streaming consumers: Scale partitions per topic, 1 consumer per partition
Context Budget
- Expected context usage: 5-8KB per pipeline design session
- When to context-optimize: When reviewing >3 pipeline files or analyzing multi-hour streaming config
- Context recovery: Use
ctx_execute for pipeline output, ctx_execute_file for large config files
Examples
Example 1: Batch ELT with dbt
User request: "Set up a dbt pipeline that ingests orders from PostgreSQL daily and produces a mart table with daily revenue."
Skill execution:
- Define raw layer:
raw.orders (append-only from source)
- Define staging model: deduplicate, rename columns, cast types
- Define intermediate model: join orders to products for revenue calc
- Define mart model:
fact_daily_revenue with date, product_id, revenue
- Add quality gates: non-null order_id, positive revenue, referential integrity
- Add freshness contract: max_staleness_hours=24, alert on stale
- Add incremental loading with watermark on
updated_at
Result: fact_daily_revenue table produced daily, with quality gate blocking bad data, freshness monitoring, and lineage metadata in catalog.
Example 2: Handling Schema Drift from Source
User request: "The upstream orders API suddenly added 3 new columns without warning."
Skill execution:
- Check if columns are additive (nullable, no NOT NULL constraint) → safe
- dbt
on_schema_change='sync_all_columns' syncs them automatically
- Review downstream models — no breaking changes since new columns are unused
- Log the drift event for audit trail
- Update documentation in data catalog with new columns
Result: Pipeline continues running. New columns available for use. Team is notified.
Example 3: CDC Pipeline with Idempotency Fix
User request: "Our CDC pipeline is creating duplicate rows every time it restarts."
Skill execution:
- Check merge key —
unique_key='order_id' is set correctly
- Check watermark lookback — set to 1 hour to capture late-arriving events
- Check CDC deduplication — add
row_number() OVER (PARTITION BY event_id) in staging
- Verify: re-run twice, compare row counts and checksums
- Add idempotency test to CI pipeline
Result: Duplicate rows eliminated. Idempotency test prevents regression.
Anti-Patterns
- Never use
INSERT without a deduplication strategy because duplicate records accumulate silently across pipeline reruns, and downstream analysts will eventually aggregate over inflated counts without knowing they are double-counting, producing reports that overstate revenue, events, or user activity by an arbitrary factor — always MERGE/UPSERT or deduplicate in staging.
- Never hard-delete records from a data warehouse because permanent deletion destroys the audit trail required for compliance, makes point-in-time analysis impossible, and breaks any downstream model that references the deleted key — use a soft-delete with
is_deleted flag and deleted_at timestamp instead.
- Never load data to a destination before running quality checks because loading dirty data and planning to clean it later means it never gets cleaned, analysts build dashboards on corrupt records, and fixing the data retroactively requires a full historical reload that can take hours or days.
- Never use a full table refresh where incremental loading is viable because full refreshes fail at scale when the source table exceeds memory or query timeout limits, and a failed mid-refresh leaves the destination in a partially overwritten state that analysts are still querying.
- Never deploy a pipeline without a freshness SLO because without a defined maximum acceptable data latency, stakeholders have no signal when data is stale — a pipeline that ran at 2 AM but silently stopped producing records by 6 AM will go undetected until an analyst acts on hours-old data.
- Never alter a non-nullable column type directly because an in-place type change locks the table during migration, breaks all downstream consumers the moment the pipeline restarts, and cannot be rolled back without restoring from backup — use the additive migration protocol: add new column, backfill, swap references, then drop old column after a deprecation window.
- Never silence a data quality failure with a try/except without alerting because swallowing a quality exception allows a corrupt batch to land in the destination silently, leaving analysts to discover data anomalies hours later when the business impact has already propagated to reports and dashboards.
Failure Modes
| Situation | Response |
|---|
| Duplicate records in destination | Root cause: missing MERGE key or late-arriving CDC events. Add deduplication in staging. Run idempotency test. |
| Schema drift from source system | Alert on unexpected column additions/removals. Use on_schema_change='sync_all_columns' in dbt as a safety net. Validate in quality gate. |
| Pipeline backpressure (Kafka lag growing) | Scale consumer replicas or increase parallelism. Add a lag alert at 60s behind production topic. |
| Late-arriving data causes missed records | Extend the watermark lookback window. Add a late-data reconciliation job that runs 6 hours after the primary job. |
| Quality check false positive blocks pipeline | Review the expectation definition. If the data is valid, update the contract. Never bypass the gate. |
| Destination table lock contention | Switch from statement-level locking to row-level upsert. Use partitioned loads with partition swap. |
References
Internal Dependencies
data-analyst — Consumes pipeline output for analysis and reporting
ml-engineer — Uses engineered features and clean data for model training
database-migrations — Handles OLTP schema changes (partner skill for OLAP)
observability-specialist — Wires pipeline metrics into monitoring dashboards
infra-architect — Provisions warehouse resources, IAM policies, network access
External Standards
Related Skills
data-analyst — Follows data-engineer in the Data vertical chain
ml-engineer — Uses clean data from pipelines for feature engineering
database-migrations — Partner skill for OLTP schema changes
Changelog
| Version | Date | Changes |
|---|
| 2.0.0 | 2026-07-09 | Upgraded to Gold Standard v2.0: added frontmatter version/category/dependencies, Blocking Violations table, Verification with commands/quality gates, Performance & Cost section, Examples, References, Changelog. Reorganized to 12-section template. |
Integration with Mega-Mind
The data engineer skill is the infrastructure layer in the Data vertical:
data-engineer (build reliable pipelines) → data-analyst (analyze trustworthy data) → ml-engineer (model and predict)
- Pair with
observability-specialist to wire pipeline metrics into monitoring dashboards and SLO alerts
- Use
database-migrations for schema changes in application databases (OLTP); this skill handles OLAP/warehouse schema changes
- Coordinate with
infra-architect for warehouse provisioning, IAM policies, and network access patterns
- When a pipeline is mature and analysts need to explore it, hand off to
data-analyst for EDA and reporting