智能运维项目的数据场景下的实现代码讲解1【2026.8.17】 总体介绍“RocketMQ接收 JSON → 消息解析与校验 → Redis/DM 状态判断 → 规则计算 → 结果上报”的调用链定位具体类并标出每个文件的职责。整个流程的主调用链如下RocketMQ JSON → TaskExecutionConsumer → TaskMessageFactory → TaskExecutionOrchestrator → DefaultTaskBusinessHandler → PrometheusComputeEngine → ComputeStrategy / CompareStrategy → RocketMqTaskResultPublisher核心文件1. TaskExecutionConsumer.javaRocketMQ消费者入口接收原始JSON并调用后续流程。2. TaskMessageFactory.java将JSON反序列化为任务对象 这里 校验模型、规则图、compute、compare和阈值。3. TaskExecutionMessage.javaRocketMQ外层消息结构包括执行ID、重试次数、过期时间和完整任务。4. TaskDefinition.javaJSON中的任务、模型、规则图、节点、指标和数据源结构。5. TaskExecutionOrchestrator.java整个流程的总调度器状态检查、加锁、执行计算、上报结果、更新成功状态。 失败重试逻辑在这里 。6. DefaultTaskBusinessHandler.java执行模型 递归处理AND/OR/NOT规则图 并判断节点是 TRIGGERED 还是 NORMAL 。7. PrometheusComputeEngine.java按监控对象查询指标、执行compute计算和compare判断。8. PrometheusQueryService.java构造PromQL并调用Prometheus query_range 接口。9. ComputeStrategy.java实现 avg/max/min/span/vary/rate 。10. CompareStrategy.java实现 gt/lt/in 阈值比较。11. RocketMqTaskResultPublisher.java将结果上报到 znyw-task-result:triggered 或 :normal 。执行状态相关- TaskExecutionRepository.java DM中的 PENDING/RUNNING/RETRYING/SUCCESS/FAILED 状态。- TaskExecutionLockService.java Redis锁和重试次数。- application-dev.yml RocketMQ Topic、重试、Prometheus和计算配置。各个文件介绍TaskExecutionConsumer.javaRocketMQ消费者入口接收原始JSON并调用后续流程。TaskExecutionConsumer.java 是整个任务执行链路的消息入口不包含具体计算逻辑。类注解第15-26行 - Service 交给Spring管理。- ConditionalOnProperty 只有配置了>先重新自己手敲了这个文件一遍。file TaskExecutionConsumer.java import java.util.Collections; //交给Spring管理使该消费者可以注入消息解析、执行编排和日志服务。 Service //只有显示启用rocketmq时才创建消费者避免本地未启动rocketmq时初始化失败 ConditionalOnProperty( prefixdata-scenaria.execution, namerocketmq-enabled, havingValuetrue ) //声明消费topic、消费者组和tag过滤规则具体值可以通过配置文件或环境变量覆盖 //前面加的是注解 RocketMQMessageListener( //运维任务执行消息所在的topic未配置时使用znyw-task-execute. topic${data-scenario.execution.topic:znyw-task-execute}, //统一消费者组中的实例以集群方式分担消息 consumerGroup${data-scenario.execution.consumer-group:znyw-task-consumer}, //默认接受全部tag也可以通过配置限制执行消费指定类型的任务 selectorExpression${data-scenario.execution.tag-expression:*}, //onMessage抛出异常后由rocketmq最多重新消费16次 maxReconsumeTimes16 ) public class TaskExecutionConsumer implements RocketMQListenerString{ //将rocketmq原始的json转换为TaskExecutionMessage并校验任务、模型和规则定义。 private final TaskMessageFactory messageFactory; //整个任务执行流程的编排器负责状态检查、加锁、计算、结果上报和失败重试 private final TaskExecutionOrchestrator orchestrator; //执行日志服务、入口消息无法解析或校验失败时用于记录INALID_MESSAGE日志 private final TaskExecutionLogService logService; //通过构造器注入消费者所需的依赖 //param messageFactory 原始json消息解析及任务定义校验组件 //param orchestrator 任务执行流程编排组件 //param logService 任务执行日志记录组件 public TaskExecutionConsumer( TaskMessageFactory messageFactory, TaskExcutionOrchestrator orchestrator, TaskExecutionLogService logService ){ this.messageFactorymessageFactory; this.orchestratororchestrator; this.logServicelogService; } //接收并处理rocketmq消息 //rocketmq Spring将消息体转换为字符串后调回本方法。方法正常返回表示本次消费成功。 //方法抛出异常表示消费失败rocketmq将根据重试配置重新投递消息。 //param rawMessage 。。RocketMQ中收到的原始json字符串。 Override public void onMessage(String rawMessage) { try{ //第一步解析json生成统一的任务执行消息并校验任务定义是否合法。 TaskExecutionMessage messagemessageFactory.fromRawMessage(rawMessage); //第二步进入任务主流程后续状态判断、计算和结果上报均由编排器负责 orchestrator.consume(message); }catch(JsonProcessingExecution exception){ //json本身存在语法问题例如缺少引号、括号不匹配或字段类型无法转换。 logService.write( //日志事件类型 event:INVALID_MESSAGE, //日志级别 level:ERROR, //json尚未成功解析无法提供有效的TaskExecutionMessage。 message:null, description:RocketMQ消息JSON格式不合法, //保留原始消息便于定位发送端的json格式问题。 Collections.singletonMap(key:rawMessage,rawMessage) ); //不能吞掉异常否则rocketmq会把非法消息误判为消费成功。 throw new IllegalArgumentException(message:RocketMQ消息JSON格式不合法,exception); }catch(IllegalArgumentException exception){ // JSON 语法正确但业务字段不合法例如缺少 taskId、规则图错误或阈值配置无效。 logService.write( INVALID_MESSAGE, ERROR, null, exception.getMessage(), Collections.singletonMap(rawMessage, rawMessage) ); // 将原异常继续交给 RocketMQ 消费容器触发消费失败及其重试机制。 throw exception; } } }