Перейти к содержанию
Learning Platform
Глоссарий
Troubleshooting

Troubleshooting — Apache Flink Internals

База знаний типичных ошибок курса Apache Flink Internals.

Категория

Показано 36 из 36 ошибок

Симптомы

  • Метрика checkpointAlignmentTimeNanos на JobMaster UI растёт с миллисекунд до десятков секунд / минут. Checkpoint duration ~ alignment time + sync + async; в первой колонке UI alignment > 90% от total. Параллельно backPressuredTimeMsPerSecond > 500 у downstream subtask. Checkpoints начинают expire'ить (checkpoint.timeout, default 10 min), джоб уходит в restart.

Причина

Aligned checkpoint требует, чтобы оператор дождался barrier на ВСЕХ input channels. Под backpressure barrier попадает в очередь сетевых buffers (по тысяче на channel), и downstream видит его только когда обработает все накопленные events. Skew по ключу: один subtask hot, остальные простаивают — barrier на быстрых channels приходит сразу, на slow ждёт минутами. Slow sink (например, JDBC sink с RDBMS, который тормозит) распространяет backpressure вверх по DAG — alignment страдает у всех stateful operators. Network memory под завязку (taskmanager.memory.network), сетевые buffers перенасыщены — barrier застревает в очереди не из-за rate, а из-за объёма.

Решение

  1. Включить unaligned checkpoints: execution.checkpointing.unaligned.enabled=true. Barrier перепрыгивает накопленные buffers; in-flight данные пишутся в checkpoint. Trade-off — больший checkpoint size.
  2. Buffer debloating: taskmanager.network.memory.buffer-debloat.enabled=true + target-ms=1000. Уменьшает effective buffer queue под backpressure -> alignment быстрее.
  3. Разобраться с источником backpressure: профилировать sink (flame graph), увеличить sink parallelism, перейти на async sink (KafkaSink с at-least-once если приемлемо, или batch JDBC sink).
  4. Перебалансировать ключ keyBy (composite key, salt-key) если виноват skew.
  5. Уменьшить checkpoint interval (если state небольшой) — короткий interval = меньше накопления -> меньше alignment.

Симптомы

  • Checkpoint expire (Checkpoint <ID> expired before completing) в логе JobMaster. Subtask показывает Aligning для всех buffer-входов, не двигается. taskmanager.memory.network под 95%+. Job уходит в restart-loop.

Причина

Aligned checkpoint держит барьер и БЛОКИРУЕТ accepting events с channel, у которого barrier пришёл первым. Если backpressure длительный — sender субтаски накапливают events в своих output buffers; когда network memory исчерпана — sender блокируется на allocateBuffer(). Маленький checkpoint.timeout (например, 1 минута) при больших checkpoint sizes — physical-side тоже не успевает finish (async snapshot, upload в S3). Slow DFS (S3) — async-phase растягивается.

Решение

  1. Unaligned checkpoints — единственное настоящее решение (см. предыдущий issue).
  2. Увеличить execution.checkpointing.timeout (default 10 min) — но это симптоматическое.
  3. Увеличить taskmanager.memory.network.max если по факту network memory bottleneck.
  4. Включить buffer debloating чтобы alignment был быстрее.
  5. Снизить execution.checkpointing.interval — меньше накопится между checkpoints.

Симптомы

  • Метрика lastCheckpointSize стабильно растёт между checkpoints (не от изменения state, а от прироста). При retainedCheckpoints=N через дни bucket S3 раздут в TB. CleanUp старых checkpoint работает (retainedCheckpoints), но размер каждого индивидуального — больше и больше. Incremental checkpoint должен был наоборот быть малым.

Причина

RocksDB compaction stalled — много старых SST не объединились в один большой. Каждый incremental checkpoint берёт diff SST -> diff растёт пропорционально несмерженным SST. Write-heavy workload с маленьким max_background_jobs (RocksDB option) — compaction не успевает за писом. TTL не configured — стейт растёт; даже после compaction финальный размер большой. Слишком частые checkpoints — incremental захватывает SST, которые не успели compact'нуться, в каждый checkpoint.

Решение

  1. Включить generic incremental checkpoint (changelog state backend): state.backend.changelog.enabled=true. Размер checkpoint станет ограничен размером changelog, materialization идёт в фоне.
  2. Тюнить RocksDB compaction: state.backend.rocksdb.thread.num=4 (или больше), state.backend.rocksdb.compaction.style=LEVEL (default), state.backend.rocksdb.writebuffer.size=128MB.
  3. Добавить state TTL: StateTtlConfig.newBuilder(Time.days(7)).cleanupFullSnapshot().build() — но проверить нет ли ttl-related performance.
  4. Снизить state.checkpoints.num-retained (default 1, многие ставят 10+ для safety) — освободит storage.
  5. Включить background cleanup tombstones: state.backend.rocksdb.compaction.level.use-dynamic-size=true.

Симптомы

  • Метрика processing-rate (records/sec) внезапно падает в 5-10x. JFR показывает время в JNI calls (RocksDB Write). Логи RocksDB (если LOG_LEVEL=INFO) показывают Stalling writes because we have N immutable memtables. Backpressure появляется upstream stateful operator.

Причина

L0 SST files превысили level0_slowdown_writes_trigger (default 20) -> write rate замедляется. Если превысили level0_stop_writes_trigger (default 36) — writes полностью блокируются. Compaction thread pool слишком маленький (max_background_jobs default 2) — за write-heavy не успевает. Slow disk — compaction чтает и пишет много на одном диске, IO throttling. Большой write rate в один CF (один state primitive) — single-CF compaction не parallelize'ится.

