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

如何优雅处理Gzip压缩大CSV:按Key提取前N行并聚合

简洁实现方案

问题回顾

我们需要处理一个超大Gzip压缩CSV文件:

  • 文件无法完整载入内存
  • 行按首列key有序排列
  • 仅保留每个key的前N行进行处理
  • 每组前N行可载入内存

原代码存在嵌套循环复杂、状态管理混乱的问题,以下是更简洁的实现:

原代码(供参考)

import gzip


def process_rows(key, rows):    
    print(f'Processed rows for key {key}')


def main(file_path, N=1000):    
    with gzip.GzipFile(filename=file_path) as file:
        curr_key = None
        rows_to_process = []

        for line in file:
            line = line.decode().strip()

            if len(line) == 0:
                continue

            fields = line.split(',')
            [key, field2, field3] = fields

            if curr_key is not None:
                if curr_key != key or (len(rows_to_process) > 0 and len(rows_to_process) % N == 0):                    
                    process_rows(key, rows_to_process)

                    # Find next key if needed
                    while curr_key == key:
                        line = next(file, None)

                        if line is None:
                            return  # End of file, exit

                        line = line.decode().strip()

                        if len(line) < 1:
                            continue

                        fields = line.split(',')
                        [key, field2, field3] = fields

                    print('Found next key', key)

                    # Reset rows to process
                    rows_to_process = []

            curr_key = key
            rows_to_process.append([key, field2, field3])

        # Flush trailing data
        if (len(rows_to_process) > 0):
            process_rows(key, rows_to_process)

优化后的代码

import gzip

def process_rows(key, rows):
    print(f'Processed {len(rows)} rows for key {key}')

def parse_line(line):
    line = line.decode().strip()
    if not line:
        return None
    fields = line.split(',')
    # 确保字段数量符合预期,可根据实际结构调整
    if len(fields) < 3:
        return None
    return fields[0], fields[1], fields[2]

def main(file_path, N=1000):
    with gzip.GzipFile(filename=file_path) as file:
        # 提前生成有效行的迭代器,过滤空行和无效行
        valid_lines = (line_data for line in file if (line_data := parse_line(line)) is not None)
        
        curr_key = None
        current_batch = []
        
        for key, field2, field3 in valid_lines:
            if key != curr_key:
                # 处理上一个key的剩余数据
                if curr_key is not None and current_batch:
                    process_rows(curr_key, current_batch)
                # 切换到新key,重置批次
                curr_key = key
                current_batch = []
            
            # 只收集当前key的前N行
            if len(current_batch) < N:
                current_batch.append([key, field2, field3])
        
        # 处理最后一个key的数据
        if curr_key is not None and current_batch:
            process_rows(curr_key, current_batch)

优化说明

  • 拆分逻辑:把行解析单独抽成parse_line函数,统一处理空行和无效行,主逻辑更简洁
  • 迭代器过滤:用生成器提前过滤无效行,避免循环内重复判断
  • 简化状态管理:仅通过curr_key判断key切换,收集行时直接判断是否小于N,超过自动跳过,无需嵌套循环跳过剩余行
  • 边界处理明确:单独处理最后一个key的剩余数据,逻辑直观
  • 可读性提升:去掉复杂嵌套和冗余判断,代码流程线性化,更容易理解和维护

内容的提问来源于stack exchange,提问作者revy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 20:25:22