【问题标题】:How to implement OR join in hadoop(scalding/cascading)如何在hadoop中实现或加入(烫伤/级联)
【发布时间】:2012-09-24 22:13:25
【问题描述】:

只需将连接字段作为reducer key发送,就可以很容易地通过单个键连接数据集。 但是通过多个键连接记录,其中至少一个应该是相同的,这对我来说并不容易。

例如我有日志,我想按用户参数对它们进行分组,我想通过 (ipAddress, sessionId,visitorCockies) 加入它们

因此,如果 log1.ip == log2.ip OR log1.session = log2.session OR log1.cockie = log2.coockie,则 log1 应与 log2 分组。也许可以创建复合键或像 minHash 这样的概率方法......

有可能吗?

【问题讨论】:

    标签: scala join hadoop cascading scalding


    【解决方案1】:

    问题在于 MapReduce 连接通常是通过为在某些字段上匹配的记录提供相同的 reduce 键来实现的,以便将它们发送到同一个 reducer。所以任何解决这个问题的方法都会有点麻烦,但有可能......

    这是我的建议:对于每个输入记录,生成三个副本,每个副本都有一个新的“关键”字段,该字段以它来自的字段为前缀。例如,假设您有以下输入:

    (ip=1.2.3.4, session=ABC, cookie=123)
    (ip=3.4.5.6, session=DEF, cookie=456)
    

    然后你会生成

    (ip=1.2.3.4, session=ABC, cookie=123, key=ip_1.2.3.4)
    (ip=1.2.3.4, session=ABC, cookie=123, key=session_ABC)
    (ip=1.2.3.4, session=ABC, cookie=123, key=cookie_123)
    (ip=3.4.5.6, session=DEF, cookie=456, key=ip_3.4.5.6)
    (ip=3.4.5.6, session=DEF, cookie=456, key=session_DEF)
    (ip=3.4.5.6, session=DEF, cookie=456, key=cookie_456)
    

    然后你可以简单地在这个新领域进行分组。

    我对 scalding/cascading 不太熟悉(尽管我一直想了解更多有关它的信息),但这肯定符合 Hadoop 中通常的连接方式。

    【讨论】:

    • 这样我会得到 3 个不同的重叠组(每个键对应地相等)所以我需要一种方法将它们合并到一个组中
    • @yura 如果左记录可以与多个右记录连接,那么通常连接会使它们不合并(因此重复值)。这样做的原因是合并会导致未定义的元组大小(表宽度),您可能会得到 1,2 或 3 个正确的记录。因此这个解决方案是正确的,(但缺乏细节也没有基本的实现;)。
    【解决方案2】:

    按照上面 Joe 的描述创建单独的连接后,您需要删除重复的连接。如果您在“OR-join”中使用的所有字段中它们都相等,则数据中的两个元组是重复的。因此,如果您之后对代表所有相关字段的键进行自然连接,您会将所有重复项组合在一起。因此,您可以将它们替换为单个出现的相应元组。

    让我们看一个示例:假设您有包含字段 (A,B,C,D) 的元组,并且您感兴趣的字段是 A、B 和 C。您首先要对A、B、C 分别。对于每一个,您都将加入初始元组流。用 (A0, B0, C0, D0) 表示第一个流,用 (A1, B1, C1, D1) 表示第二个流。结果将是元组(A0、B0、C0、D0、A1、B1、C1、D1)。对于每个元组,您将创建一个元组 (A0A1B0B1C0C1, A0, B0, C0, D0, A1, B1, C1, D1),因此所有重复项将在后续的 reducer 中组合在一起。对于每个组,只返回一个包含的元组。

    【讨论】:

      【解决方案3】:

      您能详细介绍一下“通过多个键连接记录”吗?

      如果您知道工作流程中可以连接特定键的点,那么最好的方法可能是定义一个具有多个连接的流,而不是尝试操作复杂的数据结构以将 N 个键解析为一个一步。

      这是一个示例应用程序,它展示了如何在级联中处理不同类型的连接:https://github.com/Cascading/CoPA

      【讨论】:

        【解决方案4】:

        对于级联,我最终创建了一个过滤器,用于检查 OR 中任何条件的输出是否为真。级联过滤器输出可以选择使用的真/假值。

        【讨论】:

          【解决方案5】:

          提示:使用类型别名使您的 Scalding 代码易于阅读

          注意 0:这个解决方案特别好,因为它总是只有 1 个映射作业,即使有更多的键可以加入。

          注意 1:假设每个管道没有重复的键,否则您必须使 'key 也有一个索引来自哪个日志,并且 mapTo 将是一个 flatMapTo 和有点复杂。

          注意 2:为简单起见,这将丢弃连接字段,为了保留它们,您需要一个丑陋的大元组(ip1、ip2、session1、session2、...等)。如果你真的想要,我可以写一个保留它们的例子。

          注意 3:如果你真的想合并重复的值,你可以用 groupBy 每个 logEntry1 和 logEntry2,产生一个 logEntryList,然后 cat(如评论中所述,这是加入不正常)。这将创建另外 2 个映射作业。

          type String2 = (String, String)
          type String3 = (String, String, String)
          
          def addKey(log: Pipe): Pipe = log.flatMap[String3, String](('ip, 'session, 'cookie) -> 'key)(
            _.productIterator.toList.zipWithIndex.map {
              case (key: String, index: Int) => index.toString + key
            }
          )
          
          (addKey(log1) ++ addKey(log2)).groupBy('key)(_.toList[String]('logEntry -> 'group))
          .mapTo[Iterable[String], String2]('group -> ('logEntry1, 'logEntry2))(list => (list.head, list.last))
          

          【讨论】:

            猜你喜欢
            • 2014-06-19
            • 1970-01-01
            • 1970-01-01
            • 2016-08-10
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 2014-06-07
            相关资源
            最近更新 更多