正如@anrope 已经提到的,apache_beam.io.fileio 似乎是用于写入文件的最新 Python API。 WordCount 示例目前已过时,因为它使用 WriteToText 类,该类继承自现已弃用的 apache_beam.io.filebasedsink / apache_beam.io.iobase
要添加到现有答案,这是我的管道,我在运行时动态命名输出文件。我的管道接受 N 个输入文件并创建 N 个输出文件,这些文件根据其对应的输入文件名命名。
with beam.Pipeline(options=pipeline_options) as p:
(p
| 'CreateFiles' >> beam.Create(input_file_paths)
| 'MatchFiles' >> MatchAll()
| 'OpenFiles' >> ReadMatches()
| 'LoadData' >> beam.Map(custom_data_loader)
| 'Transform' >> beam.Map(custom_data_transform)
| 'Write' >> custom_writer
)
当我加载数据时,我创建了一个元组记录(file_name, data) 的 PCollection。我的所有转换都应用于data,但我将file_name 传递到管道的末尾以生成输出文件名。
def custom_data_loader(f: beam.io.fileio.ReadableFile):
file_name = f.metadata.path.split('/')[-1]
data = custom_read_function(f.open())
return file_name, data
def custom_data_transform(record):
file_name, data = record
data = custom_transform_function(data) # not defined
return file_name, data
我将文件保存为:
def file_naming(record):
file_name, data = record
file_name = custom_naming_function(file_name) # not defined
return file_name
def return_destination(*args):
"""Optional: Return only the last arg (destination) to avoid sharding name format"""
return args[-1]
custom_writer = WriteToFiles(
path='path/to/output',
file_naming=return_destination,
destination=file_naming,
sink=TextSink()
)
用您自己的逻辑替换所有 custom_* 函数。