Flink任务未从Kafka Source输出预期告警的排查求助
问题描述
我正在开发一个Flink任务,从Kafka Source读取数据,通过CEP模式处理后将告警信息打印到控制台。执行以下命令:
./bin/flink run /Users/spartacus/icu-alarm/target/flink-kafka-stroke-risk-1.0-SNAPSHOT-jar-with-dependencies.jar > out.txt
out.txt文件中仅显示:
Job has been submitted with JobID 3811369c6f7f14d0eca0a66072550414
预期应输出类似Stroke Risk Alert: Patient ID - XYZ, Risk Level - 5的告警信息,但实际未打印。相关代码如下:
package hes.cs63.CEPMonitor; import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.api.common.serialization.SimpleStringSchema; import org.apache.flink.cep.CEP; import org.apache.flink.cep.pattern.Pattern; import org.apache.flink.cep.PatternStream; import org.apache.flink.cep.PatternSelectFunction; import org.apache.flink.cep.pattern.conditions.SimpleCondition; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.datastream.DataStreamSource; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.connector.kafka.source.KafkaSource; import org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema; import org.apache.kafka.common.serialization.StringDeserializer; import java.util.Arrays; import java.util.List; import java.util.Map; import java.util.Properties; public class StrokeRiskAlarm { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); Properties properties = new Properties(); properties.setProperty("bootstrap.servers", "localhost:9092"); properties.setProperty("group.id", "test-consumer-group"); KafkaSource<String> source = KafkaSource.<String>builder() .setBootstrapServers("localhost:9092") .setGroupId("test-consumer-group") .setTopics(Arrays.asList("geCEP")) .setDeserializer(KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class)) //.setStartingOffsets(OffsetsInitializer.earliest()) .setValueOnlyDeserializer(new SimpleStringSchema()) .build(); DataStreamSource<String> patientData = env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Source"); // Use KafkaSource Pattern<String, ?> highRiskPattern = Pattern.<String>begin("first") .where(new SimpleCondition<String>() { @Override public boolean filter(String value) { return getValue(value) > 3; } }); PatternStream<String> patternStream = CEP.pattern( patientData, highRiskPattern ); DataStream<String> strokeRiskAlerts = patternStream.select(new PatternSelectFunction<String, String>() { @Override public String select(Map<String, List<String>> pattern) throws Exception { String userId = pattern.get("first").get(0).split(",")[0].trim(); int risk = getTotalRisk(pattern.get("first").get(0), pattern.get("middle").get(0), pattern.get("last").get(0)); return "Stroke Risk Alert: Patient ID - " + userId + ", Risk Level - " + risk; } }); strokeRiskAlerts.print(); env.execute("Stroke Risk Alert Job"); } private static int getValue(String value) { // Parse the input value and extract the relevant data for risk calculation String[] parts = value.split(","); String type = parts[2].trim(); double measurementValue = Double.parseDouble(parts[3].trim()); // Implement your logic to check high stroke risk for different measurements if (type.equals("HR")) { // HeartMeasurement risk calculation logic int risk = 0; risk += measurementValue <= 50 ? 1 : 0; risk += measurementValue <= 40 ? 1 : 0; risk += measurementValue >= 91 ? 1 : 0; risk += measurementValue >= 110 ? 1 : 0; risk += measurementValue >= 131 ? 1 : 0; return risk; } else if (type.equals("SBP")) { // BloodPressureMeasurement risk calculation logic int risk = 0; risk += measurementValue <= 110 ? 1 : 0; risk += measurementValue <= 100 ? 1 : 0; risk += measurementValue <= 90 ? 1 : 0; risk += measurementValue >= 220 ? 3 : 0; return risk; } else if (type.equals("TEMP")) { // TempMeasurement risk calculation logic int risk = 0; risk += measurementValue <= 36 ? 1 : 0; risk += measurementValue <= 35 ? 2 : 0; risk += measurementValue >= 38.1 ? 1 : 0; risk += measurementValue >= 39.1 ? 1 : 0; return risk; } // Default risk calculation return 0; } private static int getTotalRisk(String firstValue, String middleValue, String lastValue) { // Calculate the total risk based on the values of the first, middle, and last measurements int firstRisk = getValue(firstValue); int middleRisk = getValue(middleValue); int lastRisk = getValue(lastValue); return firstRisk + middleRisk + lastRisk; } }
问题原因与排查方案
一、核心问题分析
CEP模式与select方法不匹配
代码定义的CEP模式是仅包含first节点的单事件匹配,但在select方法中强行获取不存在的middle和last节点数据,会直接抛出NullPointerException。Flink任务出现未捕获异常时,对应算子会静默停止,客户端不会输出错误信息。Kafka Source配置冲突
同时设置了两种反序列化逻辑:setDeserializer和setValueOnlyDeserializer,会导致反序列化逻辑冲突,可能无法正确解析Kafka消息。未配置起始偏移量
注释掉了setStartingOffsets(OffsetsInitializer.earliest()),默认会从Kafka topic的最新偏移量开始消费,若topic中没有新消息,任务不会读取到任何数据。
二、排查与修复步骤
1. 修复代码逻辑
- 调整
select方法匹配单事件模式:
若原本需要多事件连续匹配,需修正CEP Pattern定义,比如:@Override public String select(Map<String, List<String>> pattern) throws Exception { String userId = pattern.get("first").get(0).split(",")[0].trim(); int risk = getValue(pattern.get("first").get(0)); return "Stroke Risk Alert: Patient ID - " + userId + ", Risk Level - " + risk; }Pattern<String, ?> highRiskPattern = Pattern.<String>begin("first") .where(new SimpleCondition<String>() { @Override public boolean filter(String value) { return getValue(value) > 1; } }) .next("middle") .where(new SimpleCondition<String>() { @Override public boolean filter(String value) { return getValue(value) > 1; } }) .next("last") .where(new SimpleCondition<String>() { @Override public boolean filter(String value) { return getValue(value) > 1; } }); - 移除重复的反序列化配置,保留其中一种:
.setValueOnlyDeserializer(new SimpleStringSchema()) - 开启起始偏移量配置,确保能读取历史消息:
.setStartingOffsets(OffsetsInitializer.earliest())
2. 查看Flink集群日志
客户端仅输出任务提交信息,运行时错误和输出在集群日志中:
- 查看JobManager日志:
${FLINK_HOME}/log/flink-*-jobmanager-*.log,检查是否有NullPointerException等异常。 - 查看TaskManager日志:
${FLINK_HOME}/log/flink-*-taskmanager-*.log,确认算子运行状态和数据处理记录。
3. 验证Kafka数据与消费情况
- 用命令行查看topic消息,确认格式符合预期(应是
PatientID,Timestamp,Type,Value,比如XYZ,2024-05-20T10:00:00,HR,140):./bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic geCEP --from-beginning - 检查消费者组消费状态,确认偏移量是否增长:
./bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group test-consumer-group
4. 调整输出方式
Flink的print()默认输出到TaskManager的stdout,若想在客户端获取输出:
- 设置并行度为1,确保输出集中:
strokeRiskAlerts.print().setParallelism(1)。 - 或者将输出写入文件系统:
strokeRiskAlerts.writeAsText("file:///Users/spartacus/flink-output/alerts.txt").setParallelism(1);
内容的提问来源于stack exchange,提问作者Ishaan Adarsh
相关产品推荐
相关产品推荐

