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

Troubleshooting — Apache Flink

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

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

Симптомы

  • TaskManager Pod рестартится с reason OOMKilled (kubectl describe pod)
  • В логах перед рестартом: 'Container memory limit exceeded' или Java OutOfMemoryError
  • Job уходит в restart loop, checkpoint падает
  • Метрики JVM heap не доходят до limit — kill пришёл от kernel cgroup

Причина

Неправильно поделена памят: process.size != sum(heap + managed + network + jvm-overhead + jvm-metaspace + framework.off-heap) RocksDB неконтролируемо берёт memory вне managed memory (mismatched native libs, large value sizes) Большой block cache в RocksDB и одновременно много column families (state per operator) JVM Direct memory leak через Netty buffers / большой network exclusive buffers per gate Под-tuned framework.off-heap.size в новых релизах 2.x

Решение

  1. Выставить taskmanager.memory.process.size явно (=k8s limit), Flink сам распределит подкомпоненты
  2. Включить state.backend.rocksdb.memory.managed=true — RocksDB будет жить внутри managed memory
  3. Увеличить taskmanager.memory.jvm-overhead.fraction до 0.15-0.20 для RocksDB workload
  4. Снизить parallelism per TM или увеличить taskmanager.memory.process.size
  5. В Flink K8s Operator выставить podTemplate с resources.limits.memory = process.size + headroom

Симптомы

  • В Flink UI: 'Checkpoint expired before completing' или 'Checkpoint X aborted'
  • Duration метрика checkpoint > checkpointing.timeout
  • Sync phase растёт минутами, alignment time у некоторых subtasks > 1min
  • Backpressure highlighted в Web UI у operators перед checkpoint sink

Причина

Backpressure блокирует alignment — barriers застревают в очередях Unbounded state size (нет TTL, нет cleanup) — sync phase медленный Slow checkpoint storage (S3 small file throttling, NFS bottleneck) RocksDB не использует incremental — каждый раз полный snapshot SST Скэшу state на маленький local disk (ephemeral storage заполняется)

Решение

  1. Включить unaligned checkpoints: execution.checkpointing.unaligned.enabled=true для backpressure-heavy джобов
  2. Перейти на RocksDB + incremental: state.backend.incremental=true для big state
  3. Включить state TTL для keyed state, которое не должно жить вечно
  4. Увеличить checkpointing.timeout и checkpointing.interval, но сначала разобраться с root cause
  5. Использовать checkpointing.tolerable-failed-checkpoints для не-фейла job на разовых таймаутах
  6. Для S3: enable multipart upload, увеличить s3.upload.parallelism

Симптомы

  • В Flink Web UI subtasks отмечены красным 'BACKPRESSURED' или 'BUSY'
  • Throughput job ниже ожидаемого, lag в Kafka растёт
  • Checkpoint duration растёт
  • CPU/memory subtasks не загружен — но всё стоит

Причина

Slow sink (синхронный JDBC, slow remote API) Skew по ключу — один subtask получает 80% трафика GC pauses на TaskManager Async I/O с маленьким capacity, lookups блокируют Sort/Aggregate over unbounded state — операция тяжелее с ростом state

Решение

  1. Метрика backPressuredTimeMsPerSecond per subtask: ищи первый красный subtask по графу — он узкое место (но проблема в downstream, который тормозит)
  2. Flink Web UI: backpressure tab, идти по графу от sink к source — узкое место там, где красный сменяется на зелёный
  3. Для skew: pre-aggregate (LocalKeyBy + Reduce) или random-prefix + finalAggregate
  4. Для slow sink: AsyncSink с подходящим concurrency, batch writes
  5. Для GC: переход на RocksDB, G1GC tuning

Симптомы

  • Знаешь, что данные пришли, но они не входят в результат окна
  • Метрика numLateRecordsDropped > 0 (если включена)
  • Result окна оказывается заниженным/неполным

Причина

