Custom Skills для DE: SCD2, CDC, Bronze Ingestion и MERGE-паттерны

Строим переиспользуемые pipeline-генераторы для DE-задач: metadata-driven SCD2, CDC ingestion, Delta Lake MERGE и bronze-слой. Упаковываем в Custom Skills для AI-агента.

platform

Проблема 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, агент запускает проверенный движок.