91超碰碰碰碰久久久久久综合_超碰av人澡人澡人澡人澡人掠_国产黄大片在线观看画质优化_txt小说免费全本

溫馨提示×

溫馨提示×

您好,登錄后才能下訂單哦!

密碼登錄×
登錄注冊×
其他方式登錄
點擊 登錄注冊 即表示同意《億速云用戶服務條款》

Apache Flink 官方文檔--流(DataStream API)-旁路輸出

發布時間:2020-07-18 16:35:03 來源:網絡 閱讀:2181 作者:Lynn_Yuan 欄目:大數據

旁路輸出(side output)

??除了來自數據流算子的主流結果輸出之外,可以產生任意數量的流旁路輸出結果。旁路輸出結果數據類型與主流結果的數據類型以及其他旁路輸出結果數據類型可以是完全不同的。當你需要分割數據流時,這個算子非常有用。通常需要復制流,然后從每個數據流中過濾掉不需要的數據。
??當使用旁路輸出時,首先需要定義一個OutputTag來標識一個旁路輸出流。
Java

// this needs to be an anonymous inner class, so that we can analyze the type
OutputTag<String> outputTag = new OutputTag<String>("side-output") {};

Scala

val outputTag = OutputTag[String]("side-output")

??注意OutputTag是如何根據旁路輸出流包含的元素類型typed的。
??可以通過以下函數發射數據到旁路輸出。

  • ProcessFunction
  • CoProcessFunction
  • ProcessWindowFunction
  • ProcessAllWindowFunction

??可以使用Context參數(在上述函數中向用戶暴露)將數據發送到OutputTag標識的旁路輸出。以下是從ProcessFunction發出旁路輸出數據的示例:
Java:

DataStream<Integer> input = ...;

final OutputTag<String> outputTag = new OutputTag<String>("side-output"){};

SingleOutputStreamOperator<Integer> mainDataStream = input
  .process(new ProcessFunction<Integer, Integer>() {

      @Override
      public void processElement(
          Integer value,
          Context ctx,
          Collector<Integer> out) throws Exception {
        // emit data to regular output
        out.collect(value);

        // emit data to side output
        ctx.output(outputTag, "sideout-" + String.valueOf(value));
      }
    });

Scala:

val input: DataStream[Int] = ...
val outputTag = OutputTag[String]("side-output")

val mainDataStream = input
  .process(new ProcessFunction[Int, Int] {
    override def processElement(
        value: Int,
        ctx: ProcessFunction[Int, Int]#Context,
        out: Collector[Int]): Unit = {
      // emit data to regular output
      out.collect(value)

      // emit data to side output
      ctx.output(outputTag, "sideout-" + String.valueOf(value))
    }
  })

  要讀取旁路輸出流,在數據流運算后使用getSideOutput(OutputTag)。此時將會獲得鍵入旁路輸出流的結果。
Java:

final OutputTag<String> outputTag = new OutputTag<String>("side-output"){};

SingleOutputStreamOperator<Integer> mainDataStream = ...;

DataStream<String> sideOutputStream = mainDataStream.getSideOutput(outputTag);

Scala:

val outputTag = OutputTag[String]("side-output")

val mainDataStream = ...

val sideOutputStream: DataStream[String] = mainDataStream.getSideOutput(outputTag)
向AI問一下細節

免責聲明:本站發布的內容(圖片、視頻和文字)以原創、轉載和分享為主,文章觀點不代表本網站立場,如果涉及侵權請聯系站長郵箱:is@yisu.com進行舉報,并提供相關證據,一經查實,將立刻刪除涉嫌侵權內容。

AI

济宁市| 昭平县| 海阳市| 柏乡县| 普陀区| 郑州市| 射阳县| 宣化县| 钦州市| 雷州市| 上栗县| 子长县| 夏津县| 环江| 临夏市| 泰州市| 云南省| 二手房| 石门县| 瓮安县| 康平县| 双鸭山市| 禹州市| 永平县| 罗定市| 祁阳县| 奉新县| 双峰县| 云霄县| 余姚市| 图木舒克市| 丹巴县| 郑州市| 米泉市| 北碚区| 兴国县| 卢龙县| 罗江县| 克山县| 佛山市| 屏东市|