Preguntas más frecuentes

Preguntas más frecuentes sobre el uso de Kafka con Azure Databricks.

¿Por qué obtengo un error de que no se admite o no se reconoce una opción de Kafka?

Este error se produce si olvida usar el kafka. prefijo al establecer las opciones de configuración del cliente de Kafka. Todas las opciones que se pasan directamente al cliente de Kafka deben tener kafka.el prefijo :

El código siguiente muestra opciones incorrectas que faltan el kafka. prefijo:

.option("security.protocol", "SASL_SSL")
.option("sasl.mechanism", "PLAIN")

En el código siguiente se muestran las opciones correctas:

.option("kafka.security.protocol", "SASL_SSL")
.option("kafka.sasl.mechanism", "PLAIN")

Las opciones del conector de Spark Kafka (como subscribe, startingOffsets, maxOffsetsPerTrigger) no requieren el prefijo. Para obtener la lista completa de opciones, consulte Kafka.

¿Por qué aparece un error relacionado con las clases de Kafka sombreadas?

Azure Databricks requiere el uso de clases de Kafka sombreadas (prefijo con kafkashaded. o shadedmskiam.). Si ve errores, como RESTRICTED_STREAMING_OPTION_PERMISSION_ENFORCED, debe usar los nombres de clase sombreados:

  • org.apache.kafka.* Las clases requieren el kafkashaded. prefijo. Por ejemplo: kafkashaded.org.apache.kafka.common.security.plain.PlainLoginModule
  • software.amazon.msk.* Las clases requieren el shadedmskiam. prefijo. Por ejemplo: shadedmskiam.software.amazon.msk.auth.iam.IAMLoginModule

¿Por qué recibo un TimeoutException al conectarse a Kafka?

Las causas más comunes son:

  • Conectividad de red: el clúster de computación no puede comunicarse con los brokers de Kafka. Compruebe las reglas de firewall, los grupos de seguridad y las configuraciones de VPC.
  • Servidores de arranque incorrectos: compruebe que el nombre de host y el kafka.bootstrap.servers puerto son correctos.
  • Resolución DNS: compruebe que los nombres de host de los brokers de Kafka se pueden resolver desde la red de Azure Databricks.
  • Problemas de SSL/TLS: si usa SSL, compruebe que los certificados están configurados correctamente.

Para configuraciones de Private Link o de emparejamiento entre VPC, verifique que las rutas de red correctas estén configuradas.

¿Debo usar el modo de procesamiento por lotes o streaming para Kafka?

Depende de su caso de uso:

  • Modo de streaming (spark.readStream): use cuando necesite procesamiento continuo de datos o ingesta de baja latencia.
  • Modo por lotes (spark.read): utilícelo para cargas de datos puntuales, rellenos retrospectivos o depuración. Requiere tanto startingOffsets como endingOffsets.

Consulte Configuración de intervalos de desencadenador de Structured Streaming para obtener más información sobre cómo configurar intervalos de desencadenador, como AvailableNow, ProcessingTimey el modo en tiempo real.

¿Puedo leer varios temas de Kafka en una sola secuencia?

Sí, puede usar:

  • subscribe: proporcione una lista separada por comas de temas, por ejemplo .option("subscribe", "topic1,topic2").
  • subscribePattern: use un patrón regex de Java para que coincida con los nombres de tema, por ejemplo, .option("subscribePattern", "topic-.*").

¿Cómo se usa Kafka con canalizaciones de Lakeflow?

Las canalizaciones de Lakeflow tienen soporte integrado para fuentes de Kafka.

Puede definir una tabla de streaming que lee de Kafka, como en el código siguiente:

Python

import dlt

@dlt.table
def kafka_bronze():
  return (spark.readStream
    .format("kafka")
    .option("kafka.bootstrap.servers", "<server:port>")
    .option("subscribe", "<topic>")
    .load()
  )

SQL

CREATE OR REFRESH STREAMING TABLE kafka_bronze AS
SELECT * FROM STREAM read_kafka(
  bootstrapServers => '<server:port>',
  subscribe => '<topic>'
);

