Custom Skills для DE (часть 2): Spark Reviewer, SQL Reviewer, DAG Generator, Contract Validator

Строим DE-инструменты нового поколения: автоматический Spark-ревьювер anti-patterns, SQL-линтер, metadata-driven Airflow DAG generator и data contract validator.

platform

От pipeline templates к инженерным инструментам

В предыдущем уроке мы строили pipeline generators - движки для автоматической обработки данных (SCD2, CDC, Bronze). Теперь идём дальше: строим инженерные инструменты - Skills, которые помогают самим разработчикам писать лучший код и проектировать надёжные системы.

Spark Optimization Reviewer

Проблема: код работает, но упадёт на production

Junior DE пишет Spark-код, который проходит unit-тесты на маленьком датасете, но вызывает OOM, стрэглеры или бесконечные shuffle на production-данных. Эти проблемы трудно поймать в code review - reviewer должен держать в голове десятки anti-patterns.

Архитектура Spark Reviewer

Реализация статического анализа

# tools/spark_reviewer/static_analyzer.py
"""
Статический анализ PySpark-кода: поиск anti-patterns без запуска.
Работает как быстрый pre-commit check.
"""
import ast
import re
from dataclasses import dataclass
from typing import Optional
from pathlib import Path


@dataclass
class Finding:
    severity: str     # "critical", "warning", "info"
    line: int
    code: str         # код строки
    rule: str         # название правила
    message: str      # описание проблемы
    suggestion: str   # как исправить


class SparkStaticAnalyzer:
    """Анализирует PySpark-код без его выполнения."""

    def analyze_file(self, path: str) -> list[Finding]:
        source = Path(path).read_text()
        return self._analyze_source(source)

    def _analyze_source(self, source: str) -> list[Finding]:
        findings = []
        lines = source.splitlines()

        for i, line in enumerate(lines, 1):
            stripped = line.strip()

            # ── Правило 1: Python UDF ──────────────────────────────
            if re.search(r'@udf\s*\(', stripped) or \
               re.search(r'spark\.udf\.register', stripped):
                findings.append(Finding(
                    severity="warning",
                    line=i,
                    code=stripped,
                    rule="SPARK_UDF_001",
                    message="Python UDF замедляет выполнение в 5-100×",
                    suggestion="Используй built-in функции: lower(), regexp_extract(), "
                               "или Pandas UDF (@pandas_udf) для ML-логики"
                ))

            # ── Правило 2: collect() без limit ─────────────────────
            if re.search(r'\.collect\(\)', stripped) and \
               not re.search(r'\.limit\(\d+\)\.collect\(\)', stripped):
                findings.append(Finding(
                    severity="critical",
                    line=i,
                    code=stripped,
                    rule="SPARK_OOM_001",
                    message=".collect() без .limit() может убить Driver OOM",
                    suggestion="Используй .limit(N).collect() или записывай "
                               "в файл через .write.parquet()"
                ))

            # ── Правило 3: df.count() в цикле ──────────────────────
            if re.search(r'\.count\(\)', stripped):
                # Эвристика: проверяем, есть ли рядом for/while
                context_start = max(0, i - 5)
                context = "\n".join(lines[context_start:i])
                if re.search(r'\bfor\b|\bwhile\b', context):
                    findings.append(Finding(
                        severity="warning",
                        line=i,
                        code=stripped,
                        rule="SPARK_ACTION_001",
                        message=".count() в цикле: каждый вызов - отдельный Job",
                        suggestion="Вынеси .count() за пределы цикла, "
                                   "кэшируй DataFrame перед циклом"
                    ))

            # ── Правило 4: crossJoin / cartesian product ──────────
            if re.search(r'\.crossJoin\(', stripped):
                findings.append(Finding(
                    severity="critical",
                    line=i,
                    code=stripped,
                    rule="SPARK_SHUFFLE_001",
                    message="crossJoin создаёт O(N×M) записей - разрастётся до OOM",
                    suggestion="Убедись, что crossJoin действительно нужен. "
                               "Обычно это ошибка в условии JOIN"
                ))

            # ── Правило 5: select(*) ──────────────────────────────
            if re.search(r'\.select\(\s*["\']?\*["\']?\s*\)', stripped) or \
               re.search(r'select\s*\(\s*\*\s*\)', stripped):
                findings.append(Finding(
                    severity="info",
                    line=i,
                    code=stripped,
                    rule="SPARK_COLUMN_001",
                    message="select(*) читает все колонки - нарушает column pruning",
                    suggestion="Явно перечисли нужные колонки: .select('id', 'name', ...)"
                ))

        return findings

    def format_report(self, findings: list[Finding], path: str) -> str:
        if not findings:
            return f"✅ {path}: Spark anti-patterns не найдены"

        critical = [f for f in findings if f.severity == "critical"]
        warnings = [f for f in findings if f.severity == "warning"]
        infos = [f for f in findings if f.severity == "info"]

        lines = [f"🔍 Spark Review: {path}",
                 f"Critical: {len(critical)} | Warnings: {len(warnings)} | Info: {len(infos)}",
                 ""]

        for f in sorted(findings, key=lambda x: (x.severity != "critical", x.line)):
            icon = "🔴" if f.severity == "critical" else "🟡" if f.severity == "warning" else "ℹ️"
            lines.extend([
                f"{icon} [{f.rule}] Line {f.line}: {f.message}",
                f"   Code: {f.code}",
                f"   Fix:  {f.suggestion}",
                ""
            ])

        return "\n".join(lines)

