Dataflow está diseñado para ejecutar grandes canalizaciones de procesamiento de datos, ya que distribuye el trabajo en un grupo administrado de instancias de procesamiento. Comprender cómo Dataflow paraleliza el procesamiento te ayuda a diseñar canalizaciones eficientes, evitar cuellos de botella en el rendimiento y optimizar los costos de los recursos.
En esta página, se explica cómo Dataflow paraleliza el procesamiento de datos, cómo administra y escala la ejecución, los factores comunes que limitan el paralelismo y las técnicas que puedes usar para optimizar la capacidad de procesamiento de la canalización.
Modelos de paralelismo: horizontal frente a vertical
Dataflow logra el paralelismo con dos estrategias complementarias:
Paralelismo horizontal: Los datos de la canalización se particionan y procesan en varias instancias de trabajador (máquinas virtuales) de forma simultánea. Dataflow puede ajustar automáticamente el tamaño del grupo de trabajadores según la demanda de la carga de trabajo a través del ajuste de escala automático horizontal. De forma predeterminada, Dataflow establece un límite de recursos de 4,000 trabajadores por trabajo, que se puede ajustar con solicitudes de cuota.
Paralelismo vertical: Varios núcleos y subprocesos de CPU dentro de una sola instancia de trabajador procesan los datos de canalización de forma simultánea. Cada VM de trabajador ejecuta procesos de trabajador y subprocesos de arnés para usar los recursos de procesamiento disponibles. Con el escalamiento dinámico de subprocesos, Dataflow puede ajustar la cantidad de subprocesos activos por trabajador en las canalizaciones por lotes según el uso de la CPU y el espacio libre de memoria. En Dataflow Prime, el ajuste de escala automático vertical escala de forma dinámica la memoria y el procesamiento asignados a los trabajadores.
Unidades de trabajo y jerarquía de ejecución
Para distribuir el procesamiento entre trabajadores y subprocesos, Dataflow divide las canalizaciones de Apache Beam en unidades de trabajo discretas:
- PCollections y particiones: Una
PCollectionrepresenta un conjunto de datos distribuido. En el caso de los datos delimitados (canalizaciones por lotes), Dataflow divide el conjunto de datos en divisiones o fragmentos. En el caso de los datos no delimitados (canalizaciones de transmisión), los datos llegan de forma continua y se transfieren como mensajes o particiones de transmisión. - Paquetes: Dataflow agrupa los elementos en paquetes arbitrarios para que los procese un
DoFn. Un paquete es la unidad de falla y reintento: si el procesamiento de un elemento genera una excepción no controlada, se reintenta todo el paquete. Las operaciones con un alto consumo de memoria pueden aumentar la presión de memoria del trabajador y provocar errores de memoria insuficiente. - Combinación de etapas y pasos: Durante la optimización del grafo, Dataflow combina transformaciones adyacentes en etapas de ejecución combinadas para eliminar la sobrecarga de la materialización de datos intermedios. Dentro de una etapa fusionada, los elementos se procesan en un bucle de ejecución ajustado en un solo subproceso antes de pasarse a la siguiente etapa o límite de aleatorización.
Para obtener más detalles sobre la traducción de canalizaciones y la generación de grafos, consulta Ciclo de vida de la canalización.
Paralelismo y ajuste de escala automático administrados
De forma predeterminada, Dataflow administra el paralelismo de la canalización automáticamente sin necesidad de ajustar la partición de forma manual de las siguientes maneras:
- Ajuste de escala automático horizontal:
- Canalizaciones por lotes: Evalúa el trabajo restante estimado total, el trabajo pendiente de la fuente y el uso de CPU para aumentar o reducir la escala del grupo de trabajadores y completar el trabajo de forma rápida y rentable.
- Canalizaciones de transmisión: Analizan la latencia del sistema, el tamaño de la acumulación y el uso de CPU para aumentar la escala de los trabajadores durante los picos de capacidad de procesamiento y reducir la escala verticalmente durante los períodos de poco tráfico. Para obtener más detalles, consulta Cómo ajustar el ajuste de escala automático horizontal de transmisión.
- Rebalanceo dinámico de trabajos (DWR): En las canalizaciones por lotes, Dataflow supervisa el progreso de las tareas de cada trabajador. Si un trabajador termina antes o si otro se retrasa debido a la asimetría de los datos (rezagados), Dataflow divide de forma dinámica el trabajo residual sin procesar del trabajador lento y lo reasigna a un trabajador inactivo. Para obtener más información, consulta Reequilibrio dinámico del trabajo.
- Escalamiento dinámico de subprocesos: En las canalizaciones por lotes que usan el Runner portátil, ajusta automáticamente la cantidad de subprocesos de procesamiento simultáneo por trabajador según el uso de CPU y el espacio libre de memoria. Para obtener más información, consulta Escalamiento dinámico de subprocesos.
- Ajuste de escala automático vertical: En Dataflow Prime, Dataflow escala de forma dinámica la memoria del trabajador y los recursos de procesamiento para evitar errores de memoria y optimizar la utilización de los recursos. Para obtener más información, consulta Ajuste de escala automático vertical.
Factores que limitan el paralelismo
Es posible que una canalización no alcance el paralelismo esperado debido a las siguientes características de los datos o al diseño del gráfico de canalización:
Fuentes de entrada que no se pueden dividir
Si una fuente de entrada no se puede dividir en rangos independientes, Dataflow se ve obligado a leer la fuente de forma secuencial con un solo subproceso de trabajador:
- Compresión de archivos que no se pueden dividir: Los formatos como
.gz(gzip) o.bzip2(sin indexación) no se pueden leer en paralelo desde compensaciones de bytes arbitrarias. La lectura de un solo archivo comprimido grande restringe la etapa de transferencia a un solo subproceso hasta que se descomprimen y redistribuyen los datos. - Resolución: Almacena los datos en formatos de archivo divisibles (como Parquet, Avro o formatos comprimidos con Snappy) o divide los datos de entrada en varios archivos más pequeños en Cloud Storage.
Fusión de pasos y fan-out alto
La fusión de pasos mejora el rendimiento, ya que reduce la sobrecarga de la serialización, pero puede limitar inadvertidamente el paralelismo y aumentar la presión de memoria cuando un paso con bajo paralelismo produce una gran cantidad de elementos de salida (una operación de "alto fan-out"):
- Ejemplo: Una fuente lee cinco archivos y se fusiona con una transformación
FlatMapque produce 1,000,000 de elementos de salida. Si la transformaciónFlatMapse fusiona con transformaciones posteriores, los 1,000,000 de elementos siguen ejecutándose en un máximo de cinco subprocesos de trabajo, lo que limita gravemente la capacidad de procesamiento posterior. Además, si las transformaciones intermedias se expanden significativamente en la memoria antes de la confirmación, los paquetes grandes pueden agotar la memoria disponible del trabajador. - Resolución: Inserta una transformación
Redistribute(oReshuffleclásica) entre el paso de alta expansión y las transformaciones posteriores para interrumpir la fusión y redistribuir el trabajo en el grupo de trabajadores. Para depurar problemas relacionados con la memoria, consulta Soluciona problemas de errores de memoria insuficiente.
Sesgo de claves y claves activas
Las operaciones de agregación (GroupByKey, CoGroupByKey, Combine.PerKey) agrupan los elementos según su clave asociada.
- Cuello de botella de clave activa: Dataflow enruta todos los elementos con la misma clave a un solo subproceso de trabajador para la agregación. Si una sola clave contiene un gran porcentaje del conjunto de datos total, ese trabajador se convierte en rezagado y los trabajadores ascendentes pueden experimentar contrapresión. Por ejemplo, una clave
nullpredeterminada o una clave de categoría extremadamente popular. - Resolución:
- Usa Combiners (
CombineFnoCombine.PerKey) en lugar deGroupByKeysiempre que sea posible, lo que permite que Dataflow realice combinaciones locales parciales antes del ajuste de escala automático horizontal. - Agrega un prefijo o sufijo de números enteros aleatorios a las teclas de acceso rápido (salado de claves) para distribuir el espacio de claves entre los trabajadores, seguido de una agregación de segunda etapa para combinar los resultados salados.
- Usa Combiners (
Regulación del receptor downstream
Cuando se escribe el resultado de la canalización en servicios externos, como bases de datos o APIs de terceros, un alto paralelismo puede saturar el sistema de destino:
- Limitación: Cientos de subprocesos de trabajadores que emiten llamadas de escritura simultáneas pueden generar errores de límite de frecuencia, tiempos de espera de conexión o degradación de la base de datos.
- Resolución:
- Limita el paralelismo de escritura agrupando elementos con
GroupByKeyo usando receptores de lotes con paralelismo controlado. - Implementa la retirada exponencial del cliente y la lógica de reintentos en las implementaciones del receptor
DoFn.
- Limita el paralelismo de escritura agrupando elementos con
Estrategias de optimización
Para optimizar el paralelismo en tus trabajos de Dataflow, considera los siguientes enfoques:
- Evita la fusión no deseada con
Redistribute:Redistribute.arbitrarily(): Interrumpe la fusión de pasos y redistribuye los elementos de manera uniforme entre todos los trabajadores disponibles.Redistribute.byKey(): Reequilibra los pares clave-valor entre los subprocesos de trabajo y conserva la localidad de las claves.- Para ver ejemplos de implementación, consulta Evita la fusión.
- Supervisa los rezagados y los cuellos de botella: Usa los detalles de ejecución de la consola Google Cloud para identificar las etapas con una gran cantidad de rezagados o un progreso detenido:
¿Qué sigue?
- Obtén más información sobre el ciclo de vida de la canalización.
- Explora el ajuste de escala automático horizontal.
- Comprende el rebalanceo dinámico de trabajos.
- Revisa las prácticas recomendadas para canalizaciones de Dataflow.
- Obtén más información para solucionar problemas de errores de memoria insuficiente.