欢迎光临
我们一直在努力

工业物联网实战:基于显式编排器模式的9阶段数据Pipeline全链路解析

前言

在工业物联网系统中,设备上报数据的处理链路是整个平台的"心脏"。每秒可能有数千条传感器数据涌入系统,需要经过解析、校验、持久化、缓存、告警、归档等多个环节。如何保证这条链路在高并发下稳定运行?如何在出问题时快速定位故障节点?如何在不重启服务的情况下动态调整处理逻辑?

本文将结合一个真实的工业物联网项目,详细拆解一套基于显式编排器模式构建的 9 阶段数据 Pipeline,覆盖架构设计、核心源码、三级容错机制、全链路追踪以及运行时动态管控等关键能力。


一、为什么需要 Pipeline?

1.1 一个直觉类比:快递分拣中心

想象一个快递分拣中心,一个包裹从进来到发出,要经过这样的流程:

收件 → 扫码录入 → 安检 → 称重 → 分拣到区域 → 装车发走

每一步都有专人负责,安检不过的包裹直接退回,不会浪费后续资源。

工业物联网的数据处理一模一样。传感器通过 MQTT 发来的每一条数据,也要经过"解析 → 校验 → 标准化 → 存储 → 缓存 → 告警 → 归档"等多道工序——这就是 Pipeline(流水线)。

1.2 不拆 Pipeline 会怎样?

最朴素的做法是写一个巨大的 processMessage() 方法,把所有逻辑堆在一起。这在 Demo 阶段没问题,但在生产环境下会面临三个致命问题:

第一,改不动。 如果某天需要修改告警逻辑,你得在 500 行的大方法里精准定位到告警那段代码,一不小心还可能误改其他环节。Pipeline 的做法是:每个节点是一个独立的类,改告警就只动告警的类,其他文件碰都不碰。

第二,没法灵活开关。 假设线上告警模块出了 bug,需要临时关闭。如果全写在一起,就得改代码、重新打包、重新部署。Pipeline 的做法是:数据库里维护每个节点的开关状态,调一个 REST 接口就能关闭某个节点,不用重启服务。

第三,出了问题查不到。 2000+ TPS 意味着每秒 2000 多条数据在跑,某条数据处理失败了,你怎么知道它是在哪一步挂的?Pipeline 的做法是:每条数据分配一个全局唯一的 traceId,每个节点执行完都记录"成功/失败/耗时",一条 SQL 就能定位问题。


二、整体架构设计

2.1 9 阶段流水线

传感器设备(温度/湿度/振动/…)

│ MQTT 消息(JSON 格式)

┌──────────────────────────────────────────────────────────────────┐
│ 9 阶段 Data Pipeline │
│ │
│ n0 消息入口 ──→ n1 消息解析 ──→ n2 数据校验 ──→ n3 身份标准化 │
│ │ │
│ n4 写入InfluxDB ←───────────────────────────────────┘ │
│ │ │
│ n5 写入Redis ──→ n6 告警处理 ──→ n7 串口联动/副作用调度 │
│ │ │
│ n8 历史归档(MongoDB) ←────────────────────────────────┘ │
└──────────────────────────────────────────────────────────────────┘


InfluxDB(时序存储) + Redis(实时缓存) + MongoDB(冷数据归档)

各节点的职责如下表:

节点编号节点编码职责实现类
n0 MQTT_INGRESS 监听 MQTT 消息,转交给 Pipeline DeviceDataEventListener
n1 PARSE_PAYLOAD 将 JSON 字符串解析为 Java 对象 DeviceReportParser
n2 VALIDATE_REPORT 校验设备编号、数据完整性 DeviceReportValidator
n3 NORMALIZE_IDENTITY 统一设备身份标识 DeviceIdentityNormalizer
n4 WRITE_INFLUX 写入 InfluxDB 时序数据库 DeviceInfluxWriter
n5 WRITE_REALTIME_REDIS 写入 Redis 实时缓存 DeviceRedisWriter
n6 PROCESS_ALARM 告警判定与告警日志持久化 DeviceAlarmProcessor
n7 DISPATCH_SERIALPORT 串口指令下发、消息路由 DeviceAlarmSideEffectDispatcher
n8 ARCHIVE_MONGO 定时归档至 MongoDB DeviceMongoArchiveWriter