Skill для OpenCode: автоматический ревью

# .opencode/skills/spark-review.yaml
name: spark-review
description: |
  Анализирует PySpark-код на anti-patterns и предлагает оптимизации.
  Запускай перед каждым code review.
trigger: "review|проверь|оптимизируй|антипаттерн"
steps:
  - name: static_analysis
    tool: bash
    command: "python3 tools/spark_reviewer/static_analyzer.py {file}"
  - name: get_plan
    tool: bash
    command: |
      python3 -c "
      from pyspark.sql import SparkSession
      spark = SparkSession.builder.remote('sc://localhost:15002').getOrCreate()
      exec(open('{file}').read())
      # Если в файле есть df - покажем план
      if 'df' in dir():
          df.explain(mode='extended')
      "
  - name: llm_review
    llm: true
    prompt: |
      Ты - Senior Data Engineer с 10+ лет опыта Spark оптимизации.

      Проверь этот PySpark-код и Physical Plan на проблемы производительности.

      Static Analysis результаты:
      {static_analysis.output}

      Physical Plan:
      {get_plan.output}

      Дай структурированный отчёт:
      1. КРИТИЧЕСКИЕ ПРОБЛЕМЫ (риск OOM, бесконечное выполнение)
      2. ПРЕДУПРЕЖДЕНИЯ (suboptimal, но не критично)
      3. ПРЕДЛОЖЕНИЯ (улучшения без явных проблем)

      Для каждой проблемы: конкретное место в коде + исправленный вариант.
      Используй правила из docs/spark_guidelines.md нашего проекта.

Пример работы Spark Reviewer

# Входной файл: etl_bad.py (типичный код джуна)

@udf("string")
def clean_name(s):
    return s.strip().lower() if s else None

orders = spark.table("nessie.staging.orders")
users = spark.table("nessie.staging.users")

# ❌ Python UDF + collect + crossJoin - тройной anti-pattern
result = orders.withColumn("clean_user", clean_name("user_name")) \
    .crossJoin(users) \
    .collect()

# После review бот предложит:
# result = orders.withColumn("clean_user", lower(trim(col("user_name")))) \
#     .join(broadcast(users), "user_id") \
#     .write.parquet("/output/")

SQL Reviewer

Проблема: нечитаемый и неэффективный SQL

SQL - основной язык DE, но без стандартов запросы становятся нечитаемыми:

-- ❌ Типичный плохой SQL
SELECT * FROM orders o, users u WHERE o.uid=u.id AND o.status='completed'
AND o.amount>0 AND u.country IN (SELECT country FROM allowed_countries)
AND (SELECT COUNT(*) FROM order_items oi WHERE oi.order_id=o.id)>0

Проблемы:

  • Implicit join (comma-separated) вместо явного JOIN
  • SELECT * без column pruning
  • Коррелированный подзапрос в WHERE (N+1 problem)
  • Нет форматирования и CTE для читаемости