Watermark уже прошёл event-time события — окно закрыто Allowed lateness не задана или меньше реальной задержки Side output для late events не сконфигурирован, события молча отбрасываются WatermarkStrategy агрессивная (forMonotonousTimestamps вместо forBoundedOutOfOrderness)

Решение

  1. Увеличить bounded out-of-orderness в WatermarkStrategy под P99 реальной задержки
  2. Добавить allowedLateness на окно — окно перевычисляется при поступлении late event
  3. Добавить side output: .sideOutputLateData(lateTag) -> отдельный handler
  4. Pre-monitor: метрики 'numLateRecordsDropped' и 'currentEmitWatermark' vs 'currentInputWatermark'

Симптомы

  • Окна не закрываются, событие-time-timers не срабатывают
  • Метрика currentEmitWatermark стоит на одном значении
  • Один subtask в KafkaSource показывает recordsConsumedRate=0
  • Kafka topic имеет partition без новых сообщений

Причина

В topic есть пустые partition (не используются producer) Минимум watermark по всем upstream subtasks/каналам — если один не двигается, общий не двигается

Решение

  1. WatermarkStrategy.forBoundedOutOfOrderness(...).withIdleness(Duration.ofMinutes(1)) — idle splits исключаются из min watermark
  2. Уменьшить число partition в topic, если они изначально лишние
  3. Watermark Alignment (Flink 1.15+): pipeline.watermark-alignment.max-drift — выравнивает быстрые/медленные partitions
  4. Проверить partitioning стратегию producer: если по бизнес-ключу — это нормально, что некоторые пусты

Симптомы

  • После failover/restart в downstream consumer видны дублирующиеся сообщения
  • Idempotent consumer показывает повторы по message key
  • Checkpoint восстановился, но events с lastCheckpoint..now повторно отправлены

Причина

Sink настроен с DeliveryGuarantee.AT_LEAST_ONCE, не EXACTLY_ONCE EXACTLY_ONCE настроен, но consumer читает с isolation.level=read_uncommitted (default!) transactional.id.prefix не уникален — несколько jobs конфликтуют по transactional id transaction.timeout.ms < checkpointing.interval — broker abort транзакции

Решение

  1. KafkaSink.builder().setDeliveryGuarantee(EXACTLY_ONCE).setTransactionalIdPrefix('app-name-prod-')
  2. На downstream consumer обязательно: isolation.level=read_committed
  3. Уникальный transactionalIdPrefix per job (включить env+team+name)
  4. transaction.timeout.ms > checkpointing.interval + tolerance (на Kafka broker: transaction.max.timeout.ms)
  5. В Flink K8s Operator transactional.id выводить из FlinkDeployment name через env

Симптомы

  • После deploy state stateful operators пустой, метрики state size упали в 0
  • Job стартует, но 'losing' значения, ранее накопленные
  • В логе JobManager: 'Failed to rollback to checkpoint' или skip warnings о operator state
  • Restore с savepoint происходит, но конкретный оператор начинает с нуля

Причина

Кто-то изменил/добавил/удалил .uid() на stateful операторе Авто-generated UID изменился из-за рефакторинга цепочки операторов Operator chain переупорядочен — авто-UID-ы сместились

Решение

  1. Обязательно ставить .uid('stable-name') на ВСЕ stateful операторы — раз и навсегда
  2. При restore с savepoint использовать --allowNonRestoredState ТОЛЬКО осознанно (теряем state)
  3. Использовать State Processor API для миграции state на новый UID если случилось
  4. Запретить deploy без .uid() через линтер/CI-check

Симптомы

  • Низкий throughput даже на простом pipeline (десятки тысяч/сек вместо сотен)
  • В логах WARN: 'Class X cannot be used as a POJO type because [...], must be processed as GenericType'
  • CPU TaskManager-ов на 100% в Kryo serializer
  • Метрика записей/sec не масштабируется с parallelism

Причина

POJO с private fields без getter/setter (или final fields) — Flink не распознаёт, fallback на Kryo Использование record-классов до Flink 1.18 (не было POJO support) Generic Map<String, Object>, Collection без typed parameters Использование java.util.LinkedHashMap, Date — fallback на Kryo

