You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 23:30:51