Архитектура SQL Reviewer

Реализация SQL Rule Engine

# tools/sql_reviewer/rule_engine.py
"""
Кастомные SQL-правила поверх sqlfluff.
Настраивается через sql_style_guide.yaml.
"""
import re
from dataclasses import dataclass


@dataclass
class SQLFinding:
    rule: str
    severity: str
    line: Optional[int]
    message: str
    suggestion: str
    fixed_sql: Optional[str] = None


class SQLRuleEngine:
    """Применяет правила SQL style guide."""

    def analyze(self, sql: str) -> list[SQLFinding]:
        findings = []
        findings.extend(self._check_implicit_joins(sql))
        findings.extend(self._check_select_star(sql))
        findings.extend(self._check_correlated_subqueries(sql))
        findings.extend(self._check_missing_cte(sql))
        findings.extend(self._check_window_vs_self_join(sql))
        return findings

    def _check_implicit_joins(self, sql: str) -> list[SQLFinding]:
        """Implicit JOIN через запятую - архаизм, ломает читаемость."""
        pattern = r'FROM\s+\w+\s*,\s*\w+'
        if re.search(pattern, sql, re.IGNORECASE):
            return [SQLFinding(
                rule="SQL_JOIN_001",
                severity="warning",
                line=None,
                message="Implicit JOIN через запятую устарел и трудночитаем",
                suggestion="Используй явный INNER JOIN ... ON condition"
            )]
        return []

    def _check_select_star(self, sql: str) -> list[SQLFinding]:
        """SELECT * - нарушает column pruning в Spark."""
        if re.search(r'SELECT\s+\*', sql, re.IGNORECASE):
            return [SQLFinding(
                rule="SQL_COLUMN_001",
                severity="warning",
                line=None,
                message="SELECT * читает все колонки, нарушает column pruning",
                suggestion="Перечисли нужные колонки явно"
            )]
        return []

    def _check_correlated_subqueries(self, sql: str) -> list[SQLFinding]:
        """Коррелированные подзапросы = N+1 problem в SQL."""
        # Эвристика: SELECT в WHERE с ссылкой на внешнюю таблицу
        pattern = r'WHERE\s+.*\(\s*SELECT\s+.*FROM\s+.*WHERE\s+.*\.\w+\s*=\s*\w+\.\w+'
        if re.search(pattern, sql, re.IGNORECASE | re.DOTALL):
            return [SQLFinding(
                rule="SQL_SUBQUERY_001",
                severity="critical",
                line=None,
                message="Коррелированный подзапрос в WHERE: N+1 problem, "
                        "выполняется для каждой строки",
                suggestion="Замени на JOIN или оконную функцию (EXISTS, SEMI JOIN)"
            )]
        return []

    def _check_missing_cte(self, sql: str) -> list[SQLFinding]:
        """Сложные запросы без CTE - нечитаемы."""
        # Эвристика: вложенность > 2 уровней без WITH
        nested_depth = sql.count("(SELECT")
        if nested_depth > 2 and "WITH" not in sql.upper():
            return [SQLFinding(
                rule="SQL_CTE_001",
                severity="info",
                line=None,
                message=f"Запрос имеет {nested_depth} вложенных SELECT без CTE",
                suggestion="Используй WITH ... AS () для декомпозиции и читаемости"
            )]
        return []

    def _check_window_vs_self_join(self, sql: str) -> list[SQLFinding]:
        """Self-join для нумерации строк - используй ROW_NUMBER()."""
        if re.search(r'JOIN\s+\w+\s+\w+\s+ON\s+\w+\.\w+\s*=\s*\w+\.\w+', sql,
                     re.IGNORECASE) and "ROW_NUMBER" not in sql.upper():
            # Простая эвристика - не ловит все случаи
            pass
        return []

SQL Style Guide (конфиг для sqlfluff)

# .sqlfluff (проектный конфиг)
[sqlfluff]
dialect = sparksql
templater = jinja
max_line_length = 100

[sqlfluff:rules:layout.indent]
indent_unit = space
tab_space_size = 4

[sqlfluff:rules:capitalisation.keywords]
capitalisation_policy = upper

[sqlfluff:rules:capitalisation.identifiers]
capitalisation_policy = lower

