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.
"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
| Level | Scope | Example on this pipeline | Speed / count |
|---|---|---|---|
| Unit | One function in isolation | parse_event(json)→Row; one Spark transform on a 5-row fixture | ms · many |
| Integration | Components wired together | Spark reads local Parquet → writes a partition → dbt model joins/aggregates it | sec · some |
| E2E | Whole DAG, real env | Push synthetic batch through Kafka on staging; assert final BI row counts + partition lands within SLA | min · few |
| Data Quality | Live data assertions | not_null/unique on grain, RI, freshness, volume band — continuously in prod | scheduled · 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, объёмы — непрерывно в проде | по расписанию · постоянно |
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
- 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] скрывает баги, которые предполагают все данные в одной партиции.
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; пиши тест, который упражняет замыкание.
"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. Разница → либо регрессия, либо ты намеренно переблагословляешь. Лучше всего для сложных трансформаций, где ручное написание ожидаемого выхода нереально.
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 tests | Great Expectations (GX) | |
|---|---|---|
| What | generic (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 |
| Strengths | lives with transform code, runs in dbt build, version-controlled, low ceremony | rich statistical/distribution checks, profiling, cross-source, auditable reports |
| Weaknesses | SQL-warehouse-centric; weak on stats/reporting | heavier to set up & operate; another system |
| Use when | default in-warehouse quality gate | stats/profiling/cross-system, formal DQ reporting (also Soda, Deequ) |
| dbt tests | Great 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 · объём/аномалии
"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.
"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
The pragmatic "exactly-once" is at-least-once delivery + idempotent sink.
Прагматичный «exactly-once» — это at-least-once доставка + идемпотентный приёмник.
"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-ы/переобработка — первоклассный кейс → тестировать идемпотентность, а не только «запустилось один раз».
- Изменения схемы/контракта — это деплои, которые тихо ломают консьюмеров → проверки контракта + эволюции схемы.
"Test against prod data without breaking prod?" → ephemeral schemas, zero-copy/prod clones, sampled fixtures, and write-audit-publish.
«Тестировать против прод-данных без поломки прода?» → эфемерные схемы, zero-copy/прод-клоны, сэмплированные фикстуры и write-audit-publish.
"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.
"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).
"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 · Мокирование внешних систем
▾| System | Unit (mock/fake) | Integration (real-ish) |
|---|---|---|
| Object store / FS | moto, local temp dir | MinIO / s3mock |
| Kafka | fake producer behind interface | embedded Kafka / Testcontainers |
| REST API | responses/requests-mock, WireMock | WireMock server / sandbox |
| Warehouse/DB | DuckDB / SQLite stand-in | Testcontainers Postgres/Spark |
| Система | Unit (mock/fake) | Integration (реально-подобная) |
|---|---|---|
| Object store / FS | moto, локальный temp dir | MinIO / s3mock |
| Kafka | fake producer за интерфейсом | embedded Kafka / Testcontainers |
| REST API | responses/requests-mock, WireMock | WireMock server / sandbox |
| Warehouse/DB | DuckDB / SQLite замена | Testcontainers Postgres/Spark |
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-протокол, семантика транзакций).
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.
"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).
"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 · маленькие файлы/взрыв партиций
Quarantine and fix; never @Ignore forever. Muted tests erode trust in the whole suite.
Изолировать и починить; никогда не @Ignore навсегда. Заглушённые тесты разрушают доверие ко всему набору.
12 · Self-Quiz — tap a card to flip12 · Самопроверка — нажми на карточку, чтобы перевернуть
▾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типизированные фикстуры
"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 словарь.»