En esta página, se describen las características de rendimiento de los trabajos de transmisión de Dataflow que leen desde Apache Kafka y escriben en BigQuery. Se proporcionan los resultados de las pruebas comparativas para las canalizaciones de solo asignación, que realizan transformaciones por mensaje sin hacer un seguimiento del estado ni agrupar elementos en la transmisión.
Muchas cargas de trabajo de integración de datos, incluidas ETL, la validación de campos y la asignación de esquemas, pertenecen a la categoría de solo asignación. Si tu canalización sigue este patrón, puedes usar estas comparativas para evaluar tu trabajo de Dataflow en comparación con una configuración de referencia de buen rendimiento.
Metodología de prueba
Las comparativas se realizaron con los siguientes recursos:
Un clúster de Managed Service para Apache Kafka. Los mensajes se generaron con la plantilla del generador de datos de transmisión.
- Tasa de mensajes: Aproximadamente 1,000,000 mensajes por segundo
- Carga de entrada: 1 GiB/s
- Formato de mensaje: Texto JSON generado de forma aleatoria con un esquema fijo
- Tamaño del mensaje: Aproximadamente 1 KiB por mensaje
- Particiones de Kafka: 1,000
Una tabla estándar de BigQuery.
Una canalización de transmisión de Dataflow que usó la plantilla de Apache Kafka a BigQuery. Esta canalización realiza el análisis y la asignación de esquemas mínimos requeridos. No se usó ninguna función definida por el usuario (UDF) personalizada.
Después de que se estabilizó el ajuste de escala horizontal y la canalización alcanzó un estado estable, se permitió que las canalizaciones se ejecutaran durante aproximadamente un día, después del cual se recopilaron y analizaron los resultados.
Canalización de Dataflow
Esta comparativa usa una canalización de solo asignación que realiza una asignación y conversión simples de mensajes JSON. La canalización se probó con el modo de "exactamente una vez" y el modo de "al menos una vez". El procesamiento de "al menos una vez" proporciona una mejor capacidad de procesamiento. Sin embargo, solo se debe usar cuando los registros duplicados son aceptables o el receptor descendente controla la deduplicación.
Configuración del trabajo
En la siguiente tabla, se muestra cómo se configuraron los trabajos de Dataflow.
| Configuración | Valor |
|---|---|
| Tipo de máquina del trabajador | e2-standard-2 |
| CPU virtuales para trabajadores | 2 |
| RAM de la máquina del trabajador | 8 GB |
| Persistent Disk de la máquina del trabajador | Persistent Disk estándar (HDD), 30 GB |
| Cantidad máxima de trabajadores | 120 |
| Streaming Engine | Sí |
| Ajuste de escala automático horizontal | Sí |
| Modelo de facturación | Facturación basada en recursos |
| ¿Está habilitada la API de Storage Write (gRPC)? | Sí |
| Transmisiones de la API de Storage Write (gRPC) | 400 |
| Frecuencia de activación de la API de Storage Write (gRPC) | 5 segundos |
| Formato de mensaje | JSON |
| Modo de autenticación de Kafka |
Credenciales predeterminadas de la aplicación (ADC). Para obtener más información, consulta Tipos de autenticación para agentes de Kafka. |
Se recomienda la API de BigQuery Storage Write (gRPC) para las canalizaciones de transmisión. Cuando usas el modo de "exactamente una vez" con la API de Storage Write (gRPC), puedes ajustar la siguiente configuración:
Cantidad de transmisiones de escritura. Para garantizar un paralelismo de claves suficiente en la etapa de escritura, establece la cantidad de transmisiones de la API de Storage Write (gRPC) en un valor mayor que la cantidad de CPU del trabajador, mientras sigues las recomendaciones de capacidad de procesamiento por transmisión.
Frecuencia de activación. Un valor de segundo de un solo dígito es adecuado para canalizaciones de alta capacidad de procesamiento.
Para obtener más información, consulta Escribe desde Dataflow a BigQuery.
También se debe tener en cuenta la cantidad de particiones de Apache Kafka. Para garantizar un paralelismo de claves suficiente en la etapa de lectura, la cantidad de particiones debe ser al menos igual a la cantidad total de CPU virtuales del trabajador. Para obtener más información, consulta Lee desde Apache Kafka a Dataflow.
Resultados de comparativas
En esta sección, se describen los resultados de las pruebas comparativas.
Capacidad de procesamiento y uso de recursos
En la siguiente tabla, se muestran los resultados de las pruebas de capacidad de procesamiento y uso de recursos de la canalización.
| Resultado | Exactamente una vez | Al menos una vez |
|---|---|---|
| Capacidad de procesamiento de entrada por trabajador | Media: 15 MB/s, n=3 | Media: 18 MB/s, n=3 |
| Uso promedio de CPU en todos los trabajadores | Media: 70%, n=3 | Media: 75%, n=3 |
| Cantidad de nodos trabajadores | Media: 63, n=3 | Media: 53, n=3 |
| Unidades de procesamiento de Streaming Engine por hora | Media: 58, n=3 | Media: 0, n=3 |
El algoritmo de ajuste de escala automático puede afectar el nivel de uso de CPU objetivo. Para lograr un uso de CPU objetivo más alto o más bajo, puedes establecer el rango de ajuste de escala automático o la sugerencia de uso del trabajador. Los objetivos de uso más altos pueden generar costos más bajos, pero también una peor latencia de cola, en especial para cargas variables.
Latencia
En la siguiente tabla, se muestran los resultados de las comparativas de latencia de la canalización para el modo de "exactamente una vez", sin incluir la etapa de entrada.
| Latencia total de extremo a extremo de la etapa, sin incluir la etapa de entrada | Exactamente una vez |
|---|---|
| P50 | Media: 1,200 ms, n=3 |
| P95 | Media: 3,000 ms, n=3 |
| P99 | Media: 5,400 ms, n=3 |
Las pruebas midieron la latencia de extremo a extremo por etapa (la
job/streaming_engine/stage_end_to_end_latencies
métrica) en tres ejecuciones de prueba de larga duración. Esta métrica mide cuánto tiempo pasa Streaming Engine en cada etapa de la canalización. Abarca todos los pasos internos de la canalización, como los siguientes:
- Mezclar y poner en cola mensajes para el procesamiento
- El tiempo de procesamiento real; por ejemplo, convertir mensajes en objetos de fila
- Escribir estado persistente, así como el tiempo que se pasa en la cola para escribir estado persistente
Debido a una limitación de la métrica, no se informa la latencia de la etapa de entrada. Por lo tanto, no se incluye en el total.
Las comparativas que se muestran aquí representan una línea de base. La latencia es muy sensible a la complejidad de la canalización. Las UDF personalizadas, las transformaciones adicionales y la lógica de ventanas compleja pueden aumentar la latencia.
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:
- Abre la calculadora de precios.
- Haz clic en Agregar a la estimación.
- Selecciona Dataflow.
- En Tipo de servicio, selecciona "Dataflow Classic".
- Selecciona Configuración avanzada para mostrar el conjunto completo de opciones.
- Elige la ubicación en la que se ejecuta el trabajo.
- En Tipo de trabajo, selecciona "Transmisión".
- Selecciona Habilitar Streaming Engine.
- Ingresa información sobre las horas de ejecución del trabajo, los nodos trabajadores, las máquinas del trabajador y el almacenamiento de Persistent Disk.
- Ingresa la cantidad estimada de unidades de procesamiento de Streaming Engine.
El uso de recursos y el costo se ajustan 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 las comparativas.
Por ejemplo, supongamos que ejecutas una canalización de solo asignación en el modo de "exactamente una vez", con una tasa de datos de entrada de 100 MiB/s. Según los resultados de las comparativas para una canalización de 1 GiB/s, puedes estimar los requisitos de recursos de la siguiente manera:
- Factor de escala: (100 MiB/s) / (1 GiB/s) = 0.1
- Nodos trabajadores proyectados: 63 trabajadores × 0.1 = 6.3 trabajadores
- Cantidad proyectada de unidades de procesamiento de Streaming Engine por hora: 58 × 0.1 = 5.8 unidades por hora
Este valor solo debe usarse como una estimación inicial. La capacidad de procesamiento y el costo reales pueden variar de manera significativa, según factores como el tipo de máquina, la distribución del tamaño de los mensajes, 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
En esta sección, se muestran los
gcloud dataflow flex-template run
comandos que se usaron para ejecutar la canalización de solo asignación.
Modo de "exactamente una vez"
gcloud dataflow flex-template run JOB_NAME \
--project=PROJECT_ID \
--template-file-gcs-location=gs://dataflow-templates-us-central1/latest/flex/Kafka_to_BigQuery_Flex \
--enable-streaming-engine \
--parameters \
readBootstrapServerAndTopic="KAFKA_BOOTSTRAP_ADDRESS;KAFKA_TOPIC",\
kafkaReadAuthenticationMode=APPLICATION_DEFAULT_CREDENTIALS,\
messageFormat=JSON,\
writeMode=SINGLE_TABLE_NAME,\
outputTableSpec="PROJECT_ID:BQ_DATASET.BQ_TABLE_NAME",\
useBigQueryDLQ=true,\
outputDeadletterTable="PROJECT_ID:BQ_DATASET.BQ_TABLE_NAME_dlq",\
numStorageWriteApiStreams=400
Modo de "al menos una vez"
gcloud dataflow flex-template run JOB_NAME \
--project=PROJECT_ID \
--template-file-gcs-location=gs://dataflow-templates-us-central1/latest/flex/Kafka_to_BigQuery_Flex \
--enable-streaming-engine \
--additional-experiments=streaming_mode_at_least_once \
--parameters \
readBootstrapServerAndTopic="KAFKA_BOOTSTRAP_ADDRESS;KAFKA_TOPIC",\
kafkaReadAuthenticationMode=APPLICATION_DEFAULT_CREDENTIALS,\
messageFormat=JSON,\
writeMode=SINGLE_TABLE_NAME,\
outputTableSpec="PROJECT_ID:BQ_DATASET.BQ_TABLE_NAME",\
useBigQueryDLQ=true,\
outputDeadletterTable="PROJECT_ID:BQ_DATASET.BQ_TABLE_NAME_dlq",\
numStorageWriteApiStreams=400,\
useStorageWriteApiAtLeastOnce=true
Reemplaza lo siguiente:
JOB_NAME: El nombre del trabajo de DataflowPROJECT_ID: El ID del proyectoKAFKA_BOOTSTRAP_ADDRESS: La dirección de arranque del clúster de Apache KafkaKAFKA_TOPIC: El nombre del tema de KafkaBQ_DATASET: El nombre del conjunto de datos de BigQueryBQ_TABLE_NAME: El nombre de la tabla de BigQuery
Genere datos a partir de las pruebas
Para generar datos de prueba, usa el siguiente comando para ejecutar la plantilla del generador de datos de transmisión:
gcloud dataflow flex-template run JOB_NAME \
--project=PROJECT_ID \
--template-file-gcs-location=gs://dataflow-templates-us-central1/latest/flex/Streaming_Data_Generator \
--max-workers=140 \
--parameters \
schemaLocation=SCHEMA_LOCATION,\
qps=1000000,\
sinkType=KAFKA,\
bootstrapServer=KAFKA_BOOTSTRAP_ADDRESS,\
kafkaTopic=KAFKA_TOPIC,\
outputType=JSON
Reemplaza lo siguiente:
JOB_NAME: El nombre del trabajo de DataflowPROJECT_ID: El ID del proyectoSCHEMA_LOCATION: La ruta de acceso a un archivo de esquema en Cloud StorageKAFKA_BOOTSTRAP_ADDRESS: La dirección de arranque del clúster de Apache KafkaKAFKA_TOPIC: El nombre del tema de Kafka
La plantilla del generador de datos de transmisión usa un archivo JSON Data Generator para definir el esquema de mensajes. Las pruebas comparativas usaron un esquema de mensajes similar al siguiente:
{ "logStreamId": "{{integer(1000001,2000000)}}", "message": "{{alphaNumeric(962)}}" }
Próximos pasos
- Usa la interfaz de supervisión de trabajos de Dataflow
- Prácticas recomendadas para la optimización de costos de Dataflow
- Soluciona problemas de trabajos de transmisión lentos o atascados
- Lee desde Apache Kafka a Dataflow
- Escribe desde Dataflow a BigQuery