代码之家  ›  专栏  ›  技术社区  ›  Carl Rynegardh

Flink Scala-扩展WindowFunction

  •  1
  • Carl Rynegardh  · 技术社区  · 8 年前

    我在想怎么写我自己的 WindowFunction 但是有问题,我不知道为什么。我遇到的问题是apply函数,因为它无法识别 MyWindowFunction 作为有效输入,所以我无法编译。我正在传输的数据包含 (timestamp,x,y) 其中,对于测试,x和y分别为0和1。 extractTupleWithoutTs 只返回一个元组 (x,y) . 我已经成功地用简单的求和和和归约函数运行了代码。感谢您的帮助:)使用Flink 1.3

    进口:

    import org.apache.flink.streaming.api.TimeCharacteristic
    import org.apache.flink.streaming.api.functions.AssignerWithPeriodicWatermarks
    import org.apache.flink.streaming.api.scala.function.WindowFunction
    import org.apache.flink.streaming.api.watermark.Watermark
    import org.apache.flink.streaming.api.windowing.windows.TimeWindow
    import org.apache.flink.util.Collector
    

    其余代码:

    val env = StreamExecutionEnvironment.getExecutionEnvironment
    env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime)
    val text = env.socketTextStream("localhost", 9999).assignTimestampsAndWatermarks(new TsExtractor)
    val tuple = text.map( str => extractTupleWithoutTs(str))
    val counts = tuple.keyBy(0).timeWindow(Time.seconds(5)).apply(new MyWindowFunction())
    counts.print()
    env.execute("Window Stream")
    

    MyWindow函数,基本上是从示例中复制粘贴并更改类型。

    class MyWindowFunction extends WindowFunction[(Int, Int), Int, Int, TimeWindow] {
      def apply(key: Int, window: TimeWindow, input: Iterable[(Int, Int)], out: Collector[Int]): () = {
        var count = 0
        for (in <- input) {
          count = count + 1
        }
        out.collect(count)
      }
    }
    
    1 回复  |  直到 8 年前
        1
  •  3
  •   Fabian Hueske    8 年前

    问题是 WindowFunction ,即键的类型。该键是用中的索引声明的 keyBy 方法( keyBy(0) ). 因此,无法在编译时确定键的类型。如果将键声明为字符串,也会出现同样的问题,即。, keyBy("f0") .

    有两种解决方法:

    1. 使用 KeySelector keyBy公司 提取密钥(类似于 keyBy(_._1) ). 返回类型 按键选择器 函数在编译时已知,因此可以使用正确类型的 窗口函数 Int 钥匙
    2. 更改的第三个类型参数的类型 窗口函数 到 org.apache.flink.api.java.tuple.Tuple ,即。, WindowFunction[(Int, Int), Int, org.apache.flink.api.java.tuple.Tuple, TimeWindow] . Tuple 是由提取的密钥的通用持有人 keyBy公司 . 在你的情况下,这将是一个 org.apache.flink.api.java.tuple.Tuple1 . 在里面 WindowFunction.apply() 你可以投 元组 到 Tuple1 并通过访问键字段 Tuple1.f0 .