【发布时间】: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