Schema Evolution: стратегии добавления и удаления колонок без поломки

Schema Evolution: стратегии добавления и удаления колонок без поломки

streaming

В идеальном мире схема данных фиксируется на старте проекта и никогда не меняется. В реальном мире бэкенд-команда добавляет новое поле без предупреждения, аналитики переименовывают колонку «для единообразия», а бизнес требует срочно убрать персональные данные из хранилища. Каждое из этих изменений - потенциальная катастрофа для data pipeline, если команда не подготовила архитектуру к эволюции схем.

Schema Evolution - это не просто технический навык. Это искусство управлять неизбежными изменениями так, чтобы downstream-системы продолжали работать, исторические данные оставались читаемыми, а новые инженеры в команде понимали, что происходило со схемой год назад.


1. Schema Drift: почему схемы всегда меняются

1.1 Природа изменений в реальных системах

Любая живая бизнес-система эволюционирует. За год типичная production-платформа переживает:

  • Новые бизнес-требования: маркетинг хочет отслеживать utm_source - поле добавляется в события
  • Смена API-версии: платёжный провайдер перешёл с v2 на v3, структура JSON изменилась
  • CDC-обновления: разработчики добавили колонку discount_code в таблицу orders в PostgreSQL
  • Нормализация данных: команда договорилась использовать customer_id вместо user_id
  • Compliance-требования: регулятор потребовал удалить поле ip_address из хранилища
  • Оптимизация хранения: тип STRING заменяется на INTEGER для экономии места

Всё это - Schema Drift: неконтролируемое или слабо контролируемое изменение структуры данных в источниках. Проблема в том, что data pipeline написан под конкретную схему. Когда схема меняется, pipeline либо падает с ошибкой, либо - что хуже - молча пропускает новые поля, записывая null вместо реальных значений.

1.2 Классификация изменений по степени совместимости

Backward Compatible изменения безопасны для читателей, написанных под старую схему. Если они видят новую nullable колонку - просто получают null и продолжают работать. Старые Parquet-файлы без этой колонки тоже читаются корректно.

Breaking Changes - это изменения, которые ломают существующие pipeline без исправления кода. Переименование user_id в customer_id: Spark job, читающий df["user_id"], упадёт с AnalysisException. BI-дашборд, использующий SELECT user_id, вернёт ошибку.

1.3 Последствия неуправляемого Schema Drift

Каскадный эффект: одно изменение в источнике ломает весь downstream. Особенно опасны «тихие» сбои - когда pipeline не падает с ошибкой, но пишет NULL вместо реальных значений. Бизнес-дашборд показывает нулевую скидку не потому что скидок нет, а потому что поле discount_code не прошло через pipeline.


2. Schema Enforcement: защита от загрязнения данных

2.1 Что такое Schema Enforcement

Schema Enforcement - это механизм строгой проверки соответствия записываемых данных схеме целевой таблицы. Iceberg и Delta Lake по умолчанию отклоняют любую попытку записать DataFrame, схема которого не совпадает с таблицей.

Это поведение намеренно жёсткое: лучше упасть с понятной ошибкой, чем молча принять «кривые» данные и засорить хранилище.

from pyspark.sql import SparkSession
from pyspark.sql.types import (
    StructType, StructField, StringType, DecimalType, TimestampType
)
from pyspark.sql import functions as F

