- En la verificación de Bufstream 0.1.0~0.1.3, un sistema de streaming compatible con Kafka, se encontraron 2 problemas de disponibilidad y 3 de seguridad en Bufstream; los 5 quedaron corregidos en la versión 0.1.3
- Las pruebas se basaron en Java Kafka Client 3.8.0 y en pruebas Jepsen existentes para Kafka/Redpanda, y usaron configuraciones orientadas a la seguridad como
acks = all,enable.idempotence = true,enable.auto.commit = falseyread_committed - Los problemas de Bufstream incluyeron detención de consumidores y productores, respuesta incorrecta con offset
0, pérdida de commits de transacciones y pérdida de escrituras reconocidas por un bug en el filtrado del tamaño de respuesta del API fetch - Durante la investigación también salieron a la luz problemas en el cliente Java de Kafka y en el protocolo transaccional de Kafka, como bloqueo indefinido de
Consumer.close(), offsets de consumidor impredecibles y problemas de aborted read·lost write·torn transaction - Jepsen considera que, como el protocolo transaccional de Kafka no garantiza explícitamente el orden de las solicitudes del cliente ni los números de transacción, la seguridad transaccional de Kafka y de sistemas compatibles con Kafka puede romperse al usar el cliente Java oficial
Arquitectura de Bufstream y alcance de la verificación
- Kafka es un sistema de streaming que ofrece logs append-only replicados y fragmentados; Bufstream es una implementación alternativa a Kafka que prioriza la gobernanza de datos y la eficiencia de costos en entornos de nube
- Bufstream ofrece topics y partitions como Kafka, y funciona con clientes estándar de Kafka
- el producer agrega records con
producer.send() - el consumer se vincula a una partition con
consumer.assign()oconsumer.subscribe(), y luego lee records conconsumer.poll() - los consumer groups se reparten el procesamiento de records de un conjunto de topics
- el producer agrega records con
- Al integrarse con Buf Schema Registry, puede inspeccionar records de Protocol Buffers para ofrecer validación de records, control de acceso a nivel de campo y conversión de formatos de datos con otros sistemas
- A diferencia de Kafka, que usa disco local y su propio protocolo de replicación, Bufstream escribe los datos directamente en object storage
- busca reducir costos aprovechando la estructura de costos del tráfico de replicación del object storage
- los nodos de Bufstream pueden funcionar como VM stateless con autoescalado
- Bufstream está compuesto por tres subsistemas
- agent: servicio stateless que expone la API de Kafka
- object store: almacena chunks de records y los entrega a los lectores
- coordination service: actualmente usa etcd y define qué chunk fue committed y el orden de los records
- A octubre de 2024, Bufstream solo se había desplegado para algunos clientes, y aunque la documentación lo presentaba como un “drop-in replacement de Apache Kafka” y destacaba compatibilidad con transacciones de Kafka y exactly-once semantics, no hacía muchas afirmaciones concretas sobre seguridad
Configuración del cliente y supuestos transaccionales
- Jepsen ajustó la configuración del cliente para obtener un comportamiento más seguro, como en pruebas previas de sistemas compatibles con Kafka
-
Configuración del producer
- se usa el valor predeterminado
acks = all - en Bufstream,
acks = 0puede reconocer escrituras sin esperar al storage, lo que puede hacer que se pierdan escrituras ya committed acks = 1yacks = allbloquean hasta que Bufstream confirma persistencia durable- para evitar appends duplicados por los reintentos automáticos del producer de Kafka, se usó el valor predeterminado
enable.idempotence = true
- se usa el valor predeterminado
-
Configuración del consumer
- como la documentación indica que el auto-commit puede llevar a pérdida de datos, en general se usa
enable.auto.commit = false - cuando no hay offset committed, el valor predeterminado de
auto.offset.resetempieza desde el offset más reciente, por lo que no garantiza entrega at-least-once - se usa
auto.offset.reset = earliestpara que el consumer pueda observar todo el log - una transacción de Kafka está compuesta por el conjunto de records enviados por el producer y un mapa de offsets máximos por partition que el consumer consultó con poll
- solo cuando la transacción hace commit, los records enviados se vuelven durables y eventualmente visibles para un consumer
read_committed, y los offsets committed también avanzan al menos hasta el offset especificado en la transacción - si la transacción no hace commit, los offsets committed no avanzan, y la visibilidad de las escrituras puede variar según la configuración del consumer
- cuando un consumer
read_uncommittedlee valores de una transacción abortada, eso se clasifica como aborted read(G1a) - la documentación de Kafka sostiene que
read_committedevita G1a y en cierta medida garantiza la propiedad de que o bien todas las escrituras de una transacción son visibles, o no lo es ninguna; sin embargo, en las pruebas de Jepsen sobre Kafka, Redpanda y Bufstream se observaron write cycle (fenómeno similar a G0) y algunas formas de G1c
- como la documentación indica que el auto-commit puede llevar a pérdida de datos, en general se usa
Diseño de las pruebas
- Jepsen probó Bufstream desde la versión 0.1.0 hasta la 0.1.3, además de varios builds release candidate
- El test harness usó Bufstream test harness, Jepsen testing library y Java Kafka Client 3.8.0
-
Entorno de ejecución
- se usaron entre 3 y 5 nodos Debian Bookworm tanto en contenedores LXC como en VM de EC2
- se usó 1 nodo para etcd, 1 nodo para Minio y el resto como agentes de Bufstream
- el producer, consumer y admin client se inicializaron poniendo solo un nodo en
bootstrap_servers, pero no se bloqueó el smart client discovery
-
Principales configuraciones de seguridad
- auto-commit false
acks = all- reintentos 1,000
- idempotence enabled
- isolation level
read_committed auto_offset_reset = earliest- creación automática de topics del lado del servidor deshabilitada
- la inyección de fallas incluyó pausa de proceso (
SIGSTOP), crash (SIGKILL), clock skew (clock_settime) y partición de red (iptables) - como Bufstream está dividido en agent, object store y coordination service, se creó una nueva herramienta de Jepsen que permite inyectar fallas dirigidas solo a subsistemas específicos
- por ejemplo, se iban cambiando con el tiempo combinaciones como hacer crash solo a nodos de Bufstream o pausar solo al coordinador etcd
Workload de cola y workload de abort
- El workload de cola analiza la seguridad de acuerdo con el modelo de datos de Kafka
- Cada proceso lógico ejecuta un producer, un consumer y un cliente de administración
- Una clave numérica identifica un topic-partition específico
- Las claves se eligen con frecuencia exponencial, de modo que algunas se acceden con frecuencia y otras rara vez
- Se usan tres operaciones básicas
crash: finaliza un proceso lógico y lo reemplaza con un cliente nuevosubscribeoassign: cambia el conjunto de topics o partitions que el consumer harápolltxn,poll,send: ejecuta una secuencia de microoperacionespollosend
- En el workload no transaccional, cada
sendopollincluye exactamente una microoperación - En el workload transaccional, varias microoperaciones se envuelven en una transacción de Kafka
- El análisis construye un mapeo de offset a valor por clave y luego busca errores
- Si aparecen varios valores en el mismo offset, hay un offset inconsistente
- Si el mismo valor aparece en varios offsets, hay un error de duplicado
- Si un registro reconocido no se observa en absoluto, se considera lost o unseen
- Si
polldevuelve un valor enviado por una operación abortada, hay una lectura abortada - También se verifica si una transacción observa sus propias escrituras
- Después de la prueba principal, se resuelven las fallas y se pasa a la etapa de lecturas finales
- Cada proceso lee todos los topic-partition desde el offset 0 y hace
pollhasta el mayor written offset conocido - Si las lecturas finales agotan el tiempo y un registro reconocido sigue sin observarse, se clasifica como unseen
- Cada proceso lee todos los topic-partition desde el offset 0 y hace
- El workload de abort se agregó para rastrear el comportamiento del offset de
polldespués de abortar una transacción- El topic se limita a una sola partition, proceso, producer y consumer
- Después de que una transacción hace
pollde un registro, aborta intencionalmente, y luego el offset depollposterior se clasifica como advance, rewind, rewind-further u other
5 problemas encontrados en Bufstream
-
Consumidores atascados (#1)
- Desde la 0.1.0 hasta la 0.1.3-rc.8, la etapa final de lectura se detenía con frecuencia
consumer.poll()devolvía de inmediato un resultado vacío, pero en el log todavía quedaban miles de records confirmados- Este estado duraba desde decenas de segundos hasta más de 1 hora
- En una prueba, durante los primeros 120 segundos se enviaron 691 records confirmados, y al comenzar las lecturas finales 40 no fueron observados por ningún poller
- Después,
consumer.poll()no devolvió resultados durante más de 1 hora y la prueba terminó por timeout - La causa fue que un nodo de Bufstream reiniciado podía devolver un valor en caché obsoleto del last stable offset y del high watermark
- Algunas bibliotecas cliente concluían que no había records posteriores y quedaban en stall; Bufstream aplicó un parche en la 0.1.3-rc.6 para refrescar la caché al iniciar
-
Productores y consumidores atascados (#2)
- Incluso en la 0.1.3-rc.6, siguieron observándose problemas de escrituras no vistas después de pausas, crashes y particiones del coordinator, storage y nodos de Bufstream
- En algunos casos, tras una pausa del coordinator, el cliente entraba en un estado en el que esperaba
InitProducerIdhasta agotar el tiempo, aunque todos los nodos de Bufstream estuvieran en ejecución - En otros casos,
listOffsetsfallaba connode ... being disconnectedotimed out waiting for a node assignment, ypollse completaba pero no devolvía resultados - Al matar y reiniciar el nodo de Bufstream, el problema se resolvía
- La causa estaba relacionada con leases de etcd
- El agente de Bufstream usa etcd leases para rastrear agentes activos
- Debido a pausas o particiones breves, etcd eliminaba las claves ligadas al lease del agente, pero la actualización de borrado podía no llegar al agente
- El agente quedaba sin saber que había perdido su lease
- El equipo de Bufstream agregó lógica adicional de polling, y en la 0.1.3-rc.8 el problema de escrituras no vistas quedó mayormente resuelto
-
Offsets cero espurios (#3)
- Desde la 0.1.0 hasta la 0.1.3-rc.2, un valor enviado podía recibir el offset
0y luego aparecer en un offset real más alto - Esto ocurría incluso si el offset
0ya había sido asignado mucho antes - Solo el sender observaba el offset 0; el poller observaba un offset más alto
- En una prueba de 2 minutos con un solo nodo de Bufstream y una pausa del proceso etcd, 6 escrituras recibieron el offset
0y luego aparecieron en offsets más altos - La causa fue que faltaba un campo necesario en la respuesta de error de Bufstream
- Bufstream enviaba una solicitud a etcd para confirmar el commit del log y etcd la procesaba, pero por una pausa o partición Bufstream podía agotar el tiempo esperando la respuesta
- Bufstream enviaba un código de error al cliente, pero no establecía el offset del record enviado en la señal de error
-1 - El cliente Java de Kafka interpretaba esto como una respuesta exitosa con offset
0 - Franz-go, usado por la suite de pruebas de Bufstream, interpretaba este message como error, por lo que el problema no salió a la luz en las pruebas
- Bufstream lo corrigió en la 0.1.3-rc.6, y Jepsen no ha vuelto a observarlo desde entonces
- Desde la 0.1.0 hasta la 0.1.3-rc.2, un valor enviado podía recibir el offset
-
Pérdida de escrituras en transacciones (#4)
- En la 0.1.2, ocurría con frecuencia una pérdida de escrituras en la que algunos records de transacciones confirmadas desaparecían y no volvían a observarse
- En una prueba, durante 100 segundos y a lo largo de 6,761 transacciones de escritura, se perdieron 240 records escritos por transacciones confirmadas
- En el ejemplo, el valor
141de la key5fue devuelto como escrito correctamente en el offset274, pero todos losconsumer.poll()saltaron ese offset - La causa fue un bug en el mecanismo de seguridad de concurrencia agregado en la 0.1.2
- Este mecanismo asigna un número único a cada transacción dentro del epoch del productor para mitigar la falta de idempotencia del protocolo de transacciones de Kafka
- Debido a un bug en la lógica de seguimiento del número de transacción, cuando se confirmaban varias transacciones a través de múltiples epochs, algunos commits se ignoraban incorrectamente
- Una transacción que parecía haberse confirmado podía en realidad abortarse, o viceversa
- Jepsen descubrió este bug gracias a que configuró el timeout de transacción en 1 segundo
- Bufstream identificó el problema a las pocas horas del release 0.1.2, impidió que los clientes hicieran upgrade, y los clientes no hicieron upgrade a la 0.1.2
- La corrección se incluyó en la 0.1.3-rc2
-
Pérdida de escrituras por el filtrado del lado del servidor (#5)
- En la 0.1.3-rc.8, después de fallas pequeñas como pausas del proceso de Bufstream o del coordinator, o particiones entre ambos, aparecía con frecuencia una ventana breve de pérdida de escrituras
- La pérdida de datos ocurría independientemente de si se usaban transacciones o no
- En una prueba de 5 minutos, de 16,770 records, 22 fueron reconocidos pero ningún consumidor pudo hacerles poll
- Algunos records eran visibles para los pollers durante un tiempo, pero más tarde desaparecían de
poll - La causa fue la lógica de límite de tamaño de respuesta de la API fetch agregada en la 0.1.3-rc.8 para rodear un bug de una popular interfaz web de Kafka
- Un bug en la lógica de filtrado ocultaba records a consumidores con retraso, haciendo que pareciera una pérdida de escrituras
- Bufstream lo corrigió en la 0.1.3-rc.12
Problemas del cliente Java de Kafka y del protocolo de Kafka
-
KIP-588:
ProducerFencedExceptioninduce a confusión- Durante las pruebas, apareció con frecuencia el error
ProducerFencedException: There is a newer producer with the same transactionalId which fences the current one. - Este error apareció incluso en pruebas donde todos los producers recibían un transactional ID único, por lo que tomó tiempo identificar la causa
- KIP-588 señala que también puede lanzarse
ProducerFencedExceptionen caso de transaction timeout - El cliente Java de Kafka usa un
TimeoutExceptionespecífico para la mayoría de los timeouts, pero en este caso lanzaProducerFencedException - Aunque en realidad no existía un producer en conflicto, el mensaje de error decía que había una segunda instancia de producer
- KIP-588 lleva dos años abierto, y Jepsen recomienda que el equipo de Kafka cambie el mensaje de error
- Durante las pruebas, apareció con frecuencia el error
-
KAFKA-17734:
Consumer.close()puede bloquearse indefinidamente- Tanto en las pruebas de Bufstream como en las de Kafka, las pruebas se detenían cada pocas horas debido a un bug del cliente Java
Consumer.close()se bloquea en network IO de forma predeterminada- El parámetro de timeout de
close()debería evitar el bloqueo indefinido, pero no funcionó - Tampoco resultó efectivo llamar a
consumer.wakeup()desde un thread separado para interrumpir al consumer atascado en IO - Jepsen considera que los programas de larga ejecución deben poder liberar recursos como clients, connections, threads y memoria en un tiempo razonable incluso ante errores de red, y registró KAFKA-17734
-
KAFKA-17582: offset del consumer impredecible tras fallo de transaction
- La documentación oficial de Kafka casi no explica cómo debería quedar el consumer offset cuando falla un transaction commit
- La documentación de diseño de Kafka de Confluent dice que, si una transaction se aborta, la posición del consumer vuelve al valor anterior, pero el cliente Java real no siempre se comporta así
- En los resultados del workload Abort, incluso en un clúster sano, el comportamiento después del abort fue difícil de predecir
- La mayoría de los pares de transaction avanzaban hacia un offset posterior
- Algunos retrocedían a un offset anterior
- Todos los retrocesos estuvieron relacionados con eventos de rebalance, y en todos los avances no hubo rebalance
- Según la respuesta del lado de Kafka, este comportamiento es intencional
- El consumer sigue avanzando
- Si ocurre un rebalance, puede retroceder a un punto arbitrario según el committed offset
- El usuario debe hacer rewind manual de la posición del consumer cuando una transaction se aborta
- Jepsen abrió KAFKA-17582 y propuso documentar este comportamiento y considerar cambiar el comportamiento predeterminado a rewind en caso de transaction abort
- El workload Queue también fue modificado para hacer rewind explícito del consumer
-
KAFKA-17754: pérdida de escritura, lectura abortada, transaction desgarrada
- En Bufstream 0.1.0~0.1.3 se observaron lecturas abortadas, escrituras perdidas y violaciones de atomicidad solo con pause del proceso de Bufstream, pause del coordinator, crash y network partition
- El análisis condujo a un defecto fundamental en el protocolo de transactions de Kafka
- En el ejemplo, el client ejecutó una transaction con el transactional ID único
jt1234y enviócommitted = falseenEndTxnpara abortarla, pero 15 llamadas apoll()observaron escrituras de la transaction abortada - Otras escrituras de esa misma transaction no fueron observadas por ningún poller
- Al revisar juntos la captura de paquetes y el log de Bufstream, la causa fue un mensaje de commit retrasado
- Un
EndTxnde commit enviado varias transactions antes fue procesado tarde en un node - El client ya había avanzado a las siguientes transactions
- El commit retrasado se aplicó a la transaction actual, por lo que solo se committed la parte inicial de la transaction, mientras que el resto fue tratado como una transaction separada y abortado
- El protocolo de Kafka está diseñado para que el client pueda enviar requests por varias conexiones TCP y a varios nodes, pero no tiene números de secuencia que determinen el orden de las requests de un mismo client
- Tampoco existe un concepto de número de transaction, por lo que cuando el server recibe un mensaje de commit o abort no puede saber qué transaction intentaba terminar el client
- Como resultado, pueden darse las siguientes situaciones
- una transaction que parecía committed en realidad se aborta
- una transaction abortada en realidad se committed
- solo se conservan algunas escrituras de una transaction y otras se pierden, produciendo una transaction desgarrada
- El cliente oficial Java de Kafka trata los timeouts como retryable y puede enviar automáticamente varios mensajes
EndTxn, por lo que el problema puede ocurrir aunque el usuario llame a commit o abort solo una vez por transaction - Jepsen también observó en Kafka lecturas abortadas y transactions desgarradas provocadas por pause del proceso, y abrió KAFKA-17754
- Los ingenieros de Kafka consideran que KIP-890 podría corregir este problema
- KIP-890 cambia el protocolo de transactions elevando el producer epoch en cada transaction
- Como el server rechaza mensajes de epochs anteriores, puede evitar que mensajes de commit de transactions pasadas se filtren hacia transactions posteriores
- En la versión 0.1.3, Bufstream agregó un mecanismo para reducir la frecuencia usando la revisión de etcd como reloj lógico, pero eso no impide el reordenamiento entre el client y Bufstream
- Jepsen siguió observando lecturas abortadas, escrituras perdidas y transactions desgarradas también en la 0.1.3, y considera que se necesita una solución del lado del client
Resumen general de resultados
- Los 5 problemas propios de Bufstream ya fueron corregidos
- #1: el
consumerse quedaba bloqueado por unhighest stable offsetrezagado, no requería falla, corregido en 0.1.3-rc.6 - #2: el
producer/consumerse quedaba bloqueado por el vencimiento de un lease de etcd, requería pausa, corregido en 0.1.3-rc.8 - #3: offsets cero espurios, requería pausa, corregido en 0.1.3-rc.6
- #4: pérdida de escrituras de transacciones, no requería falla, corregido en 0.1.3-rc.2
- #5: pérdida de escrituras por filtrado del lado del servidor, requería pausa, corregido en 0.1.3-rc.12
- #1: el
- Los problemas relacionados con Kafka siguen pendientes
- KIP-588: mensaje de error incorrecto en
transaction timeout, sin resolver - KAFKA-17734:
ConsumerClient.close()puede bloquearse indefinidamente, sin resolver - KAFKA-17582: tras fallar una transacción, el offset del
consumerse vuelve impredecible, sin resolver - KAFKA-17754: pérdida de escrituras, lectura abortada, transacción desgarrada, sin resolver
- KIP-588: mensaje de error incorrecto en
- Jepsen advierte que la verificación experimental de seguridad puede demostrar la existencia de bugs, pero no su ausencia
- En particular, considera que KAFKA-17754 dificulta determinar si hay otros casos de pérdida de escrituras en Bufstream
Recomendaciones para usuarios y operación de Bufstream
- Los usuarios que usan transacciones de Bufstream con el cliente oficial Java de Kafka deben considerar que actualmente las transacciones podrían no ser seguras
- una transacción abortada podría en realidad terminar
commit - una transacción confirmada podría en realidad terminar
abort - una transacción podría quedar partida a la mitad y conservar solo parte de sus efectos
- una transacción abortada podría en realidad terminar
- Bufstream cree que el cliente Franz-go es menos vulnerable a este problema, pero Jepsen no probó Franz-go con técnicas como las de este trabajo
- Otros clientes podrían ser vulnerables o no
- Los usuarios de versiones anteriores a Bufstream 0.1.3 pueden sufrir los siguientes problemas
producer.send()puede devolver incorrectamente el offset0en lugar del offset real- problema de disponibilidad metaestable donde el cliente queda bloqueado
- Jepsen recomienda actualizar a 0.1.3
- Evalúa que la arquitectura general de Bufstream parece sólida
- el método de determinar el orden de fragmentos de datos inmutables con un servicio de coordinación como etcd es un enfoque relativamente simple con precedentes en sistemas OLTP y de streaming
- Desde el punto de vista operativo, recomienda dos mejoras
- si al iniciar falla una solicitud de archivo compartido del almacenamiento, el clúster puede colapsar; recomendó agregar reintentos, y Bufstream añadió una capa de reintentos
- cuando una dependencia no está disponible, recomienda que el agente siga en ejecución en vez de morir de inmediato, entregando backpressure y estado del sistema, y recuperándose de forma más gradual
- A partir de la 0.1.3, Bufstream agregó lógica adicional de reintentos para etcd, pero todavía requiere supervisión constante para mantenerse en línea
- Los usuarios deben asegurarse de tener un supervisor de procesos y probar que no se rinda y siga funcionando incluso durante outages prolongados
Necesidad de documentar y modificar el protocolo de transacciones de Kafka
- La documentación oficial de Kafka casi no habla de transacciones, por lo que los usuarios tienen que combinar varias fuentes ambiguas y contradictorias
- Jepsen recomendó al equipo de Kafka crear un documento central que aclare la semántica de las transacciones, y mencionó KAFKA-17671
- Ese documento debería especificar como mínimo lo siguiente
- cuándo un
consumerobserva offsets monótonamente crecientes - cuándo un
consumerpuede saltarse registros ya reconocidos - si un rebalance puede afectar a una transacción a mitad de camino
- cuándo el offset de escritura del
produceraumenta de forma monótona - cuándo son legales G0, G1a, G1b, G1c, fractured read y leer las escrituras de la propia transacción
- qué significado tienen el valor devuelto por
poll()y el offset después de una transacción abortada - cómo deben manejarse los errores de transacción, los errores durante
aborty los errores durante el rewind
- cuándo un
- Jepsen señala que la documentación de Confluent repite que la configuración predeterminada de Kafka ofrece entrega at-least-once, pero esto no parece ser cierto
auto.offset.reset = latestpuede hacer que registros no procesados parezcan “committed”- la documentación de Confluent sobre gestión de offsets también habla del riesgo de perder progreso de mensajes si ocurre un crash con
auto-commitpredeterminado - la documentación que dice que el
consumerhace rewind cuando se aborta una transacción tampoco coincide con la realidad
- Jepsen considera que el protocolo de transacciones de Kafka debe corregirse de raíz
- el protocolo asume implícitamente una entrega ordenada y confiable, pero existen pausas de procesos, redes no confiables, latencia distinta de cero y entrega desordenada entre varios sockets TCP
- el protocolo de Kafka distribuye mensajes entre varios nodos y sockets TCP, y el cliente reintenta mensajes automáticamente
- no existen números de secuencia para restaurar el orden de los mensajes del mismo cliente ni un número de transacción para verificar el destino de la transacción
- KIP-890 busca garantizar un orden más estricto aumentando el epoch en cada
commitde transacción - Las bibliotecas cliente también podrían ayudar reinicializando el
producery aumentando el epoch cuando un mensaje no es reconocido - Java Kafka Client 3.8.0 es vulnerable a este problema
- Jepsen considera que Franz-go puede mitigar o evitar este problema al reinicializarse cuando ocurre un timeout, pero no investigó otras bibliotecas cliente
Trabajo futuro
- Muchos usuarios dependen de la “semántica de exactamente una vez” del API de Kafka Streams, en lugar de manejar directamente las transacciones, por lo que en el futuro se podría investigar la corrección de las aplicaciones de Streams
- Jepsen también encontró escrituras no vistas en Kafka mientras investigaba KAFKA-17754, pero no pudo analizarlas por limitaciones de tiempo
- Las escrituras no vistas pueden ser una señal de transacciones colgadas, consumidores atascados o pérdida de datos
- También queda la duda de si un mensaje
Produceretrasado podría entrar en una transacción futura y violar las garantías transaccionales - También se sospecha que Kafka Java Client reutiliza el número de secuencia cuando ocurre un request timeout, lo que podría hacer que una escritura fuera reconocida pero descartada silenciosamente
- Cuando ocurre un evento de rebalanceo, la posición del consumidor puede moverse hacia adelante o hacia atrás, pero las reglas no están claras
- Jepsen quiere verificar el comportamiento previsto por Kafka si este llega a documentarse
- Jepsen explica que, al ser un proceso aleatorio, le cuesta explorar anomalías poco frecuentes
- Los problemas que ocurren una sola vez son muy difíciles de depurar y reproducir
- Bufstream también usa Antithesis, que ejecuta todo el sistema distribuido en un hipervisor determinista y una red simulada
- Combinar la generación de cargas de trabajo y la verificación del historial de Jepsen con el entorno determinista y reproducible de Antithesis puede mejorar la reproducibilidad de las pruebas
1 comentarios
Opiniones en Hacker News
Si al investigar issues como KAFKA-17754 también encontraron escrituras invisibles en Kafka, parece que ya es hora de que Jepsen vuelva a meterse a fondo con Kafka.
La última investigación fue en 2013 (https://aphyr.com/posts/293-call-me-maybe-kafka, Kafka 0.8 beta), y ahora parece que apenas están empezando a descubrir varios problemas en Kafka mismo.
Algo como “una escritura puede ser confirmada y luego descartarse silenciosamente” da bastante miedo.
Me sorprende mucho la parte de que, con el valor predeterminado
enable.auto.commit=true, un consumidor de Kafka puede confirmar offsets independientemente de si la aplicación los procesó realmente.Nunca entendí así el auto commit, y si ese es el valor predeterminado, me parece absurdo.
La documentación no es clarísima, pero en general yo la leía como que los offsets se confirmaban solo cuando el procesamiento terminaba.
Entendía que ajustar el intervalo de auto commit ayudaba a reducir la ventana de procesamiento duplicado, no la pérdida de mensajes, como uno esperaría en un modelo at-least-once.
Si no haces commit explícito, Kafka no tiene forma de saber si procesaste el mensaje.
Kafka asume que el mensaje que te entregó fue procesado de inmediato.
El auto commit es parecido a entregar un cono de helado y darte vuelta enseguida asumiendo que la otra persona se lo comió. Alguien podría dejarlo caer apenas lo recibe y no comer ni un bocado.
Si quieres esa garantía, debes enviar un acknowledgment explícito.
Por ejemplo, si lo único que haces es escribir el mensaje en una base de datos, el mensaje se considera confirmado en el momento en que entra al callback del handler del cliente.
Pero probablemente quieras que se confirme recién después de que la inserción en la DB haya tenido éxito.
Si la DB queda inaccesible por la red, Kubernetes, configuración del firewall, etc., y mientras tanto un ingeniero intenta reiniciar y el cliente se cae, es fácil terminar con mensajes sin procesar.
Otro sistema puede determinar si hubo una falla, y con esta función se puede mover la posición límite para reducir el reprocesamiento.
Aun así, si el timing coincide y ocurre una falla, después de reiniciar hay que asumir que podrías recibir de nuevo parte de lo que ya procesaste.
El problema es cuando no existe ese procesamiento antes del auto commit.
Al leerlo, parece que la intención es que el commit ocurra bastante después del procesamiento, pero también se siente contradictorio que, siendo auto commit, solo deba confirmar los elementos de unos milisegundos antes del momento del auto commit.
polly luego procesa los mensajes de forma durable.El punto confuso es que la verificación de auto commit no ocurre de forma asíncrona después de un timeout, sino en el momento de la siguiente llamada a
poll.Por lo tanto, solo debería poder perder escrituras si, antes de volver a llamar a
poll, no procesas los mensajes de forma durable y solo los guardas, por ejemplo si usas procesamiento asíncrono, demoras, colas, etc.Esto se basa en el comportamiento documentado de la biblioteca cliente de Java (https://kafka.apache.org/32/javadoc/org/apache/kafka/clients...); otra cosa es si la implementación actual realmente se comporta así.
El protocolo de Kafka queda a medio camino entre alto y bajo nivel, y no hace especialmente bien ninguna de las dos cosas.
El auto commit es una función de alto nivel que ayuda a crear aplicaciones simples fácilmente, pero por supuesto puede fallar si no se usa de la manera esperada.
Hoy creo que los usuarios finales deberían usar implementaciones de más alto nivel que manejen bien los detalles, en vez de usar directamente el cliente de Kafka: para datos, motores de procesamiento de streams; para aplicaciones, algo como motores de ejecución persistente.
Al ver la página del producto (https://buf.build/product/bufstream), me pregunto cómo encajan la afirmación de que “solo se ejecuta dentro de tu VPC de AWS o GCP y no se comunica con el exterior” con el cobro basado en uso de “US$0.002 por GiB antes de compresión”.
No creo que vayan a operar todo el negocio con un sistema de honor.
Claro que hay riesgo de abuso, pero puede ser un compromiso valioso para atraer a ciertos clientes.
Si el código fuente no está disponible, no deberías creer jamás una afirmación como “no se comunica con el exterior”.
“El protocolo de transacciones de Kafka está fundamentalmente roto y debe revisarse” suena doloroso.
Aun así, como siempre, la investigación y el artículo son excelentes.
Me pregunto si Kyle ha revisado NATS JetStream. Me daría curiosidad saber qué opina.
Varias personas sugirieron que sería… cómo decirlo… interesante :-)
No logro encontrar el proyecto bufstream en GitHub; me pregunto dónde está.
Aunque, curiosamente, tampoco tiene licencia.
Después de leer los posts y la documentación relacionados, parece que la “entrega exactamente una vez” de Kafka se define como una propiedad de una operación leer-procesar-escribir en la que un worker lee del tópico 1 y escribe al tópico 2, y ambos tópicos están dentro del mismo sistema lógico de Kafka.
Si eso es correcto, me parece que sería mejor llamarlo transacción.
Pero hay dos formas de ver “exactamente una vez”.
Una es que, como en una transacción de base de datos, los efectos no deben duplicarse ni desaparecer.
La otra se parece más a una propiedad del grafo de flujo de datos sobre las relaciones entre mensajes a través de tópicos-particiones, y está un poco más cerca de la consistencia de ACID.
Del mismo modo que un sistema de transacciones serializables garantiza cierta consistencia a nivel de dominio, se puede usar transacciones para alcanzar esa propiedad de flujo de datos.
Por ejemplo, la serializabilidad garantiza que los invariantes preservados al observar cada transacción por separado también se preserven en un historial concurrente.
Se puede ver a Kafka como intentando llegar a una “semántica exactamente una vez” de esa manera.
No confundir con https://www.warpstream.com/.
Fe de erratas: “Transactions may observe none, part, or all” debería ser “Consumers may observe none, part, or all”.
La semántica de los consumidores fuera de una transacción es más difusa.
Todas las lecturas de esta carga de trabajo ocurren en un contexto transaccional y pasan por la ruta de commit de offsets transaccional.
Me pregunto para qué se usa este software. ¿Instrumentación? ¿Caja negra?
Claro, lágrimas de alegría. Que Jepsen te preste atención ya es un logro en sí mismo.