【发布时间】: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() 方法正在打印源流,对于接收器,我尝试了 process 和 select 方法,但它们都不能打印结果。此外,我的过滤器或匹配方法的日志根本不会打印。我的印象是甚至没有使用过滤器功能。 Event 和 Alarm 对象是简单的 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