【问题标题】:Apache Flink pattern detection does not find any matchApache Flink 模式检测未找到任何匹配项
【发布时间】:2021-06-22 18:29:34
【问题描述】:

我正在尝试使用 Apache Flink CEP(复杂事件处理)库来捕捉模式。我从以下结构开始,我希望看到 2 个匹配 id [1,2] 和 [3,4]。但是我没有看到任何结果。

public class StreamingJob {

    private static Logger logger = LoggerFactory.getLogger(StreamingJob.class);

    public static void main(String[] args) throws Exception {
        // set up the streaming execution environment
        final StreamExecutionEnvironment env = StreamExecutionEnvironment
                .getExecutionEnvironment();

        ArrayList <Event> strings = new ArrayList <>();
        strings.add(new Event(0L, "room", 9));
        strings.add(new Event(1L, "room", 10));
        strings.add(new Event(2L, "garden", 11));
        strings.add(new Event(3L, "room", 12));
        strings.add(new Event(4L, "garden", 13));
        strings.add(new Event(5L, "room", 14));
        strings.add(new Event(6L, "room", 15));


        KeyedStream <Event, String> source = env.fromCollection(strings).keyBy(Event::getName);

        source.print("###-source");

        Pattern <Event, ?> pattern = Pattern. <Event>begin("room")
                .where(new SimpleCondition <Event>() {
                    @Override
                    public boolean filter(Event value) {
                        logger.info("### value: {}", value);
                        return value.getName().equals("room");
                    }
                })
                .next("garden")
                .where(new SimpleCondition <Event>() {
                    @Override
                    public boolean filter(Event value) {
                        logger.info("### value: {}", value);
                        return value.getName().equals("garden");
                    }
                });

        PatternStream <Event> patternStream = CEP.pattern(source, pattern);


        // process 
        DataStream <Alarm> result = patternStream.process(
                new PatternProcessFunction <Event, Alarm>() {
                    @Override
                    public void processMatch(
                            Map <String, List <Event>> pattern,
                            Context ctx,
                            Collector <Alarm> out) throws Exception {
                        logger.info("### pattern: {}", pattern);
                        logger.info("### ctx: {}", ctx);
                        out.collect(new Alarm(pattern.toString()));
                    }
                });

        result.print("###");

        // or select function
        patternStream.select(new PatternSelectFunction <Event, Alarm>() {
            @Override
            public Alarm select(Map <String, List <Event>> pattern) throws Exception {
                logger.info("###");
                return new Alarm(pattern.toString());
            }
        }).print("###");

        // execute program
        env.execute("Flink Streaming Java API Skeleton");
    }
}

source.print() 方法正在打印源流,对于接收器,我尝试了 processselect 方法,但它们都不能打印结果。此外,我的过滤器或匹配方法的日志根本不会打印。我的印象是甚至没有使用过滤器功能。 EventAlarm 对象是简单的 pojo,如下所示:


public class Event implements Serializable {
    Long id;
    String name;
    Integer temperature;
    //...
}

public class Alarm implements Serializable {
    String text;
    // ...
}


我还尝试更改它以过滤其他字段,例如,使用温度字段,我想捕获 odd-even 数字序列但仍然没有打印任何结果。

【问题讨论】:

    标签: apache-flink flink-streaming complex-event-processing flink-cep stream-processing


    【解决方案1】:

    CEP 依赖于能够按时间戳对事件流进行排序。这要求您提供WatermarkStrategy(如果您想使用事件时间戳)或指定您希望模式匹配以摄取顺序完成,使用处理时间语义。

    你可以通过做这个小改动来做后者:

    PatternStream <Event> patternStream = 
      CEP.pattern(source, pattern).inProcessingTime();
    

    【讨论】:

    • 这有帮助。谢谢!
    猜你喜欢
    • 2021-12-24
    • 2022-08-19
    • 1970-01-01
    • 1970-01-01
    • 2021-09-24
    • 2017-07-19
    • 2018-07-21
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多