【问题标题】:Mapreduce Table DiffMapreduce 表差异
【发布时间】:2013-05-04 18:53:04
【问题描述】:

我有两个版本(旧/新)的数据库表,大约有 100,000,000 条记录。它们在文件中:

trx-old
trx-new

结构是:

id date amount memo
1  5/1     100 slacks
2  5/1      50 wine

id 是简单的主键,其他字段是非键。我要生成三个文件:

trx-removed (ids of records present in trx-old but not in trx-new)
trx-added   (records from trx-new whose ids are not present in trx-old)
trx-changed (records from trx-new whose non-key values have changed since trx-old)

我需要每天在一个短批处理窗口中执行此操作。实际上,我需要为多个表和多个模式(为每个模式生成三个文件)执行此操作,因此实际应用程序涉及更多。但我认为这个例子抓住了问题的症结。

这感觉像是一个明显的 mapreduce 应用程序。我从未编写过 mapreduce 应用程序,我的问题是:

  1. 是否有一些 EMR 应用程序已经这样做了?
  2. 是否存在明显的 Pig 或 Cascading 解决方案?
  3. 还有其他一些与此非常接近的开源示例吗?

PS 我看到了diff between tables 的问题,但那里的解决方案看起来无法扩展。

PPS 下面是一个演示算法的 Ruby 小玩具:Ruby dbdiff

【问题讨论】:

  • 2.是的,至少对于添加和删除的部分,Pig 解决方案是明显的LEFT OUTER JOINFILTER,基于连接列是否为null。至于“更改”,我最好的猜测是内部JOIN 并根据字段是否不同进行过滤。

标签: sql hadoop mapreduce elastic-map-reduce cascading


【解决方案1】:

我认为编写自己的工作是最简单的,主要是因为当典型的 reducer 只写入一个文件时,您将希望使用 MultipleOutputs 从单个 reduce 步骤中写入三个单独的文件。您需要使用 MultipleInputs 为每个表指定一个映射器。

【讨论】:

  • 查看所有解决方案后,我想知道使用罐装 CoGroup 和一些冗余过滤是否更好,或者手动编写一个即时写入 MultipleOutputs 的解决方案是否更好。跨度>
【解决方案2】:

这似乎是级联解决的完美问题。您已经提到您从未编写过 MR 应用程序,如果您的目的是快速入门(假设您熟悉 Java),那么 Cascading 是您的方式,恕我直言。稍后我会详细介绍这一点。

可以使用 Pig 或 Hive,但如果您想对这些列执行额外的分析或更改架构,它们就不够灵活,因为您可以通过读取列标题或从您创建的映射文件以表示架构。

Cascading 你会:

  1. 设置您的传入Taps:点击 trxOld 并点击 trxNew(这些指向您的源文件)
  2. 将您的水龙头连接到Pipes:Pipe oldPipe 和 Pipe newPipe
  3. 设置您的传出Taps:点击 trxRemoved、点击 trxAdded 和点击 trxChanged
  4. 构建您的管道分析(这是有趣(伤害)发生的地方)

trx-removed: 添加 trx

Pipe trxOld = new Pipe ("old-stuff");
Pipe trxNew = new Pipe ("new-stuff");
//smallest size Pipe on the right in CoGroup
Pipe oldNnew = new CoGroup("old-N-new", trxOld, new Fields("id1"), 
                                       trxNew, new Fields("id2"), 
                                       new OuterJoin() ); 

外部连接在另一个管道(您的源数据)中缺少 id 时为我们提供了 NULLS,因此我们可以在下面的逻辑中使用 FilterNotNullFilterNull 来获取我们然后连接到 Tap 的最终管道trxRemoved 并相应地点击 trxAdded。

trx 改变

在这里,我将首先使用FieldJoiner 连接您正在寻找更改的字段,然后使用ExpressionFilter 为我们提供僵尸(因为它们已更改),例如:

Pipe valueChange = new Pipe("changed");
valueChange = new Pipe(oldNnew, new Fields("oldValues", "newValues"), 
            new ExpressionFilter("oldValues.equals(newValues)", String.class),
            Fields.All);

它的作用是过滤掉具有相同值的字段并保留差异。此外,如果上面的表达式为真,它将删除该记录。最后,将您的 valueChange 管道连接到您的 Tap trxChanged,您将获得三个输出,其中包含您正在寻找的所有数据,其中包含允许一些附加分析的代码。

