Software Testing for Data EngineeringТестирование ПО для дата-инженерии

Fast-revision cheatsheet for Data Engineering interviews. Tailored for Senior/Lead DE. Шпаргалка для быстрого повторения перед Data Engineering интервью. Заточена под Senior/Lead DE.

Test PyramidTest Pyramid Spark / StreamingSpark / Стриминг Data QualityData Quality Idempotency & ContractsИдемпотентность и контракты CI for DataCI для данных
WHY THIS MATTERSПОЧЕМУ ЭТО ВАЖНО

Strong software testing methodologies and high-quality well-tested maintainable code are core expectations for Senior/Lead Data Engineering roles. Expect deep dives. Anchor answers in real work: dbt review standards (Career.io), Scala/Spark on Hadoop (Front Tier), Kafka (IU Group), Airflow backfills (McMakler), and the 6.2B-record entity-resolution defect (Meta).

Сильные методологии тестирования ПО и высококачественный хорошо протестированный поддерживаемый код — ключевые ожидания для Senior/Lead Data Engineering ролей. Ожидай глубоких вопросов. Якори ответов — из реальной работы: стандарты ревью dbt (Career.io), Scala/Spark на Hadoop (Front Tier), Kafka (IU Group), backfill-ы Airflow (McMakler), и дефект entity-resolution на 6.2B записей (Meta).

1 · The Data Test Pyramid1 · Пирамида тестов для данных

Same shape as software (many fast unit tests at the base, few slow E2E at the top) plus an orthogonal axis: data-quality / validation tests that run on live production data, not just on code in CI.

Та же форма, что и в разработке ПО (много быстрых unit-тестов в основании, мало медленных E2E наверху) плюс ортогональная ось: тесты data-quality / валидации, которые работают на живых продакшн-данных, а не только на коде в CI.

▲ slower / fewer / costlier ┌─────────┐ │ E2E │ full DAG on staging, real-ish data, SLAs ├─────────┤ │ INTEGRA │ Spark↔storage, model→model, ephemeral schema │ -TION │ local Spark / DuckDB / Testcontainers ├─────────┤ │ UNIT │ pure transforms, 1 UDF, 1 model on a fixture └─────────┘ milliseconds, no cluster ▼ faster / many / cheap ║ DATA QUALITY ║ ← runs continuously in PROD, not just CI ║ not_null · unique · RI · freshness · volume/anomaly ║
▲ медленнее / меньше / дороже ┌─────────┐ │ E2E │ полный DAG на staging, реальноподобные данные, SLA ├─────────┤ │ INTEGRA │ Spark↔хранилище, модель→модель, эфемерная схема │ -TION │ локальный Spark / DuckDB / Testcontainers ├─────────┤ │ UNIT │ чистые трансформации, 1 UDF, 1 модель на фикстуре └─────────┘ миллисекунды, без кластера ▼ быстрее / больше / дешевле ║ DATA QUALITY ║ ← работает непрерывно в PROD, не только в CI ║ not_null · unique · RI · freshness · volume/anomaly ║
INTERVIEWER WILL PROBEО ЧЁМ СПРОСИТ ИНТЕРВЬЮЕР

"Why is it sometimes an hourglass/diamond for data?" → A lot of SQL is glue with thin unit-testable logic, so teams over-invest in integration + DQ. Honest answer: push logic into testable units where possible, and accept a heavier DQ layer because data is an input you don't control.

«Почему иногда это песочные часы/ромб для данных?» → Много SQL — это клей с тонкой логикой, пригодной для unit-тестов, поэтому команды переинвестируют в интеграцию + DQ. Честный ответ: выносить логику в тестируемые юниты, где возможно, и принять более тяжёлый слой DQ, потому что данные — это вход, который ты не контролируешь.

2 · Unit vs Integration vs E2E (one pipeline)2 · Unit vs Integration vs E2E (один пайплайн)

Pipeline:Пайплайн: Kafka events → Spark cleansing → dbt aggregates → BI table