Consulte Cargar datos en pipelines para obtener más detalles sobre las fuentes de streaming en los pipelines de Lakeflow.

¿Cómo deserializar las columnas de clave y valor de Kafka?

Las columnas key y value se devuelven como tipo BINARY. Use operaciones de DataFrame para deserializarlas en función del formato de datos:

¿Por qué aparece un error de escritura idempotente?

Databricks Runtime 13.3 LTS y versiones posteriores incluyen una versión más reciente de la biblioteca kafka-clients que permite escrituras idempotentes de manera predeterminada. Si el clúster de Kafka usa la versión 2.8.0 o inferior con las ACL configuradas pero sin IDEMPOTENT_WRITE habilitado, la escritura falla con: org.apache.kafka.common.KafkaException: Cannot execute transactional method because we are in an error state.

Para resolver este error, actualice a Kafka versión 2.8.0 o posterior, o estableciendo .option("kafka.enable.idempotence", "false") al configurar el sistema de escritura de Structured Streaming.

¿Qué es KAFKA_DATA_LOSS_ERROR y cómo puedo resolverlo?

Este error se produce cuando la fuente de Kafka detecta que los offsets almacenados en el punto de control ya no están disponibles en Kafka, normalmente debido a que:

  • La transmisión se ha pausado durante más tiempo que el periodo de retención de Kafka.
  • Los datos del tópico de Kafka se eliminaron o el tópico se recreó.
  • El broker de Kafka experimentó pérdida de datos.

Para resolver este problema:

  • Si la pérdida de datos es aceptable: configura .option("failOnDataLoss", "false") para permitir que la transmisión continúe desde el offset disponible más antiguo.
  • Si la pérdida de datos no es aceptable: restablece el punto de control y vuelve a procesar a partir de los earliest offsets, o restaura los datos de Kafka que faltan.

Consulte KAFKA_DATA_LOSS condición de error para obtener más información.

¿Cómo puedo controlar la velocidad a la que se leen los datos de Kafka?

Utiliza la opción maxOffsetsPerTrigger para limitar el número de offsets (aproximadamente el número de registros) procesados por microlote. Esto ayuda a evitar lotes grandes que podrían sobrecargar el procesamiento descendente o provocar problemas de memoria al ponerse al día en un trabajo pendiente.

Python

df = (spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:port>")
  .option("subscribe", "<topic>")
  .option("maxOffsetsPerTrigger", 10000)
  .load()
)

Scala

val df = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "<server:port>")
  .option("subscribe", "<topic>")
  .option("maxOffsetsPerTrigger", 10000)
  .load()

SQL

SELECT * FROM STREAM read_kafka(
  bootstrapServers => '<server:port>',
  subscribe => '<topic>',
  maxOffsetsPerTrigger => '10000'
);

Como alternativa, use opciones como minPartitions o maxRecordsPerPartition para controlar cuántas particiones de Spark se crean para cada lote.

¿Cómo puedo supervisar el retraso de mi flujo respecto a los últimos offsets de Kafka?

Utilice las métricas avgOffsetsBehindLatest, maxOffsetsBehindLatest y minOffsetsBehindLatest disponibles en el progreso de la consulta de streaming. Estos indican cuántos offsets está rezagado tu flujo respecto al último offset disponible en todas las particiones de los temas a los que estás suscrito. Vea Supervisión de consultas de Structured Streaming en Azure Databricks.

También puede usar estimatedTotalBytesBehindLatest para calcular el total de bytes de datos que aún no se han procesado.

¿Por qué mis métricas de retardo de desplazamiento de Kafka muestran valores persistentes que no son cero después de actualizar a Databricks Runtime 17.1?

En Databricks Runtime 17.1 y versiones posteriores, los desplazamientos de Kafka más recientes se capturan después de que se complete cada microlote. En los temas que reciben continuamente datos, las métricas de trabajos pendientes pueden mostrar valores pequeños y persistentes que no son ceros. Este es el comportamiento esperado y no indica que la transmisión se esté retrasando.

