如何使用mrjob.cat自动解压输入并流式传输至MrJob映射器?
用mrjob.cat流式处理超大压缩文件(无需提前解压)
嘿,刚好我之前折腾过这种超大压缩文件的处理场景,用mrjob.cat完全能实现不用提前解压到磁盘,直接实时解压后流式传给映射器处理,完美解决你18Gb压缩、解压后3Tb+的磁盘占用问题。下面给你一步步讲清楚怎么做:
核心原理
mrjob自带的mrjob.cat模块会自动识别常见的压缩格式(比如gzip、bzip2、xz),在读取文件时调用系统对应的解压工具实时解压,然后把解压后的内容一行一行流式传给你的映射器,全程不会把整个解压后的文件写到磁盘上,内存占用也极低。
完整代码示例
这是一个可直接复用的模板,你只需要替换掉映射器里的处理逻辑就行:
from mrjob.job import MRJob from mrjob.cat import cat class LargeCompressedDataJob(MRJob): def mapper(self, _, line): # -------------------------- # 这里写你的业务处理逻辑 # 比如解析字段、统计指标等 # 示例:拆分每行的字段,统计某个字段的出现次数 if line.strip(): fields = line.split("\t") if len(fields) >= 2: yield fields[0], 1 # -------------------------- def run(self): # 初始化runner with self.make_runner() as runner: runner.initialize() # 替换成你的压缩文件路径(本地/HDFS都支持) # 本地路径示例:"/home/you/data/large_file.gz" # HDFS路径示例:"hdfs:///user/you/data/large_file.bz2" input_file_path = "/path/to/your/18gb_compressed_file.gz" # 用cat自动解压并流式读取文件 with cat(input_file_path, self.options) as file_stream: for line in file_stream: # 把解压后的每行传给mapper处理 for key, value in self.mapper(None, line.strip()): self.output(key, value) if __name__ == '__main__': LargeCompressedDataJob.run()
关键细节说明
- 自动识别压缩格式:
mrjob.cat会通过文件扩展名(比如.gz、.bz2)或者文件头来判断压缩格式,不用你手动指定。 - 支持本地/HDFS文件:不管你的压缩文件存在本地磁盘还是HDFS分布式存储,只要路径正确就能处理(HDFS路径记得加
hdfs://前缀)。 - 内存友好:因为是逐行流式处理,哪怕是3Tb的解压后文件,也只会在内存里保留当前处理的一行数据,不会占用大量内存。
运行命令示例
直接在终端运行你的脚本,传入压缩文件路径即可:
python your_job_script.py /path/to/your/18gb_compressed_file.gz
如果是HDFS上的文件:
python your_job_script.py hdfs:///user/you/data/18gb_compressed_file.bz2
处理特殊压缩格式(比如7z)
如果你的文件是mrjob.cat默认不支持的格式(比如7z),可以手动指定解压命令,比如:
# 在cat函数里添加decompress_command参数 with cat(input_file_path, self.options, decompress_command=["7z", "x", "-so"]) as file_stream: for line in file_stream: # 处理逻辑
这里-so参数是让7z把解压内容输出到标准输出,这样mrjob.cat就能流式读取。
内容的提问来源于stack exchange,提问作者Ricardo Decal
相关产品推荐
相关产品推荐