spark = SparkSession.builder \
    .appName("Schema-Evolution") \
    .config("spark.sql.extensions",
            "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
    .config("spark.sql.catalog.lakehouse", "org.apache.iceberg.spark.SparkCatalog") \
    .config("spark.sql.catalog.lakehouse.type", "hadoop") \
    .config("spark.sql.catalog.lakehouse.warehouse", "/tmp/warehouse") \
    .getOrCreate()

# Создаём Silver-таблицу с фиксированной схемой
spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.silver.orders (
        order_id    INTEGER,
        customer_id STRING,
        status      STRING,
        amount      DECIMAL(10,2),
        created_at  TIMESTAMP
    )
    USING iceberg
""")

# Пишем корректные данные - всё работает
correct_df = spark.createDataFrame(
    [(1001, "cust_A", "confirmed", 1500.0, "2026-05-23 10:00:00")],
    ["order_id", "customer_id", "status", "amount", "created_at"]
).withColumn("created_at", F.to_timestamp("created_at")) \
 .withColumn("amount", F.col("amount").cast("decimal(10,2)"))

correct_df.writeTo("lakehouse.silver.orders").append()

# Попытка записать DataFrame с ДОПОЛНИТЕЛЬНОЙ колонкой
df_with_new_col = correct_df.withColumn("discount_code", F.lit("SALE20"))

try:
    df_with_new_col.writeTo("lakehouse.silver.orders").append()
except Exception as e:
    print(f"Schema Enforcement сработал: {type(e).__name__}")
    print(f"Сообщение: {str(e)[:200]}")
    # AnalysisException: Cannot write incompatible data to table 'silver.orders'
    # Found column discount_code which doesn't match target table schema

# Попытка записать с ДРУГИМ ТИПОМ
df_wrong_type = spark.createDataFrame(
    [(1002, "cust_B", "pending", "NOT_A_NUMBER", "2026-05-23 11:00:00")],
    ["order_id", "customer_id", "status", "amount", "created_at"]
)

try:
    df_wrong_type.writeTo("lakehouse.silver.orders").append()
except Exception as e:
    print(f"Type mismatch: {type(e).__name__}")

Schema Enforcement - это первая линия обороны Silver-слоя. Он гарантирует, что в таблицу попадают только данные, соответствующие ожидаемому контракту. Без этой защиты любой upstream-баг незаметно протечёт в аналитику.

2.2 Schema on Write vs Schema on Read

В Medallion Architecture Bronze использует Schema on Read: принимаем всё, что пришло из источника. Silver переходит на Schema on Write: схема контролируется при записи. Gold - максимально строгий Schema on Write с явными контрактами для BI-потребителей.


3. Безопасное добавление новых колонок

3.1 mergeSchema: автоматическое расширение схемы

Самый частый сценарий schema evolution - источник добавил новое поле. Нам нужно принять эти данные без перезаписи всей таблицы и без падения пайплайна.

# Сценарий: PostgreSQL добавил колонку discount_code в таблицу orders
# Debezium начинает присылать её в CDC-потоке

# Первый батч (до изменения схемы): 3 поля
batch_v1 = spark.createDataFrame([
    (1001, "cust_A", "confirmed", 1500.0, "2026-05-23 10:00:00", None),
    (1002, "cust_B", "pending",    750.0, "2026-05-23 11:00:00", None),
], ["order_id", "customer_id", "status", "amount", "created_at", "discount_code"])

batch_v1.writeTo("lakehouse.bronze.orders").append()

# Второй батч (ПОСЛЕ изменения схемы): новое поле discount_code
batch_v2 = spark.createDataFrame([
    (1003, "cust_C", "confirmed", 2000.0, "2026-05-23 12:00:00", "SALE20"),
    (1004, "cust_D", "cancelled",  500.0, "2026-05-23 13:00:00", "VIP10"),
], ["order_id", "customer_id", "status", "amount", "created_at", "discount_code"])

# Без mergeSchema: AnalysisException - discount_code не в схеме таблицы
try:
    batch_v2.writeTo("lakehouse.bronze.orders").append()
except Exception as e:
    print("Ожидаемая ошибка без mergeSchema:", type(e).__name__)

# С mergeSchema: Iceberg добавляет колонку в schema таблицы атомарно
batch_v2.writeTo("lakehouse.bronze.orders") \
    .option("mergeSchema", "true") \
    .append()

print("Schema после добавления колонки:")
spark.table("lakehouse.bronze.orders").printSchema()
# root
#  |-- order_id:      integer
#  |-- customer_id:   string
#  |-- status:        string
#  |-- amount:        decimal(10,2)
#  |-- created_at:    string
#  |-- discount_code: string    ← добавлена автоматически

3.2 Что происходит со старыми данными при добавлении колонки

Критически важный вопрос: что произойдёт со строками, записанными до добавления discount_code? Ответ зависит от table format:

Старые Parquet-файлы не перезаписываются. Это экономия ресурсов - не нужно перечитывать терабайты истории. Вместо этого Iceberg использует метаданные: при чтении старых файлов, в которых нет discount_code, значение этой колонки виртуально подставляется как NULL. Читатель видит единую согласованную схему для всех файлов.

# Демонстрация backward-compatible чтения
all_orders = spark.table("lakehouse.bronze.orders")

print("Строки ДО добавления поля:")
all_orders.filter(F.col("order_id").isin(1001, 1002)) \
    .select("order_id", "discount_code") \
    .show()
# order_id | discount_code
# 1001     | null          ← NULL для старых строк
# 1002     | null

print("Строки ПОСЛЕ добавления поля:")
all_orders.filter(F.col("order_id").isin(1003, 1004)) \
    .select("order_id", "discount_code") \
    .show()
# order_id | discount_code
# 1003     | SALE20
# 1004     | VIP10

3.3 Управляемое добавление через DDL

mergeSchema автоматически добавляет колонки при записи. Но в production-среде часто предпочтительнее явное управление через DDL: инженер осознанно добавляет колонку с описанием и типом.

# Явное добавление колонки через DDL
# Плюс: полный контроль над типом, позицией, комментарием
# Плюс: изменение задокументировано в DDL-скрипте (git history)
spark.sql("""
    ALTER TABLE lakehouse.silver.orders
    ADD COLUMN discount_code STRING COMMENT 'Промо-код, применённый к заказу (nullable)'
""")

# Добавление нескольких колонок одновременно
spark.sql("""
    ALTER TABLE lakehouse.silver.orders
    ADD COLUMNS (
        discount_amount DECIMAL(8,2) COMMENT 'Сумма скидки в рублях',
        is_first_order  BOOLEAN      COMMENT 'Флаг первого заказа клиента'
    )
""")

# Проверяем обновлённую схему
spark.sql("DESCRIBE TABLE lakehouse.silver.orders").show(20)

DDL-подход предпочтительнее mergeSchema для Silver и Gold: изменение атомарно, явно задокументировано, не зависит от наличия данных с новой колонкой.

3.4 Эволюция вложенных структур (Nested Schema)

JSON-поля и Struct-типы эволюционируют сложнее, чем плоские колонки:

from pyspark.sql.types import StructType, StructField, StringType, DecimalType

# Исходная схема: metadata как struct
schema_v1 = StructType([
    StructField("order_id",   StringType(), True),
    StructField("metadata",   StructType([
        StructField("source",   StringType(), True),
        StructField("version",  StringType(), True),
    ]), True),
])

# Новая версия: metadata получил дополнительное поле experiment_id
schema_v2 = StructType([
    StructField("order_id",   StringType(), True),
    StructField("metadata",   StructType([
        StructField("source",        StringType(), True),
        StructField("version",       StringType(), True),
        StructField("experiment_id", StringType(), True),  # новое поле
    ]), True),
])

data_v2 = spark.createDataFrame([
    ("order_1005", {"source": "mobile", "version": "3.1", "experiment_id": "exp-42"}),
], schema_v2)

# mergeSchema работает и для nested struct
data_v2.writeTo("lakehouse.bronze.orders_nested") \
    .option("mergeSchema", "true") \
    .append()

# Чтение: старые строки получат metadata.experiment_id = null
# Новые строки - реальное значение
spark.table("lakehouse.bronze.orders_nested") \
    .select("order_id", "metadata.source", "metadata.experiment_id") \
    .show()

При работе с nested struct нужно быть осторожным: Iceberg умеет добавлять поля в struct, но не умеет изменять тип существующего поля внутри struct. Если нужно изменить тип - придётся создавать новую колонку и мигрировать данные явно.


4. Ломающие изменения: удаление и переименование

4.1 Почему удаление колонки опаснее добавления

Добавление nullable колонки - backward compatible. Удаление или переименование - почти всегда breaking change. Любой код, читающий удалённую колонку, упадёт:

# Ситуация: команда решила убрать deprecated поле user_id (заменили на customer_id)
# Сначала проверяем, кто зависит от этой колонки

# 1. Проверяем Silver jobs - используют ли они user_id?
# 2. Проверяем Gold marts - есть ли JOIN по user_id?
# 3. Проверяем BI-дашборды - есть ли SELECT user_id?

# Только после аудита зависимостей - физическое удаление
spark.sql("""
    ALTER TABLE lakehouse.silver.orders DROP COLUMN user_id
""")

# ПРОБЛЕМА: физическое удаление в Iceberg (Copy-on-Write)
# перезаписывает ВСЕ data files, из которых нужно убрать колонку
# Для таблицы в 10 ТБ - это долго и дорого

Физическое удаление колонки в Iceberg (Copy-on-Write) требует перезаписи всех Parquet-файлов, содержащих эту колонку. Для большой таблицы это часы работы кластера и гигантская стоимость вычислений.

4.2 Column Mapping: логическое удаление без физической перезаписи

Iceberg поддерживает Column Mapping - механизм, позволяющий переименовывать и логически скрывать колонки без перезаписи физических файлов:

# Включаем Column Mapping для таблицы
spark.sql("""
    ALTER TABLE lakehouse.silver.orders
    SET TBLPROPERTIES (
        'write.metadata.metrics.default' = 'truncate(16)',
        'format-version' = '2'
    )
""")

# Логическое переименование: физический файл хранит "user_id",
# но Iceberg metadata маппит его на "customer_id"
spark.sql("""
    ALTER TABLE lakehouse.silver.orders
    RENAME COLUMN user_id TO customer_id
""")

# Читатели, использующие customer_id, получают корректные данные
# Физические файлы НЕ перезаписываются - операция мгновенная!
# Старые читатели, использующие user_id, получат ошибку (это корректное поведение)

Column Mapping в Iceberg хранит в метаданных маппинг: logical_name → physical_field_id. При чтении Parquet-файла Iceberg знает, что физическое поле с ID 42 теперь называется customer_id. Старые файлы не трогаются, новые записываются с новым логическим именем.

4.3 Soft Deprecation: безопасная миграция через двойную запись

В production-среде мгновенное переименование - слишком рискованно. Безопаснее использовать паттерн Dual Write - временно поддерживать оба имени:

Три фазы безопасной миграции:

  1. Dual Write: добавляем новое поле customer_id как копию user_id. Оба поля заполняются одновременно. Никакой code freeze не нужен.
  2. Migration window: BI-команда и downstream pipeline переключаются с user_id на customer_id. Этот период может длиться неделями.
  3. Cleanup: только после того как все потребители переключились - удаляем user_id. Это момент, когда ломающее изменение перестаёт быть breaking.
# Фаза 1: Dual Write в Silver pipeline
def transform_to_silver(bronze_df):
    return (
        bronze_df
        .withColumn("customer_id",
                    F.coalesce(F.col("customer_id"), F.col("user_id")))
        .withColumn("user_id",
                    F.col("user_id"))  # Сохраняем старое поле пока
        # ... остальные трансформации
    )

# Фаза 3 (через 2 недели): убираем user_id из pipeline
def transform_to_silver_v2(bronze_df):
    return (
        bronze_df
        .withColumn("customer_id",
                    F.coalesce(F.col("customer_id"), F.col("user_id")))
        .drop("user_id")  # Больше не нужно
    )

4.4 Type Widening: расширение типов данных

Iceberg 2 поддерживает Type Widening - автоматическое расширение типов без перезаписи файлов:

# Безопасные расширения типов (Type Widening)
# Iceberg умеет делать их логически, без перезаписи файлов

# INT → BIGINT (расширение диапазона)
spark.sql("""
    ALTER TABLE lakehouse.silver.orders
    ALTER COLUMN order_id TYPE BIGINT
""")
# Старые файлы хранят INT - при чтении Iceberg автоматически преобразует в BIGINT
# Новые файлы записываются как BIGINT

# FLOAT → DOUBLE (расширение точности)
spark.sql("""
    ALTER TABLE lakehouse.silver.orders
    ALTER COLUMN amount TYPE DOUBLE
""")

# Проверяем изменение
spark.table("lakehouse.silver.orders").printSchema()
# order_id: bigint  ← теперь BIGINT
# amount: double    ← теперь DOUBLE

Сужение типов (BIGINT → INT, DOUBLE → FLOAT) - это Breaking Change: исторические данные могут содержать значения, не помещающиеся в суженный тип. Iceberg запрещает Type Narrowing - это защита от потери данных.


5. Schema Evolution в Bronze слое: максимальная гибкость

5.1 Bronze как сырой буфер

Bronze принимает данные из источников в их «родном» виде. Здесь Schema Enforcement минимален - мы не хотим терять события только потому, что поставщик добавил новое поле:

# Bronze pipeline с автоматическим mergeSchema
# Этот код корректно обрабатывает любые schema changes в источнике

def ingest_to_bronze(source_df, target_table: str, pipeline_run_id: str):
    """
    Universal Bronze ingestion с автоматической эволюцией схемы.
    Принимает любые новые поля от источника без ошибок.
    """
    enriched_df = (
        source_df
        .withColumn("_ingested_at",      F.current_timestamp())
        .withColumn("_pipeline_run_id",  F.lit(pipeline_run_id))
        .withColumn("_source_schema_hash",
            # Хэш схемы для отслеживания изменений
            F.sha2(F.lit(str(source_df.schema.simpleString())), 256))
        .withColumn("ingestion_date", F.current_date())
    )

    # mergeSchema=true: любые новые поля из источника автоматически добавятся
    enriched_df.writeTo(target_table) \
        .option("mergeSchema", "true") \
        .append()

    # Логируем schema change если он произошёл
    current_schema = spark.table(target_table).schema
    print(f"Текущая схема {target_table}: {len(current_schema.fields)} колонок")

# Первый запуск: 5 полей
batch1 = spark.createDataFrame([
    ("1001", "cust_A", "confirmed", "1500.0", "2026-05-23"),
], ["order_id", "customer_id", "status", "amount", "event_date"])

ingest_to_bronze(batch1, "lakehouse.bronze.orders", "run-001")
# Схема: 5 + 4 метаданных = 9 колонок

# Второй запуск: источник добавил discount_code и experiment_id
batch2 = spark.createDataFrame([
    ("1002", "cust_B", "pending", "750.0", "2026-05-23", "SALE20", "exp-42"),
], ["order_id", "customer_id", "status", "amount", "event_date",
    "discount_code", "experiment_id"])

ingest_to_bronze(batch2, "lakehouse.bronze.orders", "run-002")
# Схема автоматически расширилась: 9 → 11 колонок
# batch1 строки получат discount_code=NULL, experiment_id=NULL

5.2 _source_schema_hash: мониторинг изменений схемы

Хранение хэша схемы источника в Bronze позволяет детектировать schema changes и алертить инженеров:

# Детектирование schema changes для мониторинга и алертинга
def detect_schema_change(spark, bronze_table: str):
    """
    Сравнивает текущую схему с предыдущей.
    Если схема изменилась - логирует предупреждение.
    """
    schema_history = (
        spark.table(bronze_table)
        .groupBy("_source_schema_hash", "ingestion_date")
        .agg(F.count("*").alias("row_count"))
        .orderBy("ingestion_date")
    )

    # Если в одном дне несколько разных хэшей - схема менялась внутри дня
    daily_schemas = (
        schema_history
        .groupBy("ingestion_date")
        .agg(F.countDistinct("_source_schema_hash").alias("schema_versions"))
        .filter(F.col("schema_versions") > 1)
    )

    if daily_schemas.count() > 0:
        print("ВНИМАНИЕ: Schema change обнаружен!")
        daily_schemas.show()

detect_schema_change(spark, "lakehouse.bronze.orders")

6. Schema Evolution в Silver слое: управляемая миграция

6.1 Silver как контракт качества

Если Bronze принимает всё, то Silver - это слой гарантий. Аналитик, читающий Silver, должен доверять схеме. Это означает, что schema evolution на Silver должна быть явной и контролируемой, а не автоматической.

6.2 Миграционный скрипт для Silver

В production-среде изменение схемы Silver оформляется как миграционный скрипт, который можно проверить в code review и повторить в любой среде:

# migration_20260523_001_add_discount_fields.py
# Миграция: добавление полей для работы со скидками

import datetime
from pyspark.sql import SparkSession
from pyspark.sql import functions as F

MIGRATION_ID = "20260523_001"
MIGRATION_DESC = "Добавление полей discount_code и discount_amount в silver.orders"

def run_migration(spark):
    """
    Идемпотентная миграция схемы Silver.
    Безопасно запускать несколько раз.
    """
    print(f"Запуск миграции {MIGRATION_ID}: {MIGRATION_DESC}")

    # Проверяем: миграция уже выполнялась?
    existing_columns = [
        f.name for f in spark.table("lakehouse.silver.orders").schema.fields
    ]

    if "discount_code" in existing_columns:
        print(f"Миграция {MIGRATION_ID} уже применена. Пропускаем.")
        return

    # Шаг 1: Добавляем новые колонки
    spark.sql("""
        ALTER TABLE lakehouse.silver.orders
        ADD COLUMNS (
            discount_code   STRING      COMMENT 'Промо-код из источника (nullable)',
            discount_amount DECIMAL(8,2) COMMENT 'Сумма скидки в рублях (nullable)'
        )
    """)
    print("  Шаг 1/3: Колонки добавлены в схему таблицы")

    # Шаг 2: Бэкфилл - заполняем discount_amount для исторических строк
    # Если у нас есть данные в Bronze, можем пересчитать
    # Если нет - оставляем NULL (backward compatible)
    print("  Шаг 2/3: Исторические данные оставляем NULL (backward compatible)")

    # Шаг 3: Валидируем результат
    schema_after = spark.table("lakehouse.silver.orders").schema
    assert "discount_code" in [f.name for f in schema_after.fields], \
        "Миграция не применилась!"
    print("  Шаг 3/3: Валидация прошла")

    # Логируем факт миграции
    spark.sql(f"""
        INSERT INTO lakehouse.metadata.schema_migrations
        (migration_id, description, applied_at, status)
        VALUES ('{MIGRATION_ID}', '{MIGRATION_DESC}', current_timestamp(), 'SUCCESS')
    """)

    print(f"Миграция {MIGRATION_ID} успешно применена.")


spark = SparkSession.builder.appName("Migration").getOrCreate()
run_migration(spark)

Ключевые свойства миграционного скрипта:

  • Идемпотентность: повторный запуск не ломает ничего и не дублирует изменения
  • Проверка состояния: скрипт сначала проверяет, не применялась ли миграция уже
  • Логирование: факт применения миграции фиксируется в schema_migrations
  • Валидация: после изменения проверяем, что схема действительно изменилась

7. Schema Evolution в Structured Streaming

7.1 Почему стриминг сложнее батча

В batch-обработке можно остановить pipeline, применить миграцию и перезапустить. В streaming-системе поток не останавливается: события продолжают поступать в Kafka, пока мы мигрируем схему.

Смешанный батч (micro-batch 2) содержит события двух версий: старые без discount_code и новые с ним. Spark должен корректно обработать оба типа в одной транзакции.

7.2 Обработка mixed-schema событий в Structured Streaming

from pyspark.sql.types import (
    StructType, StructField, StringType, DecimalType, TimestampType
)

# Определяем «максимальную» схему - включает поля всех известных версий
# Все поля, которые могут отсутствовать - nullable
TARGET_SCHEMA = StructType([
    StructField("order_id",      StringType(),         False),  # обязательное
    StructField("status",        StringType(),         True),
    StructField("amount",        DecimalType(10, 2),   True),
    StructField("event_ts",      TimestampType(),      True),
    StructField("discount_code", StringType(),         True),   # v2 поле, nullable
    StructField("experiment_id", StringType(),         True),   # v3 поле, nullable
])

# Читаем JSON из Kafka, применяем целевую схему
# from_json заполняет отсутствующие поля NULL автоматически
streaming_df = (
    spark.readStream
         .format("kafka")
         .option("kafka.bootstrap.servers", "kafka:9092")
         .option("subscribe", "orders-events")
         .load()
         .select(
             F.from_json(
                 F.col("value").cast("string"),
                 TARGET_SCHEMA
             ).alias("data")
         )
         .select("data.*")
         # Событие v1: discount_code=NULL, experiment_id=NULL
         # Событие v2: discount_code="SALE20", experiment_id=NULL
         # Событие v3: discount_code="SALE20", experiment_id="exp-42"
)

# Запись в Bronze Iceberg с mergeSchema для будущих версий
query = (
    streaming_df
    .writeStream
    .format("iceberg")
    .outputMode("append")
    .option("checkpointLocation", "s3://checkpoints/orders-events/")
    .option("fanout-enabled", "true")
    # mergeSchema позволяет добавить v4 поля без остановки стрима
    .option("mergeSchema", "true")
    .toTable("lakehouse.bronze.order_events")
)

7.3 Schema Registry: контракт между producer и consumer

В enterprise-среде schema evolution управляется через Schema Registry (Confluent Schema Registry, AWS Glue Schema Registry). Registry хранит версии схем и проверяет совместимость при публикации новой версии:

Schema Registry с режимом BACKWARD автоматически отклоняет публикацию несовместимых схем. Это ключевой инструмент schema governance: ломающие изменения обнаруживаются на этапе публикации, а не при падении Spark job в 03:47 ночи.


8. Антипаттерны schema evolution

8.1 Blind autoMerge: автоматическое загрязнение схемы

# АНТИПАТТЕРН: глобальное включение autoMerge без контроля

# В SparkSession config или SQL:
spark.conf.set("spark.databricks.delta.schema.autoMerge.enabled", "true")
# Теперь ЛЮБАЯ запись в ЛЮБУЮ таблицу автоматически расширяет схему

# Проблема: разработчик случайно пишет DataFrame с колонкой-опечаткой
bad_df = spark.createDataFrame([
    (1001, "cust_A", "confirmed", 1500.0, "SALE20")
], ["order_id", "customer_id", "status", "amount", "discont_code"])
# ^ опечатка: "discont_code" вместо "discount_code"

bad_df.writeTo("lakehouse.silver.orders").append()
# Без контроля: в таблице появится и "discount_code" и "discont_code"
# Silver-таблица молча загрязнена. Аналитики увидят дубликат колонки.

Правило: никогда не включать autoMerge глобально. Используйте mergeSchema=true только для конкретных таблиц и конкретных операций записи (Bronze ingestion).

8.2 Schema Inference в production pipeline

# АНТИПАТТЕРН: schema inference из данных в production

# Плохо: Spark пытается вывести схему из первых строк файла
df_inferred = spark.read.json("s3://raw/orders/2026-05-23.json")

# Проблемы:
# 1. Схема выводится из ПОДВЫБОРКИ данных - в конце файла могут быть поля,
#    которых нет в начале. Spark не увидит их.
# 2. Типы выводятся эвристически: "1500" → LONG, не DECIMAL
# 3. Опциональные поля могут вывестись как NOT NULL если в sample нет null
# 4. При следующем запуске inference может дать другую схему!

# Хорошо: явная схема
EXPLICIT_SCHEMA = StructType([
    StructField("order_id",     StringType(),        False),
    StructField("customer_id",  StringType(),        True),
    StructField("status",       StringType(),        True),
    StructField("amount",       DecimalType(10, 2),  True),
    StructField("created_at",   TimestampType(),     True),
    StructField("discount_code", StringType(),       True),
])

df_explicit = spark.read.schema(EXPLICIT_SCHEMA).json("s3://raw/orders/2026-05-23.json")
# Схема детерминирована: тот же результат при каждом запуске

8.3 Нет мониторинга schema changes

Без мониторинга schema changes инженеры узнают о них постфактум - когда аналитики жалуются на пустые дашборды:

# Правильный подход: автоматическое обнаружение schema changes

def validate_schema_unchanged(spark, table: str, expected_schema: StructType):
    """
    Проверяет, что схема таблицы не изменилась.
    Запускается как шаг pipeline перед обработкой.
    """
    current_schema = spark.table(table).schema
    current_fields = {f.name: f.dataType for f in current_schema.fields}
    expected_fields = {f.name: f.dataType for f in expected_schema.fields}

    # Новые поля (добавленные без нашего ведома)
    new_fields = set(current_fields.keys()) - set(expected_fields.keys())
    # Удалённые поля
    removed_fields = set(expected_fields.keys()) - set(current_fields.keys())
    # Изменённые типы
    type_changes = {
        name: (expected_fields[name], current_fields[name])
        for name in expected_fields
        if name in current_fields and expected_fields[name] != current_fields[name]
    }

    if new_fields:
        print(f"ПРЕДУПРЕЖДЕНИЕ: Обнаружены новые поля в {table}: {new_fields}")
    if removed_fields:
        raise ValueError(f"КРИТИЧНО: Удалены поля из {table}: {removed_fields}")
    if type_changes:
        raise ValueError(f"КРИТИЧНО: Изменены типы в {table}: {type_changes}")

    print(f"✓ Схема {table} соответствует ожидаемой")

9. End-to-End кейс: эволюция схемы без downtime

9.1 Сценарий

Компания добавляет в платформу систему промо-кодов. PostgreSQL получает новую колонку discount_code. Нужно обновить весь стек Bronze → Silver → Gold без остановки BI-системы.

import datetime
from pyspark.sql import SparkSession, functions as F
from pyspark.sql.window import Window

spark = SparkSession.builder \
    .appName("Schema-Evolution-E2E") \
    .config("spark.sql.extensions",
            "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
    .config("spark.sql.catalog.lakehouse", "org.apache.iceberg.spark.SparkCatalog") \
    .config("spark.sql.catalog.lakehouse.type", "hadoop") \
    .config("spark.sql.catalog.lakehouse.warehouse", "/tmp/warehouse") \
    .getOrCreate()

# ─── СОСТОЯНИЕ ДО ИЗМЕНЕНИЯ ──────────────────────────────────────────────────
print("=== ФАЗА 0: Исходное состояние ===")

# Bronze: 5 колонок из старого CDC
spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.bronze.orders_evo (
        order_id     STRING,
        customer_id  STRING,
        status       STRING,
        amount       STRING,
        created_at   STRING,
        _ingested_at  TIMESTAMP,
        ingestion_date DATE
    )
    USING iceberg
    PARTITIONED BY (ingestion_date)
""")

