【问题标题】:Dataflow Flex Template validation failing with no reason givenDataflow Flex 模板验证失败,没有给出任何理由
【发布时间】:2021-06-27 08:12:11
【问题描述】:

我一直在编写数据流管道并且正在使用弹性模板。

我的代码从 avro 读取并处理它没有问题。但是当涉及到 WriteToAvro 或 WriteToText 时,数据流作业会失败,并且看起来它在模板验证时失败了。我完全没有理由这样做。

我已经尝试了很多东西。删除输出文件的参数并将其硬编码。将 WriteToAvro 切换为 WriteToText,但同样失败。

    with beam.Pipeline(options=options) as p:
        read_from_avro = p \
                         | 'ReadFromAvro' >> ReadFromAvro(input_file)

        redact_data = read_from_avro | "RedactData" >> IdentifyRedactData(project, redact_fields)

        redact_data | 'WriteToAvro' >> WriteToAvro(
                        file_path_prefix=output_file,
                        schema=s,
                        codec='deflate',
                        file_name_suffix='.avro')

join_pcollections 的输出是一个 pcollection,每个元素都是一个字典。

数据流日志给出了这个:

2021-06-27 09:04:46.728 BST Workflow failed.

2021-06-27 09:04:46.763 BST Cleaning up.

2021-06-27 09:04:46.817 BST Worker pool stopped.

有谁知道发生了什么事。仅供参考,当我删除最后一步并运行“ProcessData”步骤时,一切运行顺利。这是刚刚中断的最后一个写入步骤。

编辑以添加需求文件。

apache-beam==2.29.0
google-cloud==0.34.0
google-cloud-dlp==3.1.0
google-cloud-storage==1.35.0
google-cloud-core==1.4.1
google-cloud-datastore==1.15.0

如果我尝试使用 apache-beam[gcp]==2.29.0,构建会失败,所以我想知道这是否与此有关。

apache-beam[gcp] 2.29.0 depends on google-cloud-dlp<2 and >=0.12.0; extra == "gcp"

【问题讨论】:

    标签: python google-cloud-dataflow dataflow


    【解决方案1】:

    已修复。我认为问题源于管道选项配置不正确。我还根据 flex wordcount 示例更改了管道的运行方式。

        options = PipelineOptions(beam_args)
        options.view_as(SetupOptions).save_main_session = True
        p = beam.Pipeline(options=options)
    
        project = options.get_all_options().get('project')
    
        read_from_avro = p \
                         | 'ReadFromAvro' >> ReadFromAvro(input_file)
    
        redact_data = read_from_avro | "RedactData" >> IdentifyRedactData(project, redact_fields)
    
        redact_data | 'WriteToAvro' >> WriteToAvro(
                        file_path_prefix=output_file,
                        schema=table_schema,
                        codec='deflate')
    
        result = p.run()
        result.wait_until_finish()
    

    【讨论】:

      【解决方案2】:

      从您的工作详细信息中,您可以导航到 Cloud Logging。显示的默认日志集可能不包含错误,因此我建议更改过滤器以显示所有日志。

      【讨论】:

      • 谢谢 Kenn,我已经这样做了,并且一直在关注 Logging 服务中的日志,但没有进一步的解释。非常沮丧!
      猜你喜欢
      • 2012-10-11
      • 2020-09-11
      • 1970-01-01
      • 2022-09-29
      • 2021-03-26
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多