Решение

  1. Увеличить background jobs: state.backend.rocksdb.thread.num=4-8 (зависит от cores).
  2. Увеличить L0 thresholds: state.backend.rocksdb.compaction.level0-file-num-compaction-trigger=8 -> write throttling reach позже.
  3. Migrate hot CF на SSD (если на HDD) — disk IO главный bottleneck.
  4. Распределить state по нескольким operators (split workload) — каждый CF своя independent compaction.
  5. Compression: state.backend.rocksdb.compression.per.level=SNAPPY,SNAPPY,SNAPPY,SNAPPY,SNAPPY,ZSTD,ZSTD — меньше disk IO на cold levels.

Симптомы

  • Метрика rocksdb_block_cache_hit_ratio < 80% (ниже 50% — критично). Метрика lookupLatency для ValueState / MapState скачет в десятки миллисекунд. CPU TM низкий, disk IO высокий. Operator throughput падает не из-за compute, а из-за state.value() latency.

Причина

Размер block cache меньше, чем working set (active hot keys). Default block_cache_size = 8MB на CF — катастрофически мало для production. Slot sharing — много CF в одном TM делят shared cache через managed memory, фрагментация. Read pattern non-locality — random access по ключам (large keyspace, hot data сheet несовпадает с cache LRU). Bloom filter не включён — каждый key lookup проверяет каждый SST.

Решение

  1. Использовать shared cache из managed memory: state.backend.rocksdb.memory.managed=true (default true). Увеличить taskmanager.memory.managed.size до 4-8GB.
  2. Включить Bloom filter: state.backend.rocksdb.use-bloom-filter=true + state.backend.rocksdb.bloom-filter.bits-per-key=10. Снижает SST checks для absent keys.
  3. Partition cache по CF (если один CF hot, других нет): state.backend.rocksdb.memory.high-prio-pool-ratio=0.2 — резерв priority для index/filter blocks (всегда in-memory).
  4. Pinning index/filter blocks: state.backend.rocksdb.use-direct-reads=true + memory.write-buffer-ratio=0.5.
  5. Уменьшить block size: state.backend.rocksdb.block.blocksize=4096 (default 4KB) — лучше для random reads.

Симптомы

  • TM pod рестартится с OOMKilled (k8s OOMKill, exit code 137). JVM heap fine (heap usage < limit), но RSS растёт постоянно. Логи Flink в порядке вплоть до kill. dmesg на node показывает Memory cgroup out of memory: Killed process N (TaskManagerRunner).

Причина

RocksDB block cache + write buffers + index/filter blocks + iterators — всё off-heap (managed memory). Без bounded sharing (managed memory shared cache disabled) каждый CF аккумулирует свой cache -> unbounded. Утечка iterators: пользовательский код держит StateIterator (mapState.iterator()) без close — pinned SST не освобождается. С Flink 1.10+ есть WriteBufferManager + shared cache; до этого версии — leak by design. Slot sharing + много CF (десятки state primitives) — index/filter blocks pinned размером десятки MB на CF, суммарно gigabytes.

Решение

  1. Убедиться state.backend.rocksdb.memory.managed=true (default true с 1.10+). Это включает WriteBufferManager и shared cache, ограниченные managed memory.
  2. Закрывать StateIterators: try (StateIterator iter = mapState.iterator()) { ... } — критично для MapState.
  3. Увеличить taskmanager.memory.jvm-overhead.fraction до 0.15 (default 0.10) — больше резерв для native leak.
  4. Использовать pinning: state.backend.rocksdb.metadata.cache.high-priority=true — гарантирует index/filter blocks не выбрасываются.
  5. Monitoring: включить RocksDB metrics (state.backend.rocksdb.metrics.block-cache-usage и др.) для мониторинга по факту.

Симптомы

  • После перехода на ForstDB (disaggregated state) operator throughput упал в 5-20x. Метрика forstdb_local_cache_hit_ratio < 70%. Latency state.value().get() = 30-100ms (вместо <1ms у local RocksDB). S3 cost растёт значительно (GET requests).

Причина

Local cache size слишком мал — рабочий set не помещается. ForstDB local cache хранит блоки SST, скачанные с S3. Random access pattern по большому keyspace — каждый запрос промахивается. Pre-warming не сделан после restart — первые минуты до 100% miss. S3 latency variability — даже cached SST может потребовать metadata round-trip.

Решение

  1. Увеличить local cache: state.backend.forst.cache.size=20GB (зависит от disk). По best practice — local cache = ~10-20% от total state size.
  2. Использовать NVMe для local cache (state.backend.forst.local-dir).
  3. Pre-warming через savepoint pre-load: restore из savepoint в новом cluster и дать pipeline покрутиться без real traffic — naturally заполняет cache.
  4. Async state access: убедиться что user code использует State V2 API (StateFuture<T>) — runtime batches I/O.
  5. Hot/cold separation: ValueState с TTL для cold data -> выпадает из hot path, не съедает cache.

Симптомы

  • После migration на ForstDB latency обработки записи скачет в 50x несмотря на async state API. Profiler показывает большую часть времени в StateFuture.get() / await(). Throughput низкий, CPU TM почти idle (waiting on I/O).

Причина

Пользовательский код вызывает stateFuture.get() (sync) сразу в processElement — это блокирует mailbox thread, ломает batching async runtime. Composition нескольких state operations через .thenCompose не используется — последовательные sync вызовы убивают concurrency. User UDF использует legacy State V1 API (ValueState.value()) на ForstDB backend — Flink wraps это в blocking sync (с warning в log).

Решение

  1. Переписать UDF на State V2: ValueState.asyncValue().thenAccept(v -> { ... outputResult ... }). Никаких .get() в processElement.
  2. Если нужно дождаться несколько state operations — комбинировать через StateFuture.combine / .thenCompose.
  3. Использовать ProcessFunction.processElement(... , output) -> output только через StateFuture callback. Никогда не блокировать mailbox thread.
  4. Проверять application code на pattern: legacy ValueState.value() в processElement -> нужен refactor.
  5. Async Execution Controller config: state.backend.forst.async.max-pending-operations=1000 — увеличить batch для high-throughput pipeline.

