En este documento, se describe cómo leer datos de Apache Kafka en Dataflow y se incluyen sugerencias de rendimiento y prácticas recomendadas.
Para la mayoría de los casos de uso, considera usar el conector de E/S administrado para leer desde Kafka.
Si necesitas un ajuste de rendimiento más avanzado, considera usar el
KafkaIO conector. El conector KafkaIO está disponible para
Java
o mediante el
framework de canalizaciones multilingüe
para Python
y Go.
Paralelismo
En las siguientes secciones, se describe cómo configurar el paralelismo cuando se lee desde Kafka.
Descripción general
El paralelismo está limitado por dos factores: la
cantidad máxima de trabajadores
(max_num_workers) y la cantidad de particiones de Kafka. Dataflow usa de forma predeterminada un fan-out de paralelismo de 4 × max_num_workers. Sin embargo, el fan-out está limitado por la cantidad de particiones. Por ejemplo, si hay 100 CPU virtuales disponibles, pero la canalización solo lee de 10 particiones de Kafka, el paralelismo máximo es 10.
Para maximizar el paralelismo, se recomienda tener al menos 4 × max_num_workers particiones de Kafka. Si tu trabajo usa
Portable Runner, considera establecer un paralelismo aún mayor.
Un buen punto de partida es tener particiones iguales al doble de la cantidad de CPU virtuales de trabajador.
Redistribuye
Si no puedes aumentar la cantidad de particiones, puedes aumentar el paralelismo llamando a
KafkaIO.Read.withRedistribute. Este método agrega una transformación Redistribute a la canalización, que proporciona una sugerencia a Dataflow para redistribuir y paralelizar los datos de manera más eficiente. Te recomendamos que especifiques la cantidad óptima de fragmentos llamando a
KafkaIO.Read.withRedistributeNumKeys. Usar
KafkaIO.Read.withRedistribute solo puede generar numerosas
claves, lo que genera problemas de rendimiento. Para obtener más información, consulta
Identifica etapas con paralelismo alto.
La redistribución de los datos agrega una sobrecarga adicional para realizar el paso de Shuffle. Para obtener más información, consulta
Evita la fusión.
Para minimizar el costo del Shuffle de redistribución, llama a
KafkaIO.Read.withOffsetDeduplication. Este modo minimiza la cantidad de datos que deben conservarse como parte del Shuffle, al tiempo que proporciona un procesamiento “exactamente una vez”.
Si no se requiere el procesamiento “exactamente una vez”, puedes permitir duplicados llamando a
KafkaIO.Read.withAllowDuplicates.
En la siguiente tabla, se resumen las opciones de redistribución:
| Opción | Modo de procesamiento | Apache Beam | Configuración |
|---|---|---|---|
| Redistribuye la entrada | Exactamente una vez | v2.60 o versiones posteriores | KafkaIO.Read.withRedistribute() |
| Permite duplicados | Al menos una vez | v2.60 o versiones posteriores | KafkaIO.Read.withRedistribute().withAllowDuplicates() |
| Anulación de duplicación de compensación | Exactamente una vez | v2.69 o versiones posteriores | KafkaIO.Read.withRedistribute().withOffsetDeduplication() |
Sesgo de carga
Intenta asegurarte de que la carga entre las particiones sea relativamente uniforme y no esté sesgada. Si la carga está sesgada, puede generar un uso deficiente de los trabajadores. Los trabajadores que leen de particiones con una carga más ligera pueden estar relativamente inactivos, mientras que los trabajadores que leen de particiones con una carga pesada pueden quedar rezagados. Dataflow proporciona métricas para el backlog por partición.
Si la carga está sesgada, el balanceo de carga dinámico puede ayudar a distribuir el trabajo. Por ejemplo, Dataflow puede asignar un trabajador para leer de varias particiones de bajo volumen y asignar otro trabajador para leer de una sola partición de gran volumen. Sin embargo, dos trabajadores no pueden leer de la misma partición, por lo que una partición con mucha carga aún puede hacer que la canalización quede rezagada.
Prácticas recomendadas
En esta sección, se incluyen recomendaciones para leer desde Kafka en Dataflow.
Temas de bajo volumen
Una situación común es leer de muchos temas de bajo volumen al mismo tiempo, por ejemplo, un tema por cliente. Crear trabajos de Dataflow independientes para cada tema es ineficiente en cuanto a costos, ya que cada trabajo requiere al menos un trabajador completo. En su lugar, considera las siguientes opciones:
Combinar temas Combina temas antes de que se transfieran a Dataflow. Transferir algunos temas de gran volumen es mucho más eficiente que transferir muchos temas de bajo volumen. Cada tema de gran volumen puede ser controlado por un solo trabajo de Dataflow que utiliza por completo sus trabajadores.
Lee varios temas. Si no puedes combinar temas antes de transferirlos a Dataflow, considera crear una canalización que lea desde varios temas. Este método permite que Dataflow asigne varios temas al mismo trabajador. Existen dos formas de implementar este método:
Paso de lectura único. Crea una sola instancia del conector
KafkaIOy configúrala para leer varios temas. Luego, filtra por nombre de tema para aplicar una lógica diferente por tema. Para obtener un código de ejemplo, consulta Lee de varios temas. Considera esta opción si todos tus temas están ubicados en el mismo clúster. Una desventaja es que los problemas con un solo sink o transformación pueden hacer que todos los temas acumulen un backlog.Para casos de uso más avanzados, pasa un conjunto de objetos
KafkaSourceDescriptorque especifiquen los temas desde los que se debe leer. El uso deKafkaSourceDescriptorte permite actualizar la lista de temas más adelante si es necesario. Esta función requiere Java con Portable Runner.Varios pasos de lectura. Para leer de temas ubicados en diferentes clústeres, tu canalización puede incluir varias instancias de
KafkaIO. Mientras se ejecuta el trabajo, puedes actualizar fuentes individuales con asignaciones de transformación. Solo se admite la configuración de un tema o clúster nuevo cuando se usa Portable Runner. La observabilidad es un desafío potencial con este método, ya que debes supervisar cada transformación de lectura individual en lugar de depender de las métricas a nivel de la canalización.
Confirmación en Kafka
De forma predeterminada, el conector KafkaIO no usa compensaciones de Kafka para hacer un seguimiento del progreso y no confirma en Kafka. Si llamas a
commitOffsetsInFinalize, el conector hace todo lo posible para confirmar en Kafka después de que se confirman los registros en
Dataflow. Es posible que los registros confirmados en Dataflow no
se procesen por completo, por lo que, si
cancelas la canalización, se puede confirmar una compensación
sin que los registros se procesen por completo.
Debido a que la configuración de enable.auto.commit=True confirma las compensaciones en cuanto se leen de Kafka sin ningún procesamiento por parte de Dataflow, no se recomienda usar esta opción.
La recomendación es establecer enable.auto.commit=False y
commitOffsetsInFinalize=True. Si estableces enable.auto.commit en True, se pueden perder datos si se interrumpe la canalización durante el procesamiento. Es posible que se descarten los registros ya confirmados en Kafka.
Marcas de agua
De forma predeterminada, el conector KafkaIO usa la hora de procesamiento actual para asignar
la marca de agua de salida
y la hora del evento. Para cambiar este comportamiento, llama a
withTimestampPolicyFactory y asigna una
TimestampPolicy. Beam proporciona implementaciones de TimestampPolicy que calculan la marca de agua en función de la hora de agregado de registros de Kafka o la hora de creación del mensaje.
Consideraciones sobre el ejecutor
El conector KafkaIO tiene dos implementaciones subyacentes para las lecturas de Kafka, la
más antigua ReadFromKafkaViaUnbounded y la más reciente
ReadFromKafkaViaSDF. Dataflow elige automáticamente la mejor implementación para tu trabajo en función del lenguaje del SDK y los requisitos del trabajo. Evita solicitar de forma explícita un ejecutor o una implementación de Kafka, a menos que necesites funciones específicas que solo estén disponibles en esa implementación. Para obtener más información sobre cómo elegir un ejecutor, consulta
Usa Dataflow Portable Runner.
Si tu canalización usa withTopic o withTopics,
la implementación anterior consulta a Kafka en el momento de la construcción de la canalización para obtener las
particiones disponibles. La máquina que crea la canalización debe tener permiso para conectarse a Kafka. Si recibes un error de permiso, verifica que tengas permisos para conectarte a Kafka de forma local. Puedes evitar este problema usando
withTopicPartitions, que no se conecta a Kafka
en el momento de la construcción de la canalización.
Implementar en producción
Cuando implementes tu solución en producción, te recomendamos que uses plantillas de Flex. Cuando se usa una plantilla de Flex, la canalización se inicia desde un entorno coherente, lo que puede ayudar a mitigar los problemas de configuración local.
El registro de KafkaIO puede ser bastante detallado. Considera reducir el nivel de registro en producción de la siguiente manera:
sdkHarnessLogLevelOverrides='{"org.apache.kafka.clients.consumer.internals.SubscriptionState":"WARN"}'.
Para obtener más información, consulta Configura los niveles de registro del trabajador de canalización.
Configura Herramientas de redes
De forma predeterminada, Dataflow inicia instancias en la red de nube privada virtual (VPC) predeterminada. Según tu configuración de Kafka, es posible que debas configurar una red y una subred diferentes para Dataflow. Para obtener más información, consulta Especifica una red y una subred. Cuando configures la red, crea reglas de firewall que permitan que las máquinas de trabajador de Dataflow lleguen a los agentes de Kafka.
Si usas los Controles del servicio de VPC, coloca el clúster de Kafka dentro del perímetro de los Controles del servicio de VPC o extiende los perímetros a la VPN autorizada o a Cloud Interconnect.
Si tu clúster de Kafka se implementa fuera de Google Cloud, debes crear una conexión de red entre Dataflow y el clúster de Kafka. Existen varias opciones de herramientas de redes con diferentes compensaciones:
- Conéctate con un espacio de direcciones RFC 1918 compartido mediante una de las siguientes opciones:
- Accede al clúster de Kafka alojado de forma externa a través de direcciones IP públicas con una de las siguientes opciones:
- Internet pública
- Intercambio de tráfico directo
- Intercambio de tráfico con proveedores
La interconexión dedicada es la mejor opción para un rendimiento y una confiabilidad predecibles, pero puede tardar más en configurarse porque los terceros deben aprovisionar los circuitos nuevos. Con una topología según la IP pública, puedes comenzar con rapidez porque hay poco trabajo de herramientas de redes para hacer.
En las siguientes dos secciones, se describen estas opciones con más detalle.
Espacio de direcciones de RFC 1918 compartido
La interconexión dedicada y la VPN de IPsec te brindan acceso directo a las direcciones IP RFC 1918 en la nube privada virtual (VPC), lo que puede simplificar la configuración de Kafka. Si usas una topología basada en VPN, considera configurar una VPN de alta capacidad de procesamiento.
De forma predeterminada, Dataflow inicia instancias en la red de VPC predeterminada
. En una topología de red privada con
rutas definidas de forma explícita en Cloud Router
que conectan subredes en Google Cloud a ese clúster de Kafka, necesitas
más control sobre dónde ubicar las instancias de Dataflow. Puedes
usar Dataflow para configurar los network y subnetwork
parámetros de ejecución.
Asegúrate de que la subred correspondiente tenga suficientes direcciones IP disponibles para que Dataflow inicie instancias mientras intenta escalar horizontalmente. Además, cuando crees una red independiente para iniciar las instancias de Dataflow, asegúrate de tener una regla de firewall que habilite el tráfico de TCP entre todas las máquinas virtuales del proyecto. La red predeterminada ya tiene configurada esta regla de firewall.
Espacio de direcciones IP públicas
Esta arquitectura usa la seguridad de la capa de transporte (TLS) para proteger el tráfico entre clientes externos y Kafka, y usa texto simple en la comunicación entre agentes. Cuando el objeto de escucha de Kafka se vincula a una interfaz de red que se usa para la comunicación interna y externa, la configuración del objeto de escucha es sencilla. Sin embargo, en muchas situaciones, las direcciones anunciadas de forma externa de los agentes de Kafka en el clúster difieren de las interfaces de red internas que usa Kafka. En tales casos, puedes usar la propiedad advertised.listeners:
# Configure protocol map listener.security.protocol.map=INTERNAL:PLAINTEXT,EXTERNAL:SSL
# Use plaintext for inter-broker communication inter.broker.listener.name=INTERNAL
# Specify that Kafka listeners should bind to all local interfaces listeners=INTERNAL://0.0.0.0:9092,EXTERNAL://0.0.0.0:9093
# Separately, specify externally visible address advertised.listeners=INTERNAL://kafkabroker-n.mydomain.com:9092,EXTERNAL://kafkabroker-n.mydomain.com:9093
Los clientes externos se conectan mediante el puerto 9093 a través de un canal SSL y los clientes internos se conectan mediante el puerto 9092 a través de un canal de texto simple. Cuando especifiques una dirección en advertised.listeners, usa nombres de DNS (kafkabroker-n.mydomain.com, en esta muestra) que se resuelven en la misma instancia para el tráfico interno y externo. Es posible que el uso de direcciones IP públicas no funcione porque las direcciones no pueden resolverse para el tráfico interno.
Ajusta Kafka
La configuración del clúster de Kafka y del cliente de Kafka puede tener un gran impacto en el rendimiento. En particular, es posible que la siguiente configuración sea demasiado baja. En esta sección, se proporcionan algunos puntos de partida sugeridos, pero debes experimentar con estos valores para tu carga de trabajo en particular.
unboundedReaderMaxElements. La configuración predeterminada es 10,000. Un valor más alto, como 100,000, puede aumentar el tamaño de los paquetes, lo que puede mejorar el rendimiento de manera significativa si tu canalización incluye agregaciones. Sin embargo, un valor más alto también puede aumentar la latencia. Para establecer el valor, usasetUnboundedReaderMaxElements. Este parámetro de configuración no se aplica a Portable Runner. Para Portable Runner, usa lasdf_checkpoint_after_output_bytesopción de servicio de Dataflow.unboundedReaderMaxReadTimeMs. El valor predeterminado es 10,000 ms. Un valor más alto, como 20,000 ms, puede aumentar el tamaño del paquete, mientras que un valor más bajo, como 5,000 ms, puede reducir la latencia o el retraso. Para establecer el valor, usasetUnboundedReaderMaxReadTimeMs. Este parámetro de configuración no se aplica a Portable Runner. Para Portable Runner, usa lasdf_checkpoint_after_durationopción de servicio de Dataflow.max.poll.records. La configuración predeterminada es 500. Un valor más alto puede funcionar mejor si recupera más registros entrantes juntos, en especial cuando se usa Portable Runner. Para establecer el valor, llama awithConsumerConfigUpdates.fetch.max.bytes. Valor predeterminado 1 MB. Un valor más alto puede mejorar la capacidad de procesamiento si reduce la cantidad de solicitudes, en especial cuando se usa Portable Runner. Sin embargo, si lo configuras demasiado alto, es posible que aumente la latencia, aunque es más probable que el procesamiento descendente sea el cuello de botella principal. Un valor inicial recomendado es de 100 MB. Para establecer el valor, llama awithConsumerConfigUpdates.max.partition.fetch.bytes. Valor predeterminado 1 MB. Este parámetro establece la cantidad máxima de datos por partición que muestra el servidor. Si aumentas el valor, puedes mejorar la capacidad de procesamiento si reduces la cantidad de solicitudes, en especial cuando se usa Portable Runner. Sin embargo, si lo configuras demasiado alto, es posible que aumente la latencia, aunque es más probable que el procesamiento descendente sea el cuello de botella principal. Un valor inicial recomendado es de 100 MB. Para establecer el valor, llama awithConsumerConfigUpdates.consumerPollingTimeout. La configuración predeterminada es de 2 segundos. Si el cliente consumidor se agota el tiempo de espera antes de poder leer un registro, intenta establecer un valor más alto. Este parámetro de configuración suele ser más relevante cuando se realizan lecturas entre regiones o lecturas con una red lenta. Para establecer el valor, llama awithConsumerPollingTimeout.
Asegúrate de que receive.buffer.bytes sea lo suficientemente grande para controlar el tamaño de los mensajes. Si el valor es demasiado pequeño, es posible que los registros muestren que los consumidores se vuelven a crear de forma continua y buscan una compensación específica.
Ejemplos
En los siguientes ejemplos de código, se muestra cómo crear canalizaciones de Dataflow que leen desde Kafka. Cuando se usan las credenciales predeterminadas de la aplicación junto con el
controlador de devolución de llamada proporcionado por Google Cloud Managed Service para Apache Kafka, se requiere la versión 3.7.0 o posterior de kafka-clients.
Lee desde un solo tema
En este ejemplo, se usa el conector de E/S administrado. Muestra cómo leer desde un tema de Kafka y escribir las cargas útiles de mensajes en archivos de texto.
Java
Para autenticarte en Dataflow, configura las credenciales predeterminadas de la aplicación. Si deseas obtener más información, consulta Configura la autenticación para un entorno de desarrollo local.
Python
Para autenticarte en Dataflow, configura las credenciales predeterminadas de la aplicación. Si deseas obtener más información, consulta Configura la autenticación para un entorno de desarrollo local.
Lee desde varios temas
En este ejemplo, se usa el conector KafkaIO. Muestra cómo leer desde varios temas de Kafka y aplicar una lógica de canalización independiente para cada tema.
Para casos de uso más avanzados, pasa de forma dinámica un conjunto de
KafkaSourceDescriptor objetos, de modo que puedas actualizar
la lista de temas desde los que se debe leer. Este método requiere Java con Portable Runner.
Java
Para autenticarte en Dataflow, configura las credenciales predeterminadas de la aplicación. Si deseas obtener más información, consulta Configura la autenticación para un entorno de desarrollo local.
Python
Para autenticarte en Dataflow, configura las credenciales predeterminadas de la aplicación. Si deseas obtener más información, consulta Configura la autenticación para un entorno de desarrollo local.
¿Qué sigue?
- Características de rendimiento de las canalizaciones de Apache Kafka a BigQuery
- Escribe en Apache Kafka.
- Obtén más información sobre la E/S administrada.