【问题标题】:MongoDB unique value aggregation via map reduce通过 map reduce 进行 MongoDB 唯一值聚合
【发布时间】:2012-05-15 03:03:39
【问题描述】:

我看到很多关于 MongoDB 中聚合的 SO 问题,但是,我还没有找到完整的解决方案。

这是我的数据示例:

{
    "fruits" : {
        "apple" : "red",
        "orange" : "orange",
        "plum" : "purple"
    }
}
{
    "fruits" : {
        "apple" : "green",
        "plum" : "purple"
    }
}
{
    "fruits" : {
        "apple" : "red",
        "orange" : "yellow",
        "plum" : "purple"
    }
}

现在,我的目标是确定每种水果的每种颜色的流行度,因此输出集合如下所示:

{
    "_id" : "apple"
    "values" : {
        "red" : 2,
        "green" : 1
    }
}
{
    "_id" : "orange"
    "values" : {
        "orange" : 1,
        "yellow" : 1
    }
}
{
    "_id" : "plum"
    "values" : {
        "purple" : 3
    }
}

我尝试了各种 M/R 功能,但最终它们要么不起作用,要么花费的时间呈指数级增长。在示例(水果)的上下文中,我在大约 10,000,000 个文档中拥有大约 1,000 种不同的水果和 100,000 种颜色。我目前的工作 M/R 是这样的:

map = function() {
    if (!this.fruits) return;
    for (var fruit in this.fruits) {
        emit(fruit, {
            val_array: [
                {value: this.fruits[fruit], count: 1}
            ]
        });
    }
};

reduce = function(key, values) {
    var collection = {
        val_array: []
    };
    var found = false;
    values.forEach(function(map_obj) {
        map_obj.val_array.forEach(function(value_obj) {
            found = false;
            // if exists in collection, inc, else add
            collection.val_array.forEach(function(coll_obj) {
                if (coll_obj.value == value_obj.value) {
                    // the collection already has this object, increment it
                    coll_obj.count += value_obj.count;
                    found = true;
                    return;
                }
            });
            if (!found) {
                // the collection doesn't have this obj yet, push it
                collection.val_array.push(value_obj);
            }
        });
    });
    return collection;
};

现在,这确实有效,对于 100 条记录,只需要一秒钟左右,但时间会非线性增加,因此 100M 条记录需要非常很长时间。问题是我正在使用 collection 数组在 reduce 函数中进行穷人子聚合,因此需要我遍历 collection 和我的 map 函数中的值。现在我只需要弄清楚如何有效地做到这一点(即使它需要多次减少)。欢迎提出任何建议!


编辑由于缺少更好的发布位置,这是我的解决方案。
首先,我创建了一个名为mr.js 的文件:
map = function() {
    if (!this.fruits) return;
    var skip_fruits = {
        'Watermelon':1,
        'Grapefruit':1,
        'Tomato':1 // yes, a tomato is a fruit
    }
    for (var fruit in this.fruits) {
        if (skip_fruits[fruit]) continue;
        var obj = {};
        obj[this.fruits[fruit]] = 1;
        emit(fruit, obj);
    }
};

reduce = function(key, values) {
    var out_values = {};
    values.forEach(function(v) {
        for(var k in v) { // iterate values
            if (!out_values[k]) {
                out_values[k] = v[k]; // init missing counter
            } else {
                out_values[k] += v[k];
            }
        }
    });
    return out_values;
};

var in_coll = "fruit_repo";
var out_coll = "fruit_agg_so";
var total_docs = db[in_coll].count();
var page_size = 100000;
var pages = Math.floor(total_docs / page_size);
print('Starting incremental MR job with '+pages+' pages');
db[out_coll].drop();
for (var i=0; i<pages; i++) {
    var skip = page_size * i;
    print("Calculating page limits for "+skip+" - "+(skip+page_size-1)+"...");
    var start_date = db[in_coll].find({},{date:1}).sort({date:1}).skip(skip).limit(1)[0].date;
    var end_date = db[in_coll].find({},{date:1}).sort({date:1}).skip(skip+page_size-1).limit(1)[0].date;
    var mr_command = {
        mapreduce: in_coll,
        map: map,
        reduce: reduce,
        out: {reduce: out_coll},
        sort: {date: 1},
        query: {
            date: {
                $gte: start_date,
                $lt: end_date
            }
        },
        limit: (page_size - 1)
    };
    print("Running mapreduce for "+skip+" - "+(skip+page_size-1));
    db[in_coll].runCommand(mr_command);
}