Решение

  1. Использовать чистые POJO: public no-arg constructor, public fields ИЛИ getter/setter, no inheritance с generic
  2. Для structured data — Avro или Protobuf, регистрировать через TypeInformation
  3. Аннотация @TypeInfo с custom TypeInfoFactory для сложных типов
  4. pipeline.generic-types=false — будет throw на любом Kryo fallback (catch в dev)
  5. Метрика serialization time per operator + Flink Profiler

Симптомы

  • Один-два subtasks горят красным (BUSY) в Web UI, остальные idle
  • Метрики numRecordsIn / numRecordsOut неравномерны между subtasks одного оператора (10x+ разрыв)
  • State size одного subtask в разы больше других
  • Backpressure только на этом subtask

Причина

Бизнес skew: один ключ (например customerId=0 или анонимный пользователь) держит 50% трафика Низкая cardinality ключа Bot/spam генерирует большую долю трафика под одним ключом

Решение

  1. Local pre-aggregation: LocalKeyByMiniBatch + Reduce, чтобы пред-агрегировать в operator перед shuffle
  2. Random salt: key = realKey + '_' + (hashOfMicroBatch % N), потом final aggregate уже без салта (two-stage aggregation)
  3. В Flink SQL: table.optimizer.distinct-agg.split.enabled=true для distinct-агрегаций
  4. Отфильтровать паразитные ключи в препроцессинге (если бизнес позволяет)
  5. Для real-time мониторинга: метрика numRecordsInPerSecond per subtask, alert на >3x median

Симптомы

  • Job в статусе CREATED дольше нескольких минут
  • Logs JobManager: 'Could not fulfill resource requirements of job' / 'NoResourceAvailableException'
  • Slot allocation timeout (slot.request.timeout)
  • В K8s: Pod-ы TaskManager в Pending

Причина

Job parallelism > доступных slots (TM * slotsPerTM) В K8s: insufficient cluster resources (cpu/memory) — Pod-ы не шедулятся Slot sharing groups жёстко разделены — нужно больше slots, чем кажется В Native K8s mode: kubernetes.taskmanager.cpu / memory не помещается в node

Решение

  1. Снизить parallelism: env.setParallelism(N) или operator.setParallelism(N)
  2. Увеличить taskmanager.numberOfTaskSlots или число TM (replicas в FlinkDeployment)
  3. В K8s: проверить кластер resources (kubectl describe nodes), HPA/Cluster Autoscaler
  4. В Native K8s mode: правильно выставить kubernetes.taskmanager.cpu соответствующий доступным cores
  5. Использовать adaptive scheduler с reactive mode — job стартует с меньшим parallelism и масштабируется

Симптомы

  • Source стоит на snapshot phase часами, нет прогресса
  • Метрики currentSnapshotProgress не двигаются
  • Backpressure от snapshot reader до downstream — sink/transformer не успевают
  • Лог: 'Snapshot taking too long' warnings

Причина

Таблица в десятки/сотни миллионов строк, snapshot single-threaded по chunk Slow downstream (sink не успевает), backpressure тормозит snapshot reader Incremental Snapshot Framework не включен (default in 3.x, но не во всех connector-ах) Chunk size слишком большой — long-running queries блокируют binlog catch-up В MySQL: snapshot.locking.mode заставляет ждать FLUSH TABLES WITH READ LOCK

Решение

  1. Включить scan.incremental.snapshot.enabled=true (chunk parallel snapshot)
  2. Увеличить scan.incremental.snapshot.chunk.size для меньшего числа chunks (но более крупных)
  3. Увеличить parallelism source: source.setParallelism(N) — chunks разойдутся параллельно
  4. Ускорить downstream sink (batch writes, увеличить parallelism)
  5. Для MySQL: snapshot.locking.mode=none если можно (требует InnoDB и repeatable read)
  6. Разбить миграцию: сначала snapshot only через batch, потом стартовать CDC из binlog offset