2.2 为什么选"显式编排器"而不是经典责任链?

经典责任链模式(Chain of Responsibility)的做法是:每个 Handler 持有一个 next 引用,处理完自动传给下一个。看起来优雅,但在实际工程中有一个痛点——传递过程的控制逻辑不好插入。

你想在每个节点执行前加一个"开关检查"、执行后加一个"耗时统计"、执行中加一个"异常兜底"?在经典责任链里,这些逻辑要么散落在每个 Handler 里,要么需要一个额外的拦截器层,越搞越复杂。

显式编排器模式的做法是:有一个"总指挥"(Orchestrator),由它按顺序显式调用每个节点。所有控制逻辑(开关判断、计时、日志、异常处理)都集中在 Orchestrator 里,每个节点只需要关心自己的业务逻辑。

经典责任链: Handler1 → Handler2 → Handler3 (传递逻辑藏在 Handler 内部)

显式编排器: Orchestrator ──调用──→ Node1
──调用──→ Node2
──调用──→ Node3 (控制逻辑集中在 Orchestrator)

后者的好处是:新增节点时只需写一个新类、在 Orchestrator 里加一行调用;想在中间插入一个"开关检查",只需要改 Orchestrator 的调用逻辑,节点代码一行不动。


三、核心源码拆解

3.1 上下文对象:Pipeline 的"面单"

在流水线上传递的每条数据,都需要一个载体来记录"我是谁、我从哪来、前面节点处理的结果是什么"。这个载体就是 Context(上下文)。

@Data
public class DeviceReportContext {
// 全局唯一追踪ID,32位UUID,相当于"快递单号"
private String traceId = UUID.randomUUID().toString().replace("-", "");

// MQTT 原始信息
private String topic; // 消息主题(标识设备类型)
private String payload; // 消息体(JSON 字符串)
private int qos; // 服务质量等级

// 解析后的业务对象(n1 的产出,后续节点都读它)
private DeviceMqttVo<? extends TemperatureVo> deviceVo;

// 时间戳
private Date receivedAt; // 消息接收时间
private Date parsedAt; // 解析完成时间

// 各阶段的执行状态
private Map<String, Boolean> flags; // 布尔标记(如 parseSuccess=true)
private Map<String, Object> result; // 结果描述(如 parsePayload="success")
}

为什么用 Context 对象而不是参数传递? 因为 9 个节点如果互相传参数,方法签名会爆炸,而且后续想加新字段(比如新增一个"处理优先级")就要改所有方法的签名。用一个 Context 对象装所有东西,每个节点都能读到前面节点的产出,也能往里写自己的结果,干净利落。

3.2 入口:事件监听器

@Slf4j
@Component
public class DeviceDataEventListener {

private final DeviceReportWorkflowOrchestrator orchestrator;

// 监听 MQTT 消息事件(Spring Event 机制)
@EventListener
public void handleDeviceData(MqttMessageEvent event) {
if (Objects.isNull(event) || StringUtils.isBlank(event.getPayload())) {
return; // 空消息直接丢弃
}
// 把消息交给编排器,开启 Pipeline 之旅
orchestrator.executeFromMqtt(
event.getTopic(), event.getPayload(), event.getQos());
}
}

这里用了 Spring Event 事件驱动:MQTT 模块收到消息后发布一个 MqttMessageEvent,监听器收到事件后转交给 Pipeline。MQTT 模块不需要知道 Pipeline 的存在,Pipeline 也不需要知道消息是怎么来的——这就是事件驱动解耦。

3.3 编排器:Pipeline 的"总指挥"

这是整个 Pipeline 最核心的类。先看它的构造函数,能看出它依赖哪些组件:

@Component
public class DeviceReportWorkflowOrchestrator {

private final DeviceReportParser parser; // n1 解析器
private final DeviceReportValidator validator; // n2 校验器
private final DeviceIdentityNormalizer identityNormalizer; // n3 标准化
private final DeviceInfluxWriter influxWriter; // n4 InfluxDB写入
private final DeviceRealtimeStageProcessor realtimeProcessor; // n5+n6+n7(实时阶段)
private final IWorkflowRuntimeService runtimeService; // 运行时管理(开关)
private final IWorkflowExecutionLogService executionLogService; // 执行日志

// … 构造函数省略
}

注意分层设计: 编排器直接管理 n1~n4,然后把 n5~n7 交给 DeviceRealtimeStageProcessor(实时阶段处理器)。为什么要分两层?因为 n5~n7 需要设备级限流(同一台设备不能太频繁处理),所以单独包了一层来做限流控制。

主流程 execute() 方法

public void execute(DeviceReportContext context) {
String workflowCode = DEFAULT_WORKFLOW_CODE;
Integer versionNo = resolveVersionNo(workflowCode);

// 记录流程开始(写入 wf_instance_log 表)
log.info("[数据上报] 流程开始 | traceId={} | topic={}",
context.getTraceId(), context.getTopic());
executionLogService.startInstance(context.getTraceId(), workflowCode, versionNo, context.getTopic());

// ========== 第一道防线:全局开关 ==========
if (!runtimeService.isWorkflowRunning(workflowCode)) {
log.warn("[数据上报] 流程中止 | traceId={} | 原因=workflow已停止", context.getTraceId());
executionLogService.finishSkipped(context.getTraceId(), "workflow_stopped");
return; // 整个流程被关闭,直接返回
}

// ========== n1:解析(失败则快速终止) ==========
if (!executeParseNode(workflowCode, context)) {
executionLogService.finishFailed(context.getTraceId(), "stopped_at_parse");
return; // 解析失败,垃圾数据不值得继续处理
}

// ========== n2:校验(失败则快速终止) ==========
if (!executeValidateNode(workflowCode, context)) {
executionLogService.finishFailed(context.getTraceId(), "stopped_at_validate");
return; // 校验失败,数据不合法
}

// ========== n3~n7:后续处理(try-catch 兜底) ==========
DeviceMqttVo<? extends TemperatureVo> deviceVo = context.getDeviceVo();
try {
executeNormalizeNode(workflowCode, context); // n3 标准化
executeInfluxNode(workflowCode, context, deviceVo); // n4 写入InfluxDB
realtimeProcessor.process(deviceVo, workflowCode, context); // n5+n6+n7
executionLogService.finishSuccess(context.getTraceId());
} catch (Exception e) {
log.error("[数据上报] 流程异常 | traceId={} | error={}",
context.getTraceId(), e.getMessage(), e);
executionLogService.finishFailed(context.getTraceId(), e.getMessage());
}
}

这段代码有三个关键设计决策:

决策一:快速失败(Fail-Fast)。 parse 和 validate 失败后直接 return,不进入后续节点。原因很直觉:数据格式都不对,往 InfluxDB 里写也是垃圾数据,不如省掉这些开销。

决策二:先持久化,再做其他。 n4(写入 InfluxDB)在 n5(写 Redis)和 n6(告警)之前执行。这样即使后面 Redis 挂了、告警模块崩了,至少原始数据已经安全落地了。

决策三:try-catch 兜底。 n3 之后的节点被 try-catch 包裹。为什么 parse 和 validate 不用 try-catch?因为它们本身就是"门槛",返回 false 就代表失败,逻辑是确定的。而 n3 之后涉及外部系统调用(InfluxDB、Redis),可能抛出意料之外的异常,需要兜底。

节点执行模板

每个节点的执行都遵循同一个模板,以 parse 节点为例:

private boolean executeParseNode(String workflowCode, DeviceReportContext context) {
// ① 开关检查:这个节点启用了吗?
if (!runtimeService.isNodeEnabled(workflowCode, WorkflowNodeCode.PARSE_PAYLOAD, true)) {
executionLogService.nodeSkipped(context.getTraceId(), workflowCode,
WorkflowNodeCode.PARSE_PAYLOAD, "node_disabled");
return true; // 跳过不算失败,流程继续
}

// ② 记录开始时间(计算耗时)
long start = System.currentTimeMillis();

// ③ 执行真正的业务逻辑
boolean ok = parser.parse(context);

// ④ 记录结果
long costMs = System.currentTimeMillis() start;
if (ok) {
executionLogService.nodeSuccess(context.getTraceId(), workflowCode,
WorkflowNodeCode.PARSE_PAYLOAD, costMs, "success");
return true;
}
executionLogService.nodeFailed(context.getTraceId(), workflowCode,
WorkflowNodeCode.PARSE_PAYLOAD, costMs, "parse_failed");
return false;
}

开关检查 → 计时 → 执行 → 记录结果,每个节点都是这个套路。这种一致性不是巧合,而是刻意设计的——它让代码行为可预测,也让排查问题时思路清晰。

3.4 解析节点源码示例

@Component
public class DeviceReportParser {

// 支持的设备主题前缀
private static final String[] SUPPORTED_TOPICS = {
"client/device/opticalFiber", // 光缆传感器
"client/device/belt", // 皮带传感器
"client/device/ambient" // 环境传感器
};

public boolean parse(DeviceReportContext context) {
String topic = context.getTopic();
String payload = context.getPayload();

// 空数据快速失败
if ("{\\"data\\":[]}".equals(payload)) {
return false;
}

try {
// 只处理已知类型的设备
if (isSupportedTopic(topic)) {
DeviceMqttVo<? extends TemperatureVo> deviceVo =
JSON.parseObject(payload, new TypeReference<DeviceMqttVo<TemperatureVo>>() {});
context.setDeviceVo(deviceVo); // 解析结果写入上下文
context.setParsedAt(new Date());
return true;
}
return false; // 不支持的设备类型
} catch (Exception e) {
log.error("JSON解析异常 | traceId={} | error={}", context.getTraceId(), e.getMessage());
return false;
}
}
}

3.5 校验节点源码示例

@Component
public class DeviceReportValidator {

public boolean validate(DeviceReportContext context) {
DeviceMqttVo<?> deviceVo = context.getDeviceVo();

// 检查1:解析结果不能为空
if (Objects.isNull(deviceVo)) {
return false;
}
// 检查2:设备编号不能为空
if (StringUtils.isBlank(deviceVo.getDeviceCode())) {
return false;
}
return true;
}
}


四、三大核心机制

4.1 机制一:运行时动态启停

这是什么?

可以在不重启服务的情况下,开启/关闭整个 Pipeline,或者其中任何一个节点。

实现原理

数据库中有三张表支撑这套机制:

表名职责关键字段
wf_definition 流程定义 workflow_code, current_version, status
wf_version 流程版本(DAG 定义) graph_json, node_config_json
wf_runtime_state 运行时状态 running(0/1), switches_json

核心在 wf_runtime_state 表的 switches_json 字段,它存储了每个节点的开关状态:

{
"PARSE_PAYLOAD": true,
"VALIDATE_REPORT": true,
"NORMALIZE_IDENTITY": true,
"WRITE_INFLUX": true,
"WRITE_REALTIME_REDIS": true,
"PROCESS_ALARM": false,
"DISPATCH_SERIALPORT": true,
"ARCHIVE_MONGO": true
}

上面的配置表示:告警处理节点(PROCESS_ALARM)已关闭,其余节点正常运行。

开关读取逻辑

public boolean isNodeEnabled(String workflowCode, String nodeCode, boolean defaultEnabled) {
WorkflowRuntimeState state = getRuntimeState(workflowCode);
Map<String, Boolean> switches = parseSwitches(state.getSwitchesJson());
Boolean enabled = switches.get(nodeCode);
return enabled == null ? defaultEnabled : enabled;
}

如果数据库里没有配置这个节点的开关,就用默认值(通常是 true)。

REST 管控接口

POST /workflow/runtime/start/{code} → 启动整个流程
POST /workflow/runtime/stop/{code} → 停止整个流程
POST /workflow/runtime/node/{code}/{node}/enable → 启用某个节点
POST /workflow/runtime/node/{code}/{node}/disable → 禁用某个节点
GET /workflow/runtime/state → 查看当前运行态快照

