【问题标题】:Stream de-duplication on Dataflow | Running services on Dataflow servicesDataflow 上的流式重复数据删除 |在 Dataflow 服务上运行服务
【发布时间】:2017-01-25 15:52:28
【问题描述】:

我想以窗口方式对基于 ID 的数据流进行重复数据删除。我们收到的流有并且我们想要在 N 小时时间窗口内删除匹配的数据。一种直接的方法是使用外部密钥库(BigTable 或类似的东西),我们在其中查找密钥并在需要时写入,但我们的 qps 非常大,使得维护这样的服务非常困难。我想出的另一种方法是在一个时间窗口内分组,这样一个时间窗口内用户的所有数据都属于同一个组,然后在每个组中,我们使用一个单独的密钥存储服务来查找键重复。所以,我对这种方法有几个问题

[1] 如果我运行 groupBy 转换,是否可以保证每个组将在同一个从站中处理?如果有保证,我们可以按 userid 分组,然后在每个组内比较每个用户的 sessionid

[2] 如果可行,我的下一个问题是我们是否可以在运行该作业的每台从机中运行此类其他服务 - 在上面的示例中,我希望运行一个本地 Redis,它可以然后每个小组也可以使用它来查找或写入 ID。

这个想法似乎与 Dataflow 应该做的不同,但我相信这样的用例应该很常见 - 所以如果有更好的模型来解决这个问题,我也很期待。考虑到我们拥有的数据量,我们基本上希望尽可能避免外部查找。

