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.
От 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 принимает решения. Инструменты находят факты, разработчик делает выводы.