Flink

Junior

¿Cómo leo una exportación de métricas de Flink, y qué números importan de verdad?

La mayor parte de lo que devuelve la API de métricas de Flink es ruido en una primera pasada. Una lista corta de contadores responde casi cualquier pregunta de «¿está bien este job?».

Leer en:
busyTimeMsPerSecond numRecordsInPerSecond duración de checkpoint
Las métricas con más señal para una revisión de salud inicial - ilustrativo, no medido.

Flink expone métricas a través de una API REST (o JMX, o un reporter a tu elección) como JSON plano: una entrada por nombre de métrica por tarea/operador/job, sin priorización incorporada de cuáles importan.

La forma de una exportación de métricas

Una consulta a /jobs/<id>/vertices/<vertex-id>/metrics devuelve un array de objetos {"id": "...", "value": "..."} - plano, sin anidamiento, sin indicación de cuáles vale la pena mirar primero. Eso depende de que tú lo sepas de antemano, ya que la propia API trata un contador de pausas de GC y tu contador real de throughput como igual de importantes.

El puñado que importa en una primera pasada

busyTimeMsPerSecond - cuánto de cada segundo una tarea estuvo realmente trabajando, lo más parecido que tiene Flink a una métrica de utilización. numRecordsInPerSecond / numRecordsOutPerSecond - throughput real, y una brecha creciente entre «in» y «out» a lo largo del pipeline es un backlog formándose. La duración del checkpoint (del endpoint REST de checkpointing, no de las métricas por tarea) - una duración de checkpoint lenta o creciente suele ser la señal de alerta más temprana de un problema, mucho antes de que el throughput baje visiblemente.

Lo que técnicamente está disponible pero suele ser ruido

Las métricas de heap de JVM y GC son reales y a veces son de verdad la respuesta, pero son una herramienta de segunda pasada - revísalas una vez que busyTimeMsPerSecond o la duración del checkpoint ya te dijeron algo, no como primer paso. Empezar por ahí en cada investigación significa revisar decenas de contadores de JVM antes de llegar a los dos o tres que normalmente explican lo que está pasando.

Es la misma razón por la que las reglas de Flink de opti-pipe solo miran una lista corta y fija de métricas - no porque la exportación no tenga más, sino porque tener más no es más útil pasado cierto punto para responder «¿está bien este job?».

Mira cómo se ve esto en tu propio pipeline.

Sube un event log real de Spark, un run_results.json de dbt, o una exportación de métricas de Flink, y recibe recomendaciones concretas que requieren tu aprobación - no otra regla general.