Apache Flink 1.16.0 CEP程序退出码为0却无控制台输出求助
Flink 1.16.0 CEP程序执行成功但控制台无输出问题排查
我正在使用Apache Flink 1.16.0版本,尝试编写一个简单的CEP程序将元素打印到控制台。但不知为何,即使程序以退出码0执行完成,控制台也没有任何内容输出。
相关代码
import org.apache.flink.cep.*; import org.apache.flink.cep.functions.PatternProcessFunction; import org.apache.flink.cep.pattern.Pattern; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.util.Collector; import java.util.List; import java.util.Map; public class Main { public static void main(String[] args) throws Exception { final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); DataStream<String> input = env.fromElements("adrian", "cincu", "diana", "sichitiu"); Pattern<String, ?> pattern = Pattern.begin("start"); PatternStream<String> patternStream = CEP.pattern(input, pattern); DataStream<String> matches = patternStream.process(new PatternProcessFunction<String, String>() { @Override public void processMatch(Map<String, List<String>> map, Context context, Collector<String> collector) throws Exception { collector.collect(map.get("start").toString()); } }); matches.print(); env.execute("myjob"); } }
解决线索
问题核心在于Flink CEP对时间语义和水位线的依赖:默认情况下CEP基于事件时间处理逻辑,而你生成的输入数据流没有携带时间戳、也未生成水位线,导致CEP无法判断何时可以安全输出匹配结果,最终作业结束时也不会触发输出动作。
可以通过以下几种方式解决:
切换为批处理模式:由于你的输入是有限数据集,将作业切换为批处理模式后,Flink会自动处理完所有元素并输出结果,无需额外配置时间戳或水位线:
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setRuntimeMode(RuntimeMode.BATCH); // 添加该行代码改用处理时间语义:如果需要保持流处理模式,将时间特征改为处理时间,CEP会立即处理到达的元素:
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // Flink 1.12+推荐的设置方式 env.getConfig().setAutoWatermarkInterval(0); DataStream<String> input = env.fromElements("adrian", "cincu", "diana", "sichitiu") .assignTimestampsAndWatermarks(WatermarkStrategy.<String>forMonotonousTimestamps() .withTimestampAssigner((element, ts) -> System.currentTimeMillis()));为数据流分配时间戳和水位线:如果要保留事件时间语义,必须为每个元素分配时间戳并生成水位线,让CEP能够推进处理流程,示例代码如上。
内容的提问来源于stack exchange,提问作者Adrian Cincu
相关产品推荐
相关产品推荐