LevelScopeExample on this pipelineSpeed / count
UnitOne function in isolationparse_event(json)→Row; one Spark transform on a 5-row fixturems · many
IntegrationComponents wired togetherSpark reads local Parquet → writes a partition → dbt model joins/aggregates itsec · some
E2EWhole DAG, real envPush synthetic batch through Kafka on staging; assert final BI row counts + partition lands within SLAmin · few
Data QualityLive data assertionsnot_null/unique on grain, RI, freshness, volume band — continuously in prodscheduled · ongoing
УровеньОхватПример на этом пайплайнеСкорость / количество
UnitОдна функция изолированноparse_event(json)→Row; одна Spark-трансформация на 5-строчной фикстуремс · много
IntegrationКомпоненты связаны вместеSpark читает локальный Parquet → пишет партицию → dbt-модель джойнит/агрегируетсек · несколько
E2EВесь DAG, реальное окружениеПротолкнуть синтетический батч через Kafka на staging; проверить финальные строки в BI + партиция приехала в SLAмин · мало
Data QualityПроверки живых данныхnot_null/unique на грейн, RI, freshness, объёмы — непрерывно в продепо расписанию · постоянно
RULE OF THUMBПРАВИЛО БОЛЬШОГО ПАЛЬЦА

Keep E2E tiny (slow + flaky + expensive). Most coverage at unit; integration for the seams; DQ for the data you can't pin down in advance.

Держать E2E минимальным (медленно + капризно + дорого). Основное покрытие — unit; интеграция для стыков; DQ для данных, которые нельзя зафиксировать заранее.

3 · Testing Spark Jobs3 · Тестирование Spark-джобов

top 5топ 5

The core move: separate logic from runtime and I/O

Ключевой ход: разделить логику от runtime и I/O

  • Factor transforms into pure functions: DataFrame → DataFrame, no reads/writes inside.
  • One local SparkSession per suite (local[2], not local[1] — surfaces partition-boundary bugs), reused.
  • Build small typed input DataFrames in code; assert with an order-insensitive equality helper (chispa for PySpark, spark-testing-base for Scala).
  • Assert the schema explicitly, not just values — schema drift is a top failure mode.
  • Выносить трансформации в чистые функции: DataFrame → DataFrame, без чтения/записи внутри.
  • Одна локальная SparkSession на suite (local[2], а не local[1] — проявляет баги на границах партиций), переиспользуемая.
  • Собирать маленькие типизированные входные DataFrames в коде; проверять с помощником order-insensitive equality (chispa для PySpark, spark-testing-base для Scala).
  • Проверять схему явно, не только значения — дрейф схемы — топовый режим сбоя.
# PySpark unit test — factored transform + local session
def flag_abuse(df):
    return df.withColumn("is_abuse", col("score") > 0.8)

def test_flag_abuse(spark):           # spark = session-scoped fixture
    inp = spark.createDataFrame(
        [(1, 0.9), (2, 0.1)], "id int, score double")
    out = flag_abuse(inp)
    assert_df_equality(out, expected, ignore_row_order=True)  # chispa
GOTCHASПОДВОДНЫЕ КАМНИ
  • Row ordering + nullable flags break naive collect() equality → use a comparator, sort first.
  • Non-determinism: current_timestamp, monotonically_increasing_id, unseeded random, shuffle order. Inject a clock, seed RNGs, never assert wall-clock.
  • Floating point needs tolerance.
  • Tests that secretly hit S3/HDFS/Hive are integration tests in disguise — slow + flaky.
  • local[1] hides bugs that assume all data in one partition.
  • Порядок строк + nullable-флаги ломают наивное collect() равенство → использовать компаратор, сортировать сначала.
  • Недетерминизм: current_timestamp, monotonically_increasing_id, не-засиженный random, порядок shuffle. Инжектить часы, сидировать RNG, никогда не проверять wall-clock.
  • Floating point требует толерантности.
  • Тесты, которые тайно обращаются к S3/HDFS/Hive — это замаскированные интеграционные тесты — медленно + капризно.
  • local[1] скрывает баги, которые предполагают все данные в одной партиции.
SCALA / JVM BRIDGE (study-gap — be explicit)SCALA / JVM BRIDGE (пробел для изучения — говори явно)

ScalaTest/specs2 + spark-testing-base (DataFrameSuiteBase, SharedSparkContext). Strongly-typed Dataset[CaseClass] makes fixtures cleaner. JVM-specific target: closure serialization — capturing a non-serializable object in a Spark closure throws at runtime; write a test that exercises the closure.

ScalaTest/specs2 + spark-testing-base (DataFrameSuiteBase, SharedSparkContext). Строго типизированный Dataset[CaseClass] делает фикстуры чище. JVM-специфичная цель: сериализация замыканий — захват несериализуемого объекта в Spark-замыкании бросает исключение в runtime; пиши тест, который упражняет замыкание.

YOUR ANCHORТВОЙ ЯКОРЬ

"I owned a Scala/Spark training pipeline on Hadoop at Front Tier end-to-end, so I've lived the local-session + factored-transform pattern and the closure-serialization gotcha."

