Características de rendimiento de las canalizaciones de Kafka a Iceberg

En esta página, se describen las características de rendimiento, a partir de la versión 2.75.0 de Apache Beam, para los trabajos de transmisión de Dataflow que leen desde Apache Kafka y escriben en tablas de Apache Iceberg. Se evalúan las diferencias de rendimiento entre las escrituras directas de Apache Iceberg y las escrituras enrutadas a través de la API de BigQuery administrada, y se comparan estos resultados con las comparativas de referencia de las canalizaciones de Kafka a BigQuery. Debido a que las optimizaciones para la E/S de Apache Iceberg están en curso, estas métricas de rendimiento están sujetas a cambios.

Las comparaciones de comparativas están disponibles en tres configuraciones de asignación sin estado principales (lo que significa que leen desde la fuente, convierten el mensaje en un registro y escriben en el receptor sin hacer un seguimiento del estado ni aplicar una lógica empresarial compleja; se hace referencia a map_only o mapping en las comparativas):

  1. Kafka a BigQuery (map_only) (comparativa de referencia del rendimiento de Kafka a BigQuery)
  2. Kafka a Iceberg directo (map_only, autosharding=false)
  3. Kafka a Iceberg con la API de BigQuery administrada (map_only)

Además, en esta guía, se evalúan los patrones de transmisión directa de Apache Iceberg, como el procesamiento por lotes con estado mediante groupbykey, y se detallan las consideraciones cruciales de nivel inferior con respecto a las distribuciones de tamaño de archivo, el comportamiento de fragmentación automática y la latencia de consulta del lado de lectura.

Metodología de prueba

Las comparativas se realizaron con los siguientes recursos:

  • Clúster de Managed Service para Apache Kafka: El tráfico se generó con la plantilla de Dataflow Streaming Data Generator.
    • Capacidad de procesamiento de entrada: 1 GBps
    • Tasa de mensajes: Aproximadamente 1,000,000 mensajes por segundo
    • Formato de mensaje: Texto JSON con un esquema fijo (aproximadamente 1 KB por mensaje)
    • Particiones: 1,000 particiones de Kafka
  • Receptores de destino:
    • BigQuery: Tabla estándar (sin particiones) escrita con la API de BigQuery Storage Write.
    • Apache Iceberg: Catálogo respaldado por Cloud Storage. El receptor directo se particiona con bucket(id, 64) (agrupado en 64 fragmentos en la clave primaria) y usa el modo de distribución hash.

Después de que se estabilizó el ajuste de escala automático horizontal, cada configuración de canalización se ejecutó en estado estable durante 24 horas. Las comparativas para cada caso de canalización se ejecutaron 3 veces por separado, y todos los valores informados representan los promedios calculados en esas ejecuciones para garantizar métricas de rendimiento confiables y sostenidas.

Rendimiento de la transferencia: Cargas de trabajo de asignación

Las canalizaciones de asignación sin estado leen desde la fuente, convierten el formato del mensaje en un registro y escriben en el receptor sin hacer un seguimiento del estado en los registros. En las siguientes secciones, se analizan las arquitecturas de referencia que se ejecutan a 1 GBps.

Configuración del trabajo

Configuración Kafka a BigQuery (map_only) Kafka a Iceberg directo (autosharding=false) Kafka a Iceberg con la API de BigQuery administrada
Tipo de máquina de trabajador e2-standard-2 e2-standard-4 e2-standard-4
CPU virtuales por trabajador 2 4 4
RAM por trabajador 8 GB 16 GB 16 GB
Streaming Engine Enabled Enabled Enabled
Ajuste de escala automático horizontal Enabled Enabled Enabled
Frecuencia de activación 5 segundos 60 segundos 60 segundos

Capacidad de procesamiento y uso de recursos

La escritura directa en archivos Parquet físicos en el almacén de objetos genera una sobrecarga de E/S más alta que la transferencia de transmisión de BigQuery. En comparación con las escrituras directas de Iceberg, el enrutamiento de escrituras a través de la API de BigQuery administrada mejora el uso de CPU del trabajador (~70% en comparación con ~60%) y reduce modestamente el consumo de Streaming Engine (~180 SECU/h en comparación con ~200 SECU/h), aunque los requisitos generales de procesamiento del trabajador siguen siendo similares (~440 CPU virtuales en comparación con ~450 CPU virtuales).

Métrica Kafka a BigQuery (map_only) Kafka a Iceberg directo (autosharding=false) Kafka a Iceberg con la API de BigQuery administrada
Capacidad de procesamiento de entrada promedio por trabajador ~15 MBps ~9 MBps ~9 MBps
Uso de CPU promedio ~70% ~60% ~70%
CPU virtuales estimadas para una entrada de 1 GBps ~126 CPU virtuales ~450 CPU virtuales ~440 CPU virtuales
Trabajadores estimados para una entrada de 1 GBps ~63 trabajadores ~110 trabajadores ~110 trabajadores
SECU estimadas por hora para 1 GBps ~58 SECU/h ~200 SECU/h ~180 SECU/h

