【问题标题】:Does flink reduce on the fly in batch modeflink 是否在批处理模式下即时减少
【发布时间】: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


    【解决方案1】:

    根据 Flink 中的 documentation on the Reduce Operation,我看到以下内容:

    应用于分组数据集的 Reduce 转换可减少 使用用户定义的 reduce 函数将每个组归为单个元素。 对于每组输入元素,依次使用一个reduce函数 将一对元素组合成一个元素,直到只有一个元素 每个组的元素仍然存在。

    请注意,对于 ReduceFunction,返回对象的键控字段 应该匹配输入值。 这是因为 reduce 是隐式的 可组合 和从组合运算符发出的对象再次 传递给 reduce 运算符时按 key 分组。

    如果我没看错,Flink 会在 mapper 端执行 reduce 操作,然后在 reducer 端再次执行,因此实际发出/序列化的数据应该很小。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2017-02-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-04-14
      • 1970-01-01
      相关资源
      最近更新 更多