该文件遍历我的整个集合,一次增量地映射/减少 100k 文档(按 date 排序,必须有一个索引!),并将它们减少到单个输出集合中。它是这样使用的:mongo db_name mr.js

然后,几个小时后,我得到了一个包含所有信息的集合。为了找出颜色最多的水果,我使用 mongo shell 中的这个来打印出前 25 个:

// Show number of number of possible values per key
var keys = [];
for (var c = db.fruit_agg_so.find(); c.hasNext();) {
    var obj = c.next();
    if (!obj.value) break;
    var len=0;for(var l in obj.value){len++;}
    keys.push({key: obj['_id'], value: len});
}
keys.sort(function(a, b){
    if (a.value == b.value) return 0;
    return (a.value > b.value)? -1: 1;
});
for (var i=0; i<20; i++) {
    print(keys[i].key+':'+keys[i].value);
}

这种方法真正酷的地方在于,由于它是增量的,我可以在 mapreduce 运行时处理输出数据。

【问题讨论】:

    标签: javascript mongodb mapreduce mongodb-query


    【解决方案1】:

    看来你真的不需要val_array。为什么不使用简单的哈希?试试这个:

    map = function() {
        if (!this.fruits) return;
        for (var fruit in this.fruits) {
            emit(fruit, 
                 {this.fruits[fruit]: 1});
        }
    };
    
    reduce = function(key, values) {
      var colors = {};
    
      values.forEach(function(v) {
        for(var k in v) { // iterate colors
          if(!colors[k]) // init missing counter
            colors[k] = 0
    
          color[k] += v[k];
        }
      });
    
      return colors;
    }
    

    【讨论】:

    • 哇,我真的想多了,不是吗!这确实做到了我想要的。我用 100、1,000 和 100,000 条记录对其进行了测试,每组的运行速度约为 20k/秒(在这些大小下显然是线性的)。我现在正在运行完整的 1000 万条记录,我可以看到随着映射数据批次变大,减少它们需要更长的时间(colors 对象必须在增长):"secs_running" : 488, "msg": "m/r: (1/3) emit phase 383999/10752083 3%"
    • 顺便说一句,我不能使用emit(fruit, {this.fruits[fruit]: 1});,因为密钥是动态生成的,所以我改用了这个 JS hack:var obj = {}; obj[this.fruits[fruit]] = 1; emit(fruit, obj);
    • 我建议尝试部分工作。也就是说,分批处理 100k(或其他)文档,然后在最终作业中减少它。这可能很难实现,所以如果它是一次性的,我不会打扰。 :)
    • @SteveK:这不是黑客行为。 :)
    • 坏消息,就像我们怀疑的那样,我的数据太大了,无法用这个 M/R 处理。这项工作现在已经运行了几个小时,预计完成时间(如果工作的其余部分是线性的)是 Thu, 24 Apr 2053 08:53:10 :P 看起来我可以有效地做 100k 批次,所以我认为我会走那条路!我想我需要将数据 M/R 到不同的集合中,然后编写一个脚本来组合结果,或者我可能会分别 M/R 每个不同的水果。感谢您的帮助!
    【解决方案2】:

    很遗憾地告诉您,MongoDB MapReduce 框架非常缓慢,并且可能会持续“相当长一段时间”(我不希望他们的路线图有所改进)。

    简单地说,我的回答是我不会使用 Mongo-MapReduce 来做这件事,而是专注于在 The New Aggregation Framework 的帮助下实现它: http://docs.mongodb.org/manual/reference/aggregation/

    或在上面运行 Hadoop: http://www.slideshare.net/spf13/mongodb-and-hadoop(漂亮而简单的介绍)

    ​​>

    在使用实现的 MapReduce 功能时,我也遇到过 MongoDB 运行缓慢的问题,我的结论是,即使在执行最简单的任务时,它在性能方面甚至都无法接近上述两种解决方案.使用新的聚合框架,您可以轻松地在商用硬件上处理超过 100 万个文档/秒(甚至更多)。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2011-07-14
      • 2012-12-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2013-05-25
      相关资源
      最近更新 更多