En Databricks Runtime 17.0 y versiones posteriores, los desplazamientos de Kafka más recientes se capturan en la hora de inicio del microlote. Las métricas de trabajos pendientes pueden devolver 0 cuando las consultas de streaming consuman constantemente todos los registros disponibles al principio del microlote.

Si los valores son grandes o siguen creciendo continuamente, es posible que la transmisión no pueda mantenerse al día con los datos entrantes. Vea Supervisión de consultas de Structured Streaming en Azure Databricks.

¿Por qué la inicialización de la secuencia de Kafka es lenta?

Los flujos de Kafka requieren tiempo para:

  1. Conéctese al clúster de Kafka y capture los metadatos.
  2. Descubre las particiones de los temas.
  3. Recupera los offsets iniciales.

En el caso de los clústeres de Kafka locales o remotos, la latencia de red puede afectar significativamente al tiempo de inicialización. Si ejecuta canalizaciones desencadenadas o programadas con reinicios frecuentes, considere la posibilidad de usar el modo de streaming continuo para evitar la sobrecarga de inicialización repetida.

¿Por qué no agregar más ejecutores de Spark aumenta el rendimiento de Kafka?

Una vez que los agentes de Kafka se saturan, agregar más ejecutores de Spark aumenta el costo sin aumentar el rendimiento.

Señales de que Kafka es el cuello de botella:

  • El rendimiento se estanca a pesar de añadir más núcleos.
  • La utilización de la CPU o de la red del broker de Kafka es elevada.
  • Las tareas de Spark se completan rápidamente, pero esperan nuevos datos.

Para resolverlo, escale el clúster de Kafka agregando agentes o aumentando los recuentos de particiones para distribuir la carga.

¿Cómo puedo optimizar el costo y el uso de proceso para el streaming de Kafka?

Para los modos microlote y AvailableNow:

  • Tamaño adecuado del clúster: supervise las métricas y establezca un tamaño de clúster fijo adecuado para la carga máxima.
  • Uso maxOffsetsPerTrigger: limite los tamaños de lote para controlar el uso de recursos durante los picos de carga.
  • Evitar el escalado automático: los trabajos de streaming se ejecutan continuamente y la adición o eliminación de nodos provoca una sobrecarga de reequilibrio de tareas.
  • Reduce el sesgo de datos: Las particiones sesgadas hacen que algunas tareas procesen una cantidad de datos significativamente mayor que otras, lo que provoca retrasos que ralentizan la finalización global del lote y desperdician recursos de proceso en tareas inactivas. Use la minPartitions opción de dividir particiones grandes de Kafka en particiones spark más pequeñas para un procesamiento más equilibrado.

En el modo en tiempo real, el tamaño del proceso es especialmente importante porque las tareas pueden permanecer inactivas mientras esperan datos. Consideraciones clave:

  • Establézcalo maxPartitions para que cada tarea controle varias particiones de Kafka para reducir la sobrecarga.
  • Ajuste spark.sql.shuffle.partitions para trabajos con gran volumen de reordenación.

Consulte Dimensionamiento de la capacidad de cálculo para obtener orientación sobre cómo dimensionar los clústeres para el modo en tiempo real.

¿Por qué mi flujo no devuelve ningún registro aunque existan datos en el tema?

Las causas más comunes son:

  • Configuración startingOffsets incorrecta: el valor predeterminado es latest, que solo lee los nuevos datos que llegan después de que se inicie el flujo. Configure startingOffsets a earliest para leer los datos existentes.
  • Nombre de tema incorrecto: compruebe que se está suscribiendo al tema correcto.
  • Problemas de autenticación: Es posible que su flujo se haya conectado correctamente, pero carezca de permisos para leer el tema. Compruebe las ACL de Kafka.
  • Caducidad de los offsets: Si su flujo ha estado detenido durante mucho tiempo y los offsets del punto de control han caducado (han sido eliminados por la directiva de retención de Kafka), es posible que tenga que restablecer el punto de control o ajustar failOnDataLoss.