前言
在工业物联网系统中,设备上报数据的处理链路是整个平台的"心脏"。每秒可能有数千条传感器数据涌入系统,需要经过解析、校验、持久化、缓存、告警、归档等多个环节。如何保证这条链路在高并发下稳定运行?如何在出问题时快速定位故障节点?如何在不重启服务的情况下动态调整处理逻辑?
本文将结合一个真实的工业物联网项目,详细拆解一套基于显式编排器模式构建的 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(流程级日志):
| a1b2c3d4… | device_report | SUCCESS | 10:00:00.000 | 10:00:00.035 |
wf_node_log(节点级日志):
| 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 因为负载过高开始拒绝写入:
三级容错层层递进,从粗粒度到细粒度,从运维操作到代码兜底,共同保障系统的高可用。
五、实时阶段处理器:限流与子编排
前面提到,主编排器把 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();
}
初始化过程做三件事:
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 的设计思路不仅适用于工业物联网场景,对于任何需要"多步骤数据处理 + 高可用 + 可观测"的系统都有参考价值。






