Troubleshooting — Apache Flink
База знаний типичных ошибок курса Apache Flink.
Симптомы
- 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
Решение
- Выставить taskmanager.memory.process.size явно (=k8s limit), Flink сам распределит подкомпоненты
- Включить state.backend.rocksdb.memory.managed=true — RocksDB будет жить внутри managed memory
- Увеличить taskmanager.memory.jvm-overhead.fraction до 0.15-0.20 для RocksDB workload
- Снизить parallelism per TM или увеличить taskmanager.memory.process.size
- В 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 заполняется)
Решение
- Включить unaligned checkpoints: execution.checkpointing.unaligned.enabled=true для backpressure-heavy джобов
- Перейти на RocksDB + incremental: state.backend.incremental=true для big state
- Включить state TTL для keyed state, которое не должно жить вечно
- Увеличить checkpointing.timeout и checkpointing.interval, но сначала разобраться с root cause
- Использовать checkpointing.tolerable-failed-checkpoints для не-фейла job на разовых таймаутах
- Для 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
Решение
- Метрика backPressuredTimeMsPerSecond per subtask: ищи первый красный subtask по графу — он узкое место (но проблема в downstream, который тормозит)
- Flink Web UI: backpressure tab, идти по графу от sink к source — узкое место там, где красный сменяется на зелёный
- Для skew: pre-aggregate (LocalKeyBy + Reduce) или random-prefix + finalAggregate
- Для slow sink: AsyncSink с подходящим concurrency, batch writes
- Для GC: переход на RocksDB, G1GC tuning
Симптомы
- Знаешь, что данные пришли, но они не входят в результат окна
- Метрика numLateRecordsDropped > 0 (если включена)
- Result окна оказывается заниженным/неполным
Причина
Watermark уже прошёл event-time события — окно закрыто Allowed lateness не задана или меньше реальной задержки Side output для late events не сконфигурирован, события молча отбрасываются WatermarkStrategy агрессивная (forMonotonousTimestamps вместо forBoundedOutOfOrderness)
Решение
- Увеличить bounded out-of-orderness в WatermarkStrategy под P99 реальной задержки
- Добавить allowedLateness на окно — окно перевычисляется при поступлении late event
- Добавить side output: .sideOutputLateData(lateTag) -> отдельный handler
- Pre-monitor: метрики 'numLateRecordsDropped' и 'currentEmitWatermark' vs 'currentInputWatermark'
Симптомы
- Окна не закрываются, событие-time-timers не срабатывают
- Метрика currentEmitWatermark стоит на одном значении
- Один subtask в KafkaSource показывает recordsConsumedRate=0
- Kafka topic имеет partition без новых сообщений
Причина
В topic есть пустые partition (не используются producer) Минимум watermark по всем upstream subtasks/каналам — если один не двигается, общий не двигается
Решение
- WatermarkStrategy.forBoundedOutOfOrderness(...).withIdleness(Duration.ofMinutes(1)) — idle splits исключаются из min watermark
- Уменьшить число partition в topic, если они изначально лишние
- Watermark Alignment (Flink 1.15+): pipeline.watermark-alignment.max-drift — выравнивает быстрые/медленные partitions
- Проверить 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 транзакции
Решение
- KafkaSink.builder().setDeliveryGuarantee(EXACTLY_ONCE).setTransactionalIdPrefix('app-name-prod-')
- На downstream consumer обязательно: isolation.level=read_committed
- Уникальный transactionalIdPrefix per job (включить env+team+name)
- transaction.timeout.ms > checkpointing.interval + tolerance (на Kafka broker: transaction.max.timeout.ms)
- В 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-ы сместились
Решение
- Обязательно ставить .uid('stable-name') на ВСЕ stateful операторы — раз и навсегда
- При restore с savepoint использовать --allowNonRestoredState ТОЛЬКО осознанно (теряем state)
- Использовать State Processor API для миграции state на новый UID если случилось
- Запретить 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
Решение
- Использовать чистые POJO: public no-arg constructor, public fields ИЛИ getter/setter, no inheritance с generic
- Для structured data — Avro или Protobuf, регистрировать через TypeInformation
- Аннотация @TypeInfo с custom TypeInfoFactory для сложных типов
- pipeline.generic-types=false — будет throw на любом Kryo fallback (catch в dev)
- Метрика 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 генерирует большую долю трафика под одним ключом
Решение
- Local pre-aggregation: LocalKeyByMiniBatch + Reduce, чтобы пред-агрегировать в operator перед shuffle
- Random salt: key = realKey + '_' + (hashOfMicroBatch % N), потом final aggregate уже без салта (two-stage aggregation)
- В Flink SQL: table.optimizer.distinct-agg.split.enabled=true для distinct-агрегаций
- Отфильтровать паразитные ключи в препроцессинге (если бизнес позволяет)
- Для 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
Решение
- Снизить parallelism: env.setParallelism(N) или operator.setParallelism(N)
- Увеличить taskmanager.numberOfTaskSlots или число TM (replicas в FlinkDeployment)
- В K8s: проверить кластер resources (kubectl describe nodes), HPA/Cluster Autoscaler
- В Native K8s mode: правильно выставить kubernetes.taskmanager.cpu соответствующий доступным cores
- Использовать 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
Решение
- Включить scan.incremental.snapshot.enabled=true (chunk parallel snapshot)
- Увеличить scan.incremental.snapshot.chunk.size для меньшего числа chunks (но более крупных)
- Увеличить parallelism source: source.setParallelism(N) — chunks разойдутся параллельно
- Ускорить downstream sink (batch writes, увеличить parallelism)
- Для MySQL: snapshot.locking.mode=none если можно (требует InnoDB и repeatable read)
- Разбить миграцию: сначала 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
Решение
- Использовать Flink CDC 3.x YAML pipeline с schema evolution-aware sinks (Paimon, StarRocks, Doris)
- Avro: обеспечить backward+forward compatibility (default values для new fields, no required deletes)
- Schema Registry с compatibility check в CI
- State Processor API для миграции state с одной schema на другую
- Для 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)
Решение
- Использовать стандартный JdbcSink из flink-connector-jdbc — он управляет pool правильно
- В custom sink: реализовать close() с pool.close(), не полагаться на JVM exit
- DB-сторона: уменьшить wait_timeout / idle_in_transaction_session_timeout
- Использовать pgbouncer/proxysql для connection pool отдельно
- Метрика 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
Решение
- Мигрировать на Flink CDC 3.x с pipeline YAML и поддержкой DDL events
- Restart job с savepoint после ручной синхронизации схемы sink-а
- Использовать Kafka в середине: Debezium -> Kafka -> Flink (с schema-registry compatibility)
- Запретить 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
Решение
- В flinkConfiguration FlinkDeployment: metrics.reporter.prom.factory.class=org.apache.flink.metrics.prometheus.PrometheusReporterFactory; metrics.reporter.prom.port=9249-9250
- Image должен включать flink-metrics-prometheus в plugins/ или lib/ (в official image он в plugins/, нужно ENABLE_BUILT_IN_PLUGINS)
- В FlinkDeployment podTemplate добавить containerPort: 9249, metrics service
- ServiceMonitor с правильным selector type=flink-native-kubernetes
- Проверить вручную: 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 не успевает
Решение
- Увеличить spec.job.savepointTriggerNonce и job.savepointTimeoutSeconds: 600 (или больше для big state)
- Включить unaligned savepoints (env.getCheckpointConfig().enableUnalignedCheckpoints())
- Проверить flinkConfiguration.state.savepoints.dir — доступность storage
- Если savepoint совсем не идёт — switch на LAST_STATE upgrade mode (теряем подверженные in-flight данные, но job апгрейдится)
- 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
Решение
- Включить HA: flinkConfiguration.high-availability.type: kubernetes; high-availability.storageDir: file:///data/ha
- Постоянный state.checkpoints.dir (S3/GCS/HDFS), не ephemeral PVC
- Не менять checkpoint.dir между deploys (это инвалидирует last-state metadata)
- Для critical state джобов использовать upgradeMode: savepoint — медленнее, но гарантирует state preservation
- 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() во всех ветках — операции зависают навсегда
Решение
- Увеличить capacity: AsyncDataStream.unorderedWait(ds, fn, timeout, TimeUnit.SECONDS, capacity=1000)
- Использовать unordered вариант если порядок не важен — выше throughput
- Обязательно: try/catch в asyncInvoke с ResultFuture.completeExceptionally(e) для error path
- ResultFuture.complete(Collections.emptyList()) в timeout handler (override timeout method)
- Профилировать 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
Решение
- WatermarkStrategy.<T>forBoundedOutOfOrderness(Duration.ofSeconds(5)).withIdleness(Duration.ofMinutes(1))
- withIdleness ставится на сам source (env.fromSource), не downstream
- Уменьшить число пустых partitions если возможно
- Watermark Alignment для борьбы со skew между splits
Симптомы
- State size бесконечно растёт
- Метрика numRegisteredEventTimeTimers / numRegisteredProcessingTimeTimers растёт линейно
- Через сутки/недели OOM на TaskManager
- В heap dump видны TimerHeap большого размера
Причина
В KeyedProcessFunction регистрируется timer, но не удаляется при обновлении значения Создаётся новый timer на каждое событие без cleanup старого State TTL не настроена + timer-driven cleanup тоже отсутствует
Решение
- Перед registerEventTimeTimer(newT) делать deleteEventTimeTimer(oldT), где oldT хранится в state
- Использовать стандартные windowsi с TTL вместо custom timer logic если возможно
- State TTL с cleanupInRocksdbCompactFilter — pulls expired records
- Метрика для alert: numRegisteredEventTimeTimers должна быть bounded relative to active keys
- 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
Решение
- В FlinkDeployment podTemplate: добавить PersistentVolumeClaim для /opt/flink/data (state.backend.rocksdb.localdir)
- Использовать SSD-класса StorageClass — RocksDB чувствителен к IOPS
- Включить state TTL на keyed state с разумными лимитами
- RocksDB tuning: state.backend.rocksdb.compaction.style=LEVEL для лучшей disk utilization
- Регулярно мониторить 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 не работает
Решение
- Добавить try/catch вокруг processElement, отправлять poison messages в side output / dead letter queue
- Для CDC sources: scan.startup.mode=specific-offset чтобы переиграть с safe offset
- Restart с предыдущим savepoint (--allowNonRestoredState если меняли operator)
- Restart strategy: failure-rate с лимитом — job не зависает в loop, оставляя FAILED state для alert
- Хороший 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
Решение
- Использовать incremental checkpoints в production (savepoint всё равно full, но обычный restart использует checkpoint, не savepoint)
- Увеличить s3.upload/download parallelism: s3.connection.maximum=100, s3.path.style.access settings
- Колоцировать checkpoint storage с compute (same region, low-latency)
- Native format savepoints (с 1.15+): state.savepoints.format=native — быстрее, меньше overhead
- Для очень 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-ов
Решение
- В FlinkDeployment scale spec.taskManager.replicas согласованно с parallelism
- Увеличить slot.request.timeout (default 5min) для медленного autoscaler-а
- Заранее scale up cluster nodes перед deploy
- Использовать Adaptive Scheduler — job стартует с available slots и подхватывает новые
- Проверить 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 с багом
Решение
- withIdleness в WatermarkStrategy
- Для long windows: emit partial results через ContinuousEventTimeTrigger или ContinuousProcessingTimeTrigger
- Чёткая модель: знать P99 lag между event time и processing time; настраивать out-of-orderness под него
- Метрика currentEmitWatermark должна отставать от текущего processing time на bounded value
- Если нужны 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)
Решение
- Чётко спроектировать metrics scope: metrics.scope.operator: <host>.taskmanager.<job_name>.<operator_name> — без attempt_id
- В custom UDF: НЕ использовать high-cardinality dimensions как labels; aggregate в самом UDF
- metric.reporter.prom.filter.includes / excludes для отфильтровать ненужное
- В Prometheus: prometheus.yml metric_relabel_configs для drop проблемных метрик
- Использовать 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
Решение
- KafkaSource.builder().setProperty('partition.discovery.interval.ms', '60000') — каждую минуту enumerator проверяет
- Restart job с свежим snapshot после изменения партиций (без --allowNonRestoredState)
- Использовать pattern-subscription для динамического discovery топиков
- Метрика watermarks и lag per partition — alert на discovery