Perfil de latencia de escritura

Las escrituras directas de Iceberg muestran una latencia final grave (P99) debido a las restricciones de confirmación de metadatos del almacenamiento de objetos. El uso de la API de BigQuery administrada elimina los picos de latencia final y mantiene una latencia media baja.

Latencia de escritura de extremo a extremo Kafka a BigQuery Kafka a Iceberg directo (autosharding=false) Kafka a Iceberg con la API de BigQuery administrada
P50 (mediana) ~1,200 ms ~1,000 ms ~1,000 ms
P95 ~3,000 ms ~7,400 ms ~1,900 ms
P99 (final) ~5,400 ms ~14,000 ms ~2,700 ms

Consideraciones y opciones de diseño de fragmentación automática

En esta sección, se analizan las implicaciones de la fragmentación automática en los tamaños de archivo y la latencia de la canalización cuando se escribe en Apache Iceberg.

Por qué se eligió autosharding=false como comparativa de referencia

En las pruebas iniciales, habilitar la fragmentación automática provocó que los tamaños de archivo se colapsaran en fragmentos pequeños y fluctuaran de forma arbitraria debido a la división dinámica de fragmentos activada por aumentos de carga localizados a nivel de subprocesos, incluso con una carga de entrada agregada constante.

Para mantener diseños de archivos Parquet estables y predecibles (~800 KB en promedio) y garantizar una comparativa de referencia justa sin vaciados prematuros, se eligió autosharding=false para la configuración del receptor directo.

¿Qué sucede si inhabilitas la fragmentación automática en comparación con mantenerla habilitada?

  • Con autosharding=false (comparativa de referencia): Obtienes tamaños de archivo iniciales más grandes (~800 KB en promedio) en comparación con la fragmentación automática. Si bien esto sigue siendo pequeño en comparación con los tamaños de archivo ideales de Iceberg (128–512 MB), requiere una compactación de nivel inferior significativamente menor. Sin embargo, la compensación es una latencia final de escritura alta (P99 que alcanza ~14.0 s) debido a los cuellos de botella de metadatos del almacenamiento de objetos.
  • Si la fragmentación automática está habilitada: Dataflow escala de forma dinámica los subprocesos de escritura para absorber los picos de capacidad de procesamiento local, lo que reduce la latencia final de escritura. Sin embargo, compromete la capa de almacenamiento al producir un gran volumen de archivos Parquet pequeños y fragmentados (~100 KB o menos). Estos tamaños de archivo muestran una alta varianza y fluctúan de forma arbitraria en las ejecuciones (con un promedio de ~39 KB a ~100 KB), lo que aumenta la necesidad de un mantenimiento agresivo de la compactación de nivel inferior.

Ajuste y recomendaciones de particiones

Durante nuestra evaluación, experimentamos con varios valores de partición fijos para la tabla de destino para encontrar un equilibrio óptimo. Descubrimos que el uso de 64 buckets (por ejemplo, bucket(id, 64)) para la partición de la tabla de destino producía los tamaños de archivo objetivo y mantenía una utilización y una capacidad de procesamiento decentes. Este enfoque nos ayudó a igualar los beneficios de rendimiento de la fragmentación automática y, al mismo tiempo, evitar los problemas de fragmentación de tamaño de archivo arbitrarios vinculados al ajuste de escala completamente dinámico.

Recomendación para profesionales: Se recomienda a los clientes que realicen pruebas preliminares similares con la configuración de partición objetivo para ubicar el punto óptimo que maximice el paralelismo de la canalización sin comprometer los tamaños de archivo Parquet.

Implicaciones de lectura de nivel inferior: Tamaños de archivo y compactación

