如何将字典列表作为参数传入MRJob任务?
不依赖临时文件传递内存数据给MRJob的方案
核心思路
MRJob默认支持从标准输入(stdin)读取数据,我们可以直接将内存中的列表数据序列化后写入stdin,无需写入临时文件。下面是分文件实现的具体步骤:
1. 编写MRJob任务文件(job.py)
该文件定义MapReduce逻辑,负责解析输入的JSON数据并统计职业数量:
from mrjob.job import MRJob import json class MRJobCounter(MRJob): def mapper(self, _, line): # 解析每行的JSON格式数据 data = json.loads(line) yield data["Job"], 1 def reducer(self, key, values): yield key, sum(values) if __name__ == '__main__': MRJobCounter.run()
2. 在调用文件(main.py)中传递内存数据
在这个文件中初始化MRJob实例,将内存中的字典列表序列化后写入任务的标准输入:
from job import MRJobCounter import json # 你的目标数据列表 mylist = [ {"name": "Kayer", "Job": "Programmer"}, {"name": "Angela", "Job": "Designer"}, {"name": "Eve", "Job": "Programmer"}, {"name": "Robert", "Job": "Programmer"}, ] # 初始化任务,本地运行示例;如果用EMR,修改args为['-r', 'emr']即可 mr_job = MRJobCounter(args=['-r', 'local']) with mr_job.make_runner() as runner: # 获取任务的标准输入流,写入序列化后的每行数据 with runner.get_stdin() as stdin: for item in mylist: json.dump(item, stdin) stdin.write('\n') runner.run() # 解析并处理输出结果 for key, value in mr_job.parse_output(runner.cat_output()): print(f"{key}: {value}")
远程运行(如EMR)的适配
如果是提交到EMR集群,只需将args改为['-r', 'emr'],MRJob会自动处理stdin的管道传递,无需修改数据写入逻辑。只需确保集群环境能正确解析JSON格式的输入数据。
关键注意点
- 必须将每个字典序列化为单独一行的JSON,避免多行JSON导致mapper解析失败。
- 拆分文件后,确保调用文件能正确导入MRJob任务类(可通过调整Python路径或把文件放在同一目录下实现)。
内容的提问来源于stack exchange,提问作者Kayer
相关产品推荐
相关产品推荐

