Ciclo de vida de la canalización

En esta página, se proporciona una descripción general del ciclo de vida de la canalización desde el código de canalización hasta un trabajo de Dataflow.

En esta página, se explican los siguientes conceptos:

  • Qué es un grafo de ejecución y cómo una canalización de Apache Beam se convierte en un trabajo de Dataflow
  • Cómo maneja Dataflow los errores
  • Cómo Dataflow paraleliza y distribuye de forma automática la lógica de procesamiento de la canalización a los trabajadores que realizan el trabajo
  • Optimizaciones de trabajos que Dataflow podría realizar

Grafo de ejecución

Cuando ejecutas la canalización de Dataflow, Dataflow crea un gráfico de ejecución a partir del código que construye tu objeto Pipeline, incluidas todas las transformaciones y sus funciones de procesamiento asociadas, como los objetos DoFn. Este es el grafo de ejecución de la canalización, y la fase se denomina tiempo de construcción del grafo.

Durante la construcción del gráfico, Apache Beam ejecuta de forma local el código desde el punto de entrada principal del código de la canalización, detiene las llamadas a un paso de origen, receptor o transformación y convierte estas llamadas en nodos del gráfico. En consecuencia, un fragmento de código en el punto de entrada de una canalización (método main de Java y Go, o el nivel superior de una secuencia de comandos de Python) se ejecuta de forma local en la máquina que ejecuta la canalización. El mismo código declarado en un método de un objeto DoFn se ejecuta en los trabajadores de Dataflow.

Por ejemplo, el ejemplo de WordCount incluido en los SDKs de Apache Beam contiene una serie de transformaciones para leer, extraer, contar, dar formato y escribir las palabras individuales en una colección de texto, junto con un recuento de ocurrencias para cada palabra. En el siguiente diagrama, se muestra cómo las transformaciones en la canalización de WordCount se expanden en un grafo de ejecución:

Las transformaciones en el programa de ejemplo de WordCount expandidas en un grafo de ejecución de los pasos que debe ejecutar el servicio de Dataflow.

Figura 1: Grafo de ejecución de ejemplo de WordCount

El grafo de ejecución a menudo difiere del orden en el que especificas las transformaciones cuando creas la canalización. Esta diferencia existe debido a que el servicio de Dataflow realiza varias optimizaciones y fusiones en el grafo de ejecución antes de que se ejecute en recursos de nube administrados. El servicio de Dataflow respeta las dependencias de datos cuando ejecuta la canalización. Sin embargo, los pasos sin dependencias de datos entre ellos pueden ejecutarse en cualquier orden.

A fin de ver el gráfico de ejecución no optimizado que Dataflow generó para la canalización, selecciona el trabajo en la interfaz de supervisión de Dataflow. Para obtener más información sobre cómo ver los trabajos, consulta Usa la interfaz de supervisión de Dataflow.

Durante la construcción del gráfico, Apache Beam valida que todos los recursos a los que hace referencia la canalización (como los depósitos de Cloud Storage, las tablas de BigQuery y los temas o suscripciones de Pub/Sub) existan en realidad y sean accesibles. La validación se realiza mediante llamadas estándar a la API a los servicios respectivos, por lo que es fundamental que la cuenta de usuario que se usa para ejecutar una canalización tenga la conectividad adecuada a los servicios necesarios y esté autorizada para llamar a las APIs de los servicios. Antes de enviar la canalización al servicio de Dataflow, Apache Beam también verifica otros errores y se asegura de que el gráfico de la canalización no contenga operaciones ilegales.

El gráfico de ejecución se traduce al formato JSON, y el gráfico de ejecución JSON se transmite al extremo del servicio de Dataflow.

Luego, el servicio de Dataflow valida el grafo de ejecución de JSON. Cuando se valida el grafo, se convierte en un trabajo en el servicio de Dataflow. Puedes ver el trabajo, su grafo de ejecución, el estado y la información de registro mediante la interfaz de supervisión de Dataflow.

Java

El servicio de Dataflow envía una respuesta a la máquina en la que se ejecutó el programa de Dataflow. Esta respuesta se encapsula en el objeto DataflowPipelineJob, que contiene el jobId de tu trabajo de Dataflow. Usa el jobId para supervisar, hacer seguimiento y solucionar problemas del trabajo mediante la interfaz de supervisión de Dataflow y la interfaz de línea de comandos de Dataflow. Si deseas obtener más información, consulta la referencia de la API para DataflowPipelineJob.

Python