old_data = spark.createDataFrame([
    ("1001", "cust_A", "confirmed", "1500.0", "2026-05-23 10:00:00"),
    ("1002", "cust_B", "pending",    "750.0", "2026-05-23 11:00:00"),
], ["order_id", "customer_id", "status", "amount", "created_at"])

(old_data
 .withColumn("_ingested_at", F.current_timestamp())
 .withColumn("ingestion_date", F.current_date())
 .writeTo("lakehouse.bronze.orders_evo").append())

print("Bronze содержит:", spark.table("lakehouse.bronze.orders_evo").count(), "строк")

# Silver: типизированные данные без discount_code
spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.silver.orders_evo (
        order_id     INTEGER,
        customer_id  STRING,
        status       STRING,
        amount       DECIMAL(10,2),
        created_at   TIMESTAMP
    )
    USING iceberg
    PARTITIONED BY (month(created_at))
""")

old_silver = (
    spark.table("lakehouse.bronze.orders_evo")
    .withColumn("order_id",   F.col("order_id").cast("integer"))
    .withColumn("amount",     F.col("amount").cast("decimal(10,2)"))
    .withColumn("created_at", F.to_timestamp("created_at"))
    .select("order_id", "customer_id", "status", "amount", "created_at")
)
old_silver.writeTo("lakehouse.silver.orders_evo").append()

# Gold: агрегат без скидок
spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.gold.revenue_evo (
        order_date    DATE,
        total_revenue DECIMAL(15,2),
        _updated_at   TIMESTAMP
    )
    USING iceberg
    PARTITIONED BY (order_date)
""")

