如何在Python Beam中读取GCS文件并添加原始行号字段
Python Beam 为GCS文件行添加原始行号并调整输出格式
要实现给GCS存储桶中文本文件的每行添加原始行号,且用逗号分隔行号与内容,核心是在处理每行数据时获取索引并调整格式化逻辑,以下是可行的实现方案:
核心代码实现
使用Beam的Map转换并启用with_indices=True参数,直接获取行索引(从0开始),加1后作为原始行号,再用逗号拼接内容:
import apache_beam as beam def main(): with beam.Pipeline() as p: ( p # 替换为你的GCS文件路径,支持单个文件或通配符(如gs://bucket/*.txt) | "读取GCS文件" >> beam.io.ReadFromText("gs://your-target-bucket/input/file.txt") # 添加行号并调整分隔符为逗号 | "添加行号" >> beam.Map( lambda idx, line: f"{idx + 1},{line}", with_indices=True ) # 输出到GCS,路径会自动生成带分片号的文件 | "写入结果" >> beam.io.WriteToText("gs://your-target-bucket/output/result") ) if __name__ == "__main__": main()
自定义DoFn版本(可选)
如果需要更复杂的行处理逻辑,可自定义DoFn配合with_indices=True:
import apache_beam as beam class AttachLineNumber(beam.DoFn): def process(self, line, index=0): # index为行的起始索引(从0开始),加1得到原始行号 yield f"{index + 1},{line}" def main(): with beam.Pipeline() as p: ( p | "读取GCS文件" >> beam.io.ReadFromText("gs://your-target-bucket/input/file.txt") | "添加行号" >> beam.ParDo(AttachLineNumber(), with_indices=True) | "写入结果" >> beam.io.WriteToText("gs://your-target-bucket/output/result") ) if __name__ == "__main__": main()
输入输出示例
- 输入文件内容(3行):
abcdefghijklmno pqrstuvwxyz 123456789
- 预期输出内容:
1,abcdefghijklmno 2,pqrstuvwxyz 3,123456789
内容的提问来源于stack exchange,提问作者Peter
相关产品推荐
相关产品推荐