El servicio de Dataflow envía una respuesta a la máquina en la que se ejecutó el programa de Dataflow. Esta respuesta se encapsula en el objeto DataflowPipelineResult, que contiene el job_id de tu trabajo de Dataflow. Usa el job_id para supervisar, hacer seguimiento y solucionar problemas del trabajo mediante la interfaz de supervisión de Dataflow y la interfaz de línea de comandos de Dataflow.

Go

El servicio de Dataflow envía una respuesta a la máquina en la que se ejecutó el programa de Dataflow. Esta respuesta se encapsula en el objeto dataflowPipelineResult, que contiene el jobID de tu trabajo de Dataflow. Usa el jobID para supervisar, hacer seguimiento y solucionar problemas del trabajo mediante la interfaz de supervisión de Dataflow y la interfaz de línea de comandos de Dataflow.

La creación del grafo también se da cuando ejecutas la canalización de forma local, pero el grafo no se traduce a JSON ni se transmite al servicio. En su lugar, el grafo se ejecuta de forma local en la misma máquina en la que se inició el programa de Dataflow. Si deseas obtener más información, consulta Configura PipelineOptions para la ejecución local.

Manejo de errores y excepciones

Tu canalización puede arrojar excepciones durante el procesamiento de datos. Algunos de estos errores son transitorios, como la dificultad temporal para acceder a un servicio externo. Otros errores son permanentes, como los errores causados por datos de entrada corruptos o no analizables, o punteros nulos durante el procesamiento.

Dataflow procesa los elementos en paquetes arbitrarios y vuelve a probar el paquete completo cuando se produce un error de cualquier elemento de ese paquete. Cuando se ejecutan en modo por lotes, los paquetes que incluyen un artículo defectuoso se reintentan cuatro veces. Si un paquete individual falla cuatro veces, la canalización fallará por completo. Cuando se ejecuta en modo de transmisión, un paquete que incluye un elemento defectuoso se reintenta de forma indefinida, lo que puede hacer que la canalización se estanque de forma permanente.

Cuando se procesa en modo por lotes, es posible que veas una gran cantidad de fallas individuales antes de que un trabajo de canalización falle por completo, lo cual ocurre cuando un paquete determinado falla después de cuatro reintentos. Por ejemplo, si la canalización intenta procesar 100 paquetes, en teoría, Dataflow podría generar varios cientos de fallas individuales hasta que un solo paquete alcance la condición de cuatro fallas para la salida.

Los errores de los trabajadores de inicio, como la falta de instalación de paquetes en los trabajadores, son transitorios. Esta situación da como resultado reintentos indefinidos y podría provocar que la canalización se detenga de forma permanente.

Paralelización y distribución

El servicio de Dataflow paraleliza y distribuye de forma automática la lógica de procesamiento de la canalización entre trabajadores y subprocesos. Dataflow usa las abstracciones en el modelo de programación para representar las funciones de procesamiento paralelo. Por ejemplo, las transformaciones ParDo hacen que Dataflow distribuya el código de procesamiento, representado por objetos DoFn, a varios trabajadores para que se ejecuten de forma simultánea.

Dataflow admite dos dimensiones complementarias de paralelismo:

Dataflow administra automáticamente el paralelismo de trabajos, controla el rebalanceo dinámico del trabajo y optimiza el gráfico de ejecución a través de optimizaciones de combinación y fusión. Factores como las fuentes de datos no divisibles, los pasos con una gran cantidad de fan-out, el sesgo de claves y los límites de receptores downstream pueden restringir el paralelismo de la canalización.

Para obtener una guía detallada sobre cómo Dataflow particiona los datos, ajusta la escala de los trabajadores y resuelve los cuellos de botella de paralelismo, consulta Comprensión del paralelismo en Dataflow.

Optimización de fusiones

Una vez que se haya validado el formulario JSON del grafo de ejecución de canalización, es posible que el servicio de Dataflow modifique el grafo para realizar optimizaciones. Las optimizaciones pueden incluir la fusión de varios pasos o transformaciones en el grafo de ejecución de tu canalización en pasos únicos. Los pasos de fusión evitan que el servicio de Dataflow tenga que materializar cada PCollection intermedia en la canalización, lo que podría ser costoso en términos de sobrecarga de procesamiento y memoria.

Aunque todas las transformaciones que especificas en la construcción de tu canalización se ejecutan en el servicio, para garantizar la ejecución más eficiente de tu canalización, las transformaciones pueden ejecutarse en un orden diferente o como parte de una transformación fusionada más grande. El servicio de Dataflow respeta las dependencias de datos entre los pasos en el grafo de ejecución, pero los demás pasos se pueden ejecutar en cualquier orden.

