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

如何将字典列表作为参数传入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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 22:27:44