(spark.table("lakehouse.silver.orders_evo")
 .filter(F.col("status") == "confirmed")
 .groupBy(F.to_date("created_at").alias("order_date"))
 .agg(F.sum("amount").alias("total_revenue"))
 .withColumn("_updated_at", F.current_timestamp())
 .writeTo("lakehouse.gold.revenue_evo").overwritePartitions())

print(f"Исходная схема Silver: "
      f"{[f.name for f in spark.table('lakehouse.silver.orders_evo').schema.fields]}")

# ─── ФАЗА 1: BRONZE ПРИНИМАЕТ НОВОЕ ПОЛЕ ─────────────────────────────────────
print("\n=== ФАЗА 1: Bronze принимает discount_code ===")

new_cdc_data = spark.createDataFrame([
    ("1003", "cust_C", "confirmed", "2000.0", "2026-05-23 12:00:00", "SALE20"),
    ("1004", "cust_D", "cancelled", "500.0",  "2026-05-23 13:00:00", "VIP10"),
], ["order_id", "customer_id", "status", "amount", "created_at", "discount_code"])

# mergeSchema: Bronze автоматически принимает новое поле
(new_cdc_data
 .withColumn("_ingested_at", F.current_timestamp())
 .withColumn("ingestion_date", F.current_date())
 .writeTo("lakehouse.bronze.orders_evo")
 .option("mergeSchema", "true")
 .append())