Ejemplo de fusión

En el siguiente diagrama, se muestra cómo el grafo de ejecución del ejemplo de WordCount incluido en el SDK de Apache Beam para Java podría optimizarse y fusionarse mediante el servicio de Dataflow a fin de obtener una ejecución eficiente:

El grafo de ejecución para el programa de ejemplo de WordCount optimizado y la fusión de pasos del servicio de Dataflow.

Figura 2: Grafo de ejecución optimizado del ejemplo de WordCount

Evita la fusión

En algunos casos, Dataflow podría suponer de forma incorrecta la mejor manera de fusionar operaciones en la canalización, lo que puede limitar la capacidad de Dataflow de usar todos los trabajadores disponibles. En esos casos, puedes darle una sugerencia a Dataflow para que redistribuya los datos con una transformación Redistribute.

Para agregar una transformación Redistribute, llama a uno de los siguientes métodos:

  • Redistribute.arbitrarily: Indica que es probable que los datos estén desequilibrados. Dataflow elige el mejor algoritmo para redistribuir los datos.

  • Redistribute.byKey: Indica que es probable que un PCollection de pares clave-valor esté desequilibrado y se debe redistribuir según las claves. Por lo general, Dataflow coloca todos los elementos de una sola clave en el mismo subproceso de trabajador. Sin embargo, no se garantiza la ubicación conjunta de las claves, y los elementos se procesan de forma independiente.

Si tu canalización contiene una transformación Redistribute, Dataflow suele evitar la fusión de los pasos anteriores y posteriores a la transformación Redistribute, y aleatoriza los datos para que los pasos posteriores a la transformación Redistribute tengan un paralelismo más óptimo.

Supervisa la fusión

Puedes acceder a tu grafo optimizado y a las etapas fusionadas en la Google Cloud consola, a través de gcloud CLI o la API.

Console

Para ver las etapas y pasos fusionados del grafo en la consola, en la pestaña Detalles de la ejecución del trabajo de Dataflow, abre la vista de grafo Flujo de trabajo de la etapa.

Si deseas ver los pasos de los componentes fusionados para una etapa, en el grafo, haz clic en la etapa fusionada. En el panel Información de la etapa, la fila Pasos de los componentes muestra las etapas fusionadas. A veces, las partes de una sola transformación compuesta se fusionan en varias etapas.

gcloud

Para acceder a tu grafo optimizado y a las etapas fusionadas a través de la gcloud CLI, ejecuta el siguiente comando de gcloud:

  gcloud dataflow jobs describe --full JOB_ID --format json

Reemplaza JOB_ID por el ID de tu trabajo de Dataflow.

Para extraer los bits relevantes, canaliza el resultado del comando gcloud a jq:

gcloud dataflow jobs describe --full JOB_ID --format json | jq '.pipelineDescription.executionPipelineStage\[\] | {"stage_id": .id, "stage_name": .name, "fused_steps": .componentTransform }'

Para ver la descripción de las etapas fusionadas en el archivo de respuesta de salida, dentro del array ComponentTransform, consulta el objeto ExecutionStageSummary.

API

Para acceder a tu grafo optimizado y a las etapas fusionadas a través de la API, llama a project.locations.jobs.get.

Para ver la descripción de las etapas fusionadas en el archivo de respuesta de salida, dentro del array ComponentTransform, consulta el objeto ExecutionStageSummary.

Optimización de combinaciones

Las operaciones de agregación son un concepto importante en el procesamiento de datos a gran escala. La agregación reúne datos que están muy separados en lo conceptual, por lo que es muy útil para la correlación. El modelo de programación de Dataflow representa las operaciones de agregación como las transformaciones GroupByKey, CoGroupByKey y Combine.

Las operaciones de agregación de Dataflow combinan datos de todo el conjunto de datos, incluidos los que se podrían distribuir entre varios trabajadores. Durante estas operaciones de agregación, a menudo es más eficiente combinar la mayor cantidad de datos posible de forma local antes de combinarlos entre instancias. Cuando aplicas una GroupByKey o alguna otra transformación de agregación, el servicio de Dataflow realiza una combinación parcial a nivel local de forma automática antes de la operación de agrupación principal.

Cuando se realiza una combinación parcial o multinivel, el servicio de Dataflow toma decisiones diferentes en función de si la canalización funciona con datos de transmisión o por lotes. Para datos limitados, el servicio favorece la eficiencia y realizará toda la combinación local posible. Para datos ilimitados, el servicio prefiere una latencia más baja y puede que no realice una combinación parcial, ya que podría aumentar la latencia.