Создание CEP с помощью Apache Flink

Я пытаюсь реализовать очень простой Apache Flink CEP для Kafka InputStream. Производитель Kafka генерирует простые двойные значения и отправляет их через тему Kafka в виде строки потребителям. В настоящий момент я кодирую CEP Consumer с помощью Flink. Пока это мой написанный код:

public static void main(String[] args) throws Exception {

        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.getConfig().disableSysoutLogging();
        env.getConfig().setRestartStrategy(RestartStrategies.fixedDelayRestart(4, 10000)); 
        env.setParallelism(3);
        Properties properties = new Properties();
        properties.setProperty("bootstrap.servers", "localhost:9092");
        properties.setProperty("group.id", "flink_consumer");

        DataStream<String> stream = env
                .addSource(new FlinkKafkaConsumer09<>("temp", new SimpleStringSchema(), properties));

        Pattern<String, ?> warning= Pattern.<String>begin("first")
                .where(new IterativeCondition<String>() {
                    private static final long serialVersionUID = 1L;
                    @Override
                    public boolean filter(String value, Context<String> ctx) throws Exception {
                        return Double.parseDouble(value) >= 89.0;
                    }
                })
                .next("second")
                .where(new IterativeCondition<String>() {
                    private static final long serialVersionUID = 1L;
                    @Override
                    public boolean filter(String value, Context<String> ctx) throws Exception {
                        return Double.parseDouble(value) >= 89.0;
                    }
                })
                .within(Time.seconds(10));  
        DataStream<String> temp = CEP.pattern(stream, warning).select(new PatternSelectFunction<String, String>() {
            private static final long serialVersionUID = 1L;

            @Override
            public String select(Map<String, List<String>> pattern) throws Exception {
                List warnung1 = pattern.get("first");
                String first = (String) warnung1.get(1);
                return first;
            }   

        });

        temp.print();
        env.execute();

    }

если я пытаюсь выполнить этот код, это сообщение об ошибке:

Исключение в потоке "main" java.lang.NoSuchFieldError: NO_INDEX в org.apache.flink.cep.PatternStream.select (PatternStream.java:102) в CEPTest.main (CEPTest.java:50)

Итак, похоже, что мой сгенерированный DataStream с шаблоном CEP неверен, но я не знаю, что не так с этим методом. Любая помощь была бы замечательной!

Изменить: я пробовал другой пример, и при каждом выполнении я получаю ту же ошибку. Я думаю, что с моими пакетами что-то не так?


person Leon    schedule 13.09.2018    source источник
comment
Убедитесь, что вы используете одну и ту же версию для всех компонентов flink. Похоже, вы используете более новую версию модуля flink-cep поверх какой-то старой версии кластера flink.   -  person Dawid Wysakowicz    schedule 14.09.2018
comment
Кроме того, если вы собираетесь использовать первый элемент списка, используйте warnung1.get (0).   -  person David Anderson    schedule 14.09.2018
comment
@DawidWysakowicz Итак, я использую Flink версии 1.12 с соединителем Flink-CEP 2.11 - как вы думаете, в этом проблема? Как я могу обновить свой Flink? Проблема в том, что я использовал существующий проект Flink из Интернета, потому что мне было трудно настраивать / выполнять файлы из учебника ci.apache.org/projects/flink/flink-docs-stable/quickstart/   -  person Leon    schedule 14.09.2018
comment
Хорошо, я сделал новый Flink Projekt, и мой код работает отлично. Спасибо за вашу помощь!   -  person Leon    schedule 14.09.2018


Ответы (1)


С Flink 1.6.0 мой код работает отлично.

person Leon    schedule 14.09.2018