Guía exhaustiva sobre los watermarks en Apache Flink: Parte 1

Fuentes: Absolutely Everything You Always Wanted to Know About Watermarks in Apache Flink - Part 1: Apache Flink

Los watermarks son uno de los conceptos más discutidos y a la vez más complejos de Apache Flink, la tecnología de referencia para procesamiento de flujos en tiempo real. Este artículo, escrito por un ingeniero de Confluent, ofrece una explicación en profundidad sobre qué son, cómo se generan y cómo se propagan dentro del dataflow de un trabajo Flink, con foco en Flink SQL y Table API y en fuentes que leen desde Apache Kafka.

Los watermarks son señales que Flink inyecta periódicamente en el flujo de registros y que viajan junto con los datos a lo largo del grafo de ejecución. Cada watermark lleva asociado un timestamp en milisegundos desde la epoch y sirve a los operadores basados en tiempo —ventanas temporales, joins temporales, MATCH_RECOGNIZE— para decidir cuándo disparar cálculos, emitir resultados o liberar estado. En el término "watermark" se mezclan dos ideas: el timestamp usado por un operador para su lógica temporal y la señal que propaga ese timestamp entre tareas. Aunque una consulta no use operaciones temporales, los watermarks definidos siguen activos y pueden afectar al consumo, por ejemplo mediante el Watermark Alignment entre particiones.

El texto detalla que cada subtask mantiene su propio tiempo de watermark, calculado como el mínimo entre sus entradas activas, y que los eventos cuyo event-time es anterior o igual al watermark se consideran tardíos y, por defecto, se descartan silenciosamente. También aborda los motivos por los que los watermarks pueden quedarse "atascados", paralizando operaciones temporales, y explica la generación en subtasks de origen: el watermark de cada partición Kafka equivale al event-time máximo observado menos el intervalo de retardo configurado menos un milisegundo.