Si bien las métricas del lado de escritura favorecen la API de BigQuery administrada para la transferencia de datos de Iceberg, la eficiencia general de la canalización depende en gran medida del rendimiento de lectura de nivel inferior:

  • Generación de archivos pequeños en la API de BigQuery administrada: La API de BigQuery administrada vacía los datos con frecuencia para garantizar una latencia de escritura baja. Este comportamiento genera un gran volumen de archivos Parquet pequeños escritos en el catálogo de Iceberg de destino.
  • Impacto en la latencia de consulta de lectura: Los motores de consultas (por ejemplo, Starburst/Trino, Apache Spark, BigQuery, Dremio) que leen tablas con millones de archivos Parquet pequeños generan una gran sobrecarga de análisis de metadatos y penalizaciones de análisis de particiones.
  • Requisitos de compactación: Para evitar la degradación del rendimiento de lectura cuando se usa la API de BigQuery administrada (o si la fragmentación automática está habilitada en las escrituras directas), ejecuta trabajos de mantenimiento de compactación de Iceberg con regularidad (por ejemplo, REWRITE DATA FILES). La sobrecarga de procesamiento para la compactación debe tenerse en cuenta en el diseño general de la arquitectura.
  • Distribución de archivos de escritura directa (autosharding=false): Las escrituras directas de Iceberg con fragmentación fija producen archivos Parquet promedio más grandes (~800 KB), lo que genera un diseño menos fragmentado para el acceso inmediato a la consulta sin demandas de compactación inmediatas (aunque aún por debajo del rango ideal).

Canalizaciones directas de Iceberg con estado (groupbykey)

Para evaluar las estrategias de procesamiento por lotes manuales, se probó la agrupación de claves con estado (groupbykey) en comparación con la canalización de Kafka a Iceberg directo (map_only, autosharding=false) de referencia. Ambas configuraciones escriben archivos Parquet directamente en el almacenamiento de objetos.

Comparación

Métrica o función Comparativa de referencia del receptor directo (autosharding=false) Receptor directo con estado (groupbykey) Impacto en el rendimiento
CPU virtuales estimadas para 1 GBps ~450 CPU virtuales ~520 CPU virtuales ~+16% de procesamiento requerido
Uso de CPU promedio ~60% ~50% ~-17% de eficiencia del trabajador
SECU/h estimadas para 1 GBps ~200 SECU/h ~300 SECU/h ~+50% de carga de Streaming Engine
Tamaño promedio de archivo ~800 KB ~100 KB Genera lotes de archivos más pequeños
Latencia de P50 ~1,000 ms ~1,200 ms ~+20% de mediana más lenta
Latencia de P95 ~7,400 ms ~5,500 ms ~-26% de latencia más baja
Latencia de P99 ~14,000 ms ~13,000 ms Cambio marginal de latencia final

Análisis de compensación

  1. Sobrecarga de Streaming Engine: Agregar un paso groupbykey con estado requiere que Beam almacene el estado intermedio en los límites de la ventana. Esto aumenta el consumo de unidades de procesamiento de Streaming Engine en ~50% (de ~200 SECU/h a ~300 SECU/h).
  2. Latencia de almacenamiento en búfer: La agregación manual de claves introduce el almacenamiento en búfer de ventanas obligatorio, lo que aumenta la latencia de escritura media (P50) a ~1,200 ms y la latencia de P95 a ~5.5 s.

Canalizaciones inversas: Transmisión de Iceberg a Kafka

Para evaluar las capacidades bidireccionales de lakehouse, también se realizaron comparativas para los datos de transmisión que fluyen en sentido inverso: leer transmisiones de solo anexos desde una tabla de Apache Iceberg y volver a publicar en Apache Kafka.

Configuración y eficiencia del trabajo

A diferencia de las canalizaciones de transferencia de datos que deben lidiar con escrituras de archivos de almacenamiento de objetos pesadas o cuellos de botella de confirmación de metadatos, la lectura y la transmisión de cambios fuera de Iceberg operan con alta eficiencia:

Métrica Iceberg a Kafka (solo anexos, exactamente una vez)
Tipo de máquina de trabajador e2-standard-4
CPU virtuales estimadas para una entrada de 1 GBps ~30 CPU virtuales
Trabajadores estimados para una entrada de 1 GBps ~7 trabajadores
SECU estimadas por hora para 1 GBps ~0.2 SECU/h

Conclusiones clave para las canalizaciones inversas

  • Sobrecarga de procesamiento significativamente menor: La lectura y la proyección de transmisiones de CDC desde Iceberg requieren muchos menos recursos de procesamiento (~30 CPU virtuales en comparación con ~450 CPU virtuales para escrituras directas) porque evita la tarea pesada de particionar, codificar y confirmar grandes volúmenes de archivos Parquet en el almacenamiento de objetos.
  • Eficiencia de recursos: El consumo o la replicación de nivel inferior basados en eventos desde un formato de lakehouse hasta las capas de transmisión son muy eficientes en comparación con las rutas de transferencia de datos de entrada.

Resumen de recomendaciones de arquitectura

