【发布时间】:2019-05-22 17:33:44
【问题描述】:
我正在 Apache Beam 中编写数据验证脚本。每当有新文件上传到 Google Cloud Storage 时,此脚本都会收到来自 PubSub 的消息,下载文件,并对文件运行一系列预定义测试。 在这些测试结束时,我需要通过电子邮件发送所有未通过测试的行的日志。
为了不多次发送电子邮件,我做了一些阅读并相信我可以使用 Beam 中的状态和计时器构造发送电子邮件一次。但是,每个文件都会有不同数量的错误,所以我如何设置它以使文件发送在发送电子邮件之前需要 X 个元素,其中每个元素都是一个错误,而不是硬编码数字。
我尝试使用带有 COUNT_STATE 的 DoFn 来计算传递给它的元素,但我得到一个关于元素是 Pcollection 而不是 K、V 元组的不同错误。
这是管道代码:
with beam.Pipeline(options=pipeline_options) as p:
# Read Lines from data
validation = (p
| "Read Element From PubSub" >> beam.io.ReadFromPubSub (topic=known_args.input_topic)
| 'Filter Messages' >> beam.ParDo(FilterMessageDoFn(known_args.project, t_options.dataset_id))
| 'After filter' >> beam.ParDo(DebugFn("DATA VALIDATION: PROCESSING FILE...", show_trace))
| 'Generate Schemas' >> beam.ParDo(GetSchemaFn(known_args.project, t_options.validation_home_path))
| 'After GetSschema' >> beam.ParDo(DebugFn("DATA VALIDATION: After OBTAINING SCHEMA...", show_trace))
| 'Validate' >> beam.ParDo(ValidateFn(known_args.project)).with_outputs(
ValidateFn.TAG_VALIDATION_GLOBAL_FAILURE,
ValidateFn.TAG_VALIDATION_CONTENT_FAILURE,
ValidateFn.TAG_VALIDATION_CONTENT_SUCCESS,
main='lines')
to_be_joined = ([validation[ValidateFn.TAG_VALIDATION_GLOBAL_FAILURE],
validation[ValidateFn.TAG_VALIDATION_CONTENT_FAILURE]]
| "Group By Key" >> beam.Flatten()
| 'Persist Global Errors to Big Query' >> beam.ParDo(PersistErrorsFn(known_args.project))
| 'Debug Errors' >> beam.ParDo(DebugFn("DATA VALIDATION: VALIDATION ERRORS", show_trace))
| 'Save Global Errors' >> beam.io.WriteToBigQuery('data_management.validation_errors',
project=known_args.project,
schema=TABLE_SCHEMA,
create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND
)
基本上,我想在写入 BigQuery 之前插入一个步骤,以发送仅在收到 VALIDATION_GLOBAL_FAILURE + VALIDATION_CONTENT_FAILURE 错误数时发送的电子邮件。
谢谢!
【问题讨论】:
-
由于这是一个流管道,您希望多久发送一次电子邮件(仅一次,每天一次,每 >Y 条消息一次)?
-
我想在每次处理文件时发送一封电子邮件。我希望每 15 分钟将文件上传到管道正在侦听的 GCS 文件夹。谢谢!
-
您是否需要确切知道电子邮件中有多少错误(并可能列出它们),或者只是发送一封电子邮件说文件 X 有 >= VALIDATION_GLOBAL_FAILURE + VALIDATION_CONTENT_FAILURE 失败?跨度>
-
我需要列出所有错误。但是,假设我没有,您将如何处理它?也许我可以在要求上有一些余地。谢谢!
标签: python python-2.7 google-cloud-dataflow apache-beam dataflow