← Files Data Engineering CopilotARCHIVED FILE

skills/data-engineering/references/data_engineering_implementation_templates.md

3.89 KB · Oct 3, 2026 · 06:38 UTC

↓ Download file

# Data Engineering Implementation Templates

## Purpose

Use for concrete SQL, PySpark, Airflow, CDC, replay, API-ingestion, data-quality, and RCA examples.

Templates are generic starting points. Adapt identifiers, ordering, schema, runtime, and platform behavior.

# Safety rules

- explicit production columns;
- deduplicate MERGE sources;
- prefer source sequence/version for CDC;
- preserve replayable raw data;
- make retries explicit;
- never embed credentials;
- include validation/rollback for material changes.

# 1. Deterministic Delta MERGE staging

```sql
CREATE OR REPLACE TEMP VIEW staged_changes AS
WITH ranked AS (
  SELECT
    business_key,
    payload_col,
    op,
    source_sequence,
    source_ts,
    ingest_ts,
    ROW_NUMBER() OVER (
      PARTITION BY business_key
      ORDER BY source_sequence DESC, source_ts DESC, ingest_ts DESC
    ) AS rn
  FROM bronze.source_cdc
  WHERE ingest_date BETWEEN DATE('${window_start}') AND DATE('${window_end}')
)
SELECT *
FROM ranked
WHERE rn = 1;
```

Validate uniqueness:

```sql
SELECT business_key, COUNT(*) AS row_count
FROM staged_changes
GROUP BY business_key
HAVING COUNT(*) > 1;
```

Expected: zero rows.

Then MERGE using the actual business rules for insert/update/delete.

Do not paste a generic MERGE until the user’s schema and delete semantics are known.

# 2. foreachBatch skeleton

```python
from pyspark.sql import DataFrame

def upsert_batch(batch_df: DataFrame, batch_id: int) -> None:
    if batch_df.isEmpty():
        return

    # 1. validate/dedupe using source ordering
    # 2. write audit metadata
    # 3. apply idempotent target change
    # 4. emit reconciliation metrics
```

```python
query = (
    source_df.writeStream
    .queryName("source_to_silver")
    .foreachBatch(upsert_batch)
    .option("checkpointLocation", "/durable/checkpoints/source_to_silver")
    .start()
)
```

Notes:
- checkpoint is unique per query;
- `batch_id` is useful for audit but does not make external side effects idempotent;
- test restart/replay behavior.

# 3. Resumable API ingestion

Operating pattern:

1. request page/cursor;
2. land raw response durably;
3. record request/run/cursor metadata;
4. commit next cursor only after durable landing;
5. resume from last committed cursor;
6. dedupe downstream using source identity/version.

For retryable HTTP statuses, use bounded retry with backoff and server-provided `Retry-After` where appropriate.

Never log credentials/tokens.

# 4. Airflow DAG pattern

```python
from airflow import DAG
from airflow.operators.python import PythonOperator

with DAG(
    dag_id="incremental_pipeline",
    schedule="@hourly",
    catchup=False,
    max_active_runs=1,
    default_args={"retries": 2},
) as dag:
    extract = PythonOperator(task_id="extract", python_callable=extract_fn)
    transform = PythonOperator(task_id="transform", python_callable=transform_fn)
    publish = PythonOperator(task_id="publish", python_callable=publish_fn)

    extract >> transform >> publish
```

Adapt imports/APIs to the actual Airflow version.

Every task should be safe for its logical data interval.

# 5. Data-quality gate

```sql
-- Null key
SELECT COUNT(*) AS null_key_count
FROM silver.entity
WHERE business_key IS NULL;

-- Duplicate grain
SELECT business_key, COUNT(*) AS row_count
FROM silver.entity
GROUP BY business_key
HAVING COUNT(*) > 1;

-- Freshness
SELECT MAX(updated_at) AS latest_update
FROM silver.entity;
```

Production gates need:
- threshold;
- severity;
- owner;
- block/quarantine/warn behavior.

# 6. Backfill record

Capture:
- replay/backfill ID;
- input scope;
- source snapshot/version;
- target scope;
- code/config version;
- concurrent-writer policy;
- input/output/rejected counts;
- start/end time;
- validation;
- rollback/repair result.

# 7. Compact RCA

- Incident/impact
- Timeline
- Dominant root cause
- Contributing factors
- Verification
- Containment
- Permanent fix
- Blast radius
- Data recovery
- Rollback
- Prevention

SHA-256: f3148c355cb9123f7f7b3f6b2edd22e340b8f6111c5a417c4129b7e6dc1d617a