[sqlfluff:rules:aliasing.table]
aliasing = explicit

[sqlfluff:rules:structure.subquery]
forbid_subquery_in = join  # подзапросы в JOIN → CTE

Skill для SQL ревью

# .opencode/skills/sql-review.yaml
name: sql-review
description: |
  Проверяет SQL на стиль, производительность и логические ошибки.
  Предлагает рефакторинг с использованием CTE и оконных функций.
trigger: "sql review|проверь sql|оптимизируй запрос"
steps:
  - name: format_check
    tool: bash
    command: "sqlfluff lint {file} --dialect sparksql"
  - name: rule_check
    tool: bash
    command: "python3 tools/sql_reviewer/rule_engine.py {file}"
  - name: read_sql
    tool: read
    path: "{file}"
  - name: llm_review
    llm: true
    prompt: |
      Ты - Senior Data Engineer, эксперт по Spark SQL.

      Проверь этот SQL запрос:
      {read_sql.output}

      sqlfluff результаты: {format_check.output}
      Rule Engine: {rule_check.output}

      Задачи:
      1. Найди логические ошибки (неверный grain, missing filter)
      2. Предложи оконные функции вместо self-join/subquery где применимо
      3. Документируй CTE (добавь комментарий к каждому)
      4. Перепиши запрос в улучшенном виде

      Формат ответа:
      ## Проблемы
      ## Улучшенный SQL
      ## Объяснение изменений

Airflow DAG Generator

Проблема: copy-paste DAG-файлы

Типичный ETL-пайплайн в Airflow: похожие DAG-файлы для разных таблиц. Каждый новый источник = копирование существующего DAG + ручная правка + ошибки.

DAG Config Schema

# configs/dags/ingestion_orders.yaml
dag_id: "ingestion_orders"
description: "Загрузка заказов из PostgreSQL в Bronze слой"
schedule: "0 2 * * *"         # Каждый день в 2:00 UTC
start_date: "2024-01-01"
max_active_runs: 1
catchup: false
tags: ["bronze", "orders", "daily"]

default_args:
  owner: "data-team"
  retries: 3
  retry_delay_minutes: 5
  execution_timeout_hours: 2
  on_failure_callback: "callbacks.slack_alert"
  on_retry_callback: "callbacks.slack_retry"

tasks:
  - id: "wait_for_source"
    type: "sql_sensor"
    sql: "SELECT COUNT(*) FROM orders WHERE DATE(updated_at) = '{{ ds }}'"
    poke_interval: 300
    timeout: 3600

  - id: "ingest_bronze"
    type: "spark_submit"
    depends_on: ["wait_for_source"]
    script: "spark_jobs/bronze_ingestion.py"
    args:
      source: "orders"
      execution_date: "{{ ds }}"
    resources:
      executor_memory: "2g"
      executor_cores: 2
      num_executors: 5

  - id: "validate_bronze"
    type: "spark_submit"
    depends_on: ["ingest_bronze"]
    script: "spark_jobs/validate_bronze.py"
    args:
      table: "nessie.bronze.orders_raw"
      date: "{{ ds }}"

  - id: "notify_success"
    type: "slack_message"
    depends_on: ["validate_bronze"]
    channel: "#data-pipelines"
    message: "✅ Orders ingested for {{ ds }}: {{ ti.xcom_pull('validate_bronze', 'row_count') }} rows"

DAG Generator Implementation

# tools/dag_generator/generator.py
"""
Metadata-driven Airflow DAG generator.
YAML конфиг → готовый Python DAG файл.
"""
import yaml
from pathlib import Path
from jinja2 import Environment, FileSystemLoader
from dataclasses import dataclass, field
from typing import Optional


