使用Apache Beam重命名CSV列并拼接生成新列的实现方法
实现方案
你现有代码的核心问题是拆分字段后仅返回了拼接后的单值,既没有处理原CSV的表头、完成列重命名,也没有保留三列的输出结构,按以下逻辑调整即可实现需求:
- 单独过滤掉原CSV自带的
first_name,last_name表头,替换成目标表头name,surname,employee_name放在输出文件最顶部 - 每行数据拆分后,保留原两个字段的值作为重命名后的
name、surname列,额外拼接两个字段值生成employee_name列,最终输出三列结构 - 统一表头和数据行的输出格式,按你需要的分隔符拼接后写入文件
可直接运行的修正代码
import apache_beam as beam p2 = beam.Pipeline() def split_row(element): # 去除首尾空白后按逗号拆分CSV字段 return element.strip().split(',') def process_row(element): name = element[0] surname = element[1] employee_name = f"{name} {surname}" # 返回三列对应值:name、surname、employee_name return [name, surname, employee_name] # 定义输出表头,用制表符分隔匹配给出的样例对齐格式 output_header = "name\tsurname\temployee_name" demodata0 = ( p2 | beam.io.ReadFromText('gs://demo/MOCK_DATA.csv') # 给每行加序号,过滤第一行原表头 | beam.WithKeys.from_count() | beam.Filter(lambda line_with_idx: line_with_idx[0] != 0) # 提取行内容做字段拆分 | beam.Map(lambda line_with_idx: split_row(line_with_idx[1])) # 生成三列结构 | beam.Map(process_row) # 转成制表符分隔的字符串,和表头格式统一 | beam.Map(lambda row: '\t'.join(row)) # 把表头插入到结果最前面 | beam.FlatMap(lambda data_rows, header: [header] + list(data_rows), header=output_header) | beam.io.WriteToText('gs://demo/temp/output2') ) p2.run()
可选优化:封装自定义PTransform
如果要符合PTransform的模块化开发要求,可以把列重命名、字段拼接的逻辑封装成独立PTransform方便复用:
class NameProcessTransform(beam.PTransform): def expand(self, input_pcoll): return ( input_pcoll | beam.Map(split_row) | beam.Map(process_row) )
在管道中直接替换对应处理步骤即可:
demodata0 = ( p2 | beam.io.ReadFromText('gs://demo/MOCK_DATA.csv') | beam.WithKeys.from_count() | beam.Filter(lambda line_with_idx: line_with_idx[0] != 0) | beam.Map(lambda line_with_idx: line_with_idx[1]) | NameProcessTransform() # 调用自定义PTransform | beam.Map(lambda row: '\t'.join(row)) | beam.FlatMap(lambda data_rows, header: [header] + list(data_rows), header=output_header) | beam.io.WriteToText('gs://demo/temp/output2') )
注意事项
- 如果你的CSV字段内容本身包含逗号,不要直接用
split(',')做字段拆分,建议使用Beam官方提供的CSV解析组件避免拆分错误 - 如果需要调整输出分隔符,把代码里的
\t(制表符)替换成你需要的符号(比如逗号)即可
内容的提问来源于stack exchange,提问作者Sandeep Mohanty
相关产品推荐
相关产品推荐