print("Bronze после Schema Evolution:")
spark.table("lakehouse.bronze.orders_evo") \
     .select("order_id", "discount_code") \
     .show()
# Старые строки: discount_code=NULL
# Новые строки:  discount_code="SALE20" / "VIP10"

# ─── ФАЗА 2: МИГРАЦИЯ SILVER ──────────────────────────────────────────────────
print("\n=== ФАЗА 2: Явная миграция Silver ===")

# Шаг 2а: DDL - добавляем колонку в Silver явно
spark.sql("""
    ALTER TABLE lakehouse.silver.orders_evo
    ADD COLUMN discount_code STRING COMMENT 'Промо-код (nullable, добавлен 2026-05-23)'
""")

# Шаг 2б: Инкрементальное обновление - MERGE в Silver с новым полем
new_silver = (
    spark.table("lakehouse.bronze.orders_evo")
    .filter(F.col("ingestion_date") == F.current_date())
    .filter(F.col("order_id").isin("1003", "1004"))
    .withColumn("order_id",   F.col("order_id").cast("integer"))
    .withColumn("amount",     F.col("amount").cast("decimal(10,2)"))
    .withColumn("created_at", F.to_timestamp("created_at"))
    .select("order_id", "customer_id", "status", "amount",
            "created_at", "discount_code")
)

