4811 字
24 分钟
车云链路双实现:从 SSH 同步直连到 Kafka 异步长任务

前面几篇面试文讲的是方法论和单点深挖()。这篇专讲一个我自己接手过的真实系统——数据采集模块车端下发链路的双实现。这是那种”不讲全了面试官追问就漏”的题,所以我把它完整拆开:两套实现各自怎么走、为什么有两套、核心的 manifest 指针模式、看门狗状态机、平台信封 codec,全部用代码注释级别的细节讲清。

背景:我接手的是数据采集模块和采集需求模块。要干的事是——云端发命令 → 车端工控机开始/停止录制自动驾驶数据 → 数据回传质检。车端是一台跑在车上的 Linux 工控机,有 5G 网络,会进隧道断网。云端要远程控制它。

历史上只有一套实现,SSH 版(云端直接 SSH 登录车端执行 shell 脚本)。后来因为一次生产事故,团队做了第二套链路版(基于 Kafka 的异步长任务),两套同时存在、用开关灰度切换。

一、为什么会有两套:一次连接池雪崩#

先讲为什么做第二套,这是整套设计的动机,也是面试第一问。

SSH 版的采集执行类,类注释写得直白——2026-05-21 线上四个 pod 集体连接池雪崩,直接根因是:这个类历史上挂了 class 级别的数据库事务注解。Spring 代理调方法时会开启数据库事务,而方法体里跑的是几十秒到几小时的 SSH 远程命令,数据库连接被借出却长时间不归还。SSH 链路一抖,连接池里的连接很快被全部借光,所有后续请求在获取连接时排队 30 秒后被中断。

修复分两步:治标是升级 HikariCP(让连接池能探活僵尸连接,前面那篇讲过);治本是把 SSH 调用搬出事务——注释里写死了规矩:本类严禁挂任何数据库事务,所有数据库写操作必须走仓储层的短事务,每次写自成一个毫秒级短事务立即释放连接

但即使修了事务,SSH 版还有三个结构性毛病改不掉:

  1. 失联就卡死:车进隧道 5G 断了,SSH 会话干等超时——启动要等 6 分钟、停止要等 120 分钟才返回,期间任务卡在”运行中”,前端一直转圈。
  2. 同步阻塞:一个 SSH 会话占着,同一台车的其他操作要排队。
  3. 结果靠同步等 跑完才返回,中间网络抖一下整条链路失败,哪怕车端其实已经启动成功了——这是”幽灵录制”(库里没记录、车端却在录)的温床。

链路版就是为解决这三个问题做的。

二、SSH 版完整链路#

通讯步骤#

用户点"开始采集"
A. 云端准备(采集主服务的启动方法)
1. 前置校验:车在不在线、任务冲突、磁盘状态
2. 生成执行实例号、落库状态=下发运行中
(用独立短事务,毫秒级释放数据库连接)
B. 建立 SSH 连接(车辆连接工具类的建连方法)—— 关键步骤
1. 从 Redis 查车在线信息
- 拿到车的内网穿透地址 + 跳板机端口映射
- 校验"在线":上报时间必须在最大在线间隔内,否则算掉线
2. 两跳 SSH 隧道(注意是两跳,不是直连):
第一跳:连跳板机(公网可达的服务器)
第二跳:从跳板机连到 127.0.0.1 的某端口(跳板机转发到车内网)
3. 强制开心跳保活,每 30 秒一次
(这是治理时加的,原来心跳间隔为零等于不发,5G 静默断连要干等 3 分钟)
4. 双板车(车上有两块计算板)要 ping 另一块板确认通,30 秒超时
C. 执行 shell 脚本(异步执行方法)
1. 从对象存储拉采集脚本,SFTP 上传到车端临时目录
2. SSH 执行脚本,脚本内部干一堆事:挂载外置盘、校验磁盘、启动录制进程
3. 同步等结果(几十秒到几分钟,SSH 通道全程占着)
4. 成功落"下发成功",失败落"下发失败"
5. SSH 输出全程记到 Redis,前端轮询看日志

SSH 版的硬约束#