@dataclass
class DAGConfig:
    dag_id: str
    description: str
    schedule: str
    start_date: str
    max_active_runs: int = 1
    catchup: bool = False
    tags: list[str] = field(default_factory=list)
    default_args: dict = field(default_factory=dict)
    tasks: list[dict] = field(default_factory=list)

    @classmethod
    def from_yaml(cls, path: str) -> "DAGConfig":
        with open(path) as f:
            data = yaml.safe_load(f)
        return cls(**data)

    def validate(self) -> list[str]:
        """Проверяет конфиг на ошибки до генерации."""
        errors = []

        # Проверка зависимостей
        task_ids = {t["id"] for t in self.tasks}
        for task in self.tasks:
            for dep in task.get("depends_on", []):
                if dep not in task_ids:
                    errors.append(f"Task '{task['id']}' зависит от несуществующей задачи '{dep}'")

        # Проверка циклических зависимостей (DFS)
        visited = set()
        rec_stack = set()

        def has_cycle(node, graph):
            visited.add(node)
            rec_stack.add(node)
            for neighbour in graph.get(node, []):
                if neighbour not in visited:
                    if has_cycle(neighbour, graph):
                        return True
                elif neighbour in rec_stack:
                    return True
            rec_stack.discard(node)
            return False

        graph = {t["id"]: t.get("depends_on", []) for t in self.tasks}
        for task_id in task_ids:
            if task_id not in visited:
                if has_cycle(task_id, graph):
                    errors.append(f"Обнаружена циклическая зависимость в task '{task_id}'")
                    break

        return errors


class DAGGenerator:
    """Генерирует Python DAG файлы из YAML конфигов."""

    def __init__(self, templates_dir: str = "tools/dag_generator/templates"):
        self.env = Environment(
            loader=FileSystemLoader(templates_dir),
            trim_blocks=True,
            lstrip_blocks=True,
        )

    def generate(self, config: DAGConfig, output_dir: str = "dags") -> Path:
        """Генерирует DAG файл. Возвращает путь к сгенерированному файлу."""
        errors = config.validate()
        if errors:
            raise ValueError(f"Ошибки в конфиге DAG:\n" + "\n".join(errors))

        template = self.env.get_template("standard_dag.py.j2")
        dag_code = template.render(cfg=config)

        output_path = Path(output_dir) / f"{config.dag_id}.py"
        output_path.write_text(dag_code)

        return output_path

Jinja2 шаблон DAG

# tools/dag_generator/templates/standard_dag.py.j2
"""
Auto-generated DAG: {{ cfg.dag_id }}
Source: configs/dags/{{ cfg.dag_id }}.yaml
DO NOT EDIT MANUALLY - изменяй конфиг и перегенерируй
"""
from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
from airflow.sensors.sql import SqlSensor

from callbacks import slack_alert, slack_retry

default_args = {
    "owner": "{{ cfg.default_args.owner | default('data-team') }}",
    "retries": {{ cfg.default_args.retries | default(3) }},
    "retry_delay": timedelta(minutes={{ cfg.default_args.retry_delay_minutes | default(5) }}),
    "execution_timeout": timedelta(hours={{ cfg.default_args.execution_timeout_hours | default(2) }}),
    "on_failure_callback": slack_alert,
    "on_retry_callback": slack_retry,
}

with DAG(
    dag_id="{{ cfg.dag_id }}",
    description="{{ cfg.description }}",
    schedule="{{ cfg.schedule }}",
    start_date=datetime({{ cfg.start_date.replace('-', ', ') }}),
    max_active_runs={{ cfg.max_active_runs }},
    catchup={{ cfg.catchup }},
    tags={{ cfg.tags | tojson }},
    default_args=default_args,
) as dag:

{% for task in cfg.tasks %}
{% if task.type == "spark_submit" %}
    {{ task.id }} = SparkSubmitOperator(
        task_id="{{ task.id }}",
        application="spark_jobs/{{ task.script }}",
        application_args=[
{% for key, val in task.args.items() %}
            "--{{ key }}", "{{ val }}",
{% endfor %}
        ],
        conf={
            "spark.executor.memory": "{{ task.resources.executor_memory | default('2g') }}",
            "spark.executor.cores": "{{ task.resources.executor_cores | default(2) }}",
            "spark.executor.instances": "{{ task.resources.num_executors | default(5) }}",
        },
    )

{% elif task.type == "sql_sensor" %}
    {{ task.id }} = SqlSensor(
        task_id="{{ task.id }}",
        conn_id="postgres_prod",
        sql="{{ task.sql }}",
        poke_interval={{ task.poke_interval | default(300) }},
        timeout={{ task.timeout | default(3600) }},
        mode="reschedule",
    )

