如何用Apache Beam的ReadFromText转换读取多行CSV文件(Python)
解决Apache Beam读取带换行字段的CSV文件问题
没问题,默认的ReadFromText确实会按物理换行符分割内容,没法识别CSV中带双引号的字段内的合法换行——这也是很多人处理CSV时踩过的坑!下面给你两种可行的解决方案,根据你的文件大小来选:
方案一:适合大文件——用带状态的DoFn合并不完整行
这个方法会在处理每行时,用状态累积那些引号未闭合的不完整行,直到凑成完整的CSV行再输出,适合超大文件(不会一次性加载整个文件到内存):
import apache_beam as beam class MergeQuotedLines(beam.DoFn): def process(self, element, accumulator=beam.DoFn.StateParam(beam.DoFn.StateSpec('partial_line', str))): # 取出之前累积的不完整行(如果有的话) partial_content = accumulator.read() or '' current_content = partial_content + element # 统计双引号数量,偶数说明引号闭合,是完整行 quote_count = current_content.count('"') if quote_count % 2 == 0: accumulator.write('') # 清空累积状态 yield current_content else: # 引号未闭合,继续累积当前内容 accumulator.write(current_content) def print_each_line(line): print(line) path = './input/testfile.csv' with beam.Pipeline() as p: (p | '读取文件行' >> beam.io.ReadFromText(path) | '合并带换行的CSV行' >> beam.ParDo(MergeQuotedLines()) | '打印每行' >> beam.Map(print_each_line) )
注意事项
Beam读取大文件时会分片处理,要是一个不完整的行刚好跨了分片,这个方法可能会出问题。你可以通过调整ReadFromText的min_bundle_size参数,让分片尽量大一些,减少跨分片的概率。
方案二:适合小文件——读取整个文件后用CSV库解析
如果你的文件不大,可以直接读取整个文件内容,用Python标准库的csv.reader来正确解析带换行的字段,这个方法更简单可靠:
import apache_beam as beam from csv import reader from io import StringIO def split_into_valid_csv_lines(file_content): # 用csv.reader自动处理带引号的换行字段 csv_reader = reader(StringIO(file_content)) for row in csv_reader: # 可选:把解析后的行重新拼接成CSV格式字符串 # 如果不需要字符串,直接处理row列表就行 yield ','.join(f'"{col}"' if ',' in col or '\n' in col else col for col in row) def print_each_line(line): print(line) path = './input/testfile.csv' with beam.Pipeline() as p: (p | '匹配文件' >> beam.io.FileIO.match(path) | '读取文件内容' >> beam.io.FileIO.read() | '分割为合法CSV行' >> beam.FlatMap(lambda file: split_into_valid_csv_lines(file.read_utf8())) | '打印每行' >> beam.Map(print_each_line) )
为什么这个方法管用?
Python的csv.reader原生支持识别双引号包裹的字段内的换行,会自动把这些多行内容合并成一个字段,完美适配你的场景。
内容的提问来源于stack exchange,提问作者Brandon
相关产品推荐
相关产品推荐