两个踩过的坑,是面试能讲的差异化:

  • 行尾必须 LF。车端 bash 要求 LF 行尾,Windows 开发机上的 CRLF 会让车端 bash 报 command not found。仓库用 .gitattributes 强制 shell 脚本为 LF,不能用会改写行尾的工具处理这些脚本。
  • 磁盘路径按车型解析。不同车型外挂盘设备名不同(有的是 /dev/sda,有的是 /dev/nvme0n1),挂载、进目录、取容量三处必须用同一个解析函数按车型解析,写死会格式化错误设备。
  • 停止采集绝不”只改库不通知车”。历史上前端”强制结束”入口直接跳过车端下发,造成”库里已结束、车端仍在录”的幽灵录制。停止一律尽力真实停车,停不下来按采集失败收尾。
  • 无论停止车端成败,下发状态一律保持成功。这是和上一条配套的红线:停止车端失败时,若把”下发状态”回写成失败,会把”采集阶段失败”误标成”下发失败”,前端会错出”重新下发”按钮——对已在收尾的任务重复下发,等于在幽灵录制上又叠一层重复下发。

三、链路版完整链路(Kafka 异步)#

通讯步骤#

用户点"开始采集"
A. 云端准备(链路服务的启动方法)—— 提交即返回,不等车端
1. 前置校验(同 SSH 版)
2. 生成确定性任务号 = "DC:start:{任务ID}"
- 幂等:同一任务重复提交得到同一任务号,任务管理器见已活跃直接忽略
- 无状态回查:监听器从任务号反解出任务ID,不用在任务表加字段
3. 构建 manifest(执行清单)—— 核心创新,下文专讲
4. 落库状态=下发运行中(短事务)
5. 任务管理器提交下发命令 → 立即返回任务号给前端
(到这里云端动作结束,前端已拿到响应,不等车端)
B. 下行通道(下行发送器 → mars 平台 → 车端 agent)
1. 云端把下行消息发到 Kafka 下行主题
- 消息包 mars 信封(因为走 mars/iot 平台中转,下文专讲)
2. mars 平台(车联网中台)把消息推到车端 agent
3. 车端 agent 收到 → 拉 manifest(按指针 URL 从对象存储下载)
4. 校验 sha256(防篡改/防传输损坏)
5. 按 manifest 的 steps 顺序执行每一步
6. 每完成一步 → 上报进度到 Kafka 上行主题
7. 全部完成 → 上报最终结果(成功/失败)到 Kafka 上行主题
C. 上行回收(上行监听器 → 消息分发器 → 业务监听器)
1. Kafka 消费者从上行主题读车端上报消息
2. 消息分发器按关联标识(=任务号)路由到对应的任务状态机
3. 状态机推进:已确认 → 执行中 → 最终结果
4. 业务监听器监听终态 → 从任务号解析出任务ID → 回写采集任务状态
5. 失联兜底:看门狗定期扫活跃任务,下文专讲

四、核心创新之一 指针模式#

链路版下行不是直接发命令,而是发一个执行清单的指针。这是和 SSH 版最大的结构性差异。

为什么不直接发命令#

Kafka 有单条消息大小限制(默认 1MB)。采集流程是多步的——挂载盘、校验磁盘、校验通道帧率、启动录制、生成摘要、上传产物——每步带具体 shell 命令,全展开可能很大。直接塞 Kafka 会被拒或被截断。

怎么做#

  1. 云端把多步流程生成成一份 JSON 清单(步骤数组 + 每步 shell 命令 + 每步超时),序列化后写入对象存储。
  2. 算 sha256 摘要。
  3. 下行消息只带指针:{schemaVersion, op, taskId, exeId, carId, vin, totalSteps, manifestOssName, manifestKey, manifestSha256, manifestUrl}——消息极小。
  4. 车端 agent 收到指针 → 下载清单 → 校验 sha256 → 按步骤执行。

这个设计带来的好处#

  • 绕开 Kafka 大消息限制:下行消息永远是几 KB 的指针。
  • 完整性校验 防传输损坏、防篡改。
  • 可审计:清单留在对象存储,任何一次下发的完整步骤都能回溯。
  • 云端编排、车端执行:云端决定流程,车端只管执行退出码——判定逻辑塞进命令(脚本退出码非零即本步失败),车端只看退出码,职责干净。
  • 凭证不下发车端:脚本下载统一用预签名 URL,访问密钥不下发车端,安全面收敛。

一个精度换可靠性的取舍#

注释里写了个诚实的设计取舍:链路版把磁盘路径用云端可推断的默认值(默认外置盘),不再像 SSH 版那样实时探测车端状态。这是用精度换可靠性——少一次实时探测就少一次失败点。后续可由车端 agent 在执行前先跑磁盘检查子命令再拼最终路径。

五、核心创新之二 看门狗状态机#

这是链路版替换 SSH 后最关键的补偿机制

为什么需要看门狗#

