Custom Skills для DE: SCD2, CDC, Bronze Ingestion и MERGE-паттерны
Строим переиспользуемые pipeline-генераторы для DE-задач: metadata-driven SCD2, CDC ingestion, Delta Lake MERGE и bronze-слой. Упаковываем в Custom Skills для AI-агента.
Проблема copy-paste ETL и цена поддержки¶
В типичной Data Platform десятки однотипных пайплайнов. Каждый разработчик пишет их по-своему:
Проблемы этого подхода:
- Баги множатся: ошибка в логике SCD2 скопирована в 10 пайплайнов
- Нет стандарта: каждый инженер по-своему решает одни и те же задачи
- Дорогой onboarding: новый DE тратит недели, понимая "наш способ" из десятков примеров
- Невозможно обновить: исправление бага требует правки в каждом пайплайне
Решение - metadata-driven pipeline generator: один движок, сотни конфигов. Движок проверен и покрыт тестами, конфиг описывает источник и целевую таблицу.
Архитектура metadata-driven pipeline¶
Добавить новую таблицу в SCD2 = написать YAML-конфиг. Никакого кода.
SCD Type 2: историзация изменений¶
Бизнес-проблема¶
Slowly Changing Dimension Type 2 (SCD2) - классическая проблема DWH. Запись изменилась (статус заказа, цена продукта, адрес пользователя) - но старые факты ссылались на старую версию. Нужно хранить полную историю.
Анатомия SCD2 записи¶
Алгоритм SCD2: три типа строк¶
Реализация SCD2Engine¶
# engines/scd2_engine.py
"""
Универсальный SCD2 Engine для Iceberg таблиц.
Использование: см. configs/scd2/*.yaml
"""
from dataclasses import dataclass
from datetime import date
from typing import Optional
import hashlib
from pyspark.sql import SparkSession, DataFrame
from pyspark.sql import functions as F
from pyspark.sql.types import StringType
@dataclass
class SCD2Config:
"""Конфигурация одной SCD2 таблицы."""
source_table: str # откуда читаем (полное имя таблицы)
target_table: str # куда пишем (Iceberg)
natural_key: list[str] # бизнес-ключи (могут быть составными)
tracked_columns: list[str] # колонки, изменение которых создаёт новую версию
execution_date: str # дата запуска (YYYY-MM-DD)
soft_delete: bool = True # закрывать ли записи, удалённые из источника
@classmethod
def from_yaml(cls, path: str) -> "SCD2Config":
import yaml
with open(path) as f:
cfg = yaml.safe_load(f)
return cls(**cfg)
class SCD2Engine:
"""
Метаданные-управляемый SCD2 pipeline.
Один движок - любое количество dimension-таблиц.
"""
EFFECTIVE_TO_MAX = "9999-12-31"
CURRENT_FLAG = "is_current"
EFFECTIVE_FROM = "effective_from"
EFFECTIVE_TO = "effective_to"
ROW_HASH = "row_hash"
def __init__(self, spark: SparkSession, config: SCD2Config):
self.spark = spark
self.cfg = config
def _compute_hash(self, df: DataFrame) -> DataFrame:
"""Вычисляет MD5-хэш tracked-колонок для детекции изменений."""
hash_expr = F.md5(
F.concat_ws("|", *[F.coalesce(F.col(c).cast(StringType()), F.lit(""))
for c in self.cfg.tracked_columns])
)
return df.withColumn(self.ROW_HASH, hash_expr)
def _read_source(self) -> DataFrame:
"""Читает snapshot источника за execution_date."""
return self.spark.table(self.cfg.source_table) \
.filter(F.to_date(F.col("updated_at")) <= self.cfg.execution_date)
def _read_current_target(self) -> Optional[DataFrame]:
"""Читает актуальные (is_current=true) записи из DWH."""
try:
return self.spark.table(self.cfg.target_table) \
.filter(F.col(self.CURRENT_FLAG) == True) # noqa: E712
except Exception:
return None # таблица ещё не существует
def _classify_changes(
self, source: DataFrame, current: Optional[DataFrame]
) -> tuple[DataFrame, DataFrame, DataFrame]:
"""
Классифицирует строки на NEW, CHANGED, DELETED.
Возвращает кортеж (new_records, changed_records, deleted_records).
"""
source_hashed = self._compute_hash(source)
key_cols = self.cfg.natural_key
if current is None:
# Первый запуск: все строки - новые
return source_hashed, self.spark.createDataFrame([], source_hashed.schema), \
self.spark.createDataFrame([], source_hashed.schema)
current_hashed = self._compute_hash(current)
# JOIN по natural key
join_cond = [source_hashed[k] == current_hashed[k] for k in key_cols]
joined = source_hashed.alias("src").join(
current_hashed.select(key_cols + [self.ROW_HASH]).alias("cur"),
on=join_cond,
how="full_outer"
)
# NEW: есть в src, нет в cur
new_records = joined.filter(
F.col(f"cur.{key_cols[0]}").isNull()
).select([F.col(f"src.{c}") for c in source_hashed.columns])
# CHANGED: hash изменился
changed_records = joined.filter(
F.col(f"cur.{key_cols[0]}").isNotNull() &
F.col(f"src.{key_cols[0]}").isNotNull() &
(F.col(f"src.{self.ROW_HASH}") != F.col(f"cur.{self.ROW_HASH}"))
).select([F.col(f"src.{c}") for c in source_hashed.columns])
# DELETED (soft): есть в cur, нет в src
deleted_records = joined.filter(
F.col(f"src.{key_cols[0]}").isNull()
).select([F.col(f"cur.{c}") for c in current_hashed.columns])
return new_records, changed_records, deleted_records
def _prepare_new_version(self, df: DataFrame) -> DataFrame:
"""Добавляет SCD2-метаданные для новой версии записи."""
return df.withColumn(self.EFFECTIVE_FROM, F.lit(self.cfg.execution_date)) \
.withColumn(self.EFFECTIVE_TO, F.lit(self.EFFECTIVE_TO_MAX)) \
.withColumn(self.CURRENT_FLAG, F.lit(True))
def run(self) -> dict:
"""Выполняет SCD2 pipeline. Возвращает метрики."""
source = self._read_source()
current = self._read_current_target()
new_recs, changed_recs, deleted_recs = self._classify_changes(source, current)
# Подготовить новые версии для INSERT (new + changed)
to_insert = self._prepare_new_version(
new_recs.unionByName(changed_recs, allowMissingColumns=True)
)
# Закрыть изменённые и удалённые записи (UPDATE effective_to)
keys_to_close = (
changed_recs.select(self.cfg.natural_key)
.union(deleted_recs.select(self.cfg.natural_key) if self.cfg.soft_delete
else self.spark.createDataFrame([], changed_recs.schema).select(self.cfg.natural_key))
)
# MERGE INTO через Spark SQL (Iceberg)
keys_to_close.createOrReplaceTempView("_scd2_keys_to_close")
self.spark.sql(f"""
MERGE INTO {self.cfg.target_table} AS target
USING _scd2_keys_to_close AS src
ON {" AND ".join(f"target.{k} = src.{k}" for k in self.cfg.natural_key)}
AND target.{self.CURRENT_FLAG} = true
WHEN MATCHED THEN UPDATE SET
target.{self.EFFECTIVE_TO} = '{self.cfg.execution_date}',
target.{self.CURRENT_FLAG} = false
""")
# INSERT новых версий
to_insert.writeTo(self.cfg.target_table).append()
return {
"new": new_recs.count(),
"changed": changed_recs.count(),
"deleted": deleted_recs.count() if self.cfg.soft_delete else 0,
}
Конфиг для конкретной таблицы¶
# configs/scd2/orders.yaml
source_table: "nessie.staging.orders_latest"
target_table: "nessie.dw.dim_orders"
natural_key: ["order_id"]
tracked_columns:
- status
- amount
- shipping_address
- payment_method
execution_date: "{{ ds }}" # Airflow шаблон
soft_delete: true
Добавить новую SCD2-таблицу = написать такой YAML. Движок не трогаем.
CDC Pipeline: обработка потока изменений¶
Что такое CDC¶
Change Data Capture (CDC) - паттерн, при котором система фиксирует каждое изменение в источнике (INSERT/UPDATE/DELETE) и транслирует его потребителям. Вместо full reload раз в день - непрерывный поток изменений.
Структура CDC-события (Debezium формат)¶
{
"before": {
"order_id": 42,
"status": "pending",
"amount": 299.99
},
"after": {
"order_id": 42,
"status": "completed",
"amount": 299.99
},
"op": "u", // "c"=create, "u"=update, "d"=delete, "r"=read(snapshot)
"ts_ms": 1705312800000,
"source": {
"table": "orders",
"lsn": 12345678
}
}
CDC Pipeline Engine¶
# engines/cdc_engine.py
"""
Batch CDC processor: обрабатывает накопленные CDC-события за период.
Идемпотентен: безопасно перезапускать.
"""
from dataclasses import dataclass
from pyspark.sql import SparkSession, DataFrame
from pyspark.sql import functions as F
from pyspark.sql.types import StructType
@dataclass
class CDCConfig:
source_table: str # staging-таблица с CDC-событиями
target_table: str # целевая Iceberg-таблица
primary_key: list[str] # ключ целевой таблицы
op_column: str # колонка с типом операции (c/u/d/r)
ts_column: str # колонка с timestamp события
execution_date: str # дата обработки (YYYY-MM-DD)
after_prefix: str = "after_" # префикс after-полей (Debezium-формат)
@classmethod
def from_yaml(cls, path: str) -> "CDCConfig":
import yaml
with open(path) as f:
return cls(**yaml.safe_load(f))
class CDCEngine:
"""
Обрабатывает CDC-события и применяет изменения к Iceberg-таблице.
Стратегия:
- Берём последнее событие для каждого ключа за период
- INSERT/UPDATE → UPSERT в target
- DELETE → hard delete или soft delete (флаг is_deleted)
"""
def __init__(self, spark: SparkSession, config: CDCConfig):
self.spark = spark
self.cfg = config
def _read_cdc_events(self) -> DataFrame:
"""Читает CDC-события за execution_date."""
return self.spark.table(self.cfg.source_table) \
.filter(F.to_date(F.col(self.cfg.ts_column)) == self.cfg.execution_date)
def _deduplicate(self, events: DataFrame) -> DataFrame:
"""
Оставляет последнее событие для каждого primary_key.
В потоке CDC один ключ может измениться несколько раз за день —
нам важно только финальное состояние.
"""
from pyspark.sql.window import Window
window = Window.partitionBy(self.cfg.primary_key) \
.orderBy(F.col(self.cfg.ts_column).desc())
return events \
.withColumn("rn", F.row_number().over(window)) \
.filter(F.col("rn") == 1) \
.drop("rn")
def _extract_after_fields(self, df: DataFrame) -> DataFrame:
"""
Извлекает after_* поля и переименовывает в обычные колонки.
Дebezium пишет: after_order_id, after_status, ...
Нам нужно: order_id, status, ...
"""
after_cols = [c for c in df.columns if c.startswith(self.cfg.after_prefix)]
renamed = df
for col_name in after_cols:
new_name = col_name[len(self.cfg.after_prefix):]
renamed = renamed.withColumnRenamed(col_name, new_name)
return renamed
def run(self) -> dict:
"""Применяет CDC-события к target таблице."""
events = self._read_cdc_events()
deduped = self._deduplicate(events)
# Разделяем на UPSERT и DELETE
upserts = deduped.filter(F.col(self.cfg.op_column).isin("c", "u", "r"))
deletes = deduped.filter(F.col(self.cfg.op_column) == "d")
upsert_records = self._extract_after_fields(upserts)
# MERGE INTO (Iceberg)
upsert_records.createOrReplaceTempView("_cdc_upserts")
join_cond = " AND ".join(
f"t.{k} = s.{k}" for k in self.cfg.primary_key
)
all_cols = [c for c in upsert_records.columns
if c not in (self.cfg.op_column, self.cfg.ts_column)]
update_set = ", ".join(f"t.{c} = s.{c}" for c in all_cols)
self.spark.sql(f"""
MERGE INTO {self.cfg.target_table} AS t
USING _cdc_upserts AS s
ON {join_cond}
WHEN MATCHED THEN UPDATE SET {update_set}
WHEN NOT MATCHED THEN INSERT *
""")
# Удаления (hard delete)
if deletes.count() > 0:
deletes.createOrReplaceTempView("_cdc_deletes")
self.spark.sql(f"""
DELETE FROM {self.cfg.target_table}
WHERE ({", ".join(self.cfg.primary_key)}) IN (
SELECT {", ".join(self.cfg.primary_key)} FROM _cdc_deletes
)
""")
return {
"upserted": upsert_records.count(),
"deleted": deletes.count(),
}
Идемпотентность в CDC¶
Ключевое требование CDC-пайплайна - идемпотентность: если пайплайн запущен дважды за один день, результат должен быть таким же, как при однократном запуске.
Delta Lake MERGE Patterns¶
Хотя в нашем стенде используется Iceberg, паттерны MERGE INTO универсальны - синтаксис аналогичен в Delta Lake, Iceberg и Hudi.
Базовый UPSERT¶
# patterns/merge_upsert.py
"""
Паттерн: UPSERT (INSERT or UPDATE).
Используется: CDC ingestion, SCD1, dimension refresh.
"""
from pyspark.sql import SparkSession, DataFrame
def upsert_to_iceberg(
spark: SparkSession,
source: DataFrame,
target_table: str,
merge_keys: list[str],
) -> None:
"""
MERGE INTO: обновляет существующие записи, вставляет новые.
Args:
source: DataFrame с новыми/изменёнными данными
target_table: Iceberg-таблица для слияния
merge_keys: ключи для JOIN (natural key таблицы)
"""
source.createOrReplaceTempView("_merge_source")
join_cond = " AND ".join(f"t.{k} = s.{k}" for k in merge_keys)
# Все колонки кроме merge keys обновляются
update_cols = [c for c in source.columns if c not in merge_keys]
update_set = ", ".join(f"t.{c} = s.{c}" for c in update_cols)
spark.sql(f"""
MERGE INTO {target_table} AS t
USING _merge_source AS s
ON {join_cond}
WHEN MATCHED THEN UPDATE SET {update_set}
WHEN NOT MATCHED THEN INSERT *
""")
Soft Delete¶
# patterns/merge_soft_delete.py
"""
Паттерн: Soft Delete - помечаем удалённые записи флагом is_deleted.
Используется: когда нельзя физически удалять данные (compliance, аудит).
"""
def soft_delete_merge(
spark: SparkSession,
active_records: DataFrame, # текущий snapshot (без удалённых)
target_table: str,
merge_keys: list[str],
deleted_flag_col: str = "is_deleted",
deleted_at_col: str = "deleted_at",
) -> None:
"""
Обновляет is_deleted=true для записей, отсутствующих в источнике.
"""
active_records.createOrReplaceTempView("_soft_delete_src")
join_cond = " AND ".join(f"t.{k} = s.{k}" for k in merge_keys)
spark.sql(f"""
MERGE INTO {target_table} AS t
USING _soft_delete_src AS s
ON {join_cond} AND t.{deleted_flag_col} = false
WHEN MATCHED THEN UPDATE SET
t.{deleted_flag_col} = false,
t.{deleted_at_col} = NULL
WHEN NOT MATCHED BY SOURCE AND t.{deleted_flag_col} = false THEN UPDATE SET
t.{deleted_flag_col} = true,
t.{deleted_at_col} = current_timestamp()
WHEN NOT MATCHED BY TARGET THEN INSERT *
""")
Insert-Only с дедупликацией¶
# patterns/merge_insert_only.py
"""
Паттерн: Insert-Only с дедупликацией.
Используется: append-only события, где UPDATE запрещён (immutable log).
"""
def insert_new_only(
spark: SparkSession,
source: DataFrame,
target_table: str,
dedup_keys: list[str],
) -> int:
"""
Вставляет только строки, которых ещё нет в target.
Предотвращает дублирование при повторных запусках.
"""
source.createOrReplaceTempView("_insert_source")
join_cond = " AND ".join(f"t.{k} = s.{k}" for k in dedup_keys)
spark.sql(f"""
MERGE INTO {target_table} AS t
USING _insert_source AS s
ON {join_cond}
WHEN NOT MATCHED THEN INSERT *
""")
return source.count()
Bronze Ingestion: первый слой Medallion Architecture¶
Medallion Architecture¶
Bronze - неизменяемый слой сырых данных. Здесь:
- Данные хранятся as-is (без трансформаций)
- Добавляются технические метаданные (
_ingested_at,_source_file,_batch_id) - Любые ошибки можно воспроизвести из Bronze
- Schema evolution: новые колонки в источнике добавляются автоматически
Generic Bronze Ingestion Engine¶
# engines/bronze_engine.py
"""
Universal Bronze Ingestion Engine.
Читает из любого источника, добавляет технические колонки,
записывает в Iceberg с поддержкой schema evolution.
"""
from dataclasses import dataclass, field
from datetime import datetime
from typing import Optional
import uuid
from pyspark.sql import SparkSession, DataFrame
from pyspark.sql import functions as F
@dataclass
class BronzeConfig:
source_format: str # "jdbc", "parquet", "json", "csv", "kafka"
source_options: dict # параметры подключения к источнику
target_table: str # Iceberg bronze-таблица
partition_by: list[str] = field(default_factory=list) # колонки партиционирования
quarantine_table: Optional[str] = None # куда писать битые записи
merge_schema: bool = True # auto schema evolution
@classmethod
def from_yaml(cls, path: str) -> "BronzeConfig":
import yaml
with open(path) as f:
return cls(**yaml.safe_load(f))
class BronzeIngestionEngine:
"""
Универсальный bronze ingestion:
- Читает из любого источника
- Добавляет технические метаданные
- Пишет в Iceberg с поддержкой schema evolution
- Изолирует битые записи в quarantine
"""
METADATA_COLS = ["_ingested_at", "_batch_id", "_source_table"]
def __init__(self, spark: SparkSession, config: BronzeConfig):
self.spark = spark
self.cfg = config
self._batch_id = str(uuid.uuid4())
def _read_source(self) -> DataFrame:
"""Читает данные из источника."""
return self.spark.read \
.format(self.cfg.source_format) \
.options(**self.cfg.source_options) \
.load()
def _add_metadata(self, df: DataFrame) -> DataFrame:
"""Добавляет технические колонки для аудита."""
return df \
.withColumn("_ingested_at", F.current_timestamp()) \
.withColumn("_batch_id", F.lit(self._batch_id)) \
.withColumn("_source_table",
F.lit(self.cfg.source_options.get("dbtable", "unknown")))
def _validate(self, df: DataFrame) -> tuple[DataFrame, DataFrame]:
"""
Базовая валидация: выделяет битые записи в quarantine.
Extends: переопределить для кастомных проверок.
"""
# Пример: строки с полностью null-ными non-metadata колонками → quarantine
data_cols = [c for c in df.columns if not c.startswith("_")]
null_count_expr = sum(F.col(c).isNull().cast("int") for c in data_cols)
all_null_condition = (null_count_expr == len(data_cols))
valid = df.filter(~all_null_condition)
quarantine = df.filter(all_null_condition) \
.withColumn("_quarantine_reason", F.lit("all_data_columns_null"))
return valid, quarantine
def run(self) -> dict:
"""Выполняет bronze ingestion. Возвращает метрики."""
raw = self._read_source()
enriched = self._add_metadata(raw)
valid, quarantine = self._validate(enriched)
# Запись в bronze (schema evolution - новые колонки добавляются автоматически)
writer = valid.writeTo(self.cfg.target_table)
if self.cfg.merge_schema:
writer = writer.option("mergeSchema", "true")
if self.cfg.partition_by:
writer = writer.partitionedBy(*self.cfg.partition_by)
try:
writer.append()
except Exception:
writer.createOrReplace() # первый запуск: таблица не существует
# Quarantine
quarantine_count = 0
if self.cfg.quarantine_table and quarantine.count() > 0:
quarantine.writeTo(self.cfg.quarantine_table).append()
quarantine_count = quarantine.count()
return {
"ingested": valid.count(),
"quarantined": quarantine_count,
"batch_id": self._batch_id,
}
Конфиг bronze-источника¶
# configs/bronze/orders_from_postgres.yaml
source_format: "jdbc"
source_options:
url: "${POSTGRES_URL}"
dbtable: "(SELECT * FROM orders WHERE updated_at::date = '{{ ds }}') t"
user: "${POSTGRES_USER}"
password: "${POSTGRES_PASSWORD}"
driver: "org.postgresql.Driver"
numPartitions: "10"
partitionColumn: "order_id"
lowerBound: "1"
upperBound: "1000000"
target_table: "nessie.bronze.orders_raw"
partition_by: [] # bronze хранится в одной партиции с метаданными
quarantine_table: "nessie.bronze.orders_quarantine"
merge_schema: true # schema evolution включена
Упаковка в Custom Skills для AI-агента¶
Каждый из этих движков превращается в Custom Skill для OpenCode - переиспользуемый инструмент агента:
# .opencode/skills/generate-scd2-config.yaml
name: generate-scd2-config
description: |
Создаёт YAML-конфиг для нового SCD2 pipeline.
Достаточно сказать: "добавь SCD2 для таблицы products".
trigger: "scd2|slowly changing|dimension"
steps:
- name: check_source_schema
tool: bash
command: |
spark-sql -e "DESCRIBE TABLE nessie.staging.{source_table}"
- name: check_existing_configs
tool: bash
command: "ls configs/scd2/"
- name: generate
llm: true
prompt: |
Создай configs/scd2/{source_table}.yaml для SCD2 pipeline.
Схема источника: {check_source_schema.output}
Паттерн из существующих конфигов: {check_existing_configs.output}
Правила:
- natural_key: уникальный бизнес-ключ таблицы (обычно *_id)
- tracked_columns: атрибуты, изменение которых создаёт новую версию
(НЕ включай технические колонки: created_at, updated_at, _ingested_at)
- soft_delete: true (сохраняем историю удалений)
- name: write_config
tool: write
path: "configs/scd2/{source_table}.yaml"
content: "{generate.output}"
- name: validate
tool: bash
command: "python3 -c \"import yaml; yaml.safe_load(open('configs/scd2/{source_table}.yaml'))\""
# .opencode/skills/run-bronze-ingestion.yaml
name: run-bronze-ingestion
description: Запускает bronze ingestion для источника, проверяет результат
trigger: "bronze|ingest|загрузи|import"
steps:
- name: check_config
tool: read
path: "configs/bronze/{source}.yaml"
- name: run_pipeline
tool: bash
command: |
python3 -c "
from pyspark.sql import SparkSession
from engines.bronze_engine import BronzeIngestionEngine, BronzeConfig
spark = SparkSession.builder.remote('sc://localhost:15002').getOrCreate()
cfg = BronzeConfig.from_yaml('configs/bronze/{source}.yaml')
engine = BronzeIngestionEngine(spark, cfg)
result = engine.run()
print(result)
"
- name: verify
tool: bash
command: |
spark-sql -e "
SELECT COUNT(*), _batch_id
FROM nessie.bronze.{source}_raw
WHERE DATE(_ingested_at) = CURRENT_DATE
GROUP BY _batch_id
"
- name: check_quarantine
tool: bash
command: |
spark-sql -e "
SELECT COUNT(*), _quarantine_reason
FROM nessie.bronze.{source}_quarantine
WHERE DATE(_ingested_at) = CURRENT_DATE
GROUP BY _quarantine_reason
"
Как агент использует эти skills¶
opencode
# > Нам нужен SCD2 для таблицы products из PostgreSQL.
# Отслеживать изменения в: name, price, category_id.
# Агент выполняет skill generate-scd2-config:
# 1. spark-sql: DESCRIBE TABLE nessie.staging.products
# 2. ls configs/scd2/ (видит orders.yaml как образец)
# 3. Генерирует products.yaml по образцу
# 4. python3 validate.py - проверяет YAML
# 5. "Конфиг создан: configs/scd2/products.yaml
# Запустить первый прогон? (python3 run_scd2.py products 2024-01-15)"
Observability: метрики и аудит для pipeline generators¶
# utils/observability.py
"""
Единое логирование и метрики для всех pipeline engines.
"""
import logging
from dataclasses import dataclass
from datetime import datetime
from pyspark.sql import SparkSession
logger = logging.getLogger("pipeline_engine")
@dataclass
class PipelineRun:
pipeline_type: str # "scd2", "cdc", "bronze"
config_name: str # имя конфига (без расширения)
execution_date: str
batch_id: str
rows_processed: int
rows_failed: int
status: str # "success", "failed", "partial"
duration_sec: float
error_message: str = ""
def log_pipeline_run(spark: SparkSession, run: PipelineRun) -> None:
"""Записывает метрику запуска в audit-таблицу."""
from pyspark.sql import Row
audit_row = Row(
pipeline_type=run.pipeline_type,
config_name=run.config_name,
execution_date=run.execution_date,
batch_id=run.batch_id,
rows_processed=run.rows_processed,
rows_failed=run.rows_failed,
status=run.status,
duration_sec=run.duration_sec,
error_message=run.error_message,
logged_at=datetime.utcnow().isoformat()
)
spark.createDataFrame([audit_row]) \
.writeTo("nessie.ops.pipeline_audit") \
.append()
Trade-offs: когда pipeline generator - overengineering¶
Правило трёх: делать абстракцию стоит, когда нужно повторить паттерн в третий раз. Первый раз - пишем напрямую. Второй - выделяем функцию. Третий - строим движок.
Практика¶
1. Запустить SCD2Engine на тестовых данных¶
# Тест SCD2: симулируем изменение данных
from pyspark.sql import SparkSession
from engines.scd2_engine import SCD2Engine, SCD2Config
spark = SparkSession.builder.remote("sc://localhost:15002").getOrCreate()
# День 1: первый прогон (все записи - новые)
cfg_day1 = SCD2Config(
source_table="nessie.staging.orders_latest",
target_table="nessie.dw.dim_orders",
natural_key=["order_id"],
tracked_columns=["status", "amount"],
execution_date="2024-01-01",
)
engine_day1 = SCD2Engine(spark, cfg_day1)
result1 = engine_day1.run()
print(f"День 1: {result1}") # {'new': N, 'changed': 0, 'deleted': 0}
# День 2: несколько заказов изменили статус
# (предварительно обновить nessie.staging.orders_latest)
cfg_day2 = SCD2Config(
source_table="nessie.staging.orders_latest",
target_table="nessie.dw.dim_orders",
natural_key=["order_id"],
tracked_columns=["status", "amount"],
execution_date="2024-01-02",
)
engine_day2 = SCD2Engine(spark, cfg_day2)
result2 = engine_day2.run()
print(f"День 2: {result2}") # {'new': 0, 'changed': M, 'deleted': 0}
# Проверка: is_current=false для старых версий
spark.sql("""
SELECT order_id, status, effective_from, effective_to, is_current
FROM nessie.dw.dim_orders
WHERE order_id IN (SELECT order_id FROM nessie.dw.dim_orders WHERE is_current = false)
ORDER BY order_id, effective_from
""").show()
2. Протестировать идемпотентность CDC¶
from engines.cdc_engine import CDCEngine, CDCConfig
cfg = CDCConfig(
source_table="nessie.bronze.orders_cdc_raw",
target_table="nessie.staging.orders_clean",
primary_key=["order_id"],
op_column="op",
ts_column="ts_ms",
execution_date="2024-01-15",
)
engine = CDCEngine(spark, cfg)
# Первый запуск
result1 = engine.run()
count1 = spark.table("nessie.staging.orders_clean").count()
# Второй запуск (тот же день - идемпотентность)
result2 = engine.run()
count2 = spark.table("nessie.staging.orders_clean").count()
assert count1 == count2, "CDC не идемпотентен! Есть дублирование."
print(f"Идемпотентность подтверждена: {count1} == {count2}")
3. Попросить агента добавить новый источник¶
opencode
# > В PostgreSQL появилась новая таблица 'shipments'.
# Нужно: bronze ingestion + SCD2 в DWH.
# Ключ: shipment_id. Отслеживать: status, carrier, tracking_number.
# Агент запустит два skill'а последовательно:
# 1. generate-scd2-config для shipments
# 2. Создаст configs/bronze/shipments_from_postgres.yaml
# 3. Запустит первый bronze ingestion
# 4. Запустит первый SCD2 прогон
# 5. Покажет результаты из audit-таблицы
Резюме¶
Metadata-driven подход + Custom Skills превращает добавление нового источника из нескольких дней в несколько минут: DE описывает источник в YAML, агент запускает проверенный движок.