Симптомы

  • Job не стартует после изменения схемы исходной БД: 'Cannot deserialize state' / 'Schema not compatible'
  • Savepoint restore падает с serializer mismatch
  • Flink CDC source: 'Failed to deserialize binlog event' после ALTER TABLE

Причина

Type schema state operator-а изменился — Flink не знает, как converted Avro/Protobuf несовместимое изменение схемы (drop required field) Flink CDC 2.x не поддерживал schema evolution (нужен 3.x с pipeline) Sink connector без поддержки schema evolution (например JDBC) — ALTER TABLE в source ломает sink

Решение

  1. Использовать Flink CDC 3.x YAML pipeline с schema evolution-aware sinks (Paimon, StarRocks, Doris)
  2. Avro: обеспечить backward+forward compatibility (default values для new fields, no required deletes)
  3. Schema Registry с compatibility check в CI
  4. State Processor API для миграции state с одной schema на другую
  5. Для JDBC sink: вручную ALTER target table перед deploy

Симптомы

  • После нескольких restart'ов job: 'Too many connections' от PostgreSQL/MySQL
  • DB-сторона: open connections растёт после каждого failover
  • Job в restart loop, downstream DB недоступна
  • JDBC sink логи: 'Could not obtain connection from datasource pool exhausted'

Причина

Sink не закрывает HikariCP/connection pool в close() при cancel Custom RichSinkFunction без правильного teardown Native JDBC connection в operator с per-record commit без batch DB не очищает idle connections достаточно быстро (wait_timeout high)

Решение

  1. Использовать стандартный JdbcSink из flink-connector-jdbc — он управляет pool правильно
  2. В custom sink: реализовать close() с pool.close(), не полагаться на JVM exit
  3. DB-сторона: уменьшить wait_timeout / idle_in_transaction_session_timeout
  4. Использовать pgbouncer/proxysql для connection pool отдельно
  5. Метрика mbean пула: HikariCP.activeConnections — alert на рост во времени

Симптомы

  • После DDL в источнике (ALTER TABLE ADD COLUMN, ALTER COLUMN TYPE) Flink CDC source падает
  • Лог: 'Failed to deserialize event' или 'No registered type for ...'
  • Job в restart loop, не может пройти через DDL event в binlog/WAL

Причина

Flink CDC 2.x не понимает DDL events — падает на schema mismatch Embedded Debezium настроен с include.schema.changes=false, но schema в state кэшируется устаревшая Не используется schema evolution pipeline

Решение

  1. Мигрировать на Flink CDC 3.x с pipeline YAML и поддержкой DDL events
  2. Restart job с savepoint после ручной синхронизации схемы sink-а
  3. Использовать Kafka в середине: Debezium -> Kafka -> Flink (с schema-registry compatibility)
  4. Запретить DDL в production источника без процедуры миграции pipeline

Симптомы

  • В Prometheus нет flink_* метрик от TaskManager/JobManager
  • ServiceMonitor скрапит endpoint, но получает 404 или connection refused
  • Метрики появляются на JM, но нет от TM (или наоборот)

Причина

metrics.reporter.prom.factory.class не настроен в config Не задеплоен flink-metrics-prometheus jar в /opt/flink/lib/ Порт metrics не expose'd в Pod (FlinkDeployment не открыл порт 9249) ServiceMonitor selector не матчит правильные Pod labels

Решение

  1. В flinkConfiguration FlinkDeployment: metrics.reporter.prom.factory.class=org.apache.flink.metrics.prometheus.PrometheusReporterFactory; metrics.reporter.prom.port=9249-9250
  2. Image должен включать flink-metrics-prometheus в plugins/ или lib/ (в official image он в plugins/, нужно ENABLE_BUILT_IN_PLUGINS)
  3. В FlinkDeployment podTemplate добавить containerPort: 9249, metrics service
  4. ServiceMonitor с правильным selector type=flink-native-kubernetes
  5. Проверить вручную: kubectl port-forward TM-pod 9249; curl localhost:9249/metrics

