【问题标题】:Apache Beam - Python : How to get the top 10 elements of a PCollection with Accumulation?Apache Beam - Python:如何通过累积获得 PCollection 的前 10 个元素?
【发布时间】:2019-10-30 04:15:45
【问题描述】:

我想像这样提取前 10 名的最高分:

Paul - 38
Michel - 27
Hugo - 27
Bob - 24
Kevin - 19
...
(10 elements)

我正在使用固定窗口和数据驱动触发器,在窗格收集 X 元素后输出早期结果。 另外,我正在使用组合器来获得前 10 名的最高分。

(inputs
         | 'Apply Window of time' >> beam.WindowInto(
                        beam.window.FixedWindows(size=5 * 60))
                        trigger=trigger.Repeatedly(trigger.AfterCount(2)),
                  accumulation_mode=trigger.AccumulationMode.ACCUMULATING)
         | beam.ParDo(PairWithOne()) # ('key', 1)
         | beam.CombinePerKey(sum)
         | 'Top 10 scores' >> beam.CombineGlobally(
                        beam.combiners.TopCombineFn(n=10,
                                                    compare=lambda a, b: a[1] < b[
                                                        1])).without_defaults())

这里的问题是第一个输出似乎是正确的,但以下输出包含重复的键:

Paul - 38
Paul - 37
Michel - 27
Paul - 36
Michel - 26
Kevin - 20
...
(10 elements)

如您所见,我得到的不是 10 个不同的 K/V 对,而是重复的密钥。

当不使用触发/累积策略时,这很有效。但如果我想有 2 小时的窗口,我想获得频繁更新...

【问题讨论】:

  • 您是否尝试过丢弃已触发的窗格? (即accumulation_mode=trigger.AccumulationMode.DISCARDING
  • 它与 DISCARDING 一起使用。但是我将此输出保存到 Cloud Firestore,并且使用该策略我不确定每次新更新都能正确更新我的数组。因为有了累积策略,我只需要用一个新的数组替换现有的数组。
  • 但我明白了……我想我应该改变保存数据的方式。
  • 否则,可能会在两者之间添加另一个 combineFn 以仅获得每个键的 Top 1,以便您获得每个用户的最新分数。编辑:我实际上是这样测试的,但我仍然得到重复的键......
  • 所以,这给了我修改 TopCombineFn 的想法,如果它们与来自新窗格的传入元素共享键,则从堆中弹出项目。它似乎对我有用,写一个(长)答案

标签: python google-cloud-dataflow apache-beam


【解决方案1】:

正如 cmets 中所讨论的,一种可能性是转换到 Discarding fired panes,可以通过 accumulation_mode=trigger.AccumulationMode.DISCARDING 设置。如果您仍想保留ACCUMULATING 模式,您可能需要修改TopCombineFn,以便同一用户的重复窗格覆盖以前的值并避免重复键。 TopDistinctFn 将以 Beam SDK 2.13.0 的代码 here 为基础。在add_input 方法中,我们将进行如下检查:

for current_top_element in enumerate(heap):
  if element[0] == current_top_element[1].value[0]:
    heap[current_top_element[0]] = heap[-1]
    heap.pop()
    heapq.heapify(heap)

基本上,我们将比较我们正在评估的元素的键 (element[0]) 与堆中的每个元素。堆元素的类型为ComparableValue,因此我们可以使用value 来取回元组(并使用value[0] 来获取密钥)。如果它们匹配,我们将希望将其从堆中弹出(因为我们正在累积总和会更大)。 Beam SDK 使用heapq 库,所以我基于this answer 删除i-th 元素(我们使用enumerate 保留索引信息)。

我添加了一些日志记录以帮助检测重复项:

logging.info("Duplicate: " + element[0] + "," + str(element[1]) + ' --- ' + current_top_element[1].value[0] + ',' + str(current_top_element[1].value[1]))

代码位于combiners 文件夹内的top.py 文件中(带有__init__.py),我将其导入:

from combiners.top import TopDistinctFn

然后,我可以在管道中使用TopDistinctFn,如下所示:

(inputs
     | 'Add User as key' >> beam.Map(lambda x: (x, 1)) # ('key', 1)
     | 'Apply Window of time' >> beam.WindowInto(
                    beam.window.FixedWindows(size=10*60),
                    trigger=beam.trigger.Repeatedly(beam.trigger.AfterCount(2)),
                    accumulation_mode=beam.trigger.AccumulationMode.ACCUMULATING)
     | 'Sum Score' >> beam.CombinePerKey(sum)   
     | 'Top 10 scores' >> beam.CombineGlobally(
                    TopDistinctFn(n=10, compare=lambda a, b: a[1] < b[1])).without_defaults()
     | 'Print results' >> beam.ParDo(PrintTop10Fn()))

完整代码可以在here找到。 generate_messages.py 是 Pub/Sub 消息生成器,top.py 包含自定义版本的 TopCombineFn 重命名为 TopDistinctFn(可能看起来很繁重,但我只添加了从第 425 行开始的几行代码)和 test_combine.py 主要管道代码。要运行它,您可以将文件放在正确的文件夹中,如果需要,安装 Beam SDK 2.13.0,修改 generate_messages.pytest_combine-py 中的项目 ID 和 Pub/Sub 主题。然后,使用python generate_messages.py 开始发布消息,并在不同的shell 中使用DirectRunner 运行管道:python test_combine.py --streaming。对于DataflowRunner,您可能需要将extra filessetup.py 文件一起添加。

例如,Bob 以 9 分领先,而在下一次更新到来时,他的得分高达 11 分。他将出现在下一次回顾中,只有更新的分数,没有重复(在我们的日志中检测到)。 9 分的条目将不会出现,并且顶部仍将根据需要有 10 个用户。 Marta 也是如此。我注意到即使不在前 10 名中,旧分数仍然出现在堆中,但我不确定垃圾收集如何与 heapq 一起工作。

INFO:root:>>> Current top 10: [('Bob', 9), ('Connor', 8), ('Eva', 7), ('Hugo', 7), ('Paul', 6), ('Kevin', 6), ('Laura', 6), ('Marta', 6), ('Diane', 4), ('Bacon', 4)]
...
INFO:root:Duplicate: Marta,8 --- Marta,6
INFO:root:Duplicate: Bob,11 --- Bob,9
INFO:root:>>> Current top 10: [('Bob', 11), ('Connor', 8), ('Marta', 8), ('Bacon', 7), ('Eva', 7), ('Hugo', 7), ('Paul', 6), ('Laura', 6), ('Diane', 6), ('Kevin', 6)]

让我知道这是否也适用于您的用例。

【讨论】:

  • 感谢您提供如此惊人而详细的答案!这似乎工作得很好。我希望这会添加到 apache-beam 的未来更新中
猜你喜欢
  • 2023-04-10
  • 2023-04-10
  • 1970-01-01
  • 2018-06-24
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-12-31
相关资源
最近更新 更多