【问题讨论】:

    标签: google-cloud-dataflow


    【解决方案1】:

    1) 在 Dataflow 模型中,不能保证同一台机器会跨窗口看到所有组的键。想象一下,一个 VM 死机或添加了新的 VM,并且工作被分配到它们之​​间以进行扩展。

    2) 欢迎您在 Dataflow 虚拟机上运行其他服务,因为它们是通用的,但请注意,您必须应对主机上其他应用程序的资源需求,这可能会导致内存不足问题。

    请注意,您可能需要查看RemoveDuplicates,如果它适合您的用例,请使用它。

    您似乎还想使用session windows 对元素进行重复数据删除。你会打电话给:

    PCollection<T> pc = ...;
    PCollection<T> windowed_pc = pc.apply(
        Window<T>into(Sessions.withGapDuration(Duration.standardMinutes(N hours))));
    

    每个新元素都会继续延长窗口的长度,因此在间隙关闭之前它不会关闭。如果您还在下游 GroupByKey 上应用 AfterCount 推测触发器和 AfterWatermark 触发器,则为 1。触发器将尽快触发,一旦它看到至少一个元素,然后在会话关闭时再次触发。在 GroupByKey 之后,您将有一个 DoFn,它根据窗格信息([3]、[4])过滤掉不是早期触发的元素。

    DoFn(T -> KV<session key, T>)
             |
            \|/
    Window.into(Session window)
             |
            \|/
    Group by key
             |
            \|/
    DoFn(Filter based upon pane information)
    

    你的描述有点不清楚,你能提供更多细节吗?

    【讨论】:

    • 感谢您的更新。我想我错过了 RemoveDuplicates 部分,因为它几乎是我想要的。我的问题几乎是关于删除可以通过会话密钥识别的重复数据,尽管我担心两件事 - 会话的持续时间,因为重复数据可能很晚,还有组本身的基数,因为我们可能有几个用户创建尽可能多的组。我对您在这里提到的触发模型并不十分熟悉,但这些似乎是不错的开始。我可以看一下,然后就此事与您联系。再次感谢。
    【解决方案2】:

    抱歉,不清楚。我尝试了你提到的设置,除了早期和晚期点火部分,它正在处理更小的样本。我有几个与扩大规模有关的后续问题。另外,我希望我能给你更多关于确切情况的信息。

    因此,我们有传入的数据流,其中的每个项目都可以通过其字段唯一标识。我们也知道重复发生的距离很远,现在,我们关心的是 6 小时内的重复。关于数据量,我们每秒至少有 10 万个事件,跨越一百万个不同的用户 - 因此在这 6 小时的窗口内,我们可以将数十亿个事件纳入管道。

    鉴于这个背景,我的问题是 [1] 对于按键发生的会话,我应该在类似

    的东西上运行它
    PCollection<KV<key, T>> windowed_pc = pc.apply(
            Window<KV<key,T>>into(Sessions.withGapDuration(Duration.standardMinutes(6 hours))));
    

    key 是我之前提到的 3 个 ID 的组合。根据 Sessions 的定义,只有在这个 KV 上运行它,我才能按密钥管理会话。这意味着 Dataflow 在任何给定时间都会有太多打开的会话等待它们关闭,我担心它是否会扩展或者我会遇到任何瓶颈。

    [2] 如上所述执行会话后,我已经根据触发删除了重复项,因为我只关心每个会话中已经销毁重复项的第一次触发。我不再需要 RemoveDuplicates 变换,我发现它是 (WithKeys, Combine.PerKey, Values) 变换的组合,基本上执行相同的操作。这是正确的假设吗?

    [3] 如果 [1] 中的解决方案有问题,另一种方法是将会话的key 减少为只是user-id, session-id,忽略sequence-id,然后运行RemoveDuplicates on每个结果窗口的顶部sequence-id。这可能会减少打开会话的数量,但仍然会留下很多打开的会话(#users * #sessions per user),很容易达到数百万。 FWIW,我不认为我们只能通过user-id 进行会话,因为那样会话可能永远不会关闭,因为同一用户的不同会话可能会不断进入,并且在这种情况下确定会话间隙变得不可行。

    希望这次我的问题更清楚一点。请让我知道我的任何方法都充分利用了 Dataflow,或者如果我遗漏了什么。

    谢谢

    【讨论】:

    • [1] 目前,Dataflow 跨机器分配工作的关键范围,因此关键的数量与机器的数量成比例。如果您遇到问题,我会采用这种方法并使用工作 ID 与我们同步。 [2] 你是对的,不需要 RemoveDuplicates,因为你已经在为第一次触发触发进行过滤。 [3] 这也可以,但在 (user-id, session-id) 键没有更多事件后,您必须等待 6 小时。尽早解雇重要吗?
    • 很高兴知道它可以扩展 - 我会试试这个。至于早期解雇,是的,这是有道理的,因为我们只关心第一个非重复的。如果 [1] 会起作用,我不想通过 [3]。感谢更新
    【解决方案3】:

    我在更大范围内尝试了这个解决方案,只要我提供足够的工人和磁盘,管道的扩展性很好,尽管我现在看到了一个不同的问题。

    在此会话化之后,我在密钥上运行Combine.perKey,然后执行ParDo,它查看c.pane().getTiming(),只拒绝EARLY 触发以外的任何内容。我尝试在这个ParDo 中计算EARLYONTIME 的触发次数,看起来准时窗格实际上比早期的窗格更精确地进行了重复数据删除。我的意思是,#early-firings 仍然有一些重复项,而#ontime-firings 则更少,并且删除了更多重复项。有什么理由会发生这种情况吗?另外,我使用Combine+ParDo 进行重复数据删除的方法是否正确,或者我可以做得更好吗?

    events.apply(
        WithKeys.<String, EventInfo>of(new SerializableFunction<EventInfo, String>() {
            @Override
            public java.lang.String apply(EventInfo input) {
                return input.getUniqueKey();
            }
        })
    )
    .apply(
        Window.named("sessioner").<KV<String, EventInfo>>into(
            Sessions.withGapDuration(mSessionGap)
        )
        .triggering(
            AfterWatermark.pastEndOfWindow()
            .withEarlyFirings(AfterPane.elementCountAtLeast(1))
        )
        .withAllowedLateness(Duration.ZERO)
        .accumulatingFiredPanes()
    );
    

    【讨论】:

      猜你喜欢
      • 2020-08-16
      • 1970-01-01
      • 2016-06-19
      • 2018-10-21
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多