You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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;
    }
}
问题原因与排查方案

一、核心问题分析

  1. CEP模式与select方法不匹配
    代码定义的CEP模式是仅包含first节点的单事件匹配,但在select方法中强行获取不存在的middle和last节点数据,会直接抛出NullPointerException。Flink任务出现未捕获异常时,对应算子会静默停止,客户端不会输出错误信息。

  2. Kafka Source配置冲突
    同时设置了两种反序列化逻辑:setDeserializer和setValueOnlyDeserializer,会导致反序列化逻辑冲突,可能无法正确解析Kafka消息。

  3. 未配置起始偏移量
    注释掉了setStartingOffsets(OffsetsInitializer.earliest()),默认会从Kafka topic的最新偏移量开始消费,若topic中没有新消息,任务不会读取到任何数据。

二、排查与修复步骤

1. 修复代码逻辑

  • 调整select方法匹配单事件模式:
    @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;
    }
    
    若原本需要多事件连续匹配,需修正CEP Pattern定义,比如:
    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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.11 02:59:58