优化100GB大文本转Parquet的PyArrow/mmap处理流程
大型文本文件转Parquet的性能优化问题
我需要处理一个100GB的大型文本文件,文件前109行之后呈现固定重复模式:一行时间戳、16行数据、一行空行。时间戳对应后续16行数据,每行数据包含10000个制表符分隔的值,用逗号作为小数分隔符,每行对应一个采集通道。
我的目标是将这些数据转换为Parquet格式,要求每行包含对应的时间戳和通道编号,且不能把整个数据集加载到内存。
目前我用mmap内存映射文件来统计总行数并提取时间戳行,再创建两个和数据行长度匹配的数组存储通道和时间戳,最后用PyArrow的open_csv流式写入Parquet。想请教有没有优化方案提升处理速度?
现有代码
def count_lines_mmap(filepath): with open(filepath, 'r+b') as f: mm = mmap.mmap(f.fileno(), 0) lines = 0 matching_lines = [] while line := mm.readline(): lines += 1 if re.match(REGEX_TIMESTAMP, line.decode()): matching_lines.append(line.decode()) mm.close() return lines, matching_lines
def process_file(input_file, output_file): # Extract wavelength information wavelength_start, wavelength_delta, num_points = extract_wavelength_info(input_file) lines, matching_lines = count_lines_mmap(FILEIN) if((lines - SKIP_ROWS) / 18 != len(matching_lines)): raise ValueError("Wrong number of timestamps found in the file.") else: print('Matching number of lines.') lines_array = pa.array(pd.to_datetime( np.repeat(matching_lines, 16), format="%d.%m.%Y %H.%M.%S.%f\n") ) channels_array = pa.array(np.tile(CHANNEL_RANGE, n_groups)) # Generate wavelength column names wavelengths = np.arange(wavelength_start, wavelength_start + num_points * wavelength_delta, wavelength_delta) wl_cols = [f"{wl:.3f}" for wl in wavelengths] def handle_invalid_row(invalid_row): return "skip" schema = pa.schema([ pa.field(i, pa.float32()) for i in wl_cols ]) writer_schema = pa.schema([ pa.field(i, pa.float32()) for i in wl_cols ] + [ pa.field("Channel", pa.uint8()), pa.field("Timestamp", pa.timestamp('ns')) ]) # Define the CSV read options read_options = csv.ReadOptions( column_names=wl_cols, skip_rows=SKIP_ROWS, block_size = 1024 * 1024 * 64 # 64MB chunks ) parse_options = csv.ParseOptions( delimiter="\t", invalid_row_handler=handle_invalid_row ) convert_options = csv.ConvertOptions( column_types=schema, decimal_point=',' ) writer = None start_index = 0 end_index = 0 with csv.open_csv( input_file, read_options=read_options, parse_options=parse_options, convert_options=convert_options, ) as reader: for next_chunk in reader: if next_chunk is None: break end_index += len(next_chunk) if writer is None: writer = pq.ParquetWriter( output_file, schema=writer_schema, #next_chunk.schema, ) next_table = pa.Table.from_batches([next_chunk]) # Add channel and timestamp columns next_table = next_table.append_column("Channel", channels_array[start_index:end_index]) next_table = next_table.append_column("Timestamp", lines_array[start_index:end_index]) writer.write_table(next_table) start_index = end_index writer.close()
优化方案
1. 消除两次文件遍历
当前代码先通过mmap全量遍历提取时间戳,再通过open_csv二次读取数据,相当于两次完整IO。改为单次流式遍历:
- 直接逐行读取文件,跳过前109行后,识别时间戳行并缓存,后续16行数据直接绑定该时间戳和对应通道号(1-16),空行直接跳过。
- 这样可以省去第一次mmap遍历的IO开销,同时避免提前生成全量时间戳/通道数组的内存占用。
2. 优化时间戳转换逻辑
替换当前依赖Pandas的时间戳转换,直接用PyArrow内置函数处理,减少跨库开销:
# 替换原lines_array生成代码 import pyarrow.compute as pc # 生成重复后的时间戳字符串数组 timestamp_strings = pa.array(np.repeat(matching_lines, 16)) # 用PyArrow原生方法转换为时间戳类型 lines_array = pc.strptime(timestamp_strings, format="%d.%m.%Y %H.%M.%S.%f\n", unit='ns')
3. 调整Parquet写入参数
- 开启快速压缩:设置
compression='snappy',在保证一定压缩比的同时大幅提升写入速度; - 增大写入批次:设置
write_batch_size=10000(或根据内存调整),减少Parquet文件的写入次数,降低IO开销。
修改后的Writer初始化示例:
writer = pq.ParquetWriter( output_file, schema=writer_schema, compression='snappy', write_batch_size=10000 )
4. 优化CSV读取的block_size
当前设置的64MB block_size可根据系统IO能力调整,尝试128MB或256MB,减少小批次读取的IO次数,但需注意不要超出系统可用内存上限。
5. 预编译正则表达式
确保REGEX_TIMESTAMP是预编译后的正则对象:
REGEX_TIMESTAMP = re.compile(r'你的时间戳正则表达式')
避免每次re.match时重复编译正则,提升匹配效率。
6. 按固定块批量处理
利用文件1+16+1行的固定模式,按18行的块大小批量读取处理,精准绑定时间戳与对应数据行,避免字节块拆分导致的逻辑复杂度,同时提升处理效率。
内容的提问来源于stack exchange,提问作者tiarap00
相关产品推荐
相关产品推荐