new_silver.createOrReplaceTempView("new_silver_orders")

spark.sql("""
    MERGE INTO lakehouse.silver.orders_evo AS target
    USING new_silver_orders AS source
    ON target.order_id = source.order_id
    WHEN MATCHED THEN UPDATE SET *
    WHEN NOT MATCHED THEN INSERT *
""")

print("Silver после миграции:")
spark.table("lakehouse.silver.orders_evo").show()
# Старые строки: discount_code=NULL (backward compatible)
# Новые строки: discount_code="SALE20" / "VIP10"

# ─── ФАЗА 3: GOLD РАСШИРЯЕТСЯ ─────────────────────────────────────────────────
print("\n=== ФАЗА 3: Gold с метрикой скидок ===")

# Gold: добавляем метрику суммы скидок
spark.sql("""
    ALTER TABLE lakehouse.gold.revenue_evo
    ADD COLUMN discounted_orders_count BIGINT
        COMMENT 'Число заказов с промо-кодом (добавлено 2026-05-23)'
""")

# Пересчитываем Gold с новой метрикой
updated_gold = (
    spark.table("lakehouse.silver.orders_evo")
    .filter(F.col("status") == "confirmed")
    .groupBy(F.to_date("created_at").alias("order_date"))
    .agg(
        F.sum("amount").alias("total_revenue"),
        # Новая метрика: сколько заказов с промо-кодом
        F.count(
            F.when(F.col("discount_code").isNotNull(), 1)
        ).alias("discounted_orders_count"),
    )
    .withColumn("_updated_at", F.current_timestamp())
)

