【发布时间】:2022-06-15 04:41:03
【问题描述】:
我正在使用 Python SDK 创建一个 Apache Beam 管道,以从 PubSub 读取数据并写入 BigQuery。我正在尝试保留来自 PubSub 的原始消息,以便如果有任何错误,我可以写出要修复的原始记录,然后重新处理。我完成这项工作的最简单方法是使用包含原始消息和工作消息的元组:
(initial_message, working_message)
然后,当我进行 Map 转换时,我转换工作消息并返回元组,保持原始消息不变:
pipeline = (pipeline
| 'Read PubSub' >> ReadFromPubSub(subscription=subscription)
| 'Decode' >> Map(lambda record: (record, record.decode('utf-8')))
| 'Example Transform' >> Map(lambda record: (record[0], some_transformation(record[1])))
)
在写入 BigQuery 之前,这似乎很有效:
write_results = (
pipeline
| 'Extract working message' >> Map(lambda record: record[1])
| 'Write to BigQuery' >> WriteToBigQuery(table=table,
project=project,
schema=schema,
create_disposition=create_disposition,
write_disposition=write_disposition,
insert_retry_strategy=insert_retry_strategy
)
write_results['FailedRows'] | 'Handle write failures' >> ?
然后如何将失败的行与原始消息相关联?
【问题讨论】:
标签: python google-bigquery streaming apache-beam google-cloud-pubsub