【发布时间】:2021-05-12 18:58:15
【问题描述】:
我目前有代码:
gs_folder = sys.argv[1]
options = PipelineOptions(
runner='DataflowRunner',
project='xxx',
job_name=xxx{uuid.uuid4()}',
region='us-central1',
temp_location='xxx')
gfs = gcs.GCSFileSystem(options)
p = beam.Pipeline(options=options)
discover_empty = p | 'Filenames' >> beam.Create(gfs.match([gs_folder])[0].metadata_list) | \
'Reshuffle' >> beam.Reshuffle() | \
'Delete empty files' >> beam.ParDo(DeleteEmpty(gfs))
p.run()
此代码基于问题here。最终发生的是我在下面收到此错误:
WARNING:apache_beam.options.pipeline_options:Discarding unparseable args: ['gs://xx/xx']
这没有多大意义,因为这是我要执行删除操作的文件夹。此外,看起来数据流作业确实成功运行,但应该删除的文件没有正确删除。我应该如何在这里传递管道选项 arg?
我还有几个关于这个过程的后续问题。看起来beam.Create() 在本地运行,然后切换到数据流。我怎样才能使管道的这一部分在数据流上运行?
【问题讨论】:
标签: python-3.x google-cloud-platform google-cloud-storage apache-beam dataflow