Процесс потоковой передачи лучей Apache с временными окнамиPython

Программы на Python
Anonymous
Процесс потоковой передачи лучей Apache с временными окнами

Сообщение Anonymous »

У меня есть конвейер потока данных, который читает сообщения из Kafka, обрабатывает их и вставляет в bigquery.

Я хочу, чтобы обработка/вставка bigquery происходила пакетами на основе времени, чтобы при каждом (1 минута), все сообщения, прочитанные из Kafka за этот интервал, будут обработаны в bigquery.
Я пытался следовать примеру документации и создал следующий конвейер:

Код: Выделить всё

import apache_beam as beam
from apache_beam.transforms import window

beam_options = SetupOptions(beam_args, streaming=True)
with beam.Pipeline(options=beam_options) as pipeline:

# read Kafka events
raw = (
pipeline
| "read kafka events" >> kafkaio.KafkaConsume(consumer_config=kafka_config)
| "extract msg" >> (beam.Map(lambda x: x[1])).with_output_types(str)
)

debug1 = (raw | "debug1" >> beam.Map(lambda x: print(f"debug1: {type(x)}, {x}")))

windows = (raw | 'apply window' >> beam.WindowInto(window.FixedWindows(60)))

debug2 = (windows | "debug2" >> beam.Map(lambda x: print(f"debug2: {type(x)}, {x}")))
Но когда я запускаю его, я вижу, что два шага отладки, debug1 и debug2, выполняются сразу один за другим, что указывает на то, что работа окон не сработала. На самом деле этого не происходит.
Я думаю, что упускаю что-то основное.
В документации здесь пример конвейера имеет шаг GroupByKey. Возможно, это необходимо при использовании окон, основанных на времени, хотя на самом деле мне не нужно группировать их по моему варианту использования.
Заранее спасибо

Подробнее здесь: https://stackoverflow.com/questions/787 ... se-windows

Вернуться в «Python»