【问题标题】:Python Hadoop streaming, secondary sorting issuesPython Hadoop 流,二次排序问题
【发布时间】:2014-06-26 17:34:34
【问题描述】:

Hadoop 新手在这里。我有一些像这样的用户事件日志,useridtimestamp 都是随机排序的:

userid  timestamp           serviceId
 aaa    2012-01-01 13:12:23 4
 aaa    2012-01-01 12:11:52 3
 ccc    2012-01-03 08:13:07 3
 bbb    2012-01-02 02:34:12 8
 aaa    2012-01-02 01:09:47 4
 ccc    2012-01-02 12:15:39 4

我想得到按 useridtimestamp 排序的中间结果,如下所示:

 aaa    2012-01-01 12:11:52 3
 aaa    2012-01-01 13:12:23 4
 aaa    2012-01-02 01:09:47 4
 bbb    2012-01-02 02:34:12 8
 ccc    2012-01-02 12:15:39 4
 ccc    2012-01-03 08:13:07 3

所以它可以很容易地被我的 Reducer 解析。

最终目标是计算用户在不同服务(serviceIds)上花费的时间。这可以使用 Python Hadoop 流实现吗?如果不是,那么您会建议什么更好的方法?非常感谢!

【问题讨论】:

    标签: algorithm sorting hadoop mapreduce hadoop-streaming


    【解决方案1】:

    在您的映射器中,您可以发出userid 作为键,timestampserviceId 作为按timestamp 排序的值(为了执行排序操作,我假设每个用户的所有行都可以放入主内存中)。

    然后,MR 框架将负责将每个用户的所有不同行发送到单个 reducer,您可以在那里轻松地执行分析。

    如果每个用户的行数过多(比如数百万行),您可以发出 userId-serviceId 作为键,并且在减少阶段之后,您将在每个 user-service 拥有一个包含在该服务上花费的时间的一行文件。如果需要,您可以使用 getmerge 加入所有这些文件

    【讨论】:

    • 太好了,谢谢!所以使用'userId'或'userId-serviceId'作为键,基本上我需要在我的Reducers中按'timestamp'对值进行排序。如果我希望 Mapper 的输出在到达 Reducer 之前已经按时间戳排序怎么办?我知道我可以在 Java 中指定一个自定义分区器(使用“userId-timestamp”作为复合键,但在“userId”上进行分区),但这在 Python 流中也可以吗?谢谢!
    • 是的,您可以将userid-timestamp 设置为key,然后使用hadoop 流的-partitioner 子句按userid 进行分区。查看官方文档中的this example
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-12-10
    • 2014-10-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-04-29
    相关资源
    最近更新 更多