Lee de Pub/Sub a Dataflow

En esta página, se describen las prácticas recomendadas para leer desde Pub/Sub en Dataflow.

Apache Beam proporciona una implementación de referencia del conector de E/S de Pub/Sub para que lo usen ejecutores que no sean de Dataflow. Sin embargo, el ejecutor de Dataflow usa su propia implementación personalizada del conector. Esta implementación aprovecha las APIs Google Cloudy los servicios internos para ofrecer marcas de agua de baja latencia, alta precisión de marcas de agua y anulación de duplicación eficiente para el procesamiento de mensajes del tipo “exactamente una vez”. El conector está disponible para Java, Python, y Go.

Procesamiento “exactamente una vez”

Pub/Sub desacopla los publicadores de eventos de los consumidores de eventos. La aplicación publica mensajes en un tema, y Pub/Sub entrega los mensajes de forma asíncrona a los suscriptores.

Pub/Sub asigna un ID de mensaje único a cada mensaje que se publica correctamente en un tema. De forma predeterminada, Pub/Sub realiza la entrega de mensajes al menos una vez. Para lograr una semántica de “al menos una vez”, si Pub/Sub no recibe la confirmación de recepción del suscriptor dentro del plazo de confirmación, vuelve a intentar la entrega del mensaje. Los reintentos también pueden ocurrir antes del plazo de confirmación o después de que se confirmó un mensaje.

Dataflow confirma recepción de mensajes después de que se procesaron de forma correcta en la primera etapa fusionada y los efectos secundarios de ese procesamiento se escribieron en el almacenamiento persistente. Para reducir la cantidad de mensajes duplicados, Dataflow extiende continuamente el plazo de confirmación mientras se procesa un lote de mensajes en esta etapa.

Debido a que Pub/Sub puede volver a entregar un mensaje, es posible que lleguen mensajes duplicados a la canalización. Si tu canalización de Dataflow usa el modo de transmisión “exactamente una vez”, Dataflow anula la duplicación de estos mensajes para lograr una semántica de “exactamente una vez” .

Si tu canalización puede tolerar algunos registros duplicados, considera usar el modo de transmisión “al menos una vez” en su lugar. Este modo puede reducir significativamente la latencia y el costo total de tu canalización. La desventaja es que los mensajes duplicados pueden procesarse dos veces. Para obtener más información, consulta Elige qué modo de transmisión usar.

Anula la duplicación por atributo de mensaje

De forma predeterminada, Dataflow anula la duplicación en función del ID del mensaje. Sin embargo, una aplicación puede enviar el mismo registro dos veces como dos mensajes de Pub/Sub distintos. Por ejemplo, los datos de origen originales pueden contener registros duplicados, o la aplicación puede publicar de forma incorrecta el mismo mensaje dos veces. Esto último puede suceder debido a reintentos, si se descartó la confirmación debido a problemas de red o a otras interrupciones. En estas situaciones, los mensajes duplicados tienen IDs de mensajes diferentes.

Según tu situación, es posible que tus datos contengan un campo único que se pueda usar para anular la duplicación. Por ejemplo, los registros pueden contener un ID de transacción único. Puedes configurar el conector de E/S de Pub/Sub para anular la duplicación de mensajes en función del valor de un atributo de mensaje, en lugar de usar el ID de mensaje de Pub/Sub. Siempre que el publicador establezca este atributo de manera coherente durante los reintentos, Dataflow podrá detectar los duplicados. Los mensajes deben publicarse en Pub/Sub en un plazo de 10 minutos entre sí para la anulación de duplicación.

Para obtener más información sobre el uso de atributos de ID, consulta los siguientes temas de referencia del SDK:

Suscripciones

Cuando configuras tu canalización, especificas un tema de Pub/Sub o una suscripción de Pub/Sub desde la que leer. Si especificas una suscripción, no uses la misma suscripción de Pub/Sub para varias canalizaciones. Si dos canalizaciones leen desde una sola suscripción, cada canalización recibirá parte de los datos de una manera no determinista, lo que puede causar la duplicación de mensajes, el retraso de la marca de agua y un ajuste de escala automático ineficiente. En su lugar, crea una suscripción independiente para cada canalización.

