Agent Context: как дать AI-агенту знания Senior Data Engineer

Превращаем AI-агента из generic code generator в эксперта по вашей платформе: architecture docs, spark_guidelines.md, data contracts, /examples, ADR и schemas.

platform

Проблема вакуума: почему LLM без контекста опасен

LLM без project-специфического контекста - образованный стажёр, который много читал, но ни разу не видел ваш стек. Он предложит:

  • df.collect() в production ETL, не зная, что в вашем кластере это запрещено (данные не помещаются в память Driver'а)
  • pandas_udf вместо встроенных функций, не зная, что ваши Executor'ы имеют лимит 4 GB
  • Delta Lake вместо Iceberg, хотя команда это решение уже отвергла полгода назад
  • Пакет pydantic==1.9 вместо pydantic==2.x, потому что LLM обучался на старых примерах

Grounding (приземление модели) - передача реальных правил, стандартов и метаданных вашей компании, чтобы LLM генерировал уместный, а не generic код.

Иерархия контекста

Контекст для агента строится слоями - от глобального к специфическому:

Каждый слой уточняет предыдущий. При конфликте - нижний слой имеет приоритет (специфичное важнее общего).

Слой 1: Архитектурная документация

Что должно быть в arch/README.md

# Data Platform Architecture

## Стек технологий

| Компонент | Технология | Версия | Где развёрнут |
|---|---|---|---|
| Compute | Apache Spark | 3.5.x | Kubernetes (Spark Operator) |
| Table Format | Apache Iceberg | 1.5.x | - |
| Catalog | Project Nessie | 0.74.x | K8s namespace `nessie` |
| Object Storage | MinIO (S3-compatible) | RELEASE.2024-01 | On-premises |
| Orchestration | Apache Airflow | 2.8.x | K8s namespace `airflow` |
| Operational DB | PostgreSQL | 15 | RDS `prod-db.internal` |
| Metrics | Prometheus + Grafana | - | Monitoring cluster |

## Topology

```mermaid
flowchart LR
    PG["PostgreSQL<br/>Source of Record"]
    SPARK["Apache Spark<br/>Processing Layer"]
    ICE["Iceberg Tables<br/>via Nessie"]
    MINIO["MinIO / S3<br/>Data Lake"]
    AIRFLOW["Apache Airflow<br/>Orchestration"]

    PG --> SPARK
    SPARK --> ICE
    ICE --> MINIO
    AIRFLOW --> SPARK

Ключевые ограничения

  • Executor Memory: максимум 8 GB на Executor, 16 ядер на Worker
  • Driver Memory: максимум 4 GB
  • Партиции: целевой размер 128–256 MB (для Parquet)
  • collect() запрещён в production-коде - данные не помещаются в Driver
  • Shuffle partitions: spark.sql.shuffle.partitions=400 в prod, 10 в dev

Запрещённые паттерны

  • Python UDF (только встроенные функции или Pandas UDF)
  • toPandas() на больших датасетах
  • Hardcoded credentials в коде (использовать Vault/K8s Secrets)
  • Прямой доступ к S3 без Iceberg API (обходит Time Travel и ACID)
    Эта документация позволяет агенту генерировать код, который соответствует реальным ограничениям кластера - не потому что он "умный", а потому что явно сказано.
    
    ## Слой 2: ADR - почему мы сделали именно так
    
    **Architecture Decision Records** - краткие документы, фиксирующие важные решения и причины их принятия. Без ADR агент будет предлагать то, что уже отвергнуто:
    
    ```markdown
    # ADR-003: Выбор Apache Iceberg вместо Delta Lake
    
    **Дата**: 2024-03-15
    **Статус**: Accepted
    **Контекст**: Нужен modern table format для Lakehouse. Рассматривались Delta Lake, Hudi, Iceberg.
    
    ## Решение: Apache Iceberg
    
    ## Причины отказа от альтернатив
    
    **Delta Lake отвергнут потому что**:
    - Тесная привязка к Databricks Runtime для advanced features
    - Лицензия BSL (ограничения на managed service)
    - Нет нативной поддержки Nessie (только Delta Sharing)
    
    **Hudi отвергнут потому что**:
    - Значительно сложнее в операционном управлении
    - Меньше community поддержки для Kubernetes-native деплоя
    
    ## Последствия
    
    - Все новые таблицы создаются через Iceberg API
    - Delta-таблицы (legacy) мигрируются по мере возможности
    - Агент не должен предлагать Delta Lake API в новом коде
    
    ## Ссылки
    - [Iceberg vs Delta vs Hudi comparison](...)
    - [RFC: Lakehouse Table Format Selection](...)
    

Когда агент видит этот ADR - он не предлагает delta.tables.DeltaTable.forPath(). Он знает, что это решение уже принято.

Слой 3: spark_guidelines.md - правила Spark Engineering

Это самый важный документ для генерации производительного Spark-кода:

# Spark Engineering Guidelines

## Правила производительности

### UDF Policy

```python
# ❌ ЗАПРЕЩЕНО: Python UDF (5-100x медленнее built-in)
@udf("string")
def clean(s):
    return s.strip().lower() if s else None

df.withColumn("name", clean("raw_name"))

# ✅ ОБЯЗАТЕЛЬНО: built-in функции
from pyspark.sql.functions import lower, trim
df.withColumn("name", lower(trim(col("raw_name"))))

# ✅ ДОПУСТИМО: Pandas UDF для ML-логики
@pandas_udf("double")
def predict(s: pd.Series) -> pd.Series:
    return model.predict(s.values)

Join Policy

# ❌ ЗАПРЕЩЕНО: shuffle join с маленькой таблицей
large_df.join(small_df, "product_id")  # shuffle обоих датасетов

# ✅ ОБЯЗАТЕЛЬНО: broadcast для таблиц < 50 MB
from pyspark.sql.functions import broadcast
large_df.join(broadcast(small_df), "product_id")

# Порог broadcast: spark.sql.autoBroadcastJoinThreshold=52428800 (50 MB)

collect() Policy

# ❌ ЗАПРЕЩЕНО в production: может убить Driver OOM
result = df.collect()  # весь датасет в память Driver

# ✅ РАЗРЕШЕНО: только для агрегатов (гарантированно маленький результат)
count = df.count()
sample = df.limit(100).collect()

# ✅ РАЗРЕШЕНО: для небольших lookup-таблиц (< 10 000 строк)
lookup = spark.table("nessie.ref.product_categories").collect()

Partitioning Standards

# Целевой размер партиции: 128-256 MB
# Для daily batch jobs:
df.repartition(400)  # prod (400 = spark.sql.shuffle.partitions)
df.repartition(10)   # dev (10 = быстрый старт)

# Для записи в Iceberg:
df.writeTo("nessie.schema.table") \
  .partitionedBy(days("event_date"))  # никогда часы/минуты для daily data

Naming Conventions

# Таблицы: snake_case, единственное число
nessie.shop.completed_order       # ✅
nessie.Shop.CompletedOrders       # ❌

# DataFrame переменные
orders_df = spark.table(...)      # ✅ суффикс _df
df = spark.table(...)             # ❌ слишком generic

# Колонки в SparkSQL: snake_case
SELECT order_id, user_id          # ✅
SELECT OrderId, UserId            # ❌

# Константы
MAX_BATCH_SIZE = 100_000          # ✅ UPPER_SNAKE_CASE

Error Handling

# ✅ Структура pipeline с логированием
import logging
from pyspark.sql import SparkSession

logger = logging.getLogger(__name__)

def run_pipeline(spark: SparkSession, execution_date: str) -> None:
    """Основная функция пайплайна. Всегда принимает SparkSession."""
    logger.info(f"Starting pipeline for {execution_date}")

    try:
        source_df = load_source(spark, execution_date)
        transformed_df = transform(source_df)
        write_output(transformed_df, execution_date)
        logger.info(f"Pipeline completed for {execution_date}")
    except Exception as e:
        logger.error(f"Pipeline failed for {execution_date}: {e}")
        raise  # Всегда re-raise для Airflow retry

Session Configuration Template

# Стандартная конфигурация SparkSession для production
def create_production_session(app_name: str) -> SparkSession:
    return SparkSession.builder \
        .appName(app_name) \
        .config("spark.sql.adaptive.enabled", "true") \
        .config("spark.sql.adaptive.coalescePartitions.enabled", "true") \
        .config("spark.sql.shuffle.partitions", "400") \
        .config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
        .config("spark.sql.catalog.nessie", "org.apache.iceberg.spark.SparkCatalog") \
        .config("spark.sql.catalog.nessie.catalog-impl", "org.apache.iceberg.nessie.NessieCatalog") \
        .config("spark.sql.catalog.nessie.uri", os.environ["NESSIE_URI"]) \
        .getOrCreate()
## Слой 4: Schemas и Data Contracts

### Описание схем таблиц

```yaml
# schemas/shop/completed_orders.yaml
table:
  name: completed_orders
  catalog: nessie
  namespace: shop
  description: "Завершённые заказы после обработки payment gateway"
  owner: "data-team@company.com"
  partition_by: [year, month]
  sort_by: [order_date]

columns:
  - name: order_id
    type: long
    nullable: false
    description: "Уникальный ID заказа из PostgreSQL orders.order_id"
    pii: false

  - name: user_id
    type: integer
    nullable: false
    description: "ID пользователя (FK к users.user_id)"
    pii: true  # требует маскирования в non-prod окружениях

  - name: amount
    type: decimal(10,2)
    nullable: false
    description: "Сумма заказа в рублях"
    business_rule: "amount > 0"

  - name: order_date
    type: date
    nullable: false
    description: "Дата заказа (партиционирующий ключ)"

  - name: year
    type: integer
    nullable: false
    description: "Год (derived from order_date, для Iceberg partition)"

  - name: month
    type: integer
    nullable: false
    description: "Месяц (derived from order_date, для Iceberg partition)"

data_quality:
  freshness_sla: "daily, T+2h"    # данные должны быть актуальны к 2:00
  completeness: "> 99%"           # не более 1% null в критичных полях
  row_count_expected: "> 1000/day"  # минимум строк в день

Data Contracts: гарантии между командами

# contracts/shop_orders_contract.yaml
contract:
  id: "SHOP-ORDERS-V2"
  version: "2.1.0"
  status: active
  effective_date: "2024-01-01"

provider:
  team: "payments-team"
  table: "nessie.shop.completed_orders"
  contact: "payments@company.com"

consumers:
  - team: "analytics-team"
    use_case: "BI dashboards"
    sla: "available by 03:00 UTC daily"

  - team: "ml-team"
    use_case: "recommendation model training"
    sla: "available by 06:00 UTC daily"

guarantees:
  schema_stability: "columns will not be removed without 30-day notice"
  backward_compatible: true
  null_policy:
    order_id: "never null"
    user_id: "never null"
    amount: "never null"

breaking_changes_process:
  notice_period_days: 30
  notification: "data-contracts@company.com"

history:
  - version: "2.0.0"
    date: "2023-06-01"
    change: "Added month column for partition optimization"
  - version: "2.1.0"
    date: "2024-01-01"
    change: "Added pii flag to user_id, affects non-prod masking"

Когда агент видит этот контракт, он понимает:

  • user_id - PII, нельзя включать в test datasets без маскирования
  • Контракт имеет версию - изменение схемы требует процесса уведомления
  • Есть SLA - pipeline должен завершиться до 03:00 UTC

Слой 5: /examples - Reference Implementations

Примеры кода - самый мощный инструмент обучения агента. LLM учится на паттернах гораздо лучше, чем на абстрактных правилах:

Структура /examples

examples/
├── patterns/
│   ├── broadcast_join.py          # Broadcast join с маленьким датасетом
│   ├── incremental_load.py        # Идемпотентная инкрементальная загрузка
│   ├── salting.py                 # Борьба с data skew через salting
│   ├── schema_evolution.py        # Добавление колонок в Iceberg без перезаписи
│   └── time_travel.py             # Time travel запросы в Iceberg
├── session/
│   ├── dev_session.py             # SparkSession для локальной разработки
│   └── prod_session.py            # SparkSession для production
├── testing/
│   ├── unit_test_example.py       # Unit тест для UDF
│   └── integration_test_example.py # Интеграционный тест с реальным Spark
└── airflow/
    ├── spark_dag_example.py       # DAG для запуска Spark через KubernetesPodOperator
    └── sensor_example.py          # Sensor для ожидания данных в Iceberg

Пример: идемпотентная инкрементальная загрузка

# examples/patterns/incremental_load.py
"""
Паттерн: Идемпотентная инкрементальная загрузка в Iceberg.
Безопасно перезапускать при сбоях - не создаёт дублей.

Используй этот паттерн для:
- Daily batch loads из PostgreSQL
- CDC-like incremental processing
- Airflow-задач с возможностью retry
"""
from pyspark.sql import SparkSession, DataFrame
from pyspark.sql.functions import col, to_date, current_timestamp


def load_incremental(
    spark: SparkSession,
    execution_date: str,       # Формат: 'YYYY-MM-DD'
    source_table: str,         # Полное имя таблицы-источника
    target_table: str,         # Полное имя Iceberg-таблицы
    date_column: str = "created_at",
) -> int:
    """
    Загружает данные за один день из source в target (Iceberg).

    Returns:
        Количество записанных строк.
    """
    # ── Шаг 1: Читаем только нужный день (partition pruning) ──
    source_df = spark.table(source_table).filter(
        to_date(col(date_column)) == execution_date
    )

    row_count = source_df.count()
    if row_count == 0:
        print(f"No data for {execution_date}, skipping write")
        return 0

    # ── Шаг 2: Удаляем уже существующие данные за этот день ──
    # (идемпотентность: при retry не будет дублей)
    spark.sql(f"""
        DELETE FROM {target_table}
        WHERE order_date = '{execution_date}'
    """)

    # ── Шаг 3: Записываем свежие данные ──
    source_df.writeTo(target_table).append()

    print(f"Loaded {row_count} rows for {execution_date}")
    return row_count


if __name__ == "__main__":
    spark = SparkSession.builder.remote("sc://localhost:15002").getOrCreate()

    rows = load_incremental(
        spark=spark,
        execution_date="2024-01-15",
        source_table="nessie.staging.orders_raw",
        target_table="nessie.shop.completed_orders",
    )
    print(f"Done: {rows} rows")

Пример: борьба с data skew через salting

# examples/patterns/salting.py
"""
Паттерн: Salting для борьбы с data skew в groupBy/join.

Проблема: один ключ группировки содержит 90% данных → один Executor
перегружен, остальные простаивают.

Решение: добавить случайный salt к ключу → данные распределились →
убрать salt на финальной агрегации.
"""
from pyspark.sql import DataFrame
from pyspark.sql.functions import col, floor, rand, concat_ws, sum as _sum


def aggregate_with_salting(
    df: DataFrame,
    group_key: str,
    value_col: str,
    num_salts: int = 10,
) -> DataFrame:
    """Агрегация с salting для skewed данных."""

    # Фаза 1: добавляем случайный salt
    salted = df.withColumn(
        "salt", floor(rand() * num_salts).cast("int")
    ).withColumn(
        "salted_key", concat_ws("_", col(group_key), col("salt"))
    )

    # Фаза 2: частичная агрегация по salted_key
    partial = salted.groupBy("salted_key", group_key) \
        .agg(_sum(value_col).alias("partial_sum"))

    # Фаза 3: финальная агрегация - убираем salt
    result = partial.groupBy(group_key) \
        .agg(_sum("partial_sum").alias(f"total_{value_col}"))

    return result

Слой 6: ADR и Runbooks

ADR для быстрого старта

# ADR-007: Запрет Python UDF в production-коде

**Дата**: 2024-02-10
**Статус**: Accepted

## Контекст
Команда столкнулась с production-инцидентом: Python UDF на датасете 500M строк
замедлил ETL с 10 минут до 3 часов. Корень проблемы - Pickle-сериализация построчно.

## Решение
**Python UDF запрещены** в production-коде. Альтернативы по приоритету:

1. Built-in Spark SQL functions (`lower()`, `regexp_extract()`, `datediff()`, ...)
2. Pandas UDF (`@pandas_udf`) для ML-логики и сложных батчевых операций
3. Scala UDF через отдельный JAR (для критичных к производительности случаев)

## Исключения
Только с явным approve от Tech Lead:
- Интеграция с legacy third-party библиотекой
- Прототипирование (только в dev-ветке)

## Последствия
- Code review автоматически отклоняет Python UDF без исключения
- OpenCode агент должен проверять и исправлять UDF при генерации кода

Runbook: Отладка Spark OOM

# Runbook: Spark Executor OOM

## Симптомы
- Executor завершился с `Java heap space` или `GC overhead limit exceeded`
- Stage перезапускается бесконечно
- `FetchFailed` из-за потери shuffle-данных

## Диагностика (шаги по порядку)

1. Spark UI → Executors tab → найти Executor с `FAILED` статусом
2. Spark UI → Stages → найти Stage с большим Shuffle Read
3. Логи: `yarn logs -applicationId APP_ID | grep "OutOfMemoryError"`
4. Проверить data skew: Task с временем 10× больше медианы

## Исправления

### Если data skew:
→ Включить AQE: `spark.sql.adaptive.skewJoin.enabled=true`
→ Использовать salting (см. examples/patterns/salting.py)

### Если маленькие партиции в памяти:
→ Уменьшить `spark.sql.shuffle.partitions`
→ Использовать `coalesce()` перед тяжёлыми операциями

### Если слишком большие объекты:
→ Проверить broadcast: отключить если broadcasted таблица > 1 GB
→ Проверить closure в UDF: не захватывать large DataFrame

## Эскалация
Если проблема не решена за 30 минут → уведомить #data-incidents в Slack

Структура AI-ready репозитория

Как агент использует контекст

Prompt injection через CLAUDE.md / .opencode

OpenCode (как и Claude Code) читает специальные файлы в начале каждой сессии - фактически это system prompt, который всегда активен:

# .opencode/AGENTS.md (или CLAUDE.md для Claude Code)

## Project Context

Это Data Platform репозиторий компании X.

## ОБЯЗАТЕЛЬНЫЕ правила при генерации PySpark кода

1. **Никогда не используй Python UDF** - только built-in функции или Pandas UDF
   (ADR-003: запрещено после production-инцидента в Q1 2024)

2. **Никогда не используй collect()** на больших датасетах
   (Executor Memory: max 8 GB, данные не помещаются в Driver)

3. **Всегда используй broadcast()** для таблиц < 50 MB при join

4. **SparkSession** создавай только через `examples/session/prod_session.py` шаблон

5. **Iceberg API**, не Delta Lake (ADR-001: выбран Iceberg)

6. **Партиционирование**: только `days()` или `months()` для temporal данных

## При генерации кода проверяй

- Соответствует ли код `docs/spark_guidelines.md`?
- Есть ли готовый паттерн в `examples/`? Используй его.
- Есть ли PII-поля согласно `contracts/`? Применяй маскирование.

## Запрос помощи

Если задача не описана в документации - сообщи об этом явно,
не изобретай решение на основе общих знаний.

Практика: до и после контекста

Без контекста

opencode  # без .opencode/AGENTS.md

> Напиши ETL для загрузки orders из PostgreSQL в Spark

# Генерирует:
# ❌ Generic код без учёта платформы
from pyspark.sql import SparkSession

spark = SparkSession.builder.master("local").getOrCreate()

orders = spark.read.jdbc(
    url="jdbc:postgresql://localhost/db",
    table="orders",
    properties={"user": "admin", "password": "admin"}  # ❌ hardcoded credentials
)

# ❌ Python UDF вместо built-in
@udf("string")
def format_status(s):
    return s.upper()

orders.withColumn("status", format_status("status")) \
    .write.saveAsTable("orders")  # ❌ не Iceberg API, не партиционировано

С контекстом

# Добавили: arch/README.md, docs/spark_guidelines.md, schemas/, examples/
opencode  # с .opencode/AGENTS.md

> Напиши ETL для загрузки orders из PostgreSQL в Iceberg
# ✅ Production-ready код на основе контекста проекта
import os
import logging
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, to_date, year, month, upper, broadcast

logger = logging.getLogger(__name__)


def create_session() -> SparkSession:
    """SparkSession по шаблону prod_session.py из examples/."""
    return SparkSession.builder \
        .appName("OrdersETL") \
        .config("spark.sql.adaptive.enabled", "true") \
        .config("spark.sql.shuffle.partitions", "400") \
        .config("spark.sql.extensions",
                "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
        .config("spark.sql.catalog.nessie",
                "org.apache.iceberg.spark.SparkCatalog") \
        .config("spark.sql.catalog.nessie.catalog-impl",
                "org.apache.iceberg.nessie.NessieCatalog") \
        .config("spark.sql.catalog.nessie.uri", os.environ["NESSIE_URI"]) \
        .getOrCreate()


def run_etl(spark: SparkSession, execution_date: str) -> None:
    logger.info(f"Loading orders for {execution_date}")

    # ✅ Credentials из env (не hardcoded)
    orders_df = spark.read.format("jdbc") \
        .option("url", os.environ["POSTGRES_URL"]) \
        .option("dbtable", f"(SELECT * FROM orders WHERE DATE(created_at) = '{execution_date}') t") \
        .option("user", os.environ["POSTGRES_USER"]) \
        .option("password", os.environ["POSTGRES_PASSWORD"]) \
        .option("driver", "org.postgresql.Driver") \
        .load()

    # ✅ Built-in functions, не Python UDF (per ADR-003)
    product_cats = spark.table("nessie.ref.product_categories")
    # product_cats < 10 MB → broadcast (per guidelines)

    enriched = orders_df \
        .join(broadcast(product_cats), "product_id", "left") \
        .withColumn("status_upper", upper(col("status"))) \
        .withColumn("order_date", to_date(col("created_at"))) \
        .withColumn("year", year(col("order_date"))) \
        .withColumn("month", month(col("order_date"))) \
        .select("order_id", "user_id", "product_id", "amount",
                "status_upper", "category_name",
                "order_date", "year", "month")
        # ✅ user_id: PII поле из контракта, НЕ маскируем в prod, но
        # не включаем в примерные датасеты (указано в комментарии)

    # ✅ Идемпотентная запись (per examples/patterns/incremental_load.py)
    spark.sql(f"""
        DELETE FROM nessie.shop.completed_orders
        WHERE order_date = '{execution_date}'
    """)

    # ✅ Iceberg writeTo API (per ADR-001), с партиционированием (per schema)
    enriched.writeTo("nessie.shop.completed_orders").append()

    logger.info(f"Loaded {enriched.count()} rows for {execution_date}")

Разница очевидна: агент с контекстом генерирует код, который соответствует стандартам команды, использует правильный стек и безопасен для production.

Anti-patterns плохого контекста

Правила хорошего контекста:

  • Актуальность: устаревший docs хуже отсутствия docs - агент будет уверенно врать
  • Согласованность: правила и примеры должны совпадать
  • Краткость: лучше 5 хороших примеров, чем 50 страниц текста
  • Do/Don't пары: показывай плохой и хороший код рядом

Knowledge Management: документация как operational dependency

В эпоху AI-агентов документация перестала быть "для людей" - она стала operational dependency: от неё зависит качество автоматически генерируемого кода.

Governance для AI-ready documentation:

  • ADR обновляется до изменения кода - агент должен знать о решении раньше, чем увидит новый код
  • Guidelines версионируются в git - можно откатить если агент начал генерировать плохой код
  • Examples проходят code review - они станут образцом для агента
  • Schemas синхронизируются с реальной схемой Iceberg автоматически (через CI)

Резюме

Качество AI-assisted DE-workflow напрямую определяется качеством контекста:

Слой контекста Что решает Ключевой файл
Архитектура Стек, топология, ограничения arch/README.md
ADR Почему именно так, что отвергнуто adr/ADR-*.md
Guidelines Как писать код, что запрещено docs/spark_guidelines.md
Schemas Структура таблиц, типы, PII schemas/**/*.yaml
Contracts SLA, ownership, breaking changes contracts/**/*.yaml
Examples Рабочие паттерны few-shot learning examples/**/*.py
Runbooks Operational knowledge runbooks/**/*.md

Без контекста агент - образованный стажёр. С контекстом - Senior DE, который знает ваши конвенции, ваш стек и ваши ограничения.