注释写得很清楚:同步 RPC(SSH)里”对端没响应”由 socket 超时直接暴露;异步消息(Kafka)里没人推就永远静默——车端进隧道断了,云端不知道,任务永远挂”运行中”。必须由云端主动判定。

状态机全貌#

CREATED ──SENT──▶ SENT ──ACK──▶ ACKED ──PROGRESS──▶ RUNNING ──RESULT_SUCCESS──▶ SUCCESS
│ │ │ │ ▲ (终态)
│ │ │ 心跳超时│ │ 进度/心跳
│ │ │ ▼ │
│ │ │ STALLED
│ 投递失败 │ 投递失败 │ 总截止 │ 总截止
▼ ▼ ▼ ▼
DELIVERY_FAILED DELIVERY_FAILED TIMEOUT TIMEOUT
(终态) (终态) (终态) (终态)
任意 ACKED/RUNNING/STALLED ──车端报错──▶ FAILED (终态)

九个状态,四个终态(成功/失败/超时/投递失败)。终态后所有事件被忽略——这是”至少一次投递”下重复消息不出错的关键(同一结果上报多次不会把已成功的任务改回运行中)。

看门狗的扫描逻辑#

看门狗是个周期调度任务,扫描所有活跃任务,做三件事:

  1. 总截止时间判定(优先级最高):当前时间 > 任务总截止时刻 → 判超时。无论任务在哪个中间态,超总时限一律超时收尾。
  2. 失联判定:任务在”已确认”或”执行中”且当前时间 > 心跳容忍截止 → 判疑似失联。这个状态可恢复——车端恢复上报进度就回”执行中”,不直接判死。
  3. 退避重发判定:任务在”退避等待”且到了重发时刻 → 触发重发回”已投递”。

ACK 超时快速重发#

一个精细设计:任务在”已投递”久未确认(最后一公里丢包),不等总截止,快速重发。看门狗注释写:阈值、白名单、预算判定收敛在任务管理器,避免死等到总截止。

为什么 STALLED 不是终态#

这是设计的关键:疑似失联(STALLED)可恢复。车进隧道 30 秒没心跳判疑似失联,出隧道恢复上报,任务继续。只有超总截止才判超时收尾。这给了网络抖动足够的容忍,不误杀。

下面这个图可以点——点任意状态看它的所有转移和触发条件:

六、核心创新之三 平台信封 codec#

链路版下行不直接发到车端,要走 mars/iot 平台中转(mars 是车联网中台)。云端内部的统一消息模型是一套业务格式,但发往 mars 平台 Kafka 下行主题的数据必须包成 mars 信封结构:

{
"productName": "autopilot",
"handleType": 1,
"deviceName": "test-1",
"data": "Zm9v"
}

data 字段的编码契约(踩过的坑都在这)#

注释里写了好几个坑,都是真实踩出来的:

  1. 类型固定为字符串,不是字节数组。为什么?避免 Jackson 默认字节数组 Base64 行为被业务方误配 ObjectMapper 关掉时,线上静默变成整数数组。显式字符串字段,协议契约由本类自身收口、不依赖框架默认。
  2. 用标准 Base64(带 padding 的 RFC 4648 字母表),与车端解码端(Java 的 Base64 解码器、Python 的 base64.b64decode)默契一致。
  3. 空负载 → data 为空字符串。Base64 编码空字节就是空串,与”无负载”语义自然吻合,避免 null 让车端解析 NPE。
  4. 统一构建入口:业务侧走公用包装方法,避免各自构造字段错位。

为什么要包信封#

mars 平台是个通用中台,接多种设备、多种产品。信封结构把”产品标识 + 操作类型 + 设备名 + 业务数据”标准化,mars 只看信封路由到对应车端,不关心业务数据内容。云端业务侧只管业务消息,转信封的事收口在一个 codec 类——这是典型的协议分层:业务层和网络传输层解耦。

七、两套对比#