Симптомы

  • В Operator-логе: 'Failed to trigger savepoint' / 'Savepoint operation timed out'
  • FlinkDeployment status.lifecycleState=UPGRADING бесконечно
  • Старый job всё ещё работает, новый не стартует
  • В Flink UI: savepoint endpoint timeouts

Причина

savepoint.timeoutSeconds в FlinkDeployment слишком короткий для big-state job Backpressure в job блокирует aligned savepoint Недоступно savepoint storage (S3 credentials истекли, PVC заполнен) Job в нестабильном состоянии (restart loop) — savepoint trigger не успевает

Решение

  1. Увеличить spec.job.savepointTriggerNonce и job.savepointTimeoutSeconds: 600 (или больше для big state)
  2. Включить unaligned savepoints (env.getCheckpointConfig().enableUnalignedCheckpoints())
  3. Проверить flinkConfiguration.state.savepoints.dir — доступность storage
  4. Если savepoint совсем не идёт — switch на LAST_STATE upgrade mode (теряем подверженные in-flight данные, но job апгрейдится)
  5. Operator 1.14+: support savepoint disposal — старые savepoints удаляются по retention

Симптомы

  • После upgrade с upgradeMode: last-state — stateful operators пустые
  • Метрики state size упали до 0
  • Job стартует, но 'забыл' предыдущие агрегации

Причина

HA не включена — last-state требует HA, иначе нет last checkpoint metadata Checkpoint storage сменился между deploy (старые checkpoint не в новой директории) В перезапуске Operator потерял reference на last checkpoint (operator restart с ConfigMap reset) spec.job.allowNonRestoredState=true случайно включён — теряет state на UID mismatch

Решение

  1. Включить HA: flinkConfiguration.high-availability.type: kubernetes; high-availability.storageDir: file:///data/ha
  2. Постоянный state.checkpoints.dir (S3/GCS/HDFS), не ephemeral PVC
  3. Не менять checkpoint.dir между deploys (это инвалидирует last-state metadata)
  4. Для critical state джобов использовать upgradeMode: savepoint — медленнее, но гарантирует state preservation
  5. Operator logs: проверить 'Found last checkpoint location: ...' перед upgrade

Симптомы

  • Async I/O оператор показывает backpressure 100%
  • Throughput через оператор близок к нулю
  • В Web UI: queue size = capacity, не уменьшается
  • Timeout-метрики: numAsyncResourcesTimedOut растёт

Причина

capacity слишком маленький под latency external system External system медленный (P99 latency > timeout, ResultFuture не complete) Connection pool в AsyncFunction исчерпан Не вызывается ResultFuture.complete() во всех ветках — операции зависают навсегда

Решение

  1. Увеличить capacity: AsyncDataStream.unorderedWait(ds, fn, timeout, TimeUnit.SECONDS, capacity=1000)
  2. Использовать unordered вариант если порядок не важен — выше throughput
  3. Обязательно: try/catch в asyncInvoke с ResultFuture.completeExceptionally(e) для error path
  4. ResultFuture.complete(Collections.emptyList()) в timeout handler (override timeout method)
  5. Профилировать external system latency, scale up или add cache layer

Симптомы

  • currentEmitWatermark не двигается
  • Окна не закрываются, события скапливаются в state
  • Source имеет несколько splits/partitions, но один не присылает данных вовсе

Причина

Один Kafka partition пустой (нет producer трафика) Source V2 split без активности — min(watermark all splits) = очень старый Idleness не настроена на WatermarkStrategy

Решение

  1. WatermarkStrategy.<T>forBoundedOutOfOrderness(Duration.ofSeconds(5)).withIdleness(Duration.ofMinutes(1))
  2. withIdleness ставится на сам source (env.fromSource), не downstream
  3. Уменьшить число пустых partitions если возможно
  4. Watermark Alignment для борьбы со skew между splits