实际应用场景

某天凌晨 2 点,告警模块因为某个边界 case 频繁误报。运维接到告警电话后,不需要叫醒开发、不需要改代码、不需要重新部署,只需要调一个接口:

POST /workflow/runtime/node/device_report/PROCESS_ALARM/disable

告警节点立刻停止执行,其他数据处理完全不受影响。等白天开发修复 bug 后,再调一个接口重新启用。


4.2 机制二:traceId 全链路追踪

这是什么?

每条进入 Pipeline 的数据都会被分配一个唯一的 traceId(32 位 UUID),后续所有日志、所有数据库记录都带着这个 ID。

生成方式

// 在 Context 创建时自动生成
private String traceId = UUID.randomUUID().toString().replace("-", "");
// 示例:a1b2c3d4e5f6789012345678abcdef01

日志记录

编排器的每一行日志都带 traceId:

log.info("[数据上报] 流程开始 | traceId={} | topic={} | qos={}",
context.getTraceId(), context.getTopic(), context.getQos());

每个节点执行完都会写一条节点日志到 wf_node_log 表:

executionLogService.nodeSuccess(traceId, workflowCode, nodeCode, costMs, "success");

数据库存储

wf_instance_log(流程级日志):

trace_idworkflow_codestatusstarted_atfinished_at
a1b2c3d4… device_report SUCCESS 10:00:00.000 10:00:00.035

wf_node_log(节点级日志):

trace_idnode_codestatuscost_ms
a1b2c3d4… PARSE_PAYLOAD SUCCESS 2
a1b2c3d4… VALIDATE_REPORT SUCCESS 1
a1b2c3d4… NORMALIZE_IDENTITY SUCCESS 0
a1b2c3d4… WRITE_INFLUX SUCCESS 15
a1b2c3d4… WRITE_REALTIME_REDIS SUCCESS 8
a1b2c3d4… PROCESS_ALARM SUCCESS 5
排障演示

假设某条数据处理异常,拿着 traceId 一条 SQL 就能还原完整执行路径:

SELECT node_code, status, cost_ms, message
FROM wf_node_log
WHERE trace_id = 'a1b2c3d4e5f6789012345678abcdef01'
ORDER BY id;

立刻就能看到:这条数据在哪个节点成功了、哪个节点失败了、每个节点花了多少毫秒。在 2000+ TPS 的场景下,这种精确到单条数据的追踪能力是排障的关键。


4.3 机制三:三级容错保障

整体结构

第一级:全局开关(最粗暴,一刀切停掉所有数据处理)
↓ 没拦住(流程是开着的)
第二级:节点开关(精准关闭某个有问题的节点,其他节点继续跑)
↓ 没拦住(节点也是开着的)
第三级:try-catch 兜底(代码层面的异常兜底,防止线程崩溃)

记录日志,标记失败,线程继续处理下一条数据

第一级:全局开关

if (!runtimeService.isWorkflowRunning(workflowCode)) {
executionLogService.finishSkipped(context.getTraceId(), "workflow_stopped");
return;
}

使用场景: 系统整体维护、数据库迁移、紧急停机等需要暂停所有数据处理的场景。

第二级:节点开关

if (!runtimeService.isNodeEnabled(workflowCode, WorkflowNodeCode.PROCESS_ALARM, true)) {
executionLogService.nodeSkipped(context.getTraceId(), workflowCode,
WorkflowNodeCode.PROCESS_ALARM, "node_disabled");
return true; // 跳过不算失败
}

使用场景: 某个节点有 bug 需要临时关闭,或者某个外部依赖(如 InfluxDB)暂时不可用时,跳过对应写入节点。

第三级:try-catch 兜底

try {
executeNormalizeNode(workflowCode, context);
executeInfluxNode(workflowCode, context, deviceVo);
realtimeProcessor.process(deviceVo, workflowCode, context);
executionLogService.finishSuccess(context.getTraceId());
} catch (Exception e) {
log.error("[数据上报] 流程异常 | traceId={} | error={}",
context.getTraceId(), e.getMessage(), e);
executionLogService.finishFailed(context.getTraceId(), e.getMessage());
}

