Tema representativo de entrevista

Entrevista de ingeniería de datos: ¿qué garantiza realmente la semántica exactly-once de Kafka?

DatosDifícil
Equipo editorial de Offer.ccPublicado Actualizado

Pregunta

Un consumidor lee eventos de órdenes desde Kafka, los transforma y escribe en otro topic. El entrevistador le pregunta cómo evita duplicados y pérdidas, y luego le pregunta por qué el mismo diseño de exactly-once no cubre automáticamente una base de datos externa. ¿Cómo respondería?

Planteamiento y contexto

Esta pregunta de ingeniería de datos y procesamiento de streams está dirigida a ingenieros de datos, ingenieros de plataformas de streaming y roles de infraestructura de datos backend. El consumidor puede fallar o sufrir un rebalanceo, y un productor puede reintentar tras perder una respuesta; los resultados regresan primero a Kafka. La respuesta debe separar el procesamiento atómico de Kafka de los efectos secundarios externos arbitrarios en lugar de tratar exactly-once como una garantía mágica para todo el sistema.

Qué evalúa el entrevistador

  • Si puede dividir exactly-once en visibilidad del output y avance atómico del offset de entrada.
  • Si puede utilizar correctamente un productor transaccional, transactional.id, read_committed y commits manuales de offsets.
  • Si puede explicar la anulación (abort), el reinicio, el fencing y la recuperación ante rebalanceos.
  • Si reconoce que una base de datos, motor de búsqueda o servicio HTTP necesita su propia transacción, clave de idempotencia o reconciliación.

Preguntas clarificadoras antes de responder

Primero confirme si la salida permanece en Kafka, si se utiliza Kafka Streams, si el grupo de consumidores realiza un rolling restart durante el despliegue y si un sistema externo debe hacer commit en la misma transacción que Kafka. Luego distinga "un output visible por input" de "un efecto secundario externo"; esto último requiere la cooperación del destino. Pregunte sobre latencia, tamaño de lote de la transacción, ventana de reintento y lag aceptable.

Marco de respuesta de 30 segundos

Pondría los registros de entrada, la salida transformada y los offsets del consumidor en una sola transacción de Kafka. El productor habilita transactional.id; el consumidor deshabilita el auto-commit y utiliza read_committed. Un commit exitoso hace que los tres sean visibles juntos, mientras que un abort oculta la salida y deja el offset antes de la transacción. Un transaction ID estable y único aplica fencing a la instancia antigua tras un reinicio. Esto cubre únicamente lecturas y escrituras de Kafka; una base de datos externa necesita el resultado y el offset en una sola transacción de almacenamiento, o bien idempotencia, un outbox y reconciliación.

Solución paso a paso

1. Establecer el límite de la garantía

El diseño de Kafka describe exactly-once de topic a topic como la actualización atómica de los registros de salida y la posición del consumidor en una sola transacción. No significa que la función se ejecute una sola vez ni que una solicitud HTTP arbitraria llegue una sola vez. Establezca este límite antes de discutir la configuración.

2. Escribir atómicamente la salida y los offsets con un productor transaccional

Deshabilite el auto-commit, procese un lote, envíe la salida y envíe los offsets de ese lote como parte de la transacción. El pseudocódigo central es:

java
producer.initTransactions();
while (running) {
  ConsumerRecords<String, Order> records = consumer.poll(timeout);
  producer.beginTransaction();
  try {
    for (ConsumerRecord<String, Order> r : records) {
      producer.send(new ProducerRecord<>("orders-enriched", r.key(), enrich(r.value())));
    }
    producer.sendOffsetsToTransaction(offsets(records), groupMetadata);
    producer.commitTransaction();
  } catch (AbortableException e) {
    producer.abortTransaction();
    consumer.seekToCommitted();
  }
}

Hacer commit de la salida y los offsets juntos evita duplicados visibles por fallas de salida-antes-del-offset y evita pérdidas por fallas de offset-antes-de-la-salida. Maneje las clases de excepción del cliente según la versión desplegada; reintentar a ciegas cada excepción no es seguro.

3. Hacer que los consumidores vean solo transacciones confirmadas (committed)

Configure isolation.level=read_committed y mantenga enable.auto.commit=false. read_uncommitted expone registros de transacciones abortadas, permitiendo que un consumidor propague salidas que debieron haberse revertido. read_committed utiliza marcadores de transacción para alinear la visibilidad con los resultados del commit.

4. Manejar reinicio, fencing y rebalanceo

Asigne a cada instancia activa de consumidor un transactional.id estable y único a nivel de clúster. Cuando una nueva instancia se registra con el mismo ID, Kafka aborta la transacción en curso de la instancia antigua y le aplica fencing, evitando que ambas hagan commit. Tras un abort, la aplicación debe recrear o retroceder explícitamente la posición del consumidor y reprocesar el lote; no debe continuar desde un cursor local que avanzó dentro de la transacción abortada. La asignación de particiones asegura que un solo miembro del grupo posea una partición a la vez.

5. Explicar por qué los sistemas externos deben cooperar