Симптомы

  • State size бесконечно растёт
  • Метрика numRegisteredEventTimeTimers / numRegisteredProcessingTimeTimers растёт линейно
  • Через сутки/недели OOM на TaskManager
  • В heap dump видны TimerHeap большого размера

Причина

В KeyedProcessFunction регистрируется timer, но не удаляется при обновлении значения Создаётся новый timer на каждое событие без cleanup старого State TTL не настроена + timer-driven cleanup тоже отсутствует

Решение

  1. Перед registerEventTimeTimer(newT) делать deleteEventTimeTimer(oldT), где oldT хранится в state
  2. Использовать стандартные windowsi с TTL вместо custom timer logic если возможно
  3. State TTL с cleanupInRocksdbCompactFilter — pulls expired records
  4. Метрика для alert: numRegisteredEventTimeTimers должна быть bounded relative to active keys
  5. Code review: any registerXxxTimer обязан иметь pair deleteXxxTimer когда state очищается

Симптомы

  • В TM-логе: 'No space left on device' / 'Failed to write SST file'
  • Job в restart loop с этой ошибкой
  • kubectl describe pod: DiskPressure
  • Метрика kubelet_volume_stats_available_bytes на ephemeral volume — близко к 0

Причина

RocksDB local working state хранится на ephemeral storage пода (default 10-20Gi) State backend grows unbounded из-за отсутствия TTL Incremental checkpoints не запускают compaction — SST накапливаются В одном TM много column families (много операторов с state), каждый создаёт свои SST

Решение

  1. В FlinkDeployment podTemplate: добавить PersistentVolumeClaim для /opt/flink/data (state.backend.rocksdb.localdir)
  2. Использовать SSD-класса StorageClass — RocksDB чувствителен к IOPS
  3. Включить state TTL на keyed state с разумными лимитами
  4. RocksDB tuning: state.backend.rocksdb.compaction.style=LEVEL для лучшей disk utilization
  5. Регулярно мониторить state.checkpoint.size — если растёт unbounded, бизнес-логика течёт

Симптомы

  • Job постоянно перезапускается с identical exception в TaskManager logs
  • restart-strategy applied (fixed-delay), но проблема воспроизводится сразу
  • Метрика jobs.uptime сбрасывается каждые N секунд

Причина

Bad data record (poison message), который всегда вызывает exception в десериализации/обработке Внешняя зависимость недоступна (DB down, Kafka cluster issue) Class not found / NoClassDefFoundError из-за версионного конфликта shaded libraries Checkpoint corrupted — restore не работает

Решение

  1. Добавить try/catch вокруг processElement, отправлять poison messages в side output / dead letter queue
  2. Для CDC sources: scan.startup.mode=specific-offset чтобы переиграть с safe offset
  3. Restart с предыдущим savepoint (--allowNonRestoredState если меняли operator)
  4. Restart strategy: failure-rate с лимитом — job не зависает в loop, оставляя FAILED state для alert
  5. Хороший error monitoring (Sentry/Bugsnag from UDF) — видеть exception trace мгновенно

Симптомы

  • Job стартует/рестартится 10-30 минут
  • В JM-логах долгое 'Restoring from savepoint' phase
  • Downloading savepoint files из S3 занимает большую часть времени
  • TM-логи: долгий 'Initializing RocksDB' / state download

Причина

Savepoint десятки/сотни GB — десериализация single-threaded на subtask Slow checkpoint storage (S3 region далеко, small parallelism) Не incremental checkpoint — каждый restore = full state download Pre-2.x RocksDB SST в savepoint без incremental — повторный download полного state

Решение

  1. Использовать incremental checkpoints в production (savepoint всё равно full, но обычный restart использует checkpoint, не savepoint)
  2. Увеличить s3.upload/download parallelism: s3.connection.maximum=100, s3.path.style.access settings
  3. Колоцировать checkpoint storage с compute (same region, low-latency)
  4. Native format savepoints (с 1.15+): state.savepoints.format=native — быстрее, меньше overhead
  5. Для очень big state — переход на kafka-source + replay из earliest вместо restore (если бизнес позволяет идемпотентный sink)