【讨论】:

  • 这看起来非常接近实际解决方案。 CoGroup 正是我需要的神奇外部连接!这让我想知道 CoGroup 的性能如何,以及它是否可以受益于我的表已预先排序这一事实。
  • MapReduce 中连接的性能应该总是很慢(技术术语)。该框架不擅长加速连接(只需询问 PIG 的 Hive),一般来说,Hadoop 应该被认为是一个 18 轮车,缓慢地拉动一吨重量,但比使用一辆勇敢的汽车要好。现在,关于预排序数据的第二点,我不知道它对优化的影响。我确实知道 CoGroup 按其自然顺序对组键进行排序,这告诉我它是内置的。我不得不相信在 GroupBy 管道中 presorted 会比 unsorted 运行得更快。
【解决方案3】:

正如@ChrisGerken 建议的那样,您必须使用MultipleOutputsMultipleInputs 才能生成多个输出文件并将自定义映射器与每种输入文件类型(旧/新)相关联。

映射器会输出:

  • key:主键(id)
  • 值:来自输入文件的记录,带有附加标志(新/旧取决于输入)

reducer 将为每个键和输出遍历所有记录 R

  • 删除文件:如果只存在带有 old 标志的记录。
  • 添加到文件:如果只存在带有新标志的记录。
  • 更改文件:如果R 中的记录不同。

由于该算法会随着 reducer 的数量而扩展,因此您很可能需要第二个作业,它将结果合并到一个文件中以作为最终输出。

【讨论】:

  • 没有与HOW相关的答案。您刚刚说过,如果第二个文件中不存在该记录,则将其视为已删除。这存在问题。
  • @smttsp:我希望这能解决问题。 reducer 将获得以下输入并将其标记为已删除:(1, {old1}),添加:(2, {new2}),更改:(3, {old3, new3})
  • 我想知道@Engineiro 的解决方案(使用 CoGroup)与该解决方案的性能相比如何。
【解决方案4】:

我想到的是:

考虑你的表格是这样的:

Table_old
1    other_columns1
2    other_columns2
3    other_columns3

Table_new 
2    other_columns2
3    other_columns3
4    other_columns4

追加table_old的元素“a”和table_new的元素“b”。

当您合并两个文件时,如果一个元素存在于第一个文件中而不是第二个文件中,则会将其删除

table_merged
1a    other_columns1
2a    other_columns2
2b    other_columns2
3a    other_columns3
3b    other_columns3
4a    other_columns4

您可以通过该文件轻松地进行操作。

另外,假设您的 id 是 n 位数字,并且您有 10 个集群+1 个主节点。您的密钥将是 id 的第一个数字,因此,您将数据均匀地划分为集群。您将进行分组+分区,以便对数据进行排序。

例子,

table_old
1...0 data
1...1 data
2...2 data

table_new
1...0 data
2...2 data
3...2 data

您的密钥是第一个数字,您根据该数字进行分组,您的分区是根据 id 的其余部分。然后你的数据会以如下方式进入你的集群

worker1
1...0b data
1...0a data
1...1a data

worker2 
2...2a data
2...2b data and so on.

请注意,a、b 不必排序。

编辑 合并将是这样的:

FileInputFormat.addInputPath(job, new Path("trx-old"));
FileInputFormat.addInputPath(job, new Path("trx-new"));

MR会得到两个输入,两个文件会被合并,

对于附加部分,您应该在 Main MR 之前再创建两个作业,这将只有 Map。第一个 Mapappend "a" 到第一个列表中的每个元素,第二个将 append "b" 到第二个列表中的元素。第三个工作(我们现在使用的那个/主地图)将只有减少工作来收集它们。所以你会有Map-Map-Reduce

可以这样添加

//you have key:Text
new Text(String.valueOf(key.toString()+"a"))

但我认为可能有不同的附加方式,其中一些可能更有效 (text hadoop)

希望对你有帮助,

【讨论】:

  • 这很有帮助,我想我明白了要点。但你假设我不具备一些 MR 知识。特别是,我不知道如何按照您早期的建议“追加”然后“合并”。
  • 我编辑了我的帖子。我可以说我的代码变得更难了,如果有任何更简单的现成代码或与之相关的现成库,你最好使用它。
猜你喜欢
  • 1970-01-01
  • 2011-01-05
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-02-01
  • 2013-09-30
  • 2010-10-25
  • 2013-05-02
相关资源
最近更新 更多