{% endif %}
{% endfor %}

    # Зависимости задач
{% for task in cfg.tasks %}
{% if task.depends_on %}
    {{ task.id }}.set_upstream([
{% for dep in task.depends_on %}
        {{ dep }},
{% endfor %}
    ])
{% endif %}
{% endfor %}

Skill для генерации DAG

# .opencode/skills/generate-airflow-dag.yaml
name: generate-airflow-dag
description: |
  Создаёт Airflow DAG из YAML-конфига.
  Проверяет зависимости, генерирует Python-файл, валидирует синтаксис.
trigger: "dag|airflow|pipeline|пайплайн"
steps:
  - name: check_existing
    tool: bash
    command: "ls configs/dags/"
  - name: generate
    tool: bash
    command: |
      python3 -c "
      from tools.dag_generator.generator import DAGGenerator, DAGConfig
      cfg = DAGConfig.from_yaml('configs/dags/{dag_name}.yaml')
      gen = DAGGenerator()
      path = gen.generate(cfg)
      print(f'Generated: {path}')
      "
  - name: validate_syntax
    tool: bash
    command: "python3 -m py_compile dags/{dag_name}.py && echo 'Syntax OK'"
  - name: airflow_check
    tool: bash
    command: "airflow dags list-import-errors 2>&1 | grep {dag_name} || echo 'No import errors'"

Data Contract Validator

Проблема: изменения схемы ломают production

Архитектура Contract Validator

Реализация Contract Validator

# tools/contract_validator/validator.py
"""
Проверяет соответствие фактической схемы Spark-таблицы контракту.
Используется в CI/CD для блокировки breaking changes.
"""
import yaml
from dataclasses import dataclass
from pyspark.sql import SparkSession
from pyspark.sql.types import StructType


@dataclass
class ColumnSpec:
    name: str
    type: str
    nullable: bool = True
    description: str = ""
    pii: bool = False


@dataclass
class ContractSpec:
    table: str
    version: str
    columns: list[ColumnSpec]
    freshness_sla: str = ""
    row_count_min: int = 0

    @classmethod
    def from_yaml(cls, path: str) -> "ContractSpec":
        with open(path) as f:
            data = yaml.safe_load(f)

        columns = [ColumnSpec(**col) for col in data.pop("columns")]
        return cls(columns=columns, **{k: v for k, v in data.items()
                                        if k in cls.__dataclass_fields__})


@dataclass
class ValidationResult:
    is_compatible: bool
    breaking_changes: list[str]
    non_breaking_changes: list[str]
    warnings: list[str]

    def is_ok_to_deploy(self) -> bool:
        return len(self.breaking_changes) == 0

    def summary(self) -> str:
        if self.is_compatible:
            icon = "✅"
            status = "Compatible"
        elif self.breaking_changes:
            icon = "❌"
            status = "BREAKING CHANGES"
        else:
            icon = "⚠️"
            status = "Non-breaking changes"

        lines = [f"{icon} Contract Validation: {status}"]

        if self.breaking_changes:
            lines.append("\n💥 BREAKING CHANGES (блокировать деплой):")
            for change in self.breaking_changes:
                lines.append(f"  - {change}")

        if self.non_breaking_changes:
            lines.append("\n⚠️ Non-breaking changes (уведомить consumers):")
            for change in self.non_breaking_changes:
                lines.append(f"  - {change}")

        if self.warnings:
            lines.append("\nℹ️ Предупреждения:")
            for w in self.warnings:
                lines.append(f"  - {w}")

        return "\n".join(lines)