Симптомы

  • Streaming SQL job с 3+ JOIN таблицами имеет state size в 10-100x больше, чем expected. EXPLAIN PLAN показывает chained join (left-deep), хотя одна из таблиц очень selective. CheckpointSize > 100GB при ожидаемых 10GB. Лаг на checkpoint грейет.

Причина

Без table statistics (CREATE TABLE ... WITH ('statistics.row-count' = ...) или Catalog statistics) CBO работает по эвристикам — часто неоптимально для streaming. JoinReorderRule не активируется для regular streaming join — нужен MultiJoin (см. issue ниже). Cardinality estimation для temporal join / lookup join некорректна — CBO не знает selectivity предиката. Predicate pushdown не сработал на каком-то source — Filter после Source, а не внутри.

Решение

  1. Запустить ANALYZE TABLE если catalog поддерживает (Hive, Iceberg). Для Paimon — статистика per-bucket автоматически.
  2. Explicit hint: SELECT /*+ JOIN_ORDER(t1, t3, t2) */ ... — заставит planner использовать заданный порядок.
  3. Включить MultiJoin: table.optimizer.multi-join-enabled=true (если не дефолт). Тогда JoinReorderRule сможет work.
  4. Изменить SQL — наиболее selective JOIN ставить первым (left-deep), CBO так выберет автоматически.
  5. EXPLAIN CHANGELOG_MODE для streaming SQL — увидеть какие join idempotent, какие нет.

Симптомы

  • Streaming SQL JOIN двух CDC-sources (например, Debezium -> Kafka). Ожидание: state линеен (lookup index). Реальность: state хранит все исторические rows, растёт без bound. EXPLAIN показывает обычный StreamPhysicalJoin вместо StreamPhysicalDeltaJoin.

Причина

Один из sources не объявлен как CDC: Kafka source с format='json' даёт append-only, не changelog. Нужен 'format'='debezium-json' или 'canal-json'. Source не SUPPORTS_FULL_SNAPSHOT — DeltaJoin требует, чтобы source мог восстановить полное состояние. Например, Kafka topic с compaction = yes; raw events = no. Версия Flink < 1.19 — DeltaJoin введён в 1.19+. table.optimizer.delta-join-enabled=false (default false в некоторых сборках).

Решение

  1. Включить DeltaJoin: table.optimizer.delta-join-enabled=true.
  2. Объявить source с changelog format: CREATE TABLE orders (...) WITH ('connector'='kafka', 'format'='debezium-json', 'scan.startup.mode'='earliest-offset').
  3. Если source append-only (event log) — DeltaJoin неприменим. Использовать regular join с state TTL.
  4. Для Paimon source: ALTER TABLE orders SET ('changelog-producer'='input') — затем Paimon отдаёт changelog stream -> DeltaJoin работает.
  5. EXPLAIN PLAN — проверить, появилась ли StreamPhysicalDeltaJoin нода.

Симптомы

  • SQL c JOIN 4+ таблиц. EXPLAIN PLAN показывает вложенные Join (Join (Join (Join (TableScan, TableScan), TableScan), TableScan)) вместо одного MultiJoin. JoinReorderRule не применяется, порядок join — точно как в SQL (что часто неоптимально).

Причина

JoinToMultiJoinRule не сработала. Условия: все joins одного типа (INNER), без OUTER, без correlated subqueries. Join condition содержит OR — MultiJoin требует CNF condition. Один из joins — temporal или interval — MultiJoin не поддерживает. table.optimizer.multi-join.enabled=false.

Решение

  1. Включить: table.optimizer.multi-join.enabled=true.
  2. Заменить OUTER JOIN на INNER если бизнес-логика позволяет — outers не reordered.
  3. Разнести temporal / interval joins в отдельные CTE — оставшаяся часть превратится в MultiJoin.
  4. Переписать сложные OR-condition в CNF (a=b AND (c=d OR c=e) -> разбить на UNION).
  5. EXPLAIN PLAN после изменений — должна появиться MultiJoin нода.

Симптомы

  • Sink ожидает upsert (например, JDBC sink с PK), а upstream — retract stream (типичный для GROUP BY с changelog). Planner добавляет StreamPhysicalSink с requireRetract=true -> downstream state увеличен, performance ниже.

Причина

Sink декларирует только UPDATE_AFTER capability (upsert), но upstream — full retract (UPDATE_BEFORE + UPDATE_AFTER). Planner добавляет RetractToUpsertConversion которая держит state. Aggregate operator с two-phase agg (table.exec.mini-batch.enabled) генерирует retract — без primary key конверсия невозможна. Missing primary key declaration: CREATE TABLE results (id BIGINT, ... PRIMARY KEY (id) NOT ENFORCED) — без NOT ENFORCED планер не знает что upsert ok.

Решение

  1. Declare PRIMARY KEY ... NOT ENFORCED на sink CREATE TABLE — planner поймёт что upsert acceptable.
  2. Использовать GROUP BY где ключ это PK sink — тогда aggregate напрямую emit'ит upsert, без retract.
  3. Включить mini-batch + local-global aggregation: table.exec.mini-batch.enabled=true, table.optimizer.agg-phase-strategy=TWO_PHASE.
  4. EXPLAIN CHANGELOG_MODE — увидеть changelog mode на каждом этапе плана. Должно быть UPSERT, не RETRACT.

Симптомы

  • Windowed aggregation не emit'ит результаты, output stream молчит. Метрика currentWatermark на window operator — застывает на старом значении. Метрика на источнике (Kafka source) — currentSourceWatermark скачет правильно. Один Kafka partition отстаёт по timestamp.

Причина

Global watermark = min(watermark всех source subtasks). Если один partition slow (less throughput / consumer лаг) — global watermark застревает на его уровне. Один partition idle (нет events) — без idleness detection блокирует watermark. Skew по producer-time: partition 0 пишет real-time, partition 1 пишется backlog (catch-up из archive) — timestamp 1 далеко в прошлом -> global wm = past. Misconfigured WatermarkStrategy: forBoundedOutOfOrderness(Duration.ofDays(1)) — wm всегда = max - 1 day, окна никогда не закроются.