updated_gold.writeTo("lakehouse.gold.revenue_evo").overwritePartitions()

print("Gold после эволюции схемы:")
spark.table("lakehouse.gold.revenue_evo").show()
# order_date  | total_revenue | discounted_orders_count | _updated_at
# 2026-05-23  | 3500.00       | 1                       | ...
# Старые BI-запросы на total_revenue не сломались!
# Новые дашборды могут использовать discounted_orders_count

# ─── ФИНАЛ: ПРОВЕРКА СОВМЕСТИМОСТИ ───────────────────────────────────────────
print("\n=== ФИНАЛ: Проверка backward compatibility ===")

# «Старый» BI-запрос (SELECT без новых полей) работает без изменений
old_bi_query = spark.sql("""
    SELECT order_date, total_revenue
    FROM lakehouse.gold.revenue_evo
    ORDER BY order_date
""")
old_bi_query.show()
print("Старые BI-запросы работают без изменений ✓")

9.2 Результат и итоги кейса

=== ФАЗА 0: Исходное состояние ===
Bronze содержит: 2 строки
Исходная схема Silver: ['order_id', 'customer_id', 'status', 'amount', 'created_at']

=== ФАЗА 1: Bronze принимает discount_code ===
Bronze после Schema Evolution:
order_id | discount_code
1001     | null           ← backward compatible
1002     | null
1003     | SALE20
1004     | VIP10

