如何优雅处理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
相关产品推荐
相关产品推荐

