【发布时间】:2018-03-20 05:44:36
【问题描述】:
我想在我的 Apache Flink 项目中使用 ProcessWindowFunction。但是我在使用进程函数时遇到了一些错误,请参见下面的代码 sn-p
错误是:
WindowedStream,Tuple,TimeWindow> 类型中的方法 process(ProcessWindowFunction,R,Tuple,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());
我的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
{
...
}
}
【问题讨论】:
-
你能仔细检查错误信息吗?似乎
process()方法的签名中缺少某些内容。
标签: apache-flink flink-streaming stream-processing