【发布时间】:2018-08-28 03:27:31
【问题描述】:
我们有一个数据流,其中每个元素都属于这种类型:
id: String
type: Type
amount: Integer
我们希望聚合此流并每周输出一次amount 的总和。
当前解决方案:
一个示例 flink 管道如下所示:
stream.keyBy(type)
.window(TumblingProcessingTimeWindows.of(Time.days(7)))
.reduce(sumAmount())
.addSink(someOutput())
输入
| id | type | amount |
| 1 | CAT | 10 |
| 2 | DOG | 20 |
| 3 | CAT | 5 |
| 4 | DOG | 15 |
| 5 | DOG | 50 |
如果窗口在记录 3 和 4 之间结束,我们的输出将是:
| TYPE | sumAmount |
| CAT | 15 | (id 1 and id 3 added together)
| DOG | 20 | (only id 2 as been 'summed')
Id 4 和 5 仍将在 flink 管道中,并将在下周输出。
因此下周我们的总产量将是:
| TYPE | sumAmount |
| CAT | 15 | (of last week)
| DOG | 20 | (of last week)
| DOG | 65 | (id 4 and id 5 added together)
新要求:
我们现在还想知道每条记录在哪周处理了每条记录。换句话说,我们的新输出应该是:
| TYPE | sumAmount | weekNumber |
| CAT | 15 | 1 |
| DOG | 20 | 1 |
| DOG | 65 | 2 |
但我们还想要这样的额外输出:
| id | weekNumber |
| 1 | 1 |
| 2 | 1 |
| 3 | 1 |
| 4 | 2 |
| 5 | 2 |
如何处理?
flink 有没有办法做到这一点?我想我们会有一个聚合函数,它可以对金额求和,但也会输出每个记录以及当前周数,但我在文档中找不到这样做的方法。
(注意:我们每周处理大约 1 亿条记录,因此理想情况下,我们只想在一周内将聚合保持在 flink 的状态,而不是所有单独的记录)
编辑:
我选择了下面安东描述的解决方案:
DataStream<Element> elements =
stream.keyBy(type)
.process(myKeyedProcessFunction());
elements.addSink(outputElements());
elements.getSideOutput(outputTag)
.addSink(outputAggregates())
KeyedProcessFunction 看起来像:
class MyKeyedProcessFunction extends KeyedProcessFunction<Type, Element, Element>
private ValueState<ZonedDateTime> state;
private ValueState<Integer> sum;
public void processElement(Element e, Context c, Collector<Element> out) {
if (state.value() == null) {
state.update(ZonedDateTime.now());
sum.update(0);
c.timerService().registerProcessingTimeTimer(nowPlus7Days);
}
element.addAggregationId(state.value());
sum.update(sum.value() + element.getAmount());
}
public void onTimer(long timestamp, OnTimerContext c, Collector<Element> out) {
state.update(null);
c.output(outputTag, sum.value());
}
}
【问题讨论】:
标签: apache-flink flink-streaming