Patrón de arquitectura Latencia de escritura de P99 Diseño de archivo Consideraciones de lectura de nivel inferior
Kafka a BigQuery (map_only) ~5.4s N/A Óptimo (motor de almacenamiento de BigQuery administrado)
Kafka a Iceberg con la API de BigQuery administrada ~2.7s Archivos arbitrariamente pequeños Requiere compactación periódica para lecturas de gran volumen
Kafka a Iceberg directo (autosharding=false) ~14.0s ~800 KB Bueno (tamaños de archivo iniciales más grandes, menor demanda de compactación)
Kafka a Iceberg directo (groupbykey) ~13.0s ~100 KB Moderado (mayor sobrecarga de procesamiento y estado)

Estimación de costos

Puedes estimar el costo de referencia de tu propia canalización comparable con la facturación basada en recursos mediante la Google Cloud calculadora de precios, de la siguiente manera:

  1. Abre la calculadora de precios.
  2. Haz clic en Agregar a la estimación.
  3. Selecciona Dataflow.
  4. En Tipo de servicio, selecciona "Dataflow Classic".
  5. Selecciona Configuración avanzada para mostrar el conjunto completo de opciones.
  6. Elige la ubicación en la que se ejecuta el trabajo.
  7. En Tipo de trabajo, selecciona "Transmisión".
  8. Selecciona Habilitar Streaming Engine.
  9. Ingresa información sobre las horas de ejecución del trabajo, los nodos trabajadores, las máquinas de trabajador y el almacenamiento en Persistent Disk.
  10. Ingresa la cantidad estimada de unidades de procesamiento de Streaming Engine.

El uso de recursos y el costo se escalan de forma aproximadamente lineal con la capacidad de procesamiento de entrada, aunque, para trabajos pequeños con solo unos pocos trabajadores, el costo total está dominado por los costos fijos. Como punto de partida, puedes extrapolar la cantidad de nodos trabajadores y el consumo de recursos de los resultados de la comparativa.

Por ejemplo, supongamos que ejecutas una canalización con la arquitectura de Kafka a Iceberg directo (autosharding=false), con una tasa de datos de entrada de 100 MBps. Según los resultados de la comparativa para una canalización de 1 GBps, puedes estimar los requisitos de recursos de la siguiente manera:

  • Factor de escala: (100 MBps) / (1024 MBps) = ~0.1
  • Nodos trabajadores proyectados: 110 trabajadores × 0.1 = ~11 trabajadores
  • Cantidad proyectada de unidades de procesamiento de Streaming Engine por hora: 200 × 0.1 = ~20 unidades por hora

Este valor solo debe usarse como una estimación inicial. La capacidad de procesamiento y el costo reales pueden variar significativamente, según factores como el tipo de máquina, la distribución del tamaño del mensaje, el código de usuario, el tipo de agregación, el paralelismo de claves y el tamaño de la ventana. Para obtener más información, consulta Prácticas recomendadas para la optimización de costos de Dataflow.

Ejecuta una canalización de prueba

Para implementar un trabajo de transmisión de Apache Iceberg con la plantilla de Flex de Dataflow, usa el gcloud dataflow flex-template run comando.

gcloud dataflow flex-template run JOB_NAME \
  --project=PROJECT_ID \
  --region=REGION \
  --template-file-gcs-location=gs://dataflow-templates-us-central1/latest/flex/Kafka_To_Iceberg_Yaml \
  --enable-streaming-engine \
  --parameters ^@^bootstrapServers="KAFKA_BOOTSTRAP_ADDRESS"\
@topic="KAFKA_TOPIC"\
@table="ICEBERG_TABLE_IDENTIFIER"\
@catalogName="CATALOG_NAME"\
@catalogProperties='{"type":"CATALOG_TYPE","warehouse":"gs://BUCKET_NAME/warehouse/"}'\
@triggeringFrequencySeconds=60\
@schema='SCHEMA_DEFINITION'

Reemplaza lo siguiente:

  • JOB_NAME: Es el nombre de tu trabajo de Dataflow.
  • PROJECT_ID: Es el ID de tu Google Cloud proyecto.
  • REGION: Es la Google Cloud región en la que se ejecuta tu trabajo (por ejemplo, us-central1).
  • KAFKA_BOOTSTRAP_ADDRESS: Es la dirección de arranque de tu clúster de Apache Kafka.
  • KAFKA_TOPIC: Es el nombre de tu tema de Kafka.
  • ICEBERG_TABLE_IDENTIFIER: Es el identificador de tu tabla de Iceberg de destino.
  • CATALOG_NAME: Es el nombre de tu catálogo de Iceberg.
  • CATALOG_TYPE: Es el tipo de catálogo que se usará (por ejemplo, hadoop o bigquery).
  • BUCKET_NAME: Es el nombre del bucket de Cloud Storage para la ubicación de tu almacén.
  • SCHEMA_DEFINITION: Es la definición de esquema para los datos de tu tema de Kafka (por ejemplo, {"type": "record", "name": "Record", "fields": [{"name": "id", "type": "string"}]}).