class ContractValidator:
    """Сравнивает фактическую схему таблицы с контрактом."""

    # Типы, которые совместимы при расширении (non-breaking)
    WIDENING_COMPATIBLE = {
        ("int", "long"),
        ("float", "double"),
        ("date", "timestamp"),
    }

    def __init__(self, spark: SparkSession):
        self.spark = spark

    def validate(self, contract: ContractSpec) -> ValidationResult:
        # Читаем фактическую схему таблицы
        actual_schema = self.spark.table(contract.table).schema
        actual_cols = {f.name: f for f in actual_schema.fields}

        breaking = []
        non_breaking = []
        warnings = []

        contract_cols = {col.name: col for col in contract.columns}

        # 1. Проверяем колонки из контракта
        for col_name, spec in contract_cols.items():
            if col_name not in actual_cols:
                if not spec.nullable:
                    # NOT NULL колонка исчезла → breaking
                    breaking.append(
                        f"Колонка '{col_name}' (NOT NULL) удалена из таблицы"
                    )
                else:
                    warnings.append(f"Nullable колонка '{col_name}' отсутствует в таблице")
                continue

            actual_field = actual_cols[col_name]
            actual_type = actual_field.dataType.simpleString()

            # Проверка типа
            if actual_type != spec.type:
                widening = (spec.type, actual_type) in self.WIDENING_COMPATIBLE
                if widening:
                    non_breaking.append(
                        f"Колонка '{col_name}': тип расширен {spec.type}{actual_type}"
                    )
                else:
                    breaking.append(
                        f"Колонка '{col_name}': несовместимое изменение типа "
                        f"{spec.type}{actual_type}"
                    )

            # Проверка nullable
            if not spec.nullable and actual_field.nullable:
                breaking.append(
                    f"Колонка '{col_name}': изменена с NOT NULL на NULLABLE"
                )

        # 2. Новые колонки в таблице (non-breaking)
        for col_name in actual_cols:
            if col_name not in contract_cols:
                non_breaking.append(f"Новая колонка '{col_name}' не описана в контракте")

        is_compatible = len(breaking) == 0

        return ValidationResult(
            is_compatible=is_compatible,
            breaking_changes=breaking,
            non_breaking_changes=non_breaking,
            warnings=warnings,
        )

    def check_quality(self, contract: ContractSpec, execution_date: str) -> list[str]:
        """Проверяет качество данных по условиям контракта."""
        issues = []
        df = self.spark.table(contract.table)

        # Row count check
        if contract.row_count_min > 0:
            count = df.filter(f"DATE(created_at) = '{execution_date}'").count()
            if count < contract.row_count_min:
                issues.append(
                    f"Row count {count} < minimum {contract.row_count_min} "
                    f"for {execution_date}"
                )

        return issues

Интеграция в CI/CD

# .github/workflows/data-contract-check.yml
name: Data Contract Validation

on:
  pull_request:
    paths:
      - 'schemas/**'
      - 'contracts/**'
      - 'spark_jobs/**'

jobs:
  validate-contracts:
    runs-on: ubuntu-latest
    steps:
      - uses: actions/checkout@v3

      - name: Run Contract Validator
        run: |
          python3 -c "
          from pyspark.sql import SparkSession
          from tools.contract_validator.validator import ContractValidator, ContractSpec

          spark = SparkSession.builder.remote('sc://staging-spark:15002').getOrCreate()
          validator = ContractValidator(spark)

          import glob
          exit_code = 0
          for contract_file in glob.glob('contracts/**/*.yaml', recursive=True):
              contract = ContractSpec.from_yaml(contract_file)
              result = validator.validate(contract)
              print(result.summary())
              if not result.is_ok_to_deploy():
                  exit_code = 1

          exit(exit_code)
          "

Skill для валидации контракта

# .opencode/skills/validate-contract.yaml
name: validate-contract
description: |
  Проверяет соответствие Spark-таблицы Data Contract.
  Запускай перед деплоем изменений схемы.
trigger: "validate|контракт|schema|совместимость"
steps:
  - name: find_contract
    tool: bash
    command: "find contracts/ -name '{table_name}*.yaml'"
  - name: run_validator
    tool: bash
    command: |
      python3 -c "
      from pyspark.sql import SparkSession
      from tools.contract_validator.validator import ContractValidator, ContractSpec
      spark = SparkSession.builder.remote('sc://localhost:15002').getOrCreate()
      contract = ContractSpec.from_yaml('{find_contract.output}'.strip())
      result = ContractValidator(spark).validate(contract)
      print(result.summary())
      "
  - name: notify
    llm: true
    condition: "breaking_changes in run_validator.output"
    prompt: |
      Обнаружены breaking changes в контракте для таблицы {table_name}.

      Результат валидации:
      {run_validator.output}

      Составь список affected consumers (посмотри в contracts/{table_name}*.yaml → consumers секция)
      и подготовь сообщение для уведомления в Slack с деталями изменений и рекомендуемыми действиями.

