Programación y desarrollo de software

Cómo construir una canalización de Kafka que preserve el orden de las sesiones usando Go

Joshua Oluikpe explica un diseño práctico para procesar mensajes de Kafka en orden estricto dentro de cada sesión, manteniendo al mismo tiempo la capacidad de las sesiones independientes para trabajar en paralelo. El diseño se basa en el hashing consistente, un trabajador de Go por sesión, reintentos dentro de la misma ruta y la gestión de puntos de compromiso contiguos para evitar la pérdida de mensajes tras un fallo.

2026-10-07
6 min de lectura
1 visitas
certi.news Editorial Team
Cómo construir una canalización de Kafka que preserve el orden de las sesiones usando Go

No basta con depender del orden de Kafka dentro de la partición (partition) cuando miles de sesiones comparten la misma partición. En sistemas de conversación y tareas secuenciales, un mensaje correctivo como «establece el destino en París» puede llegar después de un mensaje como «reserva un viaje a Londres», pero ejecutarlo antes del primer mensaje arruina el contexto. La solución que Joshua Oluikpe describe en InfoQ traslada parte de la responsabilidad del orden a la capa de aplicación escrita en Go.

Orden dentro de la sesión y paralelismo entre sesiones

El diseño se basa en dos niveles de trabajadores. El nivel de distribución recibe los registros de Kafka y los dirige según el identificador de sesión mediante hashing consistente, de modo que los mensajes de una misma sesión llegan al mismo trabajador, mientras que las sesiones diferentes pueden distribuirse entre varios trabajadores. Después, cada sesión activa obtiene una goroutine dedicada que procesa un mensaje cada vez y en el orden de llegada.

Así, una sesión lenta no se convierte en un obstáculo para otras sesiones, como ocurre en un grupo de trabajadores plano cuando el reintento de una sola sesión ocupa por completo al trabajador. Las goroutines se crean cuando son necesarias y se eliminan después de un periodo de inactividad, por lo que el consumo de recursos queda vinculado al número de sesiones activas en lugar del número total de sesiones posibles.

Reintentos sin adelantar mensajes

El reintento se realiza dentro de la goroutine de la propia sesión. Cuando se produce un error temporal, como el agotamiento del tiempo de espera de un servicio dependiente, el trabajador espera usando un retroceso exponencial con cierto grado de aleatoriedad y después vuelve a procesar el mismo mensaje antes de pasar a los mensajes posteriores. Como los mensajes siguientes permanecen en el canal detrás del mensaje fallido, el cuarto mensaje no puede adelantar al tercero.

Los errores irrecuperables, como los datos no válidos, se envían directamente a la lista de mensajes muertos (DLQ). Esta división entre errores temporales y finales elimina la necesidad de una máquina de estados independiente para gestionar los reintentos, pero también significa que los mensajes de una sesión concreta pueden acumularse durante el periodo de retroceso.

Compromiso seguro y recuperación tras un fallo

Cuando varias sesiones procesan sus mensajes en paralelo, no se puede comprometer únicamente el offset completado más alto. Si el offset 104 se completa antes que el 102 y después se compromete el 104, un fallo del consumidor podría hacer que el 102 se omitiera definitivamente. Por ello, el sistema realiza un seguimiento, para cada partición, de los offsets en proceso y de los que se han completado, y solo mueve el punto de compromiso hasta el offset completado contiguo más alto.

Como consecuencia, un reinicio puede comenzar desde un offset anterior y volver a procesar algunos mensajes. El diseño utiliza un identificador de evento estable para ayudar a eliminar duplicados, pero no afirma lograr una ejecución «exactamente una vez»; los efectos externos siguen dependiendo de la capacidad de los servicios dependientes para gestionar los duplicados de forma idempotente. En casos excepcionales, una brecha bloqueada puede omitirse después de que transcurra un tiempo de espera definido para recuperar el avance de la partición, registrándolo como un fallo operativo y dirigiendo el mensaje a la DLQ: una compensación explícita entre la continuidad y la integridad del procesamiento.

¿Qué se necesita para la producción real?

Las garantías de orden no se completan únicamente con la lógica básica. Durante la redistribución de las particiones, el consumidor deja de recibir nuevos registros, espera a que finalice el trabajo en curso y después compromete el último watermark contiguo antes de transferir la propiedad. Cuando los canales de los trabajadores se llenan, detiene temporalmente la obtención de mensajes en lugar de descartarlos, y almacena los registros extraídos previamente en almacenes específicos de cada partición para que los mensajes más recientes no los adelanten.

El sistema también necesita detectar individualmente los offsets bloqueados, mantener la trazabilidad a través de los límites asíncronos de Kafka, conservar las señales de retraso necesarias para el escalado automático durante la detención de la obtención y drenar el trabajo y escribir en la DLQ antes del cierre. El producto también aplica hashing consistente en el lado del envío para garantizar el orden de los mensajes durante la producción asíncrona.

Resultados y limitaciones

En una prueba sintética que incluyó 50,000 mensajes y 10 sesiones, con inyección de errores temporales en el 10% de los mensajes, el sistema registró 14,027 mensajes por segundo sin infracciones de orden, con una latencia mediana de 11.7 milisegundos. En una prueba de carga con 1,000 sesiones, la tasa de envío efectiva alcanzó 48,805 mensajes por segundo frente a un objetivo de 50,000, sin errores de envío, pero la latencia de cola superó el presupuesto de 200 milisegundos debido a la cuota de tasa de mensajes de API Gateway, no al mecanismo de orden de las sesiones.

El autor afirma que el sistema procesó más de 40 millones de mensajes de producción sin infracciones de orden observables, y que solo 460 mensajes se dirigieron a la DLQ, aproximadamente el 0.002%. Sin embargo, esta cifra se basa en la supervisión operativa y no en una comprobación exhaustiva de la secuencia de cada mensaje; además, la prueba de carga superior no verificó de forma independiente el orden de cada sesión. Por ello, el valor práctico del diseño consiste en convertir el orden de una suposición a nivel de Kafka en una propiedad impuesta por la aplicación, manteniendo al mismo tiempo la necesidad de realizar pruebas explícitas de secuencia y de contar con garantías de idempotencia en los servicios dependientes.

Fuente de la noticia
InfoQ - Architecture Articles
Abrir fuente original ↗
c
Autor

certi.news Editorial Team

De la misma categoría

También te puede interesar

Ver todas las noticias