I have created a FLINK CEP application. Instead of using a flink cluster, I am using the LocalExecutionEnvironment.
I cannot get CEP to work with LocalExecutionEnvironment. I am expecting the println from the pattern-filter (first, second, end), and write functions to be invoked. However, only writeWatermark of the sink is being invoked.
Can someone provide pointers to running Flink CEP in local
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment();
env.setParallelism(1);
DataStream<String> events = env.fromElements("a", "a", "b", "b", "b", "b", "a", "a", "a", "a", "b", "b", "a", "b", "b", "x");
Pattern<String, String> pattern = Pattern.<String>begin("first")
.where(new SimpleCondition<String>() {
@Override
public boolean filter(String element) throws Exception {
System.out.println("filter:first");
return (element.equals("a"));
}
}).oneOrMore()
.next("second")
.where(new SimpleCondition<String>() {
@Override
public boolean filter(String element) throws Exception {
System.out.println("filter:second");
return (element.equals("b"));
}
})
.next("end")
.where(new SimpleCondition<String>() {
@Override
public boolean filter(String element) throws Exception {
System.out.println("filter:end");
return (!element.equals("b"));
}
});
CEP.pattern(events, pattern)
.process(new PatternProcessFunction<String, String>() {
@Override
public void processMatch(Map<String, List<String>> match, Context ctx, Collector<String> out) throws Exception {
System.out.println("process");;
}
}).addSink(new SinkFunction<String>() {
@Override
public void invoke(String value) throws Exception {
System.out.println("sink:write");
}
@Override
public void invoke(String value, Context context) throws Exception {
System.out.println("sink:write");
}
@Override
public void writeWatermark(Watermark watermark) throws Exception {
System.out.println("sink:watermark");
}
});
env.execute();
new Scanner(System.in).nextLine();
}
My gradle dependencies are:
implementation 'org.apache.flink:flink-streaming-java:1.15.1'
implementation 'org.apache.flink:flink-cep:1.15.1'
implementation 'org.apache.flink:flink-clients:1.15.1'