【发布时间】:2018-09-21 02:00:45
【问题描述】:
根据 flink 流文档:
窗口函数可以是 ReduceFunction、FoldFunction 或 窗口函数。前两个可以更有效地执行(参见 State Size 部分),因为 Flink 可以增量聚合 每个窗口的元素到达时。
批处理模式是否同样适用?在下面的示例中,我正在从 cassandra 读取约 36go 的数据,但我希望减少的输出要小得多(约 0.5go)。运行此作业是否需要 flink 将整个输入存储在内存中,或者它是否足够聪明以对其进行迭代
DataSet<MyRecord> input = ...;
DataSet<MyRecord> sampled = input
.groupBy(MyRecord::getSampleKey)
.reduce(MyRecord::keepLast);
【问题讨论】:
标签: java mapreduce apache-flink