Apache Flink:ProcessWindowFunction不适用

Ksp*_*ace 3 stream-processing apache-flink flink-streaming

我想ProcessWindowFunction在我的Apache Flink项目中使用。但是使用过程函数时出现一些错误,请参见下面的代码片段

错误是:

在类型WindowedStream,元组,TimeWindow>的方法处理(ProcessWindowFunction,R,元组,TimeWindow>)是不适用的参数(JDBCExample.MyProcessWindows)

我的程序:

DataStream<Tuple2<String, JSONObject>> inputStream;

inputStream = env.addSource(new JsonArraySource());

inputStream.keyBy(0)
  .window(TumblingEventTimeWindows.of(Time.minutes(10)))
  .process(new MyProcessWindows());
Run Code Online (Sandbox Code Playgroud)

我的ProcessWindowFunction:

private class MyProcessWindows 
  extends ProcessWindowFunction<Tuple2<String, JSONObject>, Tuple2<String, String>, String, Window>
{

  public void process(
      String key, 
      Context context, 
      Iterable<Tuple2<String, JSONObject>> input, 
      Collector<Tuple2<String, String>> out) throws Exception 
  {
    ...
  }

}
Run Code Online (Sandbox Code Playgroud)

Fab*_*ske 5

问题可能是的通用类型ProcessWindowFunction。

您正在按位置(keyBy(0))引用键。因此,编译器无法推断其类型(String),您需要将其更改ProcessWindowFunction为:

private class MyProcessWindows 
    extends ProcessWindowFunction<Tuple2<String, JSONObject>, Tuple2<String, String>, Tuple, Window>
Run Code Online (Sandbox Code Playgroud)

通过替换String为Tuple您,您现在拥有一个通用的占位符,可用于Tuple1<String>在需要访问processElement()方法中的键时可以转换为的键:

public void process(
    Tuple key, 
    Context context, 
    Iterable<Tuple2<String, JSONObject>> input, 
    Collector<Tuple2<String, String>> out) throws Exception {

  String sKey = (String)((Tuple1)key).f0;
  ...
}
Run Code Online (Sandbox Code Playgroud)

您可以避免演员,如果你定义了使用正确的类型KeySelector<IN, KEY>函数来提取关键,因为返回类型KEY的KeySelector被称为编译器。