【问题标题】:Pre-shuffle aggregation in FlinkFlink 中的预混洗聚合
【发布时间】:2021-08-17 03:01:38
【问题描述】:

我们正在将 spark 作业迁移到 flink。我们在 spark 中使用了 pre-shuffle 聚合。有没有办法在火花中执行类似的操作。我们正在使用来自 apache kafka 的数据。我们正在使用键控翻转窗口来聚合数据。我们希望在执行 shuffle 之前聚合 flink 中的数据。

https://databricks.gitbooks.io/databricks-spark-knowledge-base/content/best_practices/prefer_reducebykey_over_groupbykey.html

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    是的,这是可能的,我将描述三种方式。首先是已经内置的 Flink Table API。第二种方法是您必须构建自己的预聚合运算符。第三个是动态预聚合算子,在shuffle阶段之前调整预聚合事件的数量。

    Flink 表 API

    作为it is shown here,您可以进行MiniBatch 聚合 或本地-全局聚合。第二种选择更好。您基本上告诉 Flink 为每个 5000 事件创建小批量,并在 shuffle 阶段之前预先聚合它们。

    // instantiate table environment
    TableEnvironment tEnv = ...
    
    // access flink configuration
    Configuration configuration = tEnv.getConfig().getConfiguration();
    // set low-level key-value options
    configuration.setString("table.exec.mini-batch.enabled", "true");
    configuration.setString("table.exec.mini-batch.allow-latency", "5 s");
    configuration.setString("table.exec.mini-batch.size", "5000");
    configuration.setString("table.optimizer.agg-phase-strategy", "TWO_PHASE");
    

    Flink Stream API

    这种方式比较麻烦,因为您必须使用OneInputStreamOperator 创建自己的运算符并使用doTransform() 调用它。以下是 BundleOperator 的示例。

    public abstract class AbstractMapStreamBundleOperator<K, V, IN, OUT>
      extends AbstractUdfStreamOperator<OUT, MapBundleFunction<K, V, IN, OUT>>
      implements OneInputStreamOperator<IN, OUT>, BundleTriggerCallback {
    @Override
     public void processElement(StreamRecord<IN> element) throws Exception {
      // get the key and value for the map bundle
      final IN input = element.getValue();
      final K bundleKey = getKey(input);
      final V bundleValue = this.bundle.get(bundleKey);
    
      // get a new value after adding this element to bundle
      final V newBundleValue = userFunction.addInput(bundleValue, input);
    
      // update to map bundle
      bundle.put(bundleKey, newBundleValue);
    
      numOfElements++;
      bundleTrigger.onElement(input);
     }
    
     @Override
     public void finishBundle() throws Exception {
      if (!bundle.isEmpty()) {
       numOfElements = 0;
       userFunction.finishBundle(bundle, collector);
       bundle.clear();
      }
      bundleTrigger.reset();
     }
    }
    

    回调接口定义何时触发预聚合。每次流达到if (count &gt;= maxCount) 的捆绑限制时,您的预聚合运算符都会向随机播放阶段发出事件。

    public class CountBundleTrigger<T> implements BundleTrigger<T> {
     private final long maxCount;
     private transient BundleTriggerCallback callback;
     private transient long count = 0;
    
     public CountBundleTrigger(long maxCount) {
      Preconditions.checkArgument(maxCount > 0, "maxCount must be greater than 0");
      this.maxCount = maxCount;
     }
    
     @Override
     public void registerCallback(BundleTriggerCallback callback) {
      this.callback = Preconditions.checkNotNull(callback, "callback is null");
     }
    
     @Override
     public void onElement(T element) throws Exception {
      count++;
      if (count >= maxCount) {
       callback.finishBundle();
       reset();
      }
     }
    
     @Override
     public void reset() {
      count = 0;
     }
    }
    

    然后您使用doTransform 致电您的接线员:

    myStream.map(....)
     .doTransform(metricCombiner, info, new RichMapStreamBundleOperator<>(myMapBundleFunction, bundleTrigger, keyBundleSelector))
     .map(...)
     .keyBy(...)
     .window(TumblingProcessingTimeWindows.of(Time.seconds(20)))
    

    动态预聚合

    如果您希望使用动态预聚合运算符,请查看AdCom - Adaptive Combiner for stream aggregation。它基本上是根据背压信号调整预聚合。这导致使用最大可能的洗牌阶段。

    【讨论】:

    • 看来 adcom 需要重新编译 flink 项目。
    • 是的。它以您可以调用yourStream.adcom(....) 的方式实现。那是因为JobManager调整了所有TaskManager中的AdCom参数。但是你可以在原 Flink doTransform(...) 获取代码并使用,无需动态调整。
    猜你喜欢
    • 2015-09-03
    • 1970-01-01
    • 2019-01-31
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多