Map-Reduce 是一种常见的分布式计算模型,简单来说,就是先将海量任务或大批量数据拆分为多个子任务进行处理(MAP),再把各个阶段的处理结果汇总,最终合并为一个完整结果(REDUCE)。

MongoDB 提供的 Map-Reduce 机制非常灵活,在大规模数据统计、数据分析以及复杂聚合场景中都具有较高的实用价值。
MapReduce 命令
下面是 MongoDB MapReduce 的基本语法:
>db.collection.mapReduce(
function() {emit(key,value);}, //map 函数
function(key,values) {return reduceFunction}, //reduce 函数
{
out: collection,
query: document,
sort: document,
limit: number
}
)
在使用 MapReduce 时,需要实现两个核心函数:Map 函数和 Reduce 函数。Map 函数通过 emit(key, value) 的方式遍历 collection 中的记录,并将生成的 key 和 value 传递给 Reduce 函数做进一步汇总处理。
其中,Map 函数必须调用 emit(key, value) 来输出键值对。
参数说明:
- map :映射函数(生成键值对序列,作为 reduce 函数的输入参数)。
- reduce 统计函数,reduce 函数的作用是将 key-values 转换为 key-value,也就是把 values 数组归并成一个单一的 value 值。
- out 统计结果输出到的集合(如果不指定,默认使用临时集合,客户端断开连接后会自动删除)。
- query 文档筛选条件,只有满足条件的文档才会进入 map 函数。(query、limit、sort 可以自由组合使用)
- sort 与 limit 配合使用的排序参数(也会在文档传入 map 函数前完成排序),可用于优化分组处理机制
- limit 发送到 map 函数的文档数量上限(如果没有 limit,单独使用 sort 的意义通常不大)
下面的示例会在 orders 集合中查找 status:"A" 的数据,并按照 cust_id 进行分组,同时统计 amount 的总和。
使用 MapReduce
假设有如下文档结构用于存储用户文章,文档中包含用户的 user_name 字段以及文章状态 status 字段:
>db.posts.insert({
"post_text": "菜鸟教程,最全的技术文档。",
"user_name": "mark",
"status":"active"
})
WriteResult({ "nInserted" : 1 })
>db.posts.insert({
"post_text": "菜鸟教程,最全的技术文档。",
"user_name": "mark",
"status":"active"
})
WriteResult({ "nInserted" : 1 })
>db.posts.insert({
"post_text": "菜鸟教程,最全的技术文档。",
"user_name": "mark",
"status":"active"
})
WriteResult({ "nInserted" : 1 })
>db.posts.insert({
"post_text": "菜鸟教程,最全的技术文档。",
"user_name": "mark",
"status":"active"
})
WriteResult({ "nInserted" : 1 })
>db.posts.insert({
"post_text": "菜鸟教程,最全的技术文档。",
"user_name": "mark",
"status":"disabled"
})
WriteResult({ "nInserted" : 1 })
>db.posts.insert({
"post_text": "菜鸟教程,最全的技术文档。",
"user_name": "runoob",
"status":"disabled"
})
WriteResult({ "nInserted" : 1 })
>db.posts.insert({
"post_text": "菜鸟教程,最全的技术文档。",
"user_name": "runoob",
"status":"disabled"
})
WriteResult({ "nInserted" : 1 })
>db.posts.insert({
"post_text": "菜鸟教程,最全的技术文档。",
"user_name": "runoob",
"status":"active"
})
WriteResult({ "nInserted" : 1 })
接下来,在 posts 集合中通过 mapReduce 筛选已发布的文章(status:"active"),再按照 user_name 字段进行分组,从而统计每个用户各自发布的文章数量:
>db.posts.mapReduce(
function() { emit(this.user_name,1); },
function(key, values) {return Array.sum(values)},
{
query:{status:"active"},
out:"post_total"
}
)
上述 mapReduce 的输出结果如下:
{
"result" : "post_total",
"timeMillis" : 23,
"counts" : {
"input" : 5,
"emit" : 5,
"reduce" : 1,
"output" : 2
},
"ok" : 1
}
结果说明:共有 5 个满足查询条件(status:"active")的文档, 在 map 函数中一共生成了 5 个键值对,最终再通过 reduce 函数把相同键的数据归并为 2 组。
具体参数说明:
- result:存储结果的 collection 名称,这是一个临时集合,MapReduce 连接关闭后会自动删除。
- timeMillis:执行所消耗的时间,单位为毫秒
- input:满足条件并被发送到 map 函数的文档数量
- emit:在 map 函数中 emit 被调用的次数,也可以理解为参与处理的数据总量
- output:结果集合中的文档数量(counts 对调试非常有帮助)
- ok:是否执行成功,成功时为 1
- err:如果执行失败,这里会返回失败原因,不过在实际经验中,这里的提示有时较为模糊,参考价值有限
可以使用 find 操作符来查看 mapReduce 的统计结果:
> var map=function() { emit(this.user_name,1); }
> var reduce=function(key, values) {return Array.sum(values)}
> var options={query:{status:"active"},out:"post_total"}
> db.posts.mapReduce(map,reduce,options)
{ "result" : "post_total", "ok" : 1 }
> db.post_total.find();
查询结果如下所示:
{ "_id" : "mark", "value" : 4 }
{ "_id" : "runoob", "value" : 1 }
通过类似的方法,MapReduce 还可以用于构建更大型、更复杂的 MongoDB 聚合查询与数据统计任务。
Map 函数和 Reduce 函数都可以使用 JavaScript 来实现,因此 MongoDB MapReduce 在数据处理场景下具备很强的灵活性与扩展能力。
