Bronze Layer: append-only immutable ingestion, обязательные audit-колонки (_ingest_ts, _source_system, _source_file, _event_id), почему дубликаты допустимы на этом слое

Принцип иммутабельности сырых данных, анатомия аудит-полей, легитимизация дубликатов при at-least-once доставке и отказоустойчивый ingestion-пайплайн на PySpark + Iceberg

lakehouse bronze ingestion iceberg audit

Bronze Layer: хранилище без цензуры

Bronze - это фундамент всего Data Lakehouse. Любая архитектурная ошибка на этом уровне - преждевременная фильтрация, мутация данных, потеря дублей - рушит доверие ко всем вышестоящим слоям и лишает инженеров возможности сделать полный Replay истории. Когда что-то ломается в Silver или Gold, именно Bronze позволяет пересчитать всё с нуля, не обращаясь повторно к перегруженной продакшн-базе источника.

Этот урок - про жёсткие правила, которые нельзя нарушать, и про то, почему они существуют.

1. Философия Bronze: хранилище без цензуры

Главная заповедь дата-инженера

Данные на слое Bronze должны на 100% соответствовать тому, что пришло из источника.

Это не рекомендация - это инвариант. Bronze - это цифровой слепок реальности в момент её получения. Если PostgreSQL прислал битую запись с NULL в обязательном поле - Bronze хранит эту запись с NULL. Если Kafka-топик доставил невалидный JSON - Bronze хранит этот строку в поле raw_payload. Если API вернул дубликат предыдущего события - Bronze хранит оба события.

Интуиция подсказывает: «зачем хранить мусор?» Ответ: потому что то, что сегодня выглядит как мусор, завтра может оказаться ценными данными - если изменятся бизнес-правила или выяснится, что «ошибочные» записи содержали сигнал. Решение о том, что является «мусором», принимается в Silver, а не в Bronze.

Принцип иммутабельности: Append-Only

Bronze является иммутабельным (неизменяемым) слоем. Это означает одно операционное правило без исключений:

Только INSERT / APPEND. Никогда UPDATE, DELETE или OVERWRITE.

Когда инженер выполняет UPDATE или DELETE в Bronze-таблице, он уничтожает исходное состояние данных. После этого невозможно ответить на вопрос: «Что именно прислал источник в тот момент?» Весь смысл Bronze - сохранить первоначальный снимок реальности - разрушается.

OVERWRITE особенно коварен. При батчевой загрузке инженер может написать mode("overwrite") из удобства - перезаписать партицию, чтобы убрать старые файлы. Это разрушает возможность Time Travel и Replay для перезаписанного периода.

Replayability: почему Bronze - это страховка всей платформы

Replayability - способность заново вычислить любой слой данных с нуля, используя только Bronze как первоисточник.

Представьте сценарий: спустя год работы выяснилось, что Silver-пайплайн содержал ошибку в логике дедупликации. Все KPI за год посчитаны неверно. Что делать?

  • Без Bronze: нужно повторно выгрузить данные из продакшн-базы источника за год. Это нагрузка на продакшн, потенциально невозможная задача (данные могли удалиться из источника по retention-политике)
  • С Bronze: запускаем исправленный ETL-пайплайн, который читает из Bronze и перезаписывает Silver. Источник не трогаем. Операция занимает часы вместо недель

Диаграмма: Replayability - главная ценность иммутабельного Bronze. Источники данных один раз пишут события в Bronze. Bronze хранит их вечно, неизменно. Когда Silver-логика ломается (v1 с багом), не нужно повторно обращаться к PostgreSQL или Kafka - достаточно перезапустить исправленный ETL, который читает из Bronze (v2 Replay). Это возможно только при условии, что Bronze никогда не был изменён или перезаписан.