«Я владел end-to-end пайплайном обучения на Scala/Spark на Hadoop в Front Tier, поэтому я прожил паттерн локальной сессии + факторизованных трансформаций и подводный камень сериализации замыканий.»

4 · Fixtures & Golden Datasets4 · Фикстуры и Golden-датасеты

FixturesФикстуры

Small, hand-crafted, in-code inputs with a known expected output. For unit/integration. Keep tiny so the expected output is obvious. Cover edge cases: nulls, dupes, weird encodings, empty.

Маленькие, вручную созданные, в коде входы с известным ожидаемым выходом. Для unit/integration. Держать крошечными, чтобы ожидаемый выход был очевиден. Покрывать крайние случаи: null-ы, дубли, странные кодировки, пусто.

Golden / snapshotGolden / snapshot

A stored known-good output. Run on a frozen input, diff vs the golden. Diff → either a regression or you intentionally re-bless. Best for complex transforms where hand-writing expected output is infeasible.

Сохранённый известно-правильный выход. Прогнать на замороженном входе, сравнить с golden. Разница → либо регрессия, либо ты намеренно переблагословляешь. Лучше всего для сложных трансформаций, где ручное написание ожидаемого выхода нереально.

GOTCHASПОДВОДНЫЕ КАМНИ

Goldens rot and mask bugs if blindly re-blessed. Freeze + version inputs and goldens; normalise non-deterministic columns (timestamps, ids) before diffing; review golden changes in PR like code. For volume/property coverage, prefer synthetic generators (Faker, schema-driven) over giant goldens.

Golden-ы гниют и маскируют баги, если слепо переблагословлять. Заморозить + версионировать входы и golden-ы; нормализовать недетерминированные колонки (timestamps, id-шники) перед диффом; ревьюить изменения golden-ов в PR как код. Для покрытия объёма/свойств предпочитать синтетические генераторы (Faker, schema-driven) вместо гигантских golden-ов.

5 · Data Quality: dbt tests vs Great Expectations5 · Data Quality: dbt-тесты vs Great Expectations

top 5топ 5
dbt testsGreat Expectations (GX)
Whatgeneric (not_null, unique, accepted_values, relationships) + singular SQL tests; packages add row-count/recency/dist"expectations" library + suites + checkpoints + auto HTML Data Docs; data profiling
Strengthslives with transform code, runs in dbt build, version-controlled, low ceremonyrich statistical/distribution checks, profiling, cross-source, auditable reports
WeaknessesSQL-warehouse-centric; weak on stats/reportingheavier to set up & operate; another system
Use whendefault in-warehouse quality gatestats/profiling/cross-system, formal DQ reporting (also Soda, Deequ)
dbt testsGreat Expectations (GX)
Чтоgeneric (not_null, unique, accepted_values, relationships) + singular SQL-тесты; пакеты добавляют row-count/recency/distбиблиотека «expectations» + suites + checkpoints + авто HTML Data Docs; профилирование данных
Сильные стороныживёт с кодом трансформаций, запускается в dbt build, под версионным контролем, минимум церемонийбогатые статистические/распределения проверки, профилирование, кросс-источник, аудируемые отчёты
СлабостиSQL-warehouse-центричен; слаб в статистике/отчётноститяжелее настроить и эксплуатировать; ещё одна система
Использовать когдадефолтный in-warehouse quality gateстатистика/профилирование/кросс-система, формальная отчётность DQ (также Soda, Deequ)

The 7 DQ dimensions to name out loud

7 измерений DQ, чтобы назвать вслух

completeness · uniqueness · validity (ranges/enums) · consistency (cross-col/table) · referential integrity · freshness/timeliness · volume/anomaly

полнота · уникальность · валидность (диапазоны/enums) · консистентность (cross-col/table) · referential integrity · freshness/timeliness · объём/аномалии

INTERVIEWER WILL PROBEО ЧЁМ СПРОСИТ ИНТЕРВЬЮЕР