维度SSH 版(老)链路版(新,Kafka)
通讯模型同步阻塞,云端直连车端异步消息,云端 → Kafka → mars 平台 → 车端 agent
连接方式两跳 SSH 隧道(跳板机→车内网)无长连接,Kafka 收发解耦
命令下发SFTP 上传脚本 + SSH 执行下发 manifest 指针,车端 agent 拉清单自执行
结果回收同步等 shell 跑完异步,车端上报进度/结果到 Kafka
失联处理干等超时(启动 6 分钟/停止 120 分钟)看门狗主动判疑似失联/超时,落失败
事务风险历史上挂事务导致连接池雪崩(已修)提交即返回,不持事务,天然无此风险
幂等靠 Redis SETNX 锁确定性任务号,任务管理器内置幂等
可观测SSH 输出记 Redis,前端轮询Kafka 全链路 + 状态机 + 看门狗,每步进度可见
完整性校验sha256 校验 manifest
复杂度简单直接但脆复杂(状态机/看门狗/监听器/清单/codec)但稳
依赖车端有 SSH 服务、跳板机车端有 agent、mars 平台、Kafka
灰度默认开启开关控制,默认关,与 SSH 版灰度并存

优缺点总结#

SSH 版:

  • 优点:实现简单直接、调试容易(SSH 进去就能看);不依赖额外中间件,车端只要开了 SSH 就能用;适合早期、车少、网络稳的场景。
  • 缺点:①同步阻塞占资源,长操作持数据库事务会搞垮连接池(真实事故);②失联只能干等超时;③单点阻塞,同车操作串行;④网络抖动易整链失败,产生幽灵录制;⑤无法支撑大规模车队。

链路版(Kafka):

  • 优点:①异步解耦,提交即返回;②失联有看门狗主动判死,不干等;③Kafka 天然削峰,支撑大规模车队;④状态机 + 进度上报,可观测性好;⑤幂等设计防重复。
  • 缺点:①复杂度高——状态机、看门狗、监听器、清单、上下行 codec 一大堆;②依赖重——要 mars 平台、Kafka、车端 agent,任一环挂都影响;③灰度并存期有双写双读的复杂性(两套都要落痕,否则操作历史断档);④排障难——出问题要跨 Kafka/mars/车端三处查。

八、面试怎么答这道题#

「我们采集任务下发到车端有两套并行实现,用配置开关灰度切换。老的是 SSH 版——云端两跳 SSH 隧道直连车端,SFTP 上传脚本后执行,同步等结果。新的链路版基于 Kafka——云端构建 manifest 执行清单写对象存储,下行只发指针,车端 agent 拉清单自执行、按步上报进度,终态由监听器回写。

做第二套是因为 SSH 版出了生产事故——历史上 SSH 调用挂在数据库事务里,长操作占着连接不还,四个 pod 连接池集体雪崩。但即使修了事务,SSH 还有个结构病:车进隧道 5G 断了要干等 6 分钟超时。链路版用看门狗主动判疑似失联和超时落失败,不再干等。

两套灰度并存是因为新链路复杂度高、依赖 mars 平台和车端 agent,不能一刀切。两套共用同一张表和状态字段,前端无感切换。」

三个能加分的追问点#

  1. “manifest 指针模式为什么不直接发命令?” → Kafka 有大消息限制,清单可能很大,只传指针绕开限制还能 sha256 校验完整性。云端编排、车端执行,职责干净。
  2. “看门狗的疑似失联为什么不是终态?” → 可恢复。车进隧道 30 秒没心跳判疑似失联,出隧道恢复就继续,不误杀。只有超总截止才判超时收尾。
  3. “两套灰度并存怎么保证操作历史不断?” → 两套都要落操作留痕(只记人工动作),否则前端按开关切换后,操作历史会突然断档——这是上一版留痕被回滚的头号问题。

九、这段经历对加 AI Agent 的启示#

回到之前给这套系统设想的 AI Agent 落点——这套链路版的异步长任务 + 状态机 + 看门狗基建,正是 AI Agent 落地现成的脚手架:

  • 状态机的终态幂等 调 LLM 也是异步长任务,LLM 调用慢且可能超时,看门狗判定失联落失败的机制可以直接复用。
  • manifest 指针模式 的多步执行计划(规划 → 工具调用 → 后处理)可以套同样的指针模式,下行只发计划指针,执行计划存对象存储可审计。
  • 上下行 codec 分层 和 LLM provider 之间也该有 codec 层,业务消息和传输格式解耦,换 provider 只改适配层。

所以这套链路版不只是”下发采集任务”,它是一套通用的”云端编排、车端执行、异步回收、失联兜底”长任务框架。面试时把它讲成框架而非单点功能,格局就上去了——这也是为什么设想的标定失败归因 Agent 能直接复用这套基建。

车云链路双实现:从 SSH 同步直连到 Kafka 异步长任务
https://gilgameshzzz.github.io/posts/vehicle-link-dual-implementation/
作者
Amadeus
发布于
2026-09-20
许可协议
CC BY-NC-SA 4.0