Si especificas un tema, el conector crea una suscripción temporal nueva. Esta suscripción es única por canalización.

Marcas de tiempo y marcas de agua

Todos los mensajes de Pub/Sub tienen una marca de tiempo, que representa el momento en que Pub/Sub recibe el mensaje. Tus datos también pueden tener una marca de tiempo de evento, que es el momento en que la fuente generó el registro.

Puedes configurar el conector para leer la marca de tiempo del evento desde un atributo en el mensaje de Pub/Sub. En ese caso, el conector usa la marca de tiempo del evento para la marca de agua. De lo contrario, de forma predeterminada, usa la marca de tiempo del mensaje de Pub/Sub.

Para obtener más información sobre el uso de marcas de tiempo de eventos, consulta los siguientes temas de referencia del SDK:

El conector de Pub/Sub tiene acceso a la API privada de Pub/Sub que proporciona la antigüedad del mensaje no confirmado más antiguo de una suscripción. Esta API proporciona una latencia menor que la disponible en Cloud Monitoring. Permite que Dataflow avance las marcas de agua de la canalización y emita resultados de procesamiento con ventanas con latencias bajas.

Si configuras el conector para usar marcas de tiempo de eventos, Dataflow crea una segunda suscripción de Pub/Sub, llamada suscripción de seguimiento. Dataflow usa la suscripción de seguimiento para inspeccionar los tiempos de los eventos de los mensajes que aún están en la lista de tareas pendientes. Este enfoque permite que Dataflow calcule con precisión la lista de tareas pendientes de tiempo del evento. La cuenta de servicio del trabajador debe tener al menos los siguientes permisos en el proyecto que contiene la suscripción de seguimiento:

  • pubsub.subscriptions.create
  • pubsub.subscriptions.consume
  • pubsub.subscriptions.delete

Además, necesita el permiso pubsub.topics.attachSubscription en el tema de Pub/Sub. Se recomienda crear un rol personalizado de Identity and Access Management que contenga solo estos permisos.

Para obtener más información sobre las marcas de agua, consulta la página de Stack Overflow que trata sobre cómo Dataflow calcula las marcas de agua de Pub/Sub.

Si una canalización tiene varias fuentes de Pub/Sub y una de ellas tiene un volumen muy bajo o está inactiva, retrasa el avance de toda la marca de agua, lo que aumenta la latencia general de la canalización. Si hay temporizadores o agregaciones de ventanas en la canalización en función de la marca de agua, estos también se retrasan.

Búsqueda de Pub/Sub

Con la búsqueda de Pub/Sub, los usuarios pueden volver a reproducir mensajes ya confirmados. Puedes usar la búsqueda de Pub/Sub con Dataflow para volver a procesar mensajes en una canalización.

Sin embargo, no se recomienda usar la búsqueda de Pub/Sub en una canalización en ejecución. Buscar hacia atrás en una canalización en ejecución puede generar mensajes duplicados o descartados. También invalida la lógica de marca de agua de Dataflow y entra en conflicto con el estado de una canalización que incorpora datos procesados.

Para volver a procesar mensajes con la búsqueda de Pub/Sub, se recomienda el siguiente flujo de trabajo:

  1. Crea una instantánea de la suscripción:
  2. Crea una suscripción para el tema de Pub/Sub La nueva suscripción hereda la instantánea.
  3. Drena o cancela el trabajo actual de Dataflow.
  4. Vuelve a enviar la canalización con la nueva suscripción.

Para obtener más información, consulta Reprocesamiento de mensajes con la instantánea de Pub/Sub y la búsqueda de Pub/Sub.

Paralelismo de la fuente de Pub/Sub

La fuente de Pub/Sub asigna a cada mensaje una clave determinista para el procesamiento y usa esas claves para mezclar los mensajes. Para los trabajos de Streaming Engine, se usan 1,024 claves para la mezcla. Para los trabajos que no son de Streaming Engine, la cantidad de claves es la potencia más baja de 2 mayor que (4 * maximum workers).

Para anular la cantidad predeterminada de claves de mezcla, establece la num_pubsub_keys opción de servicio:

Java

--dataflowServiceOptions=num_pubsub_keys=NUMBER_OF_KEYS

Python

--dataflow_service_options=num_pubsub_keys=NUMBER_OF_KEYS