Human-in-the-Loop: нельзя слепо доверять автоматике

Правило: автоматика снижает cognitive load, но не заменяет judgment. Reviewer принимает финальное решение - инструменты готовят факты.

Маскирование данных перед отправкой в LLM

Агент может отправлять код и SQL в внешний LLM (Claude, GPT-4). Нужно не допустить утечки PII:

# tools/pii_masker.py
"""
Маскирует чувствительные данные перед отправкой в LLM.
"""
import re


class PIIMasker:
    """Маскирует PII и корпоративные секреты в тексте."""

    PATTERNS = [
        # Email адреса
        (r'\b[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\.[A-Z|a-z]{2,}\b', '<EMAIL>'),
        # IP адреса
        (r'\b\d{1,3}\.\d{1,3}\.\d{1,3}\.\d{1,3}\b', '<IP_ADDRESS>'),
        # API ключи (паттерн: длинные hex/base64 строки)
        (r'\b[A-Za-z0-9+/]{40,}\b', '<API_KEY>'),
        # Пароли в строках подключения
        (r'password["\s:=]+["\']?[^"\s,;>]+', 'password=<REDACTED>'),
        # JWT токены
        (r'eyJ[A-Za-z0-9_-]+\.[A-Za-z0-9_-]+\.[A-Za-z0-9_-]+', '<JWT_TOKEN>'),
    ]

    # Внутренние названия (конфигурируются для каждой компании)
    INTERNAL_NAMES = [
        "prod-db.internal",
        "spark-cluster.corp.local",
    ]

    def mask(self, text: str) -> str:
        result = text
        for pattern, replacement in self.PATTERNS:
            result = re.sub(pattern, replacement, result)
        for name in self.INTERNAL_NAMES:
            result = result.replace(name, "<INTERNAL_HOST>")
        return result

Практика

1. Запустить Spark Reviewer на реальном коде

opencode

# > Проверь spark_jobs/etl_orders.py на Spark anti-patterns

# Агент:
# → static_analysis: найдёт Python UDF на строке 42, collect() на строке 87
# → get_plan: покажет Exchange узел (shuffle) в неожиданном месте
# → llm_review: объяснит почему exchange лишний и предложит broadcast
# → Напишет исправленный вариант кода

2. Создать новый Airflow DAG

opencode

# > Создай Airflow DAG для загрузки таблицы shipments из PostgreSQL.
#   Расписание: каждый день в 3:00 UTC.
#   Нужен sensor на наличие данных, потом spark submit, потом slack уведомление.

# Агент:
# → Читает configs/dags/ для примера
# → Создаёт configs/dags/ingestion_shipments.yaml
# → Генерирует dags/ingestion_shipments.py
# → Проверяет синтаксис: python3 -m py_compile
# → Показывает граф зависимостей задач

3. Провалидировать контракт перед деплоем

opencode

# > Нам нужно добавить колонку 'shipping_cost' (decimal, nullable) в orders.
#   Проверь, не сломает ли это существующие контракты.

# Агент:
# → find contracts/: находит shop_orders_v2.yaml
# → Читает контракт, проверяет consumers
# → Запускает ContractValidator с обновлённой схемой
# → Новая nullable колонка = non-breaking
# → "Безопасно добавить. Уведомить analytics-team и ml-team через data-contracts@company.com"

Best Practices: когда автоматизация - это хорошо

Инструмент Ценность Риск Mitigation
Spark Reviewer Находит критичные проблемы до code review False positives раздражают разработчиков Настрой severity, critical блокирует, info - совет
SQL Reviewer Автоформатирование + логические ошибки LLM может предложить неверный рефакторинг Human review финального SQL обязателен
DAG Generator Устраняет copy-paste, единый стандарт Сложные cases не покрыты шаблоном Поддерживай возможность писать DAG вручную
Contract Validator Блокирует breaking changes автоматически Строгая валидация замедляет разработку Non-breaking = warning, breaking = блокировка

Резюме

Четыре Custom Skills второго уровня превращают AI-агента в активного участника DE-процессов:

Ключевой принцип: автоматика снижает cognitive load, Human-in-the-Loop принимает решения. Инструменты находят факты, разработчик делает выводы.