Я хочу, чтобы обработка/вставка 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}")))
Я думаю, что упускаю что-то основное.
В документации здесь пример конвейера имеет шаг GroupByKey. Возможно, это необходимо при использовании окон, основанных на времени, хотя на самом деле мне не нужно группировать их по моему варианту использования.
Заранее спасибо
Подробнее здесь: https://stackoverflow.com/questions/787 ... se-windows