使用场景: InfluxDB 连接超时、Redis 突然不可用、JSON 序列化异常等未预期的运行时异常。即使出了这些意外,处理线程不会崩溃,错误被记录下来,线程继续处理下一条数据。

一个真实的故障场景

假设某天 InfluxDB 因为负载过高开始拒绝写入:

  • n4(WRITE_INFLUX)抛出异常
  • 第三级容错 try-catch 兜住,记录错误日志
  • 这条数据被标记为 FAILED,但不影响其他数据继续处理
  • 运维通过日志发现大量 WRITE_INFLUX 失败,定位到 InfluxDB 问题
  • 如果短期内无法修复,可以临时调接口关闭 n4 节点(第二级容错),让数据跳过 InfluxDB 直接走后续流程
  • InfluxDB 修复后,重新开启 n4 节点
  • 三级容错层层递进,从粗粒度到细粒度,从运维操作到代码兜底,共同保障系统的高可用。


    五、实时阶段处理器:限流与子编排

    前面提到,主编排器把 n5~n7 交给了 DeviceRealtimeStageProcessor。这个类有两个核心职责:设备级限流和子编排。

    5.1 为什么需要限流?

    工业场景中,某些传感器可能每秒钟上报几十条数据。如果每条数据都触发告警判断、串口指令下发等操作,不仅浪费资源,还可能导致串口设备响应不过来。

    public void process(DeviceMqttVo<?> deviceVo, String workflowCode, DeviceReportContext context) {
    String deviceCode = deviceVo.getDeviceCode();
    // 设备级限流:同一台设备在配置的最小间隔内只处理一次
    rateLimitService.executeByDevice(deviceCode, () -> {
    // n5: 写入 Redis 实时缓存
    // n6: 告警处理
    // n7: 串口联动
    });
    }

    5.2 子编排逻辑

    在限流的 Lambda 内部,n5~n7 的执行逻辑和主编排器类似:每个节点先检查开关、再计时执行、再记录日志。区别在于 n6(告警处理)依赖 n5(Redis 写入)的结果——如果 Redis 写入失败(比如设备未在系统中配置),告警处理也会被跳过。

    // n5 执行
    Object[] results = deviceRedisWriter.writeRealtimeState(deviceVo);
    boolean realtimeSuccess = results != null && results.length >= 2;

    // n6 依赖 n5 的结果
    if (realtimeSuccess) {
    // Redis 写入成功,继续处理告警
    List<InsAlarmLog> alarmLogs = deviceAlarmProcessor.process(...);
    } else {
    // Redis 写入失败(设备未配置),跳过告警
    executionLogService.nodeSkipped(..., "device_not_configured");
    }


    六、运行时管理的初始化机制

    系统启动时,WorkflowRuntimeServiceImpl 会在 @PostConstruct 阶段自动初始化默认的流程模板:

    @PostConstruct
    public void init() {
    initDefaultTemplateIfAbsent();
    }

    初始化过程做三件事:

  • 创建流程定义(wf_definition 表):记录流程名称、业务类型、当前版本号
  • 创建流程版本(wf_version 表):存储 9 个节点的 DAG 图结构(graphJson)和默认开关配置
  • 创建运行态(wf_runtime_state 表):初始状态为"运行中",所有节点开关默认为 true
  • DAG 图结构的定义如下:

    private String buildDefaultGraphJson() {
    List<Map<String, Object>> nodes = new ArrayList<>();
    nodes.add(node("n0", "MQTT_INGRESS", "消息入口"));
    nodes.add(node("n1", "PARSE_PAYLOAD", "消息解析"));
    nodes.add(node("n2", "VALIDATE_REPORT", "数据校验"));
    nodes.add(node("n3", "NORMALIZE_IDENTITY", "身份标准化"));
    nodes.add(node("n4", "WRITE_INFLUX", "写入Influx"));
    nodes.add(node("n5", "WRITE_REALTIME_REDIS", "写入实时缓存"));
    nodes.add(node("n6", "PROCESS_ALARM", "告警处理"));
    nodes.add(node("n7", "DISPATCH_SERIALPORT", "串口联动"));
    nodes.add(node("n8", "ARCHIVE_MONGO", "历史归档"));

    // 边:n0→n1→n2→…→n8(线性流水线)
    List<Map<String, String>> edges = new ArrayList<>();
    edges.add(edge("n0", "n1"));
    edges.add(edge("n1", "n2"));
    // … 依次类推
    }

    这种设计的好处是:流程的定义(DAG 结构)和运行态(开关状态)分离存储。未来如果需要发布新版本的流程(比如新增一个节点),可以创建一个新的 version 记录,而不会影响当前正在运行的版本。


    七、架构全景图

    把上面所有内容串联起来,一张图看全貌:

    ┌──────────────────────────┐
    │ MQTT Broker (EMQX) │
    └────────────┬─────────────┘
    │ 消息到达

    ┌──────────────────────────┐
    │ MqttMessageClient │
    │ 发布 MqttMessageEvent │
    └────────────┬─────────────┘
    │ Spring Event

    ┌─────────────────────────────────────┐
    │ DeviceDataEventListener (n0) │
    │ @EventListener 监听事件 │
    └─────────────────┬───────────────────┘

    ┌─────────────────▼───────────────────┐
    │ DeviceReportWorkflowOrchestrator │
    │ ┌───────────────────────────────┐ │
    │ │ 全局开关检查(第一级容错) │ │
    │ └───────────┬───────────────────┘ │
    │ ▼ │
    │ ┌─── n1 PARSE_PAYLOAD ────┐ │
    │ │ 开关检查 → JSON解析 → 日志 │ │
    │ └─── 失败则 return ─────────┘ │
    │ ▼ │
    │ ┌─── n2 VALIDATE_REPORT ──┐ │
    │ │ 开关检查 → 数据校验 → 日志 │ │
    │ └─── 失败则 return ─────────┘ │
    │ ▼ │
    │ ┌─── n3 NORMALIZE_IDENTITY ┐ │
    │ │ 开关检查 → 身份标准化 → 日志│ │
    │ └────────────┬──────────────┘ │
    │ ▼ │
    │ ┌─── n4 WRITE_INFLUX ──────┐ │
    │ │ 开关检查 → 时序写入 → 日志 │ │
    │ └────────────┬──────────────┘ │
    │ ▼ │
    │ ┌── try-catch 兜底(第三级容错)──┐ │
    │ │ │ │
    │ │ ┌─ DeviceRealtimeStage ─────┐│ │
    │ │ │ Processor ││ │
    │ │ │ ┌─ 设备级限流 ──────┐ ││ │
    │ │ │ │ │ ││ │
    │ │ │ │ n5 WRITE_REDIS │ ││ │
    │ │ │ │ n6 PROCESS_ALARM │ ││ │
    │ │ │ │ n7 DISPATCH_SERIAL │ ││ │
    │ │ │ └───────────────────┘ ││ │
    │ │ └───────────────────────────┘│ │
    │ └───────────────────────────────┘ │
    │ │ │
    │ ┌───────────▼──────────────────┐ │
    │ │ executionLogService │ │
    │ │ .finishSuccess/Failed/Skipped│ │
    │ └──────────────────────────────┘ │
    └──────────────────────────────────────┘

    ┌─────────────────▼───────────────────┐
    │ 数据存储层 │
    │ InfluxDB Redis MongoDB │
    │ (时序数据) (实时缓存) (冷数据归档) │
    └─────────────────────────────────────┘

    ┌─────────────────────────────────────┐
    │ 定时任务(每10分钟) │
    │ DeviceReportArchiveOrchestrator │
    │ n8 ARCHIVE_MONGO → 光缆数据归档 │
    └─────────────────────────────────────┘

    ┌─────────────────────────────────────┐
    │ 运行时管控 REST API │
    │ /workflow/runtime/start|stop │
    │ /workflow/runtime/node/…/enable │
    │ → wf_definition │
    │ → wf_version │
    │ → wf_runtime_state │
    └─────────────────────────────────────┘


    八、面试高频问答

    Q1:简单介绍一下你的 Pipeline 架构?

    设备上报数据的核心处理逻辑是用 Pipeline 编排引擎实现的。整个流程分 9 个节点:从 MQTT 消息入口开始,经过解析、校验、身份标准化、写入 InfluxDB、写入 Redis 实时缓存、告警处理、串口联动,最后到 MongoDB 归档。架构上用的是显式编排器模式,由一个 Orchestrator 类按顺序调用各个节点,每个节点是独立的 Spring Bean,职责单一。三个核心亮点:运行时动态启停、traceId 全链路追踪、三级容错保障 2000+ TPS 稳定运行。

    Q2:为什么不用经典责任链模式?

    经典责任链的节点传递逻辑藏在 Handler 内部,想在传递过程中加开关检查、计时、日志记录会很别扭。显式编排器把所有控制逻辑集中在 Orchestrator 里,每个节点只需要关注业务逻辑。新增节点时只需写一个新类、在 Orchestrator 里加一行调用,不影响其他节点代码。

    Q3:如果 InfluxDB 挂了怎么办?

    n4(写入 InfluxDB)在 try-catch 范围内,如果 InfluxDB 连接失败会抛异常,被第三级容错兜住。同时运维可以通过节点开关临时关闭 InfluxDB 写入节点,让数据跳过持久化直接走后续流程,等 InfluxDB 恢复后再开启。由于数据已经通过 n1~n2 的校验,跳过 n4 不会影响后续 n5~n7 的正确执行。

    Q4:traceId 用 UUID 生成,高并发下性能有没有问题?

    UUID.randomUUID() 底层用的是 SecureRandom,在极高并发下确实可能成为瓶颈。目前 2000+ TPS 下没有观察到性能问题。如果未来 TPS 进一步增长,可以考虑用 Snowflake 算法生成有序 ID,既保证唯一性又避免随机 UUID 的性能问题和 InfluxDB 写入时的索引碎片。

    Q5:节点开关每次都要查数据库吗?怎么优化?

    目前确实是每次执行都查数据库。优化方案是:启动时加载到本地缓存(ConcurrentHashMap),数据库更新时通过 Spring Event 或 Redis Pub/Sub 通知所有实例刷新缓存。这样既保证了动态生效,又避免了每次查库的开销。在多实例部署时,还可以用 Redis 的 keyspace notification 来触发缓存刷新。

    Q6:如何保证数据不丢失?

    三个层面保障:第一,MQTT 使用 QoS 1 级别,保证消息至少到达一次;第二,Pipeline 先写入 InfluxDB 做持久化,再做 Redis 缓存和告警等后续操作,即使后续步骤失败,原始数据不会丢;第三,每条数据的执行状态都有日志记录,如果某条数据 FAILED,可以通过 traceId 定位并重试。


    九、总结

    我详细拆解了一个工业物联网项目中 9 阶段数据 Pipeline 的完整设计与实现。以下式核心要点回顾:

    设计模式选型: 采用显式编排器模式而非经典责任链,将控制逻辑集中在 Orchestrator,业务逻辑分散到各节点,实现了关注点分离。

    三大核心机制: 运行时动态启停(数据库开关 + REST 接口)、traceId 全链路追踪(UUID + 双表日志)、三级容错(全局开关 → 节点开关 → try-catch 兜底),共同保障了系统在 2000+ TPS 下的稳定运行。

    分层架构: 主编排器负责 n1~n4(每条消息必经),实时阶段处理器负责 n5~n7(设备级限流),归档编排器负责 n8(定时批量),三层各司其职。

    可扩展性: 新增节点只需三步——定义节点编码、编写节点类、在 Orchestrator 中添加调用。流程定义与运行态分离存储,支持版本化管理。

    这套 Pipeline 的设计思路不仅适用于工业物联网场景,对于任何需要"多步骤数据处理 + 高可用 + 可观测"的系统都有参考价值。

    赞(0)
    未经允许不得转载:171主机测评 » 工业物联网实战:基于显式编排器模式的9阶段数据Pipeline全链路解析
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址