"What happens when a DQ test fails?" → Distinguish blocking tests (fail build, page) from warning/anomaly tests (alert + circuit-break downstream, don't always hard-fail). Tie severity to data criticality + ownership + runbook.

«Что происходит, когда DQ-тест падает?» → Различать блокирующие тесты (провал сборки, пейдж) и warning/anomaly тесты (алерт + размыкание контура ниже по потоку, не всегда хард-fail). Привязывать серьёзность к критичности данных + владение + runbook.

YOUR ANCHORТВОЙ ЯКОРЬ

"At Meta a referential-integrity + cardinality monitor would have caught the entity-resolution defect (6.2B records, 34% of threads, 100+ downstream metrics). At Career.io dbt tests were mandatory on every model PR."

«В Meta монитор referential-integrity + кардинальности поймал бы дефект entity-resolution (6.2B записей, 34% тредов, 100+ downstream метрик). В Career.io dbt-тесты были обязательны на каждый PR модели.»

6 · Idempotency Tests6 · Тесты идемпотентности

top 5топ 5

Idempotency = re-running a job for the same logical inputs yields the same final state — no duplicates, no double-counting. Mandatory because retries, backfills, and at-least-once delivery are unavoidable.

Идемпотентность = повторный запуск джоба для тех же логических входов даёт то же финальное состояние — нет дублей, нет двойного подсчёта. Обязательно, потому что ретраи, backfill-ы и at-least-once доставка неизбежны.

Achieve it

Добиться этого

  • INSERT OVERWRITE PARTITION / dbt incremental delete+insert or merge — not blind append
  • Deterministic keys; dedup on natural/surrogate key
  • MERGE INTO upsert (Iceberg/Delta) on PK
  • No now() in transforms — thread the logical run date through
  • INSERT OVERWRITE PARTITION / dbt incremental delete+insert или merge — не слепой append
  • Детерминированные ключи; дедуп на натуральный/суррогатный ключ
  • MERGE INTO upsert (Iceberg/Delta) на PK
  • Никаких now() в трансформациях — протаскивать логическую дату запуска

Test it

Тестировать это

  • Run a partition twice → assert identical (row count + checksum of sorted rows)
  • Inject a mid-write failure, re-run → no dupes, correct state
  • Backfill an existing date → overwrites cleanly, does not append
  • Прогнать партицию дважды → проверить идентичность (число строк + чексумма отсортированных строк)
  • Инжектить падение посреди записи, перезапустить → нет дублей, правильное состояние
  • Backfill существующей даты → чисто перезаписывает, не дописывает
# idempotency assertion
run_partition(ds="2026-06-14")
h1 = checksum(read(ds="2026-06-14"))
run_partition(ds="2026-06-14")   # retry / backfill
h2 = checksum(read(ds="2026-06-14"))
assert h1 == h2                # no drift, no dupes
SAY THISСКАЖИ ЭТО

The pragmatic "exactly-once" is at-least-once delivery + idempotent sink.

Прагматичный «exactly-once» — это at-least-once доставка + идемпотентный приёмник.

YOUR ANCHORТВОЙ ЯКОРЬ

"McMakler — Airflow backfills across three migrations. Idempotent partition overwrites are what made re-runs safe; I'd test by re-running a date and diffing the partition."

«McMakler — backfill-ы Airflow через три миграции. Идемпотентные перезаписи партиций — это то, что сделало перезапуски безопасными; я бы тестировал, перезапуская дату и сравнивая партицию.»

7 · CI/CD for Data7 · CI/CD для данных

top 5топ 5

On every PR

На каждый PR

  • Lint/format (sqlfluff, ruff, scalafmt) + type checks
  • Compile/parse the DAG (dbt compile, Airflow DAG import test, Spark compiles)
  • Unit tests on transforms (no warehouse)
  • Slim CI: dbt build --select state:modified+ --defer — only changed models + downstream, in an ephemeral schema
  • Integration vs local Spark/DuckDB/Testcontainers
  • Lint/format (sqlfluff, ruff, scalafmt) + проверки типов
  • Компилить/парсить DAG (dbt compile, Airflow DAG import test, Spark компилится)
  • Unit-тесты на трансформациях (без warehouse)
  • Slim CI: dbt build --select state:modified+ --defer — только изменённые модели + downstream, в эфемерной схеме
  • Интеграция vs локальный Spark/DuckDB/Testcontainers

On merge / deploy

На merge / deploy

  • Deploy DAGs, promote dbt, smoke/E2E on staging
  • Write-Audit-Publish: write to _staging / Iceberg branch → validate → atomic swap
  • Snapshot/state versioning for safe rollback (Iceberg snapshots, dbt state)
  • Деплой DAG-ов, промоут dbt, smoke/E2E на staging
  • Write-Audit-Publish: писать в _staging / Iceberg-ветку → валидировать → атомарный swap
  • Snapshot/версионирование состояния для безопасного отката (Iceberg snapshots, dbt state)

What's different from app CI

Чем отличается от app CI

  • Data is an uncontrolled input → code can be correct yet the pipeline "fails" because upstream changed → the continuous DQ layer app CI lacks.
  • Tests are slow/expensive → sampling, slim CI, ephemeral envs.
  • Backfills/reprocessing are first-class → test idempotency, not just "runs once".
  • Schema/contract changes are deploys that silently break consumers → contract + schema-evolution checks.
  • Данные — неконтролируемый вход → код может быть правильным, но пайплайн «падает», потому что upstream изменился → непрерывный слой DQ, которого нет в app CI.
  • Тесты медленные/дорогие → сэмплинг, slim CI, эфемерные окружения.
  • Backfill-ы/переобработка — первоклассный кейс → тестировать идемпотентность, а не только «запустилось один раз».
  • Изменения схемы/контракта — это деплои, которые тихо ломают консьюмеров → проверки контракта + эволюции схемы.
INTERVIEWER WILL PROBEО ЧЁМ СПРОСИТ ИНТЕРВЬЮЕР

"Test against prod data without breaking prod?" → ephemeral schemas, zero-copy/prod clones, sampled fixtures, and write-audit-publish.

«Тестировать против прод-данных без поломки прода?» → эфемерные схемы, zero-copy/прод-клоны, сэмплированные фикстуры и write-audit-publish.

YOUR ANCHORТВОЙ ЯКОРЬ

"McMakler — moved orchestration to Airflow + rebuilt the analytical layer in dbt; that's exactly where slim CI and staging-then-publish earn their keep."

«McMakler — переехал оркестрацию на Airflow + перестроил аналитический слой в dbt; это именно то место, где slim CI и staging-then-publish оправдывают себя.»

8 · Contract Testing & Schema Drift8 · Тестирование контрактов и дрейф схемы

A data contract = versioned producer↔consumer agreement: schema (fields/types/nullability), semantics, SLAs (freshness/volume), quality guarantees. Contract testing verifies producers don't break it before shipping — shifting the break left instead of a 3am downstream page.

Контракт данных = версионированное соглашение продюсер↔консьюмер: схема (поля/типы/nullable), семантика, SLA (freshness/объём), гарантии качества. Тестирование контрактов проверяет, что продюсеры не ломают его до шиппинга — сдвигая поломку влево вместо 3-часового ночного пейджа downstream.

  • Schema Registry (Confluent, Avro/Protobuf on Kafka) enforces compatibility modes: backward / forward / full — a producer literally can't publish a breaking schema.
  • Tables: dbt model contracts: enforced, schema checks in CI, consumer-driven contracts (downstream declares the columns it relies on; producer CI fails on violation).
  • Schema evolution: additive nullable = safe; rename/narrow/drop = breaking → coordinated deploy. Iceberg/Delta evolve by field ID, not position — far safer than positional CSV parsing.
  • Schema Registry (Confluent, Avro/Protobuf на Kafka) обеспечивает режимы совместимости: backward / forward / full — продюсер буквально не может опубликовать ломающую схему.
  • Таблицы: dbt-модель contracts: enforced, проверки схемы в CI, consumer-driven контракты (downstream объявляет колонки, на которые полагается; producer CI падает при нарушении).
  • Эволюция схемы: additive nullable = безопасно; rename/narrow/drop = breaking → координированный деплой. Iceberg/Delta эволюционируют по field ID, а не позиции — гораздо безопаснее, чем позиционный парсинг CSV.
INTERVIEWER WILL PROBEО ЧЁМ СПРОСИТ ИНТЕРВЬЮЕР

"How do you detect drift?" → registry compatibility checks, expect_table_columns_to_match_set, dbt source freshness + contracts, metadata diff in CI + at ingest. Feed an evolved-schema fixture and assert graceful handling (missing column = clear failure, not a silent null).

«Как обнаружить дрейф?» → проверки совместимости в registry, expect_table_columns_to_match_set, dbt source freshness + контракты, diff метаданных в CI + при ingestion. Подать фикстуру с evolved-схемой и проверить graceful обработку (отсутствующая колонка = явный провал, а не тихой null).

YOUR ANCHORТВОЙ ЯКОРЬ

"At Meta I consolidated 10 fragmented logging systems (3.6B events/day) into one model — a contract problem at its core. Each logger was an implicit contract; unifying them needed explicit schema + semantic agreements and drift detection."

«В Meta я консолидировал 10 фрагментированных систем логирования (3.6B событий/день) в одну модель — проблема контрактов в ядре. Каждый логгер был неявным контрактом; объединение их требовало явных схем + семантических соглашений и обнаружения дрейфа.»

9 · Mocking External Systems9 · Мокирование внешних систем

SystemUnit (mock/fake)Integration (real-ish)
Object store / FSmoto, local temp dirMinIO / s3mock
Kafkafake producer behind interfaceembedded Kafka / Testcontainers
REST APIresponses/requests-mock, WireMockWireMock server / sandbox
Warehouse/DBDuckDB / SQLite stand-inTestcontainers Postgres/Spark
СистемаUnit (mock/fake)Integration (реально-подобная)
Object store / FSmoto, локальный temp dirMinIO / s3mock
Kafkafake producer за интерфейсомembedded Kafka / Testcontainers
REST APIresponses/requests-mock, WireMockWireMock server / sandbox
Warehouse/DBDuckDB / SQLite заменаTestcontainers Postgres/Spark
PRINCIPLEПРИНЦИП

Mock at the boundary you own (your client wrapper), not deep library internals. Prefer fakes / Testcontainers / local engines over heavy mocking for data systems — mocking a warehouse's SQL semantics is a trap; a real ephemeral DB catches dialect bugs. The axis is fidelity vs speed: mock cheap edges, use real systems where wrong assumptions are expensive (SQL dialect, Kafka protocol, transaction semantics).

Мокировать на границе, которой владеешь (твой клиентский обёртка), а не глубоких внутренностях библиотеки. Предпочитать фейки / Testcontainers / локальные движки вместо тяжёлого мокирования для систем данных — мокировать SQL-семантику warehouse — это ловушка; реальная эфемерная БД ловит баги диалекта. Ось — точность vs скорость: мокировать дешёвые края, использовать реальные системы там, где неправильные предположения дороги (SQL-диалект, Kafka-протокол, семантика транзакций).

ALWAYS TEST THE SAD PATHВСЕГДА ТЕСТИРОВАТЬ SAD PATH

Inject timeouts, partial reads, throttling, out-of-order/late records, and assert retry/backoff behaviour — not just the happy path.

Инжектить таймауты, частичные чтения, throttling, out-of-order/поздние записи и проверять поведение retry/backoff — не только happy path.

YOUR ANCHORТВОЙ ЯКОРЬ

"Kafka ingestion at IU Group (Syntea) — I'd use Testcontainers/embedded Kafka for integration and a fake producer for unit-level consumer logic."

«Kafka ingestion в IU Group (Syntea) — я бы использовал Testcontainers/embedded Kafka для интеграции и fake producer для unit-уровневой логики консьюмера.»

10 · Testing Streaming Jobs10 · Тестирование стриминговых джобов

Streaming adds time, ordering, and state as test dimensions.

Стриминг добавляет время, порядок и состояние как измерения тестов.

  • Determinism via event-time + injected watermarks — drive tests by controlled event-time, not wall-clock; assert window outputs.
  • Spark Structured Streaming: MemoryStream source + in-memory sink; advance batches; assert per-batch.
  • Flink: OneInputStreamOperatorTestHarness, MiniClusterWithClientResource — push timestamped records, fire watermarks/timers, snapshot/restore state to test checkpoint recovery + exactly-once.
  • Cases: late/out-of-order past watermark (drop vs side-output), window boundaries, stateful aggregations across restarts (restore from savepoint, assert continuity), exactly-once vs at-least-once + idempotent sink, backpressure (load/integration).
  • Детерминизм через event-time + инжектированные watermarks — управлять тестами контролируемым event-time, а не wall-clock; проверять выходы окон.
  • Spark Structured Streaming: MemoryStream source + in-memory sink; продвигать батчи; проверять per-batch.
  • Flink: OneInputStreamOperatorTestHarness, MiniClusterWithClientResource — толкать timestamped-записи, запускать watermarks/таймеры, snapshot/restore состояние, чтобы тестировать восстановление checkpoint + exactly-once.
  • Кейсы: поздние/out-of-order за watermark (drop vs side-output), границы окон, stateful агрегации через рестарты (restore из savepoint, проверять continuity), exactly-once vs at-least-once + идемпотентный sink, backpressure (load/integration).
FLINK = STUDY-GAP (be explicit, don't overclaim)FLINK = ПРОБЕЛ ДЛЯ ИЗУЧЕНИЯ (говори явно, не переоценивай)

"I have Kafka production experience (IU Group), but Flink itself is a study area. I'd lean on the Flink test harnesses + MiniCluster and the savepoint-restore test pattern. The mental model transfers from Spark Structured Streaming's MemoryStream + watermark testing."

«У меня продакшн-опыт Kafka (IU Group), но Flink сам по себе — область для изучения. Я бы опирался на тестовые harnesses Flink + MiniCluster и паттерн теста savepoint-restore. Ментальная модель переносится из Spark Structured Streaming MemoryStream + тестирования watermark.»

11 · Differentiators (depth)11 · Дифференциаторы (глубина)

Write-Audit-Publish (WAP)

Write-Audit-Publish (WAP)

Write to staging/Iceberg branch → audit with DQ/contract checks → publish (atomic swap) only if audits pass. Blue-green for data; bad data is never visible; clean rollback by repointing snapshot.

Писать в staging/Iceberg-ветку → аудит с DQ/контрактными проверками → publish (атомарный swap) только если аудиты прошли. Blue-green для данных; плохие данные никогда не видны; чистый откат переключением snapshot.

Property-based testing

Property-based тестирование

Hypothesis / ScalaCheck. Assert properties over generated inputs: "dedup ⇒ key unique", "Σ partition aggs = global agg", "serialize∘deserialize = id". Auto-shrinks to minimal counterexample.

Hypothesis / ScalaCheck. Проверять свойства на сгенерированных входах: «dedup ⇒ ключ уникален», «Σ аггрегаций партиций = глобальная аггрегация», «serialize∘deserialize = id». Автосжатие до минимального контрпримера.

Testing vs observability

Тестирование vs observability

Tests = explicit assertions for known failure classes (CI + scheduled). Observability = continuous instrumentation surfacing unknown-unknowns / anomalies (Monte-Carlo-style). Need both.

Тесты = явные проверки для известных классов сбоев (CI + по расписанию). Observability = непрерывная инструментация, всплывающая unknown-unknowns / аномалии (Monte-Carlo-стиль). Нужны оба.

Beyond line coverage

Дальше line coverage

Mutation testing (mutate code; do tests fail?), DQ coverage of critical columns/grains, escaped-defect rate, alert precision/recall, flakiness rate, and data scenarios exercised (nulls/dupes/late/skew).

Mutation testing (мутировать код; падают ли тесты?), DQ-покрытие критичных колонок/грейнов, escaped-defect rate, precision/recall алертов, частота капризности (flakiness), и сценарии данных, которые упражнены (nulls/дубли/поздние/skew).

Distributed failure modes to "reason about"

Распределённые режимы сбоев, про которые «рассуждать»

partial failure/retries · data skew (hot key) · late/out-of-order/duplicate · schema drift · stragglers/OOM/shuffle spill · clock skew & TZ/DST · backpressure & checkpoint recovery · small-files/partition explosion

частичный сбой/ретраи · перекос данных (горячий ключ) · поздние/out-of-order/дубли · дрейф схемы · отстающие/OOM/shuffle spill · clock skew & TZ/DST · backpressure & восстановление checkpoint · маленькие файлы/взрыв партиций

FLAKY TEST = A BUGКАПРИЗНЫЙ ТЕСТ = БАГ

Quarantine and fix; never @Ignore forever. Muted tests erode trust in the whole suite.

Изолировать и починить; никогда не @Ignore навсегда. Заглушённые тесты разрушают доверие ко всему набору.

12 · Self-Quiz — tap a card to flip12 · Самопроверка — нажми на карточку, чтобы перевернуть

SparkSpark
Why run unit tests with local[2] not local[1]? Почему запускать unit-тесты с local[2], а не local[1]?
tap to revealнажми, чтобы открыть
local[1] puts all data in one partition and hides bugs in logic that assumes single-partition data (or wrong partition-boundary handling). local[2] forces real partitioning so those surface. local[1] кладёт все данные в одну партицию и скрывает баги в логике, которая предполагает single-partition данные (или неправильную обработку границ партиций). local[2] заставляет реальное партиционирование, так что эти баги всплывают.
IdempotencyИдемпотентность
Pragmatic recipe for "exactly-once"? Прагматичный рецепт для «exactly-once»?
tap to revealнажми, чтобы открыть
At-least-once delivery + idempotent sink (partition overwrite / MERGE upsert on a key). True exactly-once is expensive; idempotent effect is robust and cheap. At-least-once доставка + идемпотентный sink (перезапись партиции / MERGE upsert на ключ). Настоящий exactly-once дорог; идемпотентный эффект робастен и дёшев.
DataFrame testDataFrame-тест
#1 fix for comparing two DataFrames? №1 фикс для сравнения двух DataFrames?
tap to revealнажми, чтобы открыть
Order-insensitive comparison (sort first, compare schema + values, watch nullable flags + float tolerance). Naive collect() equality fails on row order. Order-insensitive сравнение (сначала сортировать, сравнивать схему + значения, следить за nullable-флагами + толерантностью float). Наивное collect() равенство падает на порядке строк.
CI for dataCI для данных
What is "slim CI" in dbt? Что такое «slim CI» в dbt?
tap to revealнажми, чтобы открыть
dbt build --select state:modified+ --defer — build/test only changed models + their downstream against the prod manifest. Fast PR feedback without rebuilding the world. dbt build --select state:modified+ --defer — собирать/тестировать только изменённые модели + их downstream против прод-манифеста. Быстрая обратная связь по PR без пересборки всего мира.
ContractsКонтракты
How does a Schema Registry enforce contracts? Как Schema Registry обеспечивает контракты?
tap to revealнажми, чтобы открыть
Compatibility modes (backward/forward/full). A producer cannot publish a schema that would break consumers under the configured mode — contract testing baked into the platform. Режимы совместимости (backward/forward/full). Продюсер не может опубликовать схему, которая сломала бы консьюмеров в настроенном режиме — тестирование контрактов вшито в платформу.
Schema evolutionЭволюция схемы
Why is Iceberg evolution safer than CSV? Почему эволюция Iceberg безопаснее, чем CSV?
tap to revealнажми, чтобы открыть
Iceberg resolves columns by field ID, not position — add/rename/reorder without breaking readers. Positional CSV parsing breaks the moment a column moves. Iceberg разрешает колонки по field ID, а не позиции — добавлять/переименовывать/переупорядочивать без поломки читателей. Позиционный парсинг CSV ломается в момент, когда колонка перемещается.
WAP
What is Write-Audit-Publish? Что такое Write-Audit-Publish?
tap to revealнажми, чтобы открыть
Write to staging/branch → audit with DQ/contract checks → publish (atomic swap) only if audits pass. Blue-green for data; bad data never reaches consumers; easy rollback. Писать в staging/ветку → аудит с DQ/контрактными проверками → publish (атомарный swap) только если аудиты прошли. Blue-green для данных; плохие данные никогда не достигают консьюмеров; лёгкий откат.
CoverageПокрытие
Why isn't line coverage enough? Почему line coverage недостаточно?
tap to revealнажми, чтобы открыть
It shows what ran, not whether assertions are meaningful. Use mutation testing, DQ coverage of critical grains, escaped-defect rate, and which data scenarios (nulls/dupes/late/skew) you actually exercise. Показывает, что запустилось, а не являются ли проверки осмысленными. Использовать mutation testing, DQ-покрытие критичных грейнов, escaped-defect rate и какие сценарии данных (nulls/дубли/поздние/skew) ты реально упражняешь.

13 · Study-Gap Cheat (say confidently, don't overclaim)13 · Шпаргалка по пробелам (говори уверенно, не переоценивай)

Iceberg

  • snapshots + time-travel (rollback/testing)
  • snapshots + time-travel (откат/тестирование)
  • branches/tags → WAP
  • ветки/тэги → WAP
  • hidden partitioning
  • скрытое партиционирование
  • field-ID schema evolution
  • MERGE INTO upserts

Flink

  • event-time + watermarks
  • keyed state + checkpoints/savepoints
  • exactly-once via 2-phase-commit sinks
  • exactly-once через 2-phase-commit sink-и
  • test: operator harness + MiniCluster
  • тест: operator harness + MiniCluster
  • bridge from Spark Structured Streaming
  • мост от Spark Structured Streaming

Scala / JVM testingScala / JVM тестирование

  • ScalaTest / specs2 / JUnit
  • ScalaCheck (property-based)
  • spark-testing-base (DataFrameSuiteBase)
  • Testcontainers-scala
  • Dataset[CaseClass] typed fixturesтипизированные фикстуры
FRAMING THE GAP HONESTLYЧЕСТНАЯ ФОРМУЛИРОВКА ПРОБЕЛА

"My recent stack is Python/SQL/dbt, but I have real Scala/Spark-on-Hadoop (Front Tier), Kafka (IU Group), and Airflow (McMakler). The testing principles — factor logic, inject time, idempotency, contracts, WAP — transfer directly; I'm actively bridging the Flink/Iceberg/Java vocabulary."

«Мой недавний стек — Python/SQL/dbt, но у меня реальный опыт Scala/Spark-на-Hadoop (Front Tier), Kafka (IU Group) и Airflow (McMakler). Принципы тестирования — факторизовать логику, инжектить время, идемпотентность, контракты, WAP — переносятся напрямую; я активно мосту Flink/Iceberg/Java словарь.»