Planteamiento y Cuándo Aplica
Un topic de Kafka tiene 24 particiones y recibe 120,000 registros por segundo en el pico de tráfico. Un inquilino (tenant) produce el 45% del tráfico y el productor particiona por tenant_id, por lo que una partición acumula lag continuamente mientras la mayoría de los demás consumidores permanecen inactivos. Un consumidor puede sostener 8,000 registros por segundo con el manejador (handler) actual. El negocio requiere orden únicamente dentro de un order_id, no en todos los eventos del inquilino, y no se pueden agregar particiones durante el incidente.
Explica cómo demostrarías que el sesgo de claves (key skew), y no una falla en el consumidor, el broker o los sistemas downstream, causó el síntoma. Luego, describe cómo reducirías el crecimiento del backlog hoy, rediseñarías la clave de partición y la migración, y confirmarías offsets de forma segura tras agregar procesamiento asíncrono. Finaliza con las métricas y las pruebas de fallas que demuestren la solución.
Esta es una pregunta de resolución de problemas en ingeniería de datos y plataformas de streaming. Una guía pública actual de entrevistas de Kafka presenta esencialmente la misma combinación: una partición caliente, su consumidor retrasándose, consumidores inactivos en otros lugares, un productor con particionamiento por clave y la imposibilidad de cambiar de inmediato el número de particiones. Pide a los candidatos conectar el particionamiento del productor, el paralelismo del consumidor y el ordenamiento. La documentación de Apache Kafka establece que el productor predeterminado elige una partición mediante el hash de la clave presente. El particionamiento semántico preserva la localidad y el orden dentro de la partición elegida, pero también concentra allí el tráfico de una clave. La guía de Kafka de Huawei Cloud señala de manera similar que una partición solo puede ser consumida por un consumidor a la vez y que agregar particiones temporalmente no es una vía rápida para vaciar el backlog de una partición existente.
El rendimiento, la proporción de tráfico, la capacidad del consumidor y el alcance del ordenamiento en este planteamiento son supuestos de entrevista, no cifras de producción atribuidas a una empresa en particular.
Qué Evalúa el Entrevistador
La primera señal es si puedes pasar de métricas agregadas a evidencia a nivel de partición. El lag global del topic, el uso promedio de CPU del consumidor y el número de consumidores pueden ocultar una partición caliente. Una respuesta sólida alinea en una misma línea de tiempo la tasa de producción, la tasa de consumo, la pendiente del lag, el broker líder, la distribución de claves y el tiempo de procesamiento downstream de cada partición. Eso separa el sesgo del productor de un manejador lento, rebalanceos repetidos o presión de recursos en el broker.
La segunda señal es si cuantificas el incidente:
120,000 × 45% = 54,000 records/secondSi el consumidor asignado a esa partición solo puede procesar 8,000 registros por segundo, el backlog crece a razón de:
54,000 - 8,000 = 46,000 records/second
46,000 × 600 = 27,600,000 records in 10 minutesEse cálculo explica por qué agregar consumidores ordinarios no cambia el límite máximo de esta partición y por qué cambiar el umbral de una alerta no mitiga el incidente.
La tercera señal es si identificas la verdadera invariante de ordenamiento. La clave actual expande el dominio de orden a todo un inquilino, mientras que el negocio solo necesita orden dentro de un pedido. Cambiar la clave a order_id aumenta la cardinalidad y distribuye los nuevos pedidos, pero solo si la migración evita que un pedido cruce dos particiones o topics.
La cuarta señal es la corrección de offsets. Una vez que los registros de una partición se ejecutan en un grupo de trabajadores (worker pool), el orden de finalización puede diferir del orden de los offsets. Si el offset 105 termina mientras el 104 aún se está ejecutando, confirmar hasta el 105 puede omitir el 104 tras una caída del sistema. Un diseño sólido rastrea el offset completado contiguo más alto y hace que los efectos downstream sean idempotentes, ya que los registros completados pero no confirmados pueden reproducirse.
La señal final es si estableces el límite estricto. Más particiones crean más ranuras paralelas, pero no dividen una entidad gigante que permanece vinculada a una sola clave. Más consumidores en un grupo de consumidores tradicional no permiten que dos consumidores posean la misma partición simultáneamente. Si un solo pedido excede la capacidad segura de una partición y debe permanecer estrictamente ordenado, las únicas palancas restantes son optimizar la ruta serial, aplicar throttling o redefinir secuencias de negocio independientes.
Preguntas para Aclarar Antes de Responder
- ¿Se requiere ordenamiento por inquilino, por pedido o para algún flujo de eventos más pequeño? Si el orden en todo el inquilino es obligatorio, dividir un inquilino es inválido. Si el orden por pedido es suficiente,
order_ides el límite de partición más preciso. - ¿El 45% describe cantidad de registros, bytes o costo de procesamiento? Registros grandes o escrituras downstream costosas pueden crear sesgo de costo incluso cuando el conteo de registros parece equilibrado. Inspecciona registros, bytes y tiempo del manejador.
- ¿Aumentó la tasa de producción o cayó la capacidad de consumo? Un flujo constante de 54,000 registros por segundo frente a una capacidad de 8,000 registros indica sesgo de clave y capacidad insuficiente por clave. Si la entrada no cambia pero el consumo cae de 8,000 a 2,000, investiga primero el sistema downstream, la recolección de basura (garbage collection), la red, el disco y los rebalanceos.
- ¿Dónde escribe el consumidor? Un pipeline de Kafka a Kafka puede confirmar atómicamente los offsets de salida y de entrada con transacciones de Kafka. Una base de datos, almacenamiento de objetos o API externa generalmente necesita procesamiento at-least-once más una clave de idempotencia de negocio o
topic-partition-offset. - ¿Cuánto tiempo puede permanecer activo un pedido? Los pedidos de corta duración pueden permanecer en la ruta heredada hasta su finalización mientras los nuevos pedidos usan la nueva ruta. Los pedidos de larga duración requieren una barrera explícita, secuencia o estado de enrutamiento.
- ¿Puede la respuesta al incidente limitar la tasa (throttle) o degradar el tráfico? Una cuota de inquilino, el retraso de eventos analíticos no críticos o la consolidación (coalescing) de actualizaciones de estado pueden reducir la entrada más rápido que una migración de código y particiones.
- ¿Están los consumidores excediendo
max.poll.interval.ms? Si el trabajo pesado bloquea el hilo de sondeo (poll thread), los rebalanceos amplifican el lag. Separa el sondeo del procesamiento y delimita el trabajo en vuelo antes de limitarte a aumentar el tiempo de espera.
Estructura de Respuesta en 30 Segundos
“Desglosaría el lag agregado en tasa de producción, tasa de consumo y pendiente de lag por partición, y luego los alinearía con la frecuencia de claves, bytes, tiempo del manejador, registros de rebalanceo y métricas del broker líder. El inquilino caliente produce 54,000 registros por segundo mientras que un consumidor maneja 8,000, por lo que el backlog crece aproximadamente a 46,000 por segundo; agregar consumidores ordinarios no puede acelerar esa partición. Hoy aplicaría throttling o degradaría al inquilino caliente, aislaría la partición en una instancia dedicada y aprovecharía el verdadero límite de orden serializando cada pedido mientras proceso diferentes pedidos de forma concurrente. Confirmaría solo el offset completado contiguo más alto y usaría escrituras downstream idempotentes. A largo plazo, los nuevos pedidos se moverían a un nuevo topic particionado por order_id, mientras que los pedidos existentes permanecerían en la ruta heredada hasta que terminen. Validaría la pendiente de lag por partición, p99 de extremo a extremo, duplicados y violaciones de orden bajo carga sesgada, caídas de consumidores y rebalanceos.”
Respuesta Detallada Paso a Paso
Paso 1: Demostrar qué capa creó el punto caliente
Usa una ventana de tiempo pico y alinea:
- registros producidos por segundo, bytes por segundo y crecimiento del high-watermark por partición;
- registros consumidos por segundo, offset confirmado y pendiente de lag por partición;
- frecuencia de claves, bytes y distribución estimada del costo de procesamiento;
- CPU, red, espera de disco y latencia de solicitudes en el broker líder de la partición caliente;
- intervalos de poll, tamaños de lote, latencia del manejador, recolección de basura, errores y registros de rebalanceo para el consumidor asignado;
- latencia correlacionada con la partición o throttling en la base de datos, sistema de almacenamiento o API downstream.
La regla de decisión se deduce de los números. La partición caliente recibe unos 54,000 registros por segundo. Las otras 23 particiones comparten los 66,000 restantes, con un promedio aproximado de 2,870 registros por segundo si ese remanente es razonablemente uniforme. Un consumidor de 8,000 registros tiene capacidad de sobra en una partición ordinaria pero no puede igualar la entrada caliente. Eso explica completamente una partición creciente y capacidad ociosa en otros lugares. Si la entrada de la partición caliente es normal mientras el consumo se degrada al mismo ritmo que la latencia downstream, el diseño de la clave aún no es la causa raíz comprobada.
Verifica también la ubicación en los brokers. Una partición cuyo líder se encuentra en un broker sobrecargado puede sufrir una producción y extracción más lentas. Mover el liderazgo o rebalancear réplicas puede eliminar ese cuello de botella de ubicación, pero no cambia el hecho de que un tenant_id todavía se asigna a una sola partición.
Paso 2: Reducir la pendiente del backlog antes de intentar drenarlo
El primer objetivo del incidente es:
hot-partition input rate ≤ hot-partition safe processing rateLa palanca más rápida suele ser la admisión. Aplica una cuota explícita al inquilino caliente, retrasa eventos analíticos reconstruibles, consolida actualizaciones para las cuales solo importa el estado más reciente o coloca el trabajo no crítico en una ruta degradada declarada. Cada medida necesita semánticas explícitas de pérdida, retraso y reintento. Descartar registros silenciosamente no es una estrategia de throttling.
En el lado del consumidor, una ventana de mantenimiento controlada puede detener el grupo existente y reemplazarlo con una asignación explícita, exclusiva y sin solapamiento: una instancia adecuadamente aprovisionada posee solo la partición caliente, y las instancias restantes poseen las otras particiones. Iniciar un segundo grupo de consumidores ordinario no ayuda; consume independientemente el topic completo y duplica los efectos de negocio. El aislamiento evita que la partición caliente deje sin recursos a las particiones ordinarias en el mismo proceso, aunque no eleva el manejador serial original por encima de los 8,000 registros por segundo.
Debido a que el negocio solo necesita orden dentro de un pedido, el consumidor caliente puede despachar por order_id: una cola serial por pedido activo y un grupo de trabajadores delimitado entre diferentes pedidos. El grupo debe tener un límite de trabajo en vuelo. Cuando esté lleno, pausa la partición o reduce la cantidad enviada a los trabajadores para que el lag de Kafka no se convierta en memoria de proceso ilimitada. El bucle de poll debe mantenerse receptivo; de lo contrario, exceder max.poll.interval.ms desencadena un rebalanceo y agrega otra pausa.
Paso 3: Confirmar una marca de agua de finalización contigua
La concurrencia intrapartición cambia el orden de finalización, pero no debe cambiar el orden de confirmación (commit). Mantén este estado para cada partición:
nextCommitOffset = smallest unfinished offset
completed = offsets that finished but still have a gap before them
onComplete(offset):
add offset to completed
while completed contains nextCommitOffset:
remove nextCommitOffset from completed
increment nextCommitOffset
commit nextCommitOffsetKafka confirma la siguiente posición a leer. Por lo tanto, solo después de que 104, 105 y 106 hayan terminado, la posición confirmada puede avanzar a 107. Si 105 termina mientras 104 se reintenta, el punto de confirmación permanece en 104. Una caída antes del siguiente commit reproduce algunos registros finalizados, por lo que las escrituras en la base de datos deben usar event_id u otra clave única de negocio para un upsert idempotente. Si no existe una clave de negocio, topic-partition-offset puede identificar el registro de origen.
No describas una confirmación de offset y un efecto secundario externo como naturalmente atómicos. Una aplicación de Kafka a Kafka puede colocar los registros de salida y los offsets consumidos en una sola transacción de Kafka. Con una base de datos externa, un contrato más común es consumo at-least-once más una escritura idempotente, o una transacción de base de datos que contenga tanto el registro de deduplicación como la mutación de negocio.
Paso 4: Ajustar el alcance de la partición al verdadero dominio de ordenamiento
La capacidad promedio de 24 particiones no es necesariamente insuficiente. Si 8,000 registros por segundo es un máximo seguro medido y la utilización planificada tiene un tope del 70%, la capacidad planificada por partición es de 5,600:
120,000 ÷ 5,600 ≈ 21.4Con una distribución uniforme, 24 particiones cubren el pico asumido con un margen modesto. La falla proviene de colocar el 45% del tráfico detrás de una sola clave de baja cardinalidad, no del conteo total de particiones. Una clave duradera debe representar el dominio de orden más pequeño requerido, tener suficiente cardinalidad y permanecer predeciblemente distribuida en el pico. Aquí, order_id es la opción natural. Una codificación estable de (tenant_id, order_id) también es posible si la localidad por inquilino tiene un valor operativo real.
El salting controlado es válido solo cuando los registros dentro de la clave original pueden reordenarse o una etapa downstream puede restaurar el orden. Aplicar salting aleatorio a un pedido viola la invariante de este planteamiento porque ese pedido puede llegar concurrentemente desde varias particiones. Si un solo pedido excede la capacidad de una partición, una mayor cardinalidad de clave no ayuda; optimiza o limita la ruta serial de ese pedido, o rediseña el protocolo de negocio en secuencias explícitamente independientes.
Paso 5: Migrar con enrutamiento versionado
Agregar particiones al topic existente y cambiar las claves inmediatamente genera dos riesgos. El mapeo de hash predeterminado puede mover claves existentes cuando cambia el número de particiones. Los eventos de un mismo pedido también pueden aterrizar en diferentes particiones antes y después del cambio (cutover), mientras que Kafka no garantiza el orden entre particiones.
Un diseño más seguro crea un nuevo topic particionado por order_id y versiona el enrutamiento del productor:
- los pedidos creados después del cutover usan el nuevo topic;
- los pedidos que ya existían permanecen en el topic antiguo y con la clave heredada hasta que se cierren;
- cada productor utiliza el mismo estado de enrutamiento de pedidos en lugar de comparar su reloj local con una hora de corte;
- los consumidores leen ambas rutas, pero un pedido pertenece a una sola ruta activa en cualquier momento;
- el topic antiguo se retira después de que los pedidos heredados se vacíen y se cumplan los requisitos de retención.
Si los pedidos no terminan de forma natural, crea una barrera de migración por pedido: pausa los nuevos eventos para ese pedido, espera hasta que la ruta antigua alcance una secuencia u offset final registrado, cambia la versión de enrutamiento y reanuda. Una alternativa sin pausas puede transportar números de secuencia monotónicos y fusionar ambas rutas downstream, pero eso introduce almacenamiento en búfer, tiempos de espera y recuperación de brechas (gap recovery). Solo se justifica si el requisito de negocio compensa esa complejidad.
Paso 6: Validar con tráfico sesgado y fallas
El rendimiento agregado no es evidencia suficiente. Prueba al menos:
- una distribución en la que un inquilino produce el 45% del tráfico y la distribución de pedidos se asemeja al pico real;
- entrada, consumo, pendiente de lag y lag máximo por partición;
- p50, p95 y p99 de extremo a extremo más el tiempo estimado de drenado;
- recuento de trabajadores en vuelo, antigüedad de la tarea más antigua, reintentos y cartas muertas (dead letters);
- efectos duplicados, violaciones de orden por pedido y conflictos de idempotencia;
- una caída del consumidor mientras los offsets completados contienen un salto o brecha;
- si un procesamiento prolongado desencadena un rebalanceo y cuánto tiempo toma la recuperación;
- si un pedido aparece en una sola ruta en el límite entre el topic antiguo y el nuevo.
La condición de aprobación incluye un comportamiento sostenido: la pendiente de lag de la partición caliente ya no es positiva en el pico constante; una caída puede reproducir trabajo pero no puede perder un efecto de negocio; no se observa ningún pedido fuera de secuencia; los nuevos pedidos se distribuyen entre las particiones; y los pedidos heredados se drenan según lo programado. Si el rendimiento agregado aumenta mientras un pedido grande crea repetidamente un punto caliente, el problema del dominio de orden o de admisión del negocio sigue sin resolverse.
Ejemplo de Respuesta de Alta Calidad
“No comenzaría agregando consumidores porque el planteamiento ya nos dice que una partición está retrasada mientras que los consumidores en otros lugares están inactivos. Primero demostraría el sesgo utilizando registros por segundo, bytes por segundo, pendiente de lag y frecuencia de claves por partición, descartando al mismo tiempo el broker líder, rebalanceos y latencia downstream.
El inquilino caliente produce 54,000 registros por segundo. El consumidor de una sola partición procesa 8,000, por lo que el lag crece a razón de aproximadamente 46,000 registros por segundo, o 27.6 millones en 10 minutos. Si el otro 55% se distribuye aproximadamente en 23 particiones, cada una promedia unos 2,870 registros por segundo. Eso explica tanto el déficit del consumidor caliente como la capacidad ociosa en otros lugares. En un grupo de consumidores tradicional, una partición pertenece a un consumidor a la vez, por lo que más instancias ordinarias no la aceleran.
Hoy reduciría primero la pendiente de entrada con una cuota documentada para el inquilino y movería los eventos tolerantes al retraso o consolidables a una ruta degradada. En una ventana controlada, aislaría la partición caliente en una instancia dedicada para que no prive de recursos a las particiones ordinarias. Dado que solo importa el orden por pedido, despacharía por order_id a un pool delimitado: serial dentro de un pedido, concurrente entre pedidos. No confirmaría la tarea que termine más rápido. Rastrearía el offset completado contiguo más alto y me detendría ante cualquier brecha. Una caída puede reproducir registros completados pero no confirmados, por lo que el destino utiliza event_id o una clave única de negocio para idempotencia.
A largo plazo, agregar particiones no es la solución completa. Al 70% de utilización planificada, cada partición de 8,000 registros contribuye con unos 5,600 registros por segundo, por lo que 24 particiones cargadas uniformemente pueden cubrir el pico asumido de 120,000 registros. El problema es que tenant_id fija el 45% en una sola partición. Crearía un nuevo topic con clave order_id. Los nuevos pedidos después del corte lo usan, mientras que los pedidos heredados activos permanecen en la ruta antigua hasta su finalización. El estado de enrutamiento compartido garantiza que un pedido nunca abarque ambos topics.
Antes del despliegue, reproduciría el mismo sesgo del 45% del inquilino e inspeccionaría el rendimiento por partición, la pendiente del lag, el p99 de extremo a extremo, los duplicados y las violaciones de orden. Provocaría la caída del consumidor mientras los offsets se completan fuera de orden y verificaría que el reinicio cause solo una reproducción idempotente, desencadenaría rebalanceos y mediría la recuperación, y verificaría que cada pedido en el límite de la ruta aparezca en un solo topic. Si un pedido por sí mismo excede la capacidad de una sola partición, establecería el límite estricto: ni más consumidores ni más particiones lo resuelven sin optimizar, aplicar throttling o cambiar el modelo de secuenciación de ese pedido.”
Errores Comunes
- Mirar solo el lag a nivel de topic → Un promedio oculta la pendiente de entrada y consumo de una partición → Grafica registros, bytes, lag y distribución de claves por partición.
- Agregar consumidores cada vez que sube el lag → Una partición en un grupo tradicional tiene un solo consumidor asignado a la vez → Compara primero el conteo de particiones, la asignación y la capacidad por partición.
- Agregar particiones de inmediato → El backlog existente no se redistribuye automáticamente y el mapeo de claves predeterminado puede cambiar → Detén primero el crecimiento del backlog, luego usa un topic versionado y un límite de migración.
- Aplicar salting aleatorio a la clave caliente → Un pedido puede cruzar particiones y llegar desordenado → Aplica salting solo cuando reordenar sea aceptable; usa
order_idpara el verdadero alcance de ordenamiento de este problema. - Confirmar un offset tan pronto como su trabajador termine → Un offset inferior no terminado puede ser omitido tras una caída → Confirma solo la marca de agua de finalización contigua.
- Llamar procesamiento exactly-once a una confirmación de offset → Una mutación en una base de datos externa y un offset de Kafka no suelen ser una sola transacción → Establece los límites de at-least-once, idempotencia y transacciones.
- Iniciar un segundo grupo de consumidores para ayudar → El segundo grupo lee su propia copia completa y duplica efectos → Usa asignación explícita exclusiva o un rediseño de procesamiento controlado.
- Probar solo con tráfico uniforme → Un promedio favorable no demuestra que la clave caliente haya desaparecido → Reproduce un sesgo realista e inspecciona la partición máxima, no solo la media.
- Ignorar el límite de una sola entidad → Una entidad estrictamente ordenada no se puede paralelizar entre particiones gratuitamente → Declara el compromiso estricto entre orden, throttling y procesamiento serial.
Preguntas de Seguimiento y Respuestas
Pregunta de seguimiento 1: ¿Por qué no aumentar el conteo de particiones de 24 a 48 de inmediato?
Nuevas particiones crean ranuras paralelas futuras, pero no dividen el backlog ya almacenado en una partición antigua, ni hacen que un solo tenant_id se asigne a varias particiones. Con el hashing de clave predeterminado, cambiar el conteo de particiones también puede remapear claves existentes y colocar un pedido en diferentes particiones a lo largo del corte. Corrige primero el dominio de ordenamiento y la migración, luego usa una prueba de carga sesgada para decidir si son necesarias más particiones agregadas.
Pregunta de seguimiento 2: ¿Cómo evitas la memoria ilimitada tras agregar concurrencia intrapartición?
Establece un límite máximo de registros en vuelo y una ventana máxima de offsets no confirmados por partición. Cuando se alcance cualquiera de los límites, pausa esa partición o libera menos registros al grupo de trabajadores mientras mantienes el sondeo separado del procesamiento pesado. Rastrea la antigüedad del offset no terminado más antiguo. Si un pedido permanece bloqueado, aíslalo, reinténtalo o enrútalo para intervención en lugar de permitir que cada offset posterior consuma memoria indefinidamente.
Pregunta de seguimiento 3: ¿Qué pasa si la base de datos downstream no admite un upsert idempotente?
Inserta un registro de deduplicación y realiza la mutación de negocio en una sola transacción de base de datos. Usa un event_id u topic-partition-offset de negocio como clave única; un conflicto de unicidad significa que el evento ya ha sido aplicado. Para una API externa no transaccional, usa su clave de idempotencia, una tabla outbox o un estado de operación consultable. Si no se puede construir un límite de idempotencia, no puedes garantizar que la reproducción tras una caída no tenga efectos duplicados.
Pregunta de seguimiento 4: ¿Qué cambia si el negocio requiere posteriormente un orden estricto para todo el inquilino?
Entonces tenant_id es el dominio de orden indivisible, y la concurrencia entre pedidos ya no es válida. Si un inquilino excede la capacidad de una sola partición, optimiza esa ruta serial, aplica throttling al inquilino o renegocia qué secuencias de eventos son independientes. Un flujo fragmentado (sharded stream) secuenciado globalmente con un fusionador downstream es posible, pero traslada las esperas de ordenamiento, la recuperación de brechas y el costo de disponibilidad al consumidor; no es escalabilidad gratuita.
Pregunta de seguimiento 5: ¿Cuánto tiempo tomará drenar el backlog de 27.6 millones de registros?
La entrada debe caer primero por debajo de la capacidad de procesamiento. Si el throttling reduce la entrada de la partición caliente a 3,000 registros por segundo y la optimización eleva la capacidad de procesamiento a 12,000, la tasa neta de drenado es de 9,000:
27,600,000 ÷ 9,000 ≈ 3,067 seconds ≈ 51 minutesEsa sigue siendo una estimación de tasa constante. Un plan de recuperación real agrega reintentos, throttling downstream, variación en el tamaño de los registros y margen de seguridad, y luego revisa continuamente la estimación a partir de la pendiente de lag observada.
Pregunta de seguimiento 6: ¿Cómo demuestras que ningún pedido cruza entre el topic antiguo y el nuevo?
Almacena un partitioning_version para cada pedido. Cada productor lee o almacena en caché el mismo registro de enrutamiento versionado, y la versión cambia solo después de que la barrera de migración tenga éxito. Los consumidores registran el primer topic y versión vistos para un pedido, alertan si un pedido aparece en ambas rutas activas y detienen el avance automático para ese pedido. Durante las pruebas de carga y el despliegue canary, concilia los registros del productor, los offsets en ambos topics y las secuencias de pedidos downstream en lugar de verificar únicamente que los recuentos totales de registros sean iguales.