Why is Flink CEP not working in local execution environment, created using createLocalEnvironment?

Viewed 31

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'
0 Answers
Related