Si el resultado va a PostgreSQL, una transacción de Kafka no incluye automáticamente el commit de la base de datos. Una caída donde la base de datos se ejecuta primero deja el offset de Kafka sin confirmar, por lo que un reintento necesita una clave única o una actualización de versión condicional; un commit de offset primero puede perder la escritura. Un diseño más sólido coloca el resultado y el offset en una sola transacción de base de datos, o escribe un registro reproducible de outbox/conector y deja que el destino deduplique y reconcilie. Sin la cooperación del destino, prometa ejecución at-least-once con duplicados detectables, no exactly-once de extremo a extremo.

6. Verificar con inyección de fallos, no con capturas de configuración

Provoque caídas antes y después de sendOffsetsToTransaction, pierda la respuesta del commit, active rebalanceos, aplique fencing a la instancia antigua, expire una transacción y haga que el sink no esté disponible. Consuma las salidas con read_committed y verifique el recuento de salidas visibles por ID de entrada, los offsets finales, el lag tras el reinicio y los registros de transacciones abortadas. Para un sink externo, pruebe los conflictos de claves únicas, la reproducción y la reconciliación por separado. Monitoree la tasa de aborts, el lag del consumidor, la latencia de procesamiento y el recuento de fencings.

Respuesta de muestra de alta calidad

Primero delimitaría la afirmación. Si tanto la entrada como la salida están en Kafka, usaría Kafka Streams o un bucle transaccional equivalente de consumo-transformación-producción. El consumidor deshabilita el auto-commit; el productor tiene un transactional.id estable y único; y las salidas de cada lote más los offsets se confirman con sendOffsetsToTransaction. Los consumidores posteriores usan read_committed, por lo que la salida abortada es invisible. Al reiniciar, el mismo transaction ID aplica fencing a la instancia antigua, y un abort requiere retroceder al último offset confirmado. Si el destino es PostgreSQL, no expandiría la garantía de Kafka a una afirmación de extremo a extremo: pondría el resultado y el offset en una sola transacción de base de datos, o usaría un outbox, clave de idempotencia y reconciliación. Luego inyectaría caídas, respuestas perdidas, rebalanceos y fencing, y verificaría las salidas visibles, los offsets y el estado del destino por ID de entrada.

Errores comunes

  • Llamar exactly-once a un productor idempotente → principalmente evita entradas de log duplicadas en reintentos del productor y no hace que la salida más el offset sean atómicos → agregue transacciones y commits de offsets.
  • Configurar únicamente read_committed cambia la visibilidad pero no hace commit de transacciones ni recupera caídas → implemente el ciclo de vida completo de transacciones con auto-commit desactivado.
  • Compartir un solo transactional.id entre instancias activas → las instancias se aplican fencing mutuamente y desestabilizan el rendimiento → asigne un ID único y estable por instancia activa.
  • Continuar desde la posición local tras un abort → esa posición puede estar dentro del lote no confirmado → recargue el offset confirmado o haga seek explícitamente.
  • Afirmar que Kafka hace que un efecto secundario en la base de datos ocurra una sola vez → los dos sistemas carecen de un commit atómico automático → utilice una transacción de destino, escritura idempotente, outbox o reconciliación.

Preguntas de seguimiento y respuestas

¿Exactly-once significa que la función de negocio se ejecuta una sola vez?

No. La función puede ejecutarse y luego ejecutarse nuevamente tras un fallo de commit; la garantía es que la salida visible de Kafka y el commit del offset coincidan. Mantenga la función libre de efectos secundarios externos siempre que sea posible, o haga que esos efectos sean repetibles y reconciliables.

¿Por qué read_committed todavía puede añadir latencia?

El consumidor debe omitir los registros abortados y esperar a que se completen las transacciones. Las transacciones abiertas o largas aumentan la latencia de visibilidad y el lag. Limite el tamaño del lote y el timeout, monitoree la duración de las transacciones y pondere el costo frente al procesamiento at-least-once para rutas de baja latencia.

¿Qué sucede cuando ocurre un rebalanceo a mitad de una transacción?

El nuevo propietario reanuda desde el último offset confirmado. La transacción actual debe abortarse; la instancia antigua se libera o se le aplica fencing; y la nueva instancia reprocesa el lote no confirmado. No haga commit de offsets a menos que el miembro actual todavía posea la partición.

¿Cómo enviaría la salida de Kafka a PostgreSQL?

Escriba el resultado de negocio y la posición del consumidor en una sola transacción de PostgreSQL, o escriba una fila de outbox con un ID de evento único y deje que un relay confiable la entregue. Si no pueden compartir una transacción, elija entrega at-least-once, restricciones únicas, actualizaciones de versión condicionales y reconciliación; no lo promocione como exactly-once.

¿Qué métricas demuestran que el diseño funciona?

Cuente las salidas visibles únicas de read_committed por ID de evento de entrada, luego correlacione offsets confirmados, recuentos de aborts y fencings, lag del consumidor, volumen de reproducción tras reinicios y conflictos de claves duplicadas en el sink. Mantenga muestras antes y después de la inyección de fallos; una captura de configuración no puede probar el invariante.

Fuentes públicas

Preguntas relacionadas