Replayability - это не только про исправление багов. Это про:

  • Смену бизнес-правил: изменилось определение «активного пользователя» - пересчитываем Gold без обращения к источникам
  • Аудит и регуляторные требования: проверяющий орган требует объяснить, откуда взялись данные в отчёте за прошлый год - показываем Bronze
  • Эксперименты: Data Scientist хочет проверить новую модель на исторических данных - берёт Bronze и строит альтернативный Silver

2. Паспорт данных: обязательные аудит-колонки

Зачем нужны метаданные

Без метаданных Bronze быстро превращается в анонимное болото. Через год невозможно ответить на базовые вопросы:

  • Когда именно эта строка попала в Lakehouse?
  • Из какого источника она пришла, если таблица агрегирует несколько систем?
  • Из какого конкретно файла или Kafka-партиции?
  • Как найти эту строку, если нужно воспроизвести инцидент?

Аудит-колонки - это «паспорт» каждой строки в Bronze. Они добавляются не источником, а самим ingestion-пайплайном в момент записи. Бизнес-данные приходят «как есть» в поле raw_payload, а технические метаданные навешиваются Spark-джобой.

Диаграмма: анатомия Bronze-строки. Исходная запись из PostgreSQL (5 бизнес-полей) проходит через Spark Ingestion Job и превращается в Bronze-строку с 6 полями: raw_payload хранит оригинальный JSON как строку, остальные пять - технические аудит-поля, добавленные Spark. Это разделение критично: бизнес-данные неизменны, технические метаданные помогают отследить происхождение любой строки в любой момент времени.

Анатомия аудит-полей

_ingest_ts (TIMESTAMP) - время загрузки

Это момент, когда Spark записал строку в Lakehouse, а не когда событие произошло в источнике. Разница принципиальна: событие могло произойти в 09:00, но из-за задержки в Kafka попасть в Bronze в 10:30. Поле _ingest_ts используется для:

  • Инкрементального захвата: WHERE _ingest_ts > last_processed_ts - читаем только новые записи с момента предыдущего запуска
  • SLA-мониторинга: насколько свежи данные в Bronze прямо сейчас
  • Партиционирования: PARTITIONED BY (days(_ingest_ts)) даёт партиции по дням загрузки, а не по дням события

_source_system (STRING) - идентификатор источника

Критичен при мультисорсинге - когда одна Bronze-таблица агрегирует данные из нескольких систем. Например, bronze.orders может принимать заказы из ecom_website, mobile_app и crm_oracle одновременно. Без _source_system невозможно понять, откуда пришла конкретная строка.

Рекомендованный формат: {system_name}_{transport}, например ecom_postgres_cdc, clickstream_kafka, payments_api_rest.

_source_file (STRING) - lineage на уровне файлов

Хранит путь к физическому файлу (для batch-загрузок) или topic:partition:offset (для Kafka). Это незаменимо при разборе аварий: «В какой день и из какого файла к нам прилетел этот битый JSON?» Spark предоставляет встроенную функцию input_file_name(), которая возвращает путь к текущему читаемому файлу.

_event_id (STRING) - технический UUID события

UUID, генерируемый ingestion-пайплайном для каждой строки. Важно понимать: Bronze не гарантирует уникальность _event_id - если событие доставлено дважды, оба раза генерируется новый UUID. Это поле используется Silver для отслеживания конкретной строки в Bronze (для аудита и отладки), а не для дедупликации.

Для дедупликации Silver использует бизнес-ключ из самого payload (например, order_id из raw_payload).

_batch_id (STRING) - идентификатор батча

Один UUID на весь запуск ETL-джобы. Все строки, загруженные в рамках одного запуска, получают одинаковый _batch_id. Это позволяет:

  • Найти все записи конкретного запуска и проверить их
  • В случае неудачного батча - найти и отфильтровать все его строки
  • Отследить, в какой конкретно запуск была допущена ошибка

Практика: создание Bronze-таблицы с аудит-колонками

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

spark = SparkSession.builder \
    .appName("Bronze-Layer-Best-Practices") \
    .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()