Решение

  1. Включить watermark alignment: WatermarkStrategy.<T>forBoundedOutOfOrderness(...).withWatermarkAlignment('group1', Duration.ofSeconds(20), Duration.ofSeconds(1)). Fast readers паузятся пока slow догонит.
  2. Включить idleness: WatermarkStrategy.<T>forBoundedOutOfOrderness(...).withIdleness(Duration.ofMinutes(2)). Idle partitions перестанут блокировать.
  3. Уменьшить outOfOrderness до realistic значения (секунды-минуты, не дни). Pour through histogram lateness в DataDog/Prometheus.
  4. Если backlog catch-up — отделить historical processing от real-time (отдельный job для catch-up), не mix в одном partition set.
  5. Метрики watermark per subtask: currentWatermark на каждом subtask UI — найти laggard.

Симптомы

  • Job throughput значительно ниже expected. UI показывает все subtasks OK (зелёные, isBackPressured=false). Но input rate < output rate ожидаемого. metrics.backPressuredTimeMsPerSecond никуда не двигается. Чувство, что что-то тормозит, но не видно где.

Причина

Backpressure detection в UI смотрит на ratio time spent на output buffer pool — если оператор chained (одна задача = несколько operators), внутренний backpressure скрывается. Источник sloweed by external system (slow Kafka broker), но в UI это выглядит как Source healthy — Source не получает events, backpressure не возникает. Slow sink написан асинхронно (внутренняя очередь) — оператор быстро отдаёт, очередь sink растёт без видимого backpressure. Skew — один subtask hot, остальные idle. UI показывает по-операторно (avg), individual subtask backpressure теряется.

Решение

  1. Отключить operator chaining для diagnostic: env.disableOperatorChaining() (или per-operator .disableChaining()). Backpressure станет видим на boundaries.
  2. Смотреть per-subtask: UI -> Job -> Vertices -> subtask details -> backpressure. Median может быть OK, max — критичный.
  3. Метрика inputQueueLength / outputQueueLength на subtask — non-zero значит данные накапливаются.
  4. Включить flame graph: env.java.opts -Dflink.web.flame-graph.enabled=true -> UI -> Vertex -> Flame Graph. Покажет реальную bottleneck функцию.
  5. Source-side метрики: pendingRecords (Kafka source), sourceIdleTime — если source idle, проблема upstream Kafka.

Симптомы

  • TM pod restart с exit code 137 (OOMKilled by k8s). Flink logs показывают нормальную работу до самого момента kill. JVM heap метрики OK. Stats показывают direct memory ниже limit. Но RSS process растёт неограниченно. dmesg на node показывает Memory cgroup out of memory.

Причина

JNI native allocations: RocksDB, Kafka native compression libraries (zstd, snappy), Beam Python worker — все эти не учтены в Flink memory model. JVM Metaspace растёт — много classes (особенно при PyFlink или dynamic class loading) — taskmanager.memory.jvm-metaspace.size default 256MB мал. Native thread stacks: каждый Java thread = ~1MB native. С 1000+ threads (Kafka, Netty, gRPC) — 1GB. Flink memory model резервирует jvm-overhead.fraction=0.10 — это часто недостаточно для production.

Решение

  1. Увеличить jvm-overhead: taskmanager.memory.jvm-overhead.fraction=0.15, taskmanager.memory.jvm-overhead.min=512MB.
  2. Увеличить metaspace: taskmanager.memory.jvm-metaspace.size=512MB.
  3. Профилировать native memory: jcmd <pid> VM.native_memory summary — увидеть какая категория растёт.
  4. Использовать native memory tracking: -XX:NativeMemoryTracking=detail.
  5. На k8s — поставить pod memory limit с 20% headroom над Flink total: total flink memory = pod memory * 0.85.

Симптомы

  • Job не стартует после рестарта или upgrade. Логи JobMaster: NoResourceAvailableException: Could not allocate slot for ... after slot.request.timeout. Slots на UI показываются Available, но не assigned. RM долго думает.

Причина

Suspended TM-pods после crash — RM ждёт их heartbeat timeout (default 50s) перед запросом новых. Slot sharing groups + co-location constraints — scheduler не может разместить subtask group в одном slot, инициирует cascading slot allocation. K8s scheduler медленно scheduleит TM pods (cluster под завязку, taints/tolerations не подобраны). Default slot.request.timeout=5 min не достаточно для large job (1000+ slots).

Решение

  1. Увеличить slot.request.timeout до 15 min для large jobs: slot.request.timeout=900000.
  2. Уменьшить heartbeat.timeout если TM short-lived environment: heartbeat.timeout=20000 — быстрее освобождать dead slots.
  3. Использовать adaptive scheduler — fail-fast не на all slots, gradient рестарт с меньшим parallelism.
  4. Pre-warming TM pods через standalone deployment с overprovisioning.
  5. K8s priorityClass для Flink TM — гарантия scheduling приоритета.

Симптомы

  • Job в state SCHEDULED, adaptive scheduler ждёт slots. После добавления TM pods (manual scale up) — job всё ещё не запускается. Логи: Could not assign all required slots after timeout.

Причина

Adaptive scheduler требует jobmanager.adaptive-scheduler.resource-stabilization-timeout перед re-evaluation (default 10s). Если кластер ещё не стабилизировался — wait. Минимум slots не достигнут: jobmanager.adaptive-scheduler.min-parallelism-increase=N — если scale up < N, scheduler не реагирует. Slot sharing config не правильный — нужно min(parallelism групп), не sum. Reactive mode не настроен — adaptive scheduler в обычном mode не масштабирует, только при slot.timeout.

