← Files Data Engineering CopilotARCHIVED FILE
skills/data-engineering/references/data_engineering_implementation_templates.md
3.89 KB · Oct 2, 2026 · 00:37 UTC
# 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