使用 Apache Flink 创建 CEP
Creating CEP with Apache Flink
我正在尝试为 Kafka InputStream 实现一个非常简单的 Apache Flink CEP。
Kafka 生产者生成一个简单的 Double Value 并通过 Kafka 主题将它们作为字符串发送给消费者。目前,我正在使用 Flink 编写 CEP 消费者代码。
到目前为止,这是我编写的代码:
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();
}
如果我尝试执行此代码,这是错误消息:
Exception in thread "main" java.lang.NoSuchFieldError: NO_INDEX at
org.apache.flink.cep.PatternStream.select(PatternStream.java:102) at
CEPTest.main(CEPTest.java:50)
所以看起来我用 CEP 模式生成的 DataStream 是错误的,但我不知道那个方法有什么问题。每一个帮助都会很棒!
编辑:我尝试了其他一些示例,但每次执行时我都会遇到同样的错误。所以我觉得我的包裹有问题?
使用 Flink 1.6.0,我的代码可以完美运行。
我正在尝试为 Kafka InputStream 实现一个非常简单的 Apache Flink CEP。 Kafka 生产者生成一个简单的 Double Value 并通过 Kafka 主题将它们作为字符串发送给消费者。目前,我正在使用 Flink 编写 CEP 消费者代码。 到目前为止,这是我编写的代码:
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();
}
如果我尝试执行此代码,这是错误消息:
Exception in thread "main" java.lang.NoSuchFieldError: NO_INDEX at org.apache.flink.cep.PatternStream.select(PatternStream.java:102) at CEPTest.main(CEPTest.java:50)
所以看起来我用 CEP 模式生成的 DataStream 是错误的,但我不知道那个方法有什么问题。每一个帮助都会很棒!
编辑:我尝试了其他一些示例,但每次执行时我都会遇到同样的错误。所以我觉得我的包裹有问题?
使用 Flink 1.6.0,我的代码可以完美运行。