Решение

  1. Включить reactive mode: jobmanager.scheduler=adaptive, scheduler-mode=reactive — job всегда занимает все available slots.
  2. Тюнить stabilization: jobmanager.adaptive-scheduler.resource-stabilization-timeout=30s + scaling-interval.min=30s.
  3. Снизить min-parallelism-increase=1 — реагировать на любое изменение.
  4. Проверить slot config: каждый TM имеет N slots; ResourceManager metric numAvailableTaskSlots должен показывать ожидаемое значение.
  5. Использовать Flink Kubernetes Operator с autoscaler — он умеет управлять scaling корректно.

Симптомы

  • Job постоянно рестартится (parallelism прыгает 4 -> 8 -> 4 -> 12). Autoscaler (k8s HPA или Flink K8s Operator autoscaler) видит metric, добавляет TM, Flink рестартит, lag временно растёт (catch-up), HPA реагирует обратно. Throughput нестабильный, checkpoint frequency скачет.

Причина

Слишком короткое stabilization window в autoscaler — реагирует на noise. Метрика autoscaler — CPU или backpressure — после rescale временно высокая (catch-up). Cooldown между scaling events отсутствует. Adaptive scheduler restart triggers — каждый rescale = downtime + state restore.

Решение

  1. Stabilization window в HPA: behavior.scaleDown.stabilizationWindowSeconds=300, scaleUp.stabilizationWindowSeconds=120.
  2. Cooldown между rescale операциями: scaling-interval.min=2min.
  3. Использовать Flink K8s Operator autoscaler (FLIP-271) — он understands backpressure cycle, не реагирует на catch-up noise.
  4. Hysteresis: target utilization 70% scale-down, 80% scale-up (avoid yoyoing around 75%).
  5. Disaggregated state (ForstDB) драматически снижает stop-the-world rescale time -> oscillation менее болезненно.

Симптомы

  • Kafka transaction в state PrepareCommit на брокере, висит часами. Downstream consumer в read_committed mode не видит данные (зависшая транзакция). После Flink restart те же transaction id используются — но broker fences old transactions, повтор повисает. Лог Flink: Transaction with id ... is currently in state ... timeout.

Причина

Flink JobMaster крашнул МЕЖДУ commit на checkpoint и notifyCheckpointComplete. Transaction в prepared state, но Flink не успел committed. Kafka broker transaction.timeout.ms превышен (default 1 hour). Brokerсейchen abort'ит transaction, но Flink думает что закоммитил — recovery ломается. Transactional id collision: два разных Flink джобы используют одинаковый prefix -> один fences другого. checkpoint.tolerable-failed-checkpoints большой -> Flink не fail'ит job, продолжает копить prepared transactions.

Решение

  1. Manual recovery: kafka-transactions.sh --describe-producers --topic ... — найти hanging txn; --abort-transactions.
  2. Increase Kafka transaction.timeout.ms = 2 * checkpoint.interval + safety margin (например, 30 min для 10-min checkpoint).
  3. Уникальный transactional.id.prefix для каждого Flink job: sink.setTransactionalIdPrefix('orders-v1') — обязательно change при upgrade.
  4. checkpoint.tolerable-failed-checkpoints=0 — fail-fast при first failure, не аккумулировать hanging txn.
  5. Monitoring: alert на zombie transactions старше N часов через kafka-broker metrics.

Симптомы

  • Один Flink job начинает выкидывать InvalidProducerEpochException постоянно. Sink subtasks не могут write — каждый attempt fence'ится. Лог: Producer attempted an operation with an old epoch. Other producer with same transactional ID fenced this producer.

Причина

Два разных Flink job (например, prod + staging) деплоятся с одинаковым transactional.id.prefix -> каждый fences другой. Upgrade job без смены transactional.id.prefix — старый job не полностью остановился, новый стартует с тем же prefix. Default prefix не настроен — некоторые versions Flink использовали generic prefix. Несколько Flink clusters пишут в один Kafka cluster (например, multi-region) — без unique prefix per cluster — collision.

Решение

  1. Уникальный prefix на каждый deployment: KafkaSink.builder().setTransactionalIdPrefix('orders-prod-region-eu-v3').setDeliveryGuarantee(EXACTLY_ONCE).build().
  2. Включить в prefix region/env/version: {service}-{env}-{region}-{version}.
  3. При upgrade — bump version в prefix (v3 -> v4), old prefix transactions corrupt не страшно.
  4. Naming convention документировать в playbook.
  5. Pre-deploy check: kafka-transactions.sh --list — увидеть существующие transactional ids перед deploy.

Симптомы

  • Source V2 connector (Kafka, Iceberg) запущен, но reader subtasks idle. Метрика numberOfAssignedSplits на readers = 0. Метрика SplitEnumerator's coordinator не показывает activity. SourceReader logs: Waiting for split assignment.

Причина

