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

优化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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 16:31:15