=== ФАЗА 2: Явная миграция Silver ===
Silver после миграции:
order_id | status    | amount | discount_code
1001     | confirmed | 1500.0 | null          ← старые строки не тронуты
1002     | pending   | 750.0  | null
1003     | confirmed | 2000.0 | SALE20
1004     | cancelled | 500.0  | VIP10

=== ФАЗА 3: Gold с метрикой скидок ===
Gold после эволюции:
order_date  | total_revenue | discounted_orders_count
2026-05-23  | 3500.00       | 1

=== ФИНАЛ ===
Старые BI-запросы работают без изменений ✓

BI-система ни разу не видела ошибку в процессе всей миграции. Старые запросы на total_revenue продолжали возвращать корректные данные. Новые дашборды начали использовать discounted_orders_count после завершения фазы 3.


10. Чек-лист безопасной Schema Evolution

Изменение Безопасность Рекомендуемый подход
Добавить nullable колонку ✅ Safe ALTER TABLE ADD COLUMN + mergeSchema на Bronze
Расширить тип INT → BIGINT ✅ Safe ALTER TABLE ALTER COLUMN TYPE
Переименовать колонку ⚠️ Breaking Dual Write → Migration Window → Drop old
Удалить колонку ⚠️ Breaking Audit зависимостей → Dual Write → Drop after migration
Сузить тип BIGINT → INT ❌ Dangerous Новая колонка + трансформация + rename
Изменить смысл поля ❌ Dangerous Новая колонка с новым именем
Изменить nested struct ⚠️ Breaking Добавить поле (OK), изменить тип поля (новая колонка)

Главные правила schema evolution:

  • Bronze: mergeSchema=true, принимает всё, логирует schema changes через _source_schema_hash
  • Silver: явные ALTER TABLE DDL, идемпотентные миграции в виде скриптов в git
  • Gold: версионирование (gold.revenue_v2), параллельная поддержка версий в migration window
  • Streaming: явная целевая схема через from_json(schema), никакого schema inference
  • Мониторинг: алерты при обнаружении новых/удалённых полей в Bronze-потоке
  • Никогда: не используйте autoMerge глобально, не полагайтесь на schema inference в production