SplitEnumerator выбросил exception, координатор силен (но не fail'ит job целиком). Лог JobMaster: SplitEnumerator failed: .... Static enumerator (для batch source) уже all splits ассигнил, ничего нового не shows — это OK для batch. Discovery interval не настроен (для streaming source) — scan.partition.discovery.interval.ms=0 (default) -> never re-scan. Coordinator memory leak или GC stuck — SplitEnumerator на jobmanager process, не на TM.

Решение

  1. Проверить JobMaster логи на SplitEnumerator failures — это часто config error в source.
  2. Включить partition discovery для Kafka: scan.partition.discovery.interval.ms=30000 — каждые 30s checking новые partitions.
  3. Профилировать JobMaster (не только TM!) — coordinator memory.
  4. Restart job чтобы re-init SplitEnumerator state — если разовый glitch.
  5. Reader-side метрика sourceReaderAssignedSplits должна matching enumerator assignedSplits — если рассинхрон, network issue.

Симптомы

  • Watermark alignment включён, но fast reader не паузится. Метрика watermark на slow reader не reflects на других readers. WatermarkCoordinator (JobMaster-side) метрика не показывает aggregation.

Причина

withWatermarkAlignment передал split id, разный для каждого subtask -> coordinator не группирует. Source connector не имплементирует SourceReaderBase / SourceReader interface полностью — методы pauseOrResumeSplits не overridden. Update interval слишком большой: alignment(group, drift, Duration.ofMinutes(1)) — drift накопится за минуту до коррекции. Custom WatermarkGenerator не emit'ит watermark через ctx (silently swallows).

Решение

  1. Одинаковый split group для всех subtasks одного source: WatermarkStrategy.<T>forBoundedOutOfOrderness(...).withWatermarkAlignment('kafka-group', Duration.ofSeconds(30), Duration.ofSeconds(2)).
  2. Update interval = 1-5 секунд для tight alignment.
  3. Проверить SourceReader implementation — должен реагировать на assignedSplits и pauseSplits.
  4. Watermark generator emit через ctx.emitWatermark(new Watermark(ts)).
  5. Метрика на JobMaster: watermarkAlignmentReceivedReports — non-zero значит aggregation работает.

Симптомы

  • Paimon table bucket в S3 разрастается до TB (sometimes 10x от reasonable size). Hadoop ls показывает тысячи snapshot файлов в snapshot/. Query latency для time-travel запросов растёт. Compaction stalls.

Причина

snapshot.num-retained.min/max не настроены — default min=10, max=Integer.MAX_VALUE -> snapshots живут пока их не expire'нет background. Tag retention: tags не expire вообще — каждый tag pinning snapshots под собой. Стейт changelog для CDC sink — каждый change рождает snapshot, частота commit низкая -> много snapshots. Snapshot expire не запущен — нужно явно вызвать CALL sys.expire_snapshots('db.table', 100) periodically или включить write.compaction-task-num.

Решение

  1. Настроить retention: ALTER TABLE t SET ('snapshot.num-retained.max'='100', 'snapshot.time-retained'='24h').
  2. Tag retention: ALTER TABLE t SET ('tag.num-retained-max'='10').
  3. Запланировать periodic expire через Flink batch job: CALL sys.expire_snapshots ('default.t', RETENTION_TIME => 7 days).
  4. Compaction: write.compaction-task-num=2 -> background job сливает small snapshots.
  5. Включить async commit для writer: write.compaction.async=true -> не блокирует на каждый snapshot.

Симптомы

  • SELECT с predicate на Paimon table занимает минуты ДО начала read данных (query planning slow). Manifest list файл в S3 — десятки/сотни MB. Listing операций на S3 medio-долгий. Каждый commit добавляет manifest entry, его не consolidate.

Причина

Manifest consolidation не настроен. Каждый write commit рождает manifest entry (или несколько); они аккумулируются в manifest list. Высокая частота commit (streaming write каждые 30s) — manifests множатся. manifest.target-file-size слишком маленький -> много мелких manifest вместо нескольких больших. Compaction не сливает manifest tree.

Решение

  1. Настроить compaction manifest: manifest.target-file-size=8MB (default), manifest.merge-min-count=30.
  2. Включить full compaction по schedule: CALL sys.compact ('default.t', WHERE => '...').
  3. Снизить частоту commit (commit batch более крупный): write.batch-size=10000.
  4. Использовать file index (filebroker.path-prefix) для S3 чтобы избежать full listing.
  5. Monitoring metric manifestListSize — alert при превышении 50MB.

Симптомы

  • Query SELECT col1, col2 FROM fluss_table — ожидание только 2 column scan. Реальность: bandwidth с Fluss-cluster соответствует full row scan. Метрика fluss server bytes_out сильно выше expected. Latency высокая.

Причина

Source connector не пушит SupportsProjectionPushDown — старая версия flink-fluss-connector. SELECT * во вложенном CTE или view — top-level columns pruned, но Fluss получает request на все. Predicate в WHERE содержит UDF или complex expression — Fluss не может evaluate server-side. Materialized table downstream — projection pushed только в materialized definition, не в каждый SELECT.

Решение

  1. Обновить flink-fluss-connector до версии с full pushdown support.
  2. EXPLAIN PLAN — проверить что в TableSourceScan указаны только нужные columns, не *.
  3. Избегать complex predicate, использовать column-level filtering: WHERE region = 'EU' AND ts > '2026-01-01' (поддерживаются), не WHERE my_udf(col) > 0 (не пушится).
  4. Метрика fluss client metric — pushed_filters, projected_columns — verify pushdown.
  5. Если CTE — попробовать inline (рассмотреть materialization).

Симптомы

  • Job с CEP pattern API state size растёт линейно с временем. Метрика state size per key — десятки/сотни MB. Checkpoint duration deteriorates. Job рестартится из-за OOM на heap state backend или disk full на RocksDB.

Причина

Pattern с unbounded quantifier (oneOrMore без within): Pattern.begin('a').oneOrMore() — NFA держит все possible matches пока pattern не complete (а complete никогда). Pattern с followedBy без within — все intermediate events accumulate в partial matches. Non-deterministic contiguity (followedByAny) — exponential blow-up state. Skip strategy AfterMatchSkipStrategy.noSkip — переиспользует matched events, blow-up.

Решение

  1. Bounded patterns: всегда указывать .within(Time.minutes(5)) — partial matches expire после timeout.
  2. Determinictic contiguity: предпочитать .next() (strict) или .followedBy() (relaxed), избегать .followedByAny().
  3. AfterMatchSkipStrategy.skipPastLastEvent — после match не переиспользовать events.
  4. Greedy with limit: .times(1,10) вместо .oneOrMore() — bounded max occurrences.
  5. State TTL для CEP intermediate state: configurable через CEP.pattern(input, pattern, comparator, eventTime, ttl-config).

Симптомы

  • ML_PREDICT TVF добавлен в SQL, и pipeline резко замедляется. Метрика recordsOut на этом операторе ~ 10/sec вместо 1000/sec. Logs: ModelClient request timeout after 30s или backpressure upstream.

Причина

External model endpoint medio (OpenAI API, AWS Bedrock) имеет latency 100ms-5s/request. Sync вызов в processElement -> throughput = 1 / latency. max-concurrent-requests слишком маленький (default 8) — недоиспользуется async. Network reliability — external API rate-limit'ит, retries добавляют latency. Batch size = 1 (default) — не batches inputs, каждый element = round-trip.

Решение

  1. Async ML inference: использовать AsyncWaitOperator под капотом ML_PREDICT (с Flink 2.1+ default async).
  2. Увеличить async pool: ml-predict.async.max-concurrent-requests=64 (зависит от rate limit endpoint).
  3. Batch requests: ml-predict.async.batch.size=32 — agrupar inputs в одном HTTP call.
  4. Timeout config: ml-predict.async.timeout=60s, retry-on-failure=3.
  5. Если endpoint медленный — рассмотреть self-hosted model (vllm, triton) с lower latency.

Симптомы

  • VECTOR_SEARCH TVF на векторной таблице (Milvus, Qdrant). Ожидание: HNSW/IVF index lookup за миллисекунды. Реальность: query занимает секунды, server-side CPU спайки. Vector DB logs показывают full collection scan.

Причина

Predicate в WHERE содержит фильтр по metadata, который не pushed down — Flink делает post-filter после vector search. Vector DB вынужден return top-K = whole collection. Index не построен на vector column на Vector DB side. SELECT не использует index, делает brute force. Filter expression сложный (UDF, complex AND/OR) — Vector DB не парсит, fallback на brute force. VECTOR_SEARCH connector не поддерживает SupportsFilterPushDown — старая версия.

Решение

  1. Verify index в Vector DB — Milvus describe_index, Qdrant collection info.
  2. Filter rewrite: вместо complex predicate, use Vector DB-native filter syntax (если connector forwards).
  3. EXPLAIN PLAN — должен показать FilterPushDown в VectorTableScan.
  4. Использовать pre-filter (Vector DB-side) вместо post-filter: WHERE category = 'shoes' ORDER BY vector_dist(...) — должен push category в Milvus.
  5. Обновить flink-vector-search connector — recent versions support full pushdown.

Симптомы

  • PyFlink job с UDF на Python processing 1k records/sec, аналогичный Java job — 50k records/sec. CPU на TM ~30% (не utilized). Profiler: большая часть времени в Beam fn worker boundaries (serialize/deserialize).

Причина

Java ↔ Python boundary через Beam fn API: каждый record serialize (Java side) -> IPC -> deserialize (Python side) -> UDF -> serialize -> IPC -> deserialize. Bundle size слишком маленький: python.fn-execution.bundle.size=1000 (default) — каждый bundle = 1 round-trip. Малый bundle -> много round-trips. Python single-thread bottleneck (GIL) — даже async UDF не parallelize'ится. Использование Python UDF (slow) vs Pandas UDF / Vectorized UDF (быстрее на 10-50x).

Решение

  1. Use Vectorized UDF (Pandas Function): @udf(result_type=DataTypes.BIGINT(), func_type='pandas') — batch processing, numpy/pandas internally.
  2. Увеличить bundle size: python.fn-execution.bundle.size=10000, python.fn-execution.bundle.time=1000 (ms) — больше batching.
  3. Multiple Python processes per TM: python.fn-execution.arrow.batch.size — вместе с vectorized используют Arrow.
  4. Use Table API + SQL where possible, переходить на Python UDF только когда необходимо.
  5. Pure Java реализация UDF — если throughput critical и rewrite acceptable.

Симптомы

  • Restore из savepoint после schema change для Avro state class фейлится: IncompatibleStateException: ... fields don't match или SerializerSnapshot incompatibility. Allow non-restored state не помогает — Flink fails before reaching state matching phase.

Причина

Изменение Avro schema нарушило evolution rules: rename field, change type, remove required field — все breaking changes. Default value не указан для нового optional field — старый state не имеет этого поля, новый serializer не знает что подставить. Реестр schema не используется (embedded schema in state) — каждый snapshot содержит свою schema, mismatch detects. Avro Generic vs Specific — switch между ними при upgrade ломает state.

Решение

  1. Только additive evolution: добавлять только optional fields с default value: "default": null или "default": 0.
  2. Никогда не remove required field — пометить deprecated, продолжать писать default value.
  3. Использовать Confluent Schema Registry: register schema, Flink resolves at runtime.
  4. Тест-suite на evolution: bootstrap savepoint в старой схеме, restore в новой — для каждой PR.
  5. State Processor API для migration: read old state, transform, write new savepoint с новой schema.

Симптомы

  • Throughput Flink job неожиданно низкий. Profiler: большая часть в KryoSerializer.serialize / deserialize. Class который думали POJO — не POJO (missing no-arg constructor, или final field) -> silent fallback на Kryo. State migration после change class — breaks без warning.

Причина

Class объявлен с конструктором arg-only, без public no-arg -> Flink не может instantiate как POJO. Final fields — POJO requirement: not final non-static fields with getter/setter, OR no setter но public not-final field. Inner class without static modifier — Flink не сможет instantiate (требует enclosing instance). Default Flink config — pipeline.generic-types=true (Kryo fallback enabled silently).

Решение

  1. Strict mode: pipeline.generic-types=false -> Flink fails-fast при Kryo fallback вместо silent slow.
  2. ENV.getConfig().disableGenericTypes() — programmatic equivalent.
  3. Best practice: declare PojoTypeInfo.isPojo(MyClass.class) в unit test — fails сразу если class не POJO.
  4. Если class должен быть POJO — добавить public no-arg constructor + public getters/setters + non-final fields.
  5. Если Avro нужен — TypeInfo.of(MyAvroRecord.class) явно, не полагаться на auto-detect.

Симптомы

  • Включили JFR на production TM. Throughput упал на 10-20%, GC time выросло, latency P99 скакнуло в 2x. Default JFR settings — heavy для production: stack traces на каждом allocation.

Причина

Default JFR settings — high-detail (settings=default.jfc). На high-throughput Flink — allocation rate миллионы в секунду, JFR sample на 1% — overhead. Continuous recording — file растёт unbounded, write IO impact. Stack depth deep (Flink stacks 80+ frames) — каждое sample expensive. Object allocation profiling on (default.jfc включает -XX:FlightRecorderOptions=defaultrecording=true,jfc=default).

Решение

  1. Use profile.jfc instead of default.jfc — disables allocation profiling: -XX:StartFlightRecording=settings=profile.jfc,duration=5m.
  2. Limit recording duration: 1-5 минут вместо continuous.
  3. Reduce sampling: -XX:FlightRecorderOptions=samplethreads=true,stackdepth=16.
  4. Schedule profiling: nightly cron, не on-demand на live traffic.
  5. Alternative — async-profiler с 1% sample rate (--rate 1000) — обычно overhead < 1%.

Симптомы

  • Region failover (primary down -> switch на secondary). Flink job в secondary стартует из последнего savepoint, но обнаруживает обработанные duplicate events. Downstream metric показывает inconsistency. Data loss ~ N минут (между last checkpoint и failure).

Причина

Checkpoint interval — например, 10 минут. RPO = 10 min minimum (data между checkpoints потеряно). Async checkpoint to DFS (S3) — checkpoint занимает 30s sync + minute async. Если failure between sync и async complete — checkpoint не finalized, restore из предыдущего -> ещё больший gap. Cross-region DFS replication async (S3 cross-region replication ~ minutes) — checkpoint только что committed может не быть на secondary region. Kafka replication lag (MM2) — secondary cluster offset отстаёт от primary.

Решение

  1. Чаще checkpoints: execution.checkpointing.interval=60s (если performance acceptable).
  2. Sync DFS storage в обоих regions: использовать multi-region S3 bucket (например, AWS S3 Multi-Region Access Point) или сheckpoint в две DFS параллельно.
  3. Monitoring RPO: alert если кадр между committed checkpoints > N minutes.
  4. Disaggregated state (ForstDB) на shared DFS — само хранилище уже multi-region replicated.
  5. Application-level idempotency для downstream (key dedup при failover) — для случаев когда RPO потеря неизбежна.

Симптомы

  • Failover на DR-region. Flink job стартует, читает из mirrored Kafka topics. Window operators не emit (watermark не двигается). Метрика currentWatermark на source = старая. Mirror lag (mm2 latency) = 30-60s.

Причина

MirrorMaker 2 / Cluster Linking — replication async. Если в момент failover был lag — secondary topic не имеет events до lag moment. Flink reads из secondary topic, watermark = max-event-time-seen. Если последние events не reached secondary — watermark stuck. Idleness detection на secondary — partition seemingly idle (нет events до тех пор пока mm2 не догонит). Watermark alignment между partitions усугубляет — fastest partition паузится.

Решение

  1. Включить idleness: .withIdleness(Duration.ofMinutes(2)) — partitions без events перестают блокировать.
  2. Topic rename в MM2 (typically prefix cluster_alias.topic) — убедиться Flink reads correct topic name на DR.
  3. Monitoring MM2 lag — alert при > 30s, прежде чем failover критично.
  4. Use Cluster Linking (Confluent) вместо MM2 — лучше lag (typically < 1s) и в native Kafka protocol (Flink не различает).
  5. Buffer downstream при failover (manual circuit breaker) — пауза consumers на минуту чтобы Flink выровнял watermark, потом resume.

Симптомы

  • Enabled pipeline.generic-types=false для cleanup Kryo fallback. Перестроили class на POJO-friendly. Restore из старого savepoint фейлится: SerializerSnapshot indicates incompatibility (KryoSerializerSnapshot vs PojoSerializerSnapshot).

Причина

Kryo serialized state и POJO serialized state — binary incompatible. Schema-evolution flow Flink не покрывает миграцию между serializer families. TypeInfo class changed — Kryo writes type info as part of stream, POJO writes только fields. PojoSerializerSnapshot vs KryoSerializerSnapshot — incompatible через Flink type system.

Решение

  1. State Processor API для migration: read old state с KryoSerializer, write new savepoint с PojoSerializer — single batch job.
  2. Хотя бы один deploy с обоими serializer'ами (legacy fallback) — graceful migration через двойная запись.
  3. Best practice: pipeline.generic-types=false с самого начала проекта — никаких Kryo silent fallback.
  4. Запланировать migration: stop job -> State Processor batch (read Kryo, write POJO) -> start new job -> done.
  5. Tests: roundtrip serializer compatibility — каждая PR проверяет no-Kryo-fallback.

Симптомы

  • Checkpoints начинают falling. Логи: S3 PUT request 503 SlowDown или RequestLimitExceeded. Под нагрузкой много параллельных TM writeput в один bucket prefix. Job restarts из-за tolerable-failed-checkpoints exceeded.

Причина

S3 rate limits на prefix (3500 PUT/sec по docs, в реальности меньше). Все checkpoint files в один prefix (s3://bucket/checkpoint/{job-id}/...) — concurrent writes от N TM hit same prefix. Hot prefix не distributed — S3 партиционирует по hash prefix, single-prefix bottleneck. Slow DFS — каждый file write requires network round-trip.

Решение

  1. S3 prefix entropy: state.checkpoints.create-subdirectories=true -> каждый task subdir.
  2. Request S3 prefix scaling от AWS (free, takes time) или random prefix в path: state.checkpoints.dir=s3://bucket/{random-hash}/checkpoints.
  3. Снизить parallelism upload: state.backend.fs.write-buffer-size=4MB — меньше files (но больше каждый).
  4. Use S3 Express One Zone — single-prefix throughput выше.
  5. Generic incremental checkpoint (changelog) — меньше files per checkpoint.