Go

--dataflow_service_options=num_pubsub_keys=NUMBER_OF_KEYS

Reemplaza NUMBER_OF_KEYS por la cantidad de claves. Se usa la siguiente potencia de 2 mayor o igual que el valor especificado.

Por ejemplo, puedes establecer esta opción en las siguientes situaciones:

Si estableces esta opción, considera las compensaciones que se describen en Paralelización y distribución.

No puedes cambiar la cantidad de claves como parte de una actualización de canalización. Para cambiar la cantidad de claves de un trabajo de canalización existente, debes iniciar un trabajo nuevo.

Características de Pub/Sub no compatibles

Las siguientes funciones de Pub/Sub no son compatibles con la implementación del ejecutor de Dataflow del conector de E/S de Pub/Sub.

Retirada exponencial

Cuando creas una suscripción de Pub/Sub, puedes configurarla para que use una política de reintentos de retirada exponencial. Sin embargo, la retirada exponencial no funciona con Dataflow. En su lugar, crea la suscripción con la política de reintentos Volver a intentar de inmediato.

La retirada exponencial se activa con una confirmación de recepción negativa o cuando vence el plazo de confirmación. Sin embargo, Dataflow no envía confirmaciones de recepción negativas cuando falla el código de la canalización. En su lugar, vuelve a intentar el procesamiento del mensaje de forma indefinida, mientras extiende el plazo de confirmación para el mensaje de forma continua.

Temas de mensajes no entregados

No uses temas de mensajes no entregados de Pub/Sub con Dataflow por los siguientes motivos:

  • Dataflow envía confirmaciones de recepción negativas por varios motivos internos (por ejemplo, si un trabajador se está cerrando). Como resultado, es posible que los mensajes se entreguen al tema de mensajes no entregados, incluso cuando no se producen fallas en el código de la canalización.

  • Dataflow reconoce los mensajes después de que la primera etapa fusionada procesa correctamente un paquete de mensajes. Si la canalización tiene varias etapas fusionadas y se producen fallas en cualquier momento después de la primera etapa, los mensajes ya se confirmaron y no van al tema de mensajes no entregados.

En su lugar, implementa el patrón de mensajes no entregados de forma explícita en la canalización, enrutando los mensajes con errores a un destino para su procesamiento posterior. Algunos receptores de E/S tienen compatibilidad integrada con las filas de buzón de destino. En los siguientes ejemplos, se implementan patrones de mensajes no entregados:

Entrega de Pub/Sub “exactamente una vez”

Debido a que Dataflow tiene sus propios mecanismos para el procesamiento “exactamente una vez” , no se recomienda usar la entrega “exactamente una vez” de Pub/Sub con Dataflow. Habilitar la entrega “exactamente una vez” de Pub/Sub reduce el rendimiento de la canalización, ya que limita la cantidad de mensajes disponibles para el procesamiento paralelo.

Pub/Sub: ordenamiento de mensajes

El ordenamiento de mensajes es una función de Pub/Sub que permite que un suscriptor reciba mensajes en el orden en que se publicaron.

No se recomienda usar el orden de mensajes con Dataflow por los siguientes motivos:

  • Es posible que el conector de E/S de Pub/Sub no conserve el orden de los mensajes.
  • Apache Beam no define lineamientos estrictos con respecto al orden en el que se procesan los elementos. Por lo tanto, es posible que el orden no se conserve en las transformaciones posteriores.
  • Usar el ordenamiento de mensajes de Pub/Sub con Dataflow puede aumentar la latencia y disminuir el rendimiento.

Transformaciones de mensajes únicos de Pub/Sub

Las transformaciones de mensajes únicos (SMT) te permiten manipular, validar y filtrar mensajes en función de sus atributos o datos a medida que se transmiten por el sistema. Las suscripciones que se envían a Dataflow no deben usar SMT que filtren mensajes, ya que pueden interferir con el ajuste de escala automático. Esto sucede porque el filtrado de SMT de suscripción puede hacer que la lista de tareas pendientes parezca más grande que lo que se entrega a Dataflow hasta que la SMT procese los mensajes filtrados. Las SMT de temas que filtran mensajes no causarán problemas con el ajuste de escala automático.

¿Qué sigue?