Симптомы

  • После изменения parallelism в FlinkDeployment job не стартует
  • В JM-логах: 'NoResourceAvailableException' / 'Slot request bulk timed out'
  • TaskManager Pod-ы Pending в K8s — не помещаются на nodes
  • Operator UPGRADING долго

Причина

Новый parallelism требует больше slots, чем доступно (TM * slotsPerTM < parallelism) Native K8s mode: TaskManager Pod-ы не помещаются (insufficient node resources) Cluster autoscaler не успевает добавить nodes в timeout Resource quota в namespace не разрешает больше Pod-ов

Решение

  1. В FlinkDeployment scale spec.taskManager.replicas согласованно с parallelism
  2. Увеличить slot.request.timeout (default 5min) для медленного autoscaler-а
  3. Заранее scale up cluster nodes перед deploy
  4. Использовать Adaptive Scheduler — job стартует с available slots и подхватывает новые
  5. Проверить ResourceQuota: kubectl describe quota -n flink-namespace

Симптомы

  • Окно открыто, события набираются в state, но result не emit
  • На малом потоке (dev/test) — окна закрываются после долгого простоя или вообще не закрываются
  • В production: некоторые ключи дают результаты, другие нет
  • Метрика state size окна растёт, не сбрасывается

Причина

Watermark не двигается из-за idle source partition / нет новых данных Event-time timestamp в событиях устарел — watermark позади realtime Окно длинное (день+) и watermark двигается медленнее, чем ожидается Custom WindowAssigner или Trigger с багом

Решение

  1. withIdleness в WatermarkStrategy
  2. Для long windows: emit partial results через ContinuousEventTimeTrigger или ContinuousProcessingTimeTrigger
  3. Чёткая модель: знать P99 lag между event time и processing time; настраивать out-of-orderness под него
  4. Метрика currentEmitWatermark должна отставать от текущего processing time на bounded value
  5. Если нужны periodic results — переход на processing-time окна или sliding event-time с эмиссией каждые N секунд

Симптомы

  • Prometheus storage растёт быстро, scrape тормозит
  • Grafana dashboards медленные, queries timeout
  • В Prometheus метриках видны лейблы с UUID, timestamps, transactional IDs
  • Память prometheus постоянно растёт

Причина

metrics.scope.* settings с включёнными high-cardinality scopes (например <task_attempt_id>) Custom UDF метрики с динамическими labels (user_id как label) Sink метрики per-transaction (transactional id меняется каждый checkpoint)

Решение

  1. Чётко спроектировать metrics scope: metrics.scope.operator: <host>.taskmanager.<job_name>.<operator_name> — без attempt_id
  2. В custom UDF: НЕ использовать high-cardinality dimensions как labels; aggregate в самом UDF
  3. metric.reporter.prom.filter.includes / excludes для отфильтровать ненужное
  4. В Prometheus: prometheus.yml metric_relabel_configs для drop проблемных метрик
  5. Использовать VictoriaMetrics или Mimir — лучше переносят high cardinality

Симптомы

  • В Kafka topic добавили partitions, но Flink job их не читает
  • Метрика numActivePartitions не выросла
  • Лаг на новых partitions растёт линейно, никто не консьюмит

Причина

partition.discovery.interval.ms не задан (по умолчанию выключено в Flink 1.x, default 5min в 2.x) Subscription по topic name работает, но enumerator не пере-сканирует partitions Job переподнят с savepoint, который содержит фиксированный список partition assignments

Решение

  1. KafkaSource.builder().setProperty('partition.discovery.interval.ms', '60000') — каждую минуту enumerator проверяет
  2. Restart job с свежим snapshot после изменения партиций (без --allowNonRestoredState)
  3. Использовать pattern-subscription для динамического discovery топиков
  4. Метрика watermarks и lag per partition — alert на discovery