# DDL Bronze-таблицы: raw_payload + 5 аудит-колонок
# Партиционируем по дням _ingest_ts (не по бизнес-дате!)
# Дни инжекции стабильны и предсказуемы, бизнес-даты могут прийти с опозданием
spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.bronze.orders (
        raw_payload    STRING    COMMENT 'Original payload from source, stored as-is',
        _ingest_ts     TIMESTAMP COMMENT 'When Spark wrote this record to Lakehouse',
        _source_system STRING    COMMENT 'Source system identifier, e.g. ecom_postgres_cdc',
        _source_file   STRING    COMMENT 'Physical file path or kafka topic:partition:offset',
        _event_id      STRING    COMMENT 'Technical UUID per row for traceability',
        _batch_id      STRING    COMMENT 'ETL job run UUID - same for all rows in one batch'
    ) USING iceberg
    PARTITIONED BY (days(_ingest_ts))
    TBLPROPERTIES (
        'write.format.default'                 = 'parquet',
        'write.parquet.compression-codec'      = 'zstd',
        'history.expire.max-snapshot-age-ms'   = '2592000000'
    )
""")

Обратите внимание на PARTITIONED BY (days(_ingest_ts)) - это Iceberg Hidden Partitioning. Вместо партиционирования по бизнес-дате события (которая может прийти с опозданием на несколько дней), мы партиционируем по дате фактической загрузки. Это делает партиции предсказуемыми и равномерными.

3. Легитимизация дубликатов на Bronze слое

Природа распределённых систем: at-least-once доставка

В распределённых системах доставка событий работает по одному из трёх принципов:

  • At-most-once: событие может потеряться, но никогда не дублируется
  • Exactly-once: событие доставляется ровно один раз (очень дорого)
  • At-least-once: событие доставляется минимум один раз, возможны дубли

Kafka, RabbitMQ, Kinesis и большинство других брокеров сообщений по умолчанию работают в режиме at-least-once. Это означает: если Spark-консьюмер прочитал событие и обработал его, но упал до того, как зафиксировал (коммит) смещение (offset) в Kafka, при перезапуске он снова прочитает то же событие. Bronze получит дубль.

Диаграмма: at-least-once доставка создаёт дубли в Bronze. Spark читает событие abc123 из Kafka, записывает его в Bronze, но падает до коммита смещения. При перезапуске Kafka снова отдаёт то же смещение (offset=100), и Spark снова записывает abc123 в Bronze. Результат: два идентичных ряда. Это не ошибка - это корректное поведение при at-least-once семантике. Bronze осознанно принимает дубли, чтобы не терять события.

Попытка обеспечить exactly-once на уровне Bronze-ingestion требует транзакционной координации между Kafka и Iceberg - это возможно (через Iceberg ACID + Kafka transactional producer), но дорого по инфраструктуре и задержке. Для большинства продакшн-систем правильное решение: принять дубли в Bronze и удалить их в Silver.

CDC и хронология событий: дубли - это не всегда дубли

Второй источник «дубликатов» в Bronze - это CDC (Change Data Capture) системы вроде Debezium. CDC захватывает каждое изменение строки в источнике и отправляет его как отдельное событие.

Если пользователь создал заказ и затем 4 раза его изменил (изменил количество, адрес, применил купон, подтвердил оплату), CDC отправит 5 событий - все с одним order_id. С точки зрения бизнес-ключа это «дубликаты». С точки зрения хронологии изменений - это 5 разных состояний одного объекта.

Диаграмма: CDC-хронология одного заказа. Пять CDC-событий описывают жизненный цикл одного заказа: от создания до доставки. Bronze сохраняет все пять - это хронологический лог изменений. Silver затем решает, что с ними делать: взять только последнее состояние (SCD Type 1) или хранить всю историю (SCD Type 2). Если Bronze удалил бы «дубли» при записи, вся хронология была бы утеряна - и SCD Type 2 стал бы невозможным.

Если бы Bronze применил dropDuplicates(["order_id"]) при записи, остался бы только первый или последний event. Какой именно зависит от порядка обработки - то есть результат был бы недетерминированным. Кроме того, при сбое и перезапуске результат мог бы оказаться другим.

Цена преждевременной дедупликации

Дедупликация при записи в Bronze требует либо stateful shuffle (дорого по CPU и памяти), либо lookup в существующую таблицу (дорого по I/O). Оба варианта:

  • Замедляют запись: ingestion latency вырастает в 5–20 раз
  • Повышают риск потери данных: если дедупликация произошла по неправильному ключу, потерянные данные невосстановимы
  • Делают пайплайн хрупким: любое изменение бизнес-ключа требует переписывания Bronze-джобы

Деградация производительности особенно болезненна в стриминговых пайплайнах с высокой пропускной способностью. Bronze должен работать как append-only буфер с минимальной обработкой - тогда он масштабируется линейно.

4. Практика: Ingestion Pipeline на PySpark

Диаграмма: полный Bronze ingestion pipeline. Три типа источников (batch-файлы, Kafka, CDC) объединяются через один Spark-пайплайн. Ключевые шаги: чтение сырых данных → добавление аудит-колонок → минимальная валидация формата (только «это текст?», не «это валидный заказ?») → запись в Iceberg в режиме APPEND. На всём пути нет ни фильтрации, ни трансформации бизнес-данных.

Batch ingestion: чтение файлов из Landing Zone

Типичный сценарий: внешняя система сбрасывает JSON-файлы на S3 по расписанию (раз в час или раз в день). Spark читает всё новое с момента последнего запуска.

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

spark = SparkSession.builder \
    .appName("Bronze-Batch-Ingestion") \
    .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()

# --- Шаг 1: Создаём тестовые файлы в Landing Zone ---
# Имитируем три JSON-файла от внешней системы
# Обратите внимание: среди них есть дубль (order_1 в двух файлах)
# и «битая» запись (amount: null)

spark.createDataFrame([
    ('{"id": "order_1", "amount": 100.00, "status": "NEW"}',),
    ('{"id": "order_2", "amount": 250.00, "status": "PAID"}',),
], ["value"]).coalesce(1).write.mode("overwrite").text("/tmp/landing/orders/batch_2024_01_15.json")

spark.createDataFrame([
    ('{"id": "order_1", "amount": 100.00, "status": "NEW"}',),  # дубль!
    ('{"id": "order_3", "amount": null, "status": "ERROR"}',),  # битая запись
], ["value"]).coalesce(1).write.mode("overwrite").text("/tmp/landing/orders/batch_2024_01_15_retry.json")

# --- Шаг 2: Bronze Ingestion Job ---
batch_id = str(uuid.uuid4())  # один UUID на весь батч
landing_path = "/tmp/landing/orders/*.json"

# Читаем файлы как «голый» текст - никаких парсеров, никакой схемы
# Каждая строка файла → одна строка DataFrame с полем "value"
raw_df = spark.read.text(landing_path)

# Навешиваем аудит-колонки
# input_file_name() - нативная Spark-функция, возвращает путь к текущему файлу
# expr("uuid()") - генерирует уникальный UUID для каждой строки
bronze_df = raw_df \
    .withColumnRenamed("value", "raw_payload") \
    .withColumn("_ingest_ts",     F.current_timestamp()) \
    .withColumn("_source_system", F.lit("ecom_files_s3")) \
    .withColumn("_source_file",   F.input_file_name()) \
    .withColumn("_event_id",      F.expr("uuid()")) \
    .withColumn("_batch_id",      F.lit(batch_id))

# Записываем в Bronze - только APPEND, дубли остаются
# .writeTo() + .append() - современный Iceberg API
bronze_df.writeTo("lakehouse.bronze.orders").append()

print(f"Batch {batch_id}: загружено {bronze_df.count()} строк")
print("Дубли сохранены - это нормально для Bronze")
spark.table("lakehouse.bronze.orders").show(truncate=False)

Функция input_file_name() - одна из самых важных для Bronze. Она возвращает полный путь к S3-объекту, из которого Spark читает строку в данный момент. Если через три месяца нужно будет найти, из какого конкретно файла пришли проблемные данные, это поле даст мгновенный ответ.

Kafka Streaming ingestion

Для событийных источников (Kafka, Kinesis) используется Spark Structured Streaming. Пайплайн работает непрерывно, обрабатывая микробатчи.

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

spark = SparkSession.builder \
    .appName("Bronze-Kafka-Streaming") \
    .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()

# Читаем из Kafka - каждое сообщение как бинарный payload
kafka_df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "kafka:9092") \
    .option("subscribe", "orders_raw") \
    .option("startingOffsets", "earliest") \
    .load()

# Kafka отдаёт: key, value (bytes), topic, partition, offset, timestamp, headers
# Важно сохранить topic:partition:offset как _source_file для Kafka lineage
# Это позволит точно найти сообщение в Kafka при разборе инцидентов

bronze_stream_df = kafka_df \
    .select(
        F.col("value").cast("string").alias("raw_payload"),
        F.current_timestamp().alias("_ingest_ts"),
        F.lit("kafka_orders_raw").alias("_source_system"),
        F.concat_ws(":",
            F.col("topic"),
            F.col("partition").cast("string"),
            F.col("offset").cast("string")
        ).alias("_source_file"),  # topic:partition:offset вместо пути к файлу
        F.expr("uuid()").alias("_event_id"),
        F.col("timestamp").alias("kafka_event_ts"),  # время события в Kafka (от producer)
    )

# Streaming APPEND в Bronze Iceberg
# checkpointLocation хранит смещение (offset) последнего обработанного сообщения
# При рестарте Spark продолжит с того же места, но может повторить последний микробатч
# - отсюда at-least-once и дубли в Bronze
query = bronze_stream_df.writeStream \
    .format("iceberg") \
    .outputMode("append") \
    .option("path", "lakehouse.bronze.orders") \
    .option("checkpointLocation", "/tmp/checkpoints/bronze_orders") \
    .trigger(processingTime="60 seconds") \
    .start()

query.awaitTermination()

Обратите внимание на поле kafka_event_ts - это время, когда Kafka-producer создал сообщение (business timestamp). Оно отличается от _ingest_ts (время записи в Bronze). Разница между ними - это сквозная задержка (end-to-end latency) пайплайна. Хранение обоих timestamps позволяет её мониторить.

Хранение raw_payload как STRING vs структурированная запись

Существует два подхода к хранению данных в Bronze:

Подход 1: raw_payload как STRING (рекомендуется при нестабильном источнике)

# Весь JSON/XML/CSV-payload как одна строка
# Плюсы: устойчив к schema evolution источника
# Bronze не ломается, если источник добавил/удалил поле
# Минусы: нельзя читать отдельные поля без парсинга
bronze_df = raw_df \
    .withColumnRenamed("value", "raw_payload")  # STRING

Подход 2: структурированная Bronze-запись (при стабильном источнике)

# Парсим JSON прямо при записи в Bronze
# Плюсы: можно читать отдельные поля без from_json в Silver
# Минусы: Bronze ломается при schema drift источника
from pyspark.sql.types import StructType, StructField, StringType, DoubleType

order_schema = StructType([
    StructField("id",     StringType()),
    StructField("amount", DoubleType()),
    StructField("status", StringType()),
])

bronze_df = raw_df \
    .withColumn("data", F.from_json(F.col("value"), order_schema)) \
    .select("data.*")

Рекомендация: для нестабильных источников (внешние API, Kafka без schema registry) используйте STRING. Для стабильных источников с контролируемой схемой (внутренние PostgreSQL с жёстким контрактом) - структурированная запись допустима.

5. Инкрементальный захват данных

После первого полного ingestion нужна стратегия для загрузки только новых данных. Без этого каждый запуск читает всю таблицу целиком - что неприемлемо при больших объёмах.

Watermark по _ingest_ts

Простейший паттерн: запоминаем максимальный _ingest_ts предыдущего запуска и при следующем запуске читаем только записи, добавленные после него.

from datetime import datetime, timedelta

def get_last_processed_ts(spark, checkpoint_table: str) -> datetime:
    """Читаем watermark из специальной checkpoint-таблицы."""
    if spark.catalog.tableExists(checkpoint_table):
        row = spark.table(checkpoint_table) \
            .agg(F.max("last_ingest_ts").alias("ts")) \
            .collect()[0]
        return row["ts"] or datetime(2020, 1, 1)
    return datetime(2020, 1, 1)

def save_watermark(spark, checkpoint_table: str, ts: datetime):
    """Сохраняем новый watermark после успешного батча."""
    spark.createDataFrame([(ts,)], ["last_ingest_ts"]) \
        .writeTo(checkpoint_table) \
        .append()

# --- Инкрементальный читатель Bronze ---
last_ts = get_last_processed_ts(spark, "lakehouse.meta.bronze_checkpoints")

# Читаем только новые строки, добавленные после последнего запуска
# Iceberg partition pruning автоматически откроет только нужные партиции
new_bronze_df = spark.table("lakehouse.bronze.orders") \
    .filter(F.col("_ingest_ts") > last_ts)

current_max_ts = new_bronze_df.agg(F.max("_ingest_ts")).collect()[0][0]

# Обрабатываем new_bronze_df → Silver
# ...

# Обновляем watermark только после успешного завершения
if current_max_ts:
    save_watermark(spark, "lakehouse.meta.bronze_checkpoints", current_max_ts)
    print(f"Обработано {new_bronze_df.count()} новых строк, watermark обновлён до {current_max_ts}")

Партиционирование Bronze по days(_ingest_ts) делает этот паттерн эффективным: Iceberg автоматически пропускает партиции, где _ingest_ts до watermark. Запрос с WHERE _ingest_ts > yesterday не читает данные за прошлые месяцы.

6. Анти-паттерны Bronze слоя

Диаграмма: правильный Bronze vs анти-паттерны. Левая колонка показывает ключевые свойства правильного Bronze: только APPEND, дубли сохраняются, все записи без исключений, атомарная ACID-запись через Iceberg. Правая колонка показывает пять типичных анти-паттернов, каждый из которых разрушает одно из этих свойств. Особенно опасен OVERWRITE: он убивает Replay за перезаписанный период даже при наличии правильных аудит-колонок.

Расшифровка каждого анти-паттерна:

  • UPDATE/DELETE raw data: нарушает иммутабельность, делает Replay невозможным
  • dropDuplicates() при записи: медленно (требует shuffle), уничтожает CDC-хронологию, результат недетерминирован при сбоях
  • Бизнес-фильтры (WHERE status = 'OK'): «ошибочные» записи теряются навсегда; если через год выяснится, что их нужно было проанализировать - данных нет
  • JOIN при записи: смешивает ingestion с трансформацией; если JOIN-таблица изменится, историческое поведение пайплайна изменится ретроспективно
  • OVERWRITE партиций: уничтожает возможность Time Travel и Replay для перезаписанного периода

7. Кейс: PII и GDPR - как удалять данные из иммутабельного слоя

На практике возникает конфликт: Bronze должен быть иммутабельным, но GDPR (и российский ФЗ-152) требуют «Права на забвение» - удаления персональных данных по запросу пользователя.

Типичный сценарий: безопасник обнаруживает, что в raw_payload хранятся незашифрованные паспортные данные клиентов. Он требует еженедельно запускать скрипт DELETE FROM bronze.orders WHERE user_id = ? для маскирования. Как поступить?

Диаграмма: три подхода к PII в иммутабельном Bronze. Правильный ответ на требование безопасника зависит от этапа жизненного цикла. Лучший вариант (1) - не допускать PII в Bronze вообще, токенизируя их в потоке до записи. Компромиссный вариант (2) - Crypto-Shredding: PII зашифрованы, удаление ключа делает их нечитаемыми без физического DELETE. Последний резерв (3) - Iceberg DELETE при обязательном требовании физического удаления по закону; это единственный случай, когда нарушение иммутабельности оправдано.

Правильный ответ на требование безопасника:

Еженедельный DELETE/UPDATE в Bronze - это архитектурная катастрофа. Он превращает Bronze из архива в изменяемую базу данных и разрушает Replayability. Правильные решения:

Вариант 1 (лучший): Маскирование до Lake. PII-поля заменяются токенами ещё в потоке данных - в Kafka Streams, Flink или специализированном PII-vault. В Bronze попадают только токены. Оригинальные PII хранятся в защищённом KMS, связанном с токеном. При запросе на удаление достаточно удалить запись из KMS - данные в Bronze становятся анонимными.

Вариант 2 (хороший): Crypto-Shredding. Bronze хранит PII в зашифрованном виде (AES-256 per-user key). Ключ шифрования хранится в KMS и привязан к user_id. При запросе на забвение удаляем ключ из KMS. PII в Bronze становятся нечитаемой последовательностью байт - без физического удаления записей. Иммутабельность Bronze не нарушается.

Вариант 3 (крайняя мера): Iceberg DELETE + expire_snapshots. Если закон (GDPR, ФЗ-152) требует физического удаления данных - Iceberg поддерживает точечный DELETE WHERE user_id = ?. После этого необходимо запустить expire_snapshots() для физического удаления файлов со старыми версиями. Это нарушает иммутабельность для удалённых строк, но для регуляторного «Права на забвение» (право на удаление) это единственный корректный путь. Ни в коем случае это не должно быть еженедельным батчем - только редкое исключение по конкретным user_id.

8. Чек-лист Bronze-пайплайна

Чек-лист для код-ревью

Критерий Правильно Неправильно
Режим записи mode("append") / .append() mode("overwrite")
Дубли Сохраняются в Bronze dropDuplicates() при записи
Бизнес-фильтры Отсутствуют WHERE status = 'OK'
JOIN при записи Нет JOIN со справочниками
Аудит-колонки Все 5 на месте Только raw_payload
_ingest_ts current_timestamp() Дата из payload
_source_file input_file_name() Хардкод строки
Партиционирование По days(_ingest_ts) По бизнес-дате из payload
Формат таблицы Iceberg / Delta (ACID) Голые Parquet файлы в S3
Обработка ошибок Все строки записываются Битые строки отфильтровываются

Различие между бизнес-ошибкой и техническим сбоем

Это различие определяет, кто исправляет проблему:

Бизнес-ошибка (исправляется в Silver): данные корректно прибыли в Bronze, но содержат невалидные значения по бизнес-правилам. Например, amount = -100, status = 'UNKNOWN', user_id = null. Bronze хранит их как есть. Silver фильтрует и отправляет в silver.quarantine. Исправление: обновить бизнес-правила в Silver-пайплайне.

Технический сбой инжекта (перезапускаем Bronze-джобу): данные не попали в Bronze из-за инфраструктурной проблемы. Например, сеть была недоступна, S3 вернул 503, Kafka была недоступна. Исправление: перезапустить Bronze-джобу. Благодаря at-least-once семантике и иммутабельности Bronze, перезапуск безопасен - дубли попадут в Bronze, Silver справится с ними через MERGE.

Bronze никогда не «исправляет» данные. Он только принимает и хранит.