Hcc的Blog

嵌入式 · AI · 折腾不止

0%

DDS入门指南

学习 DDS,使用 claude 实现最小例程并学习。

1. 先看一个具体问题

假设你在做一个机器人,系统里有这些进程:

  • 硬件抽象层:以 500Hz 读取关节编码器、IMU,输出机器人状态
  • 运动控制:以 40Hz 跑控制算法,需要机器人状态,输出关节指令
  • 遥控接收:接手柄/遥控器输入,输出速度指令
  • 日志记录:想看到上面所有数据

它们要互相传数据。你会怎么做?

方案一:TCP socket。 于是你要处理:谁连谁(谁是 server)、连接断了怎么重连、
数据怎么编码成字节又解回来、一份数据要发给 3 个订阅者怎么办(连 3 条连接?)、
新增一个进程要改几处配置。这些都是与业务无关的苦工,而且每一项都容易写错。

方案二:消息队列(Redis / RabbitMQ / ZeroMQ)。 好一些,但引入了 broker:
一个额外的进程,成为单点故障,且数据要多走一跳(发布端 → broker → 订阅端)。
对 500Hz 的控制回路,这一跳的延迟和抖动是要命的。而且 broker 不理解你的数据,
它只是搬运字节。

方案三:ROS 的 topic。 这就接近了 —— 事实上 ROS 2 底层用的就是 DDS。

DDS 要解决的正是这个问题:多个进程之间,高频、低延迟、无中心节点地共享有类型的数据。
它是 OMG(对象管理组织)的国际标准,工业界有多个实现(CycloneDDS、FastDDS、
RTI Connext、OpenDDS),本例程用的是 Eclipse CycloneDDS。


2. DDS 是什么:三种范式的对比

DDS 全称 Data Distribution Service,核心范式叫 DCPS —— 以数据为中心的发布订阅
(Data-Centric Publish-Subscribe)。

这个”以数据为中心”是关键,先和你熟悉的东西对比:

消息队列 RPC(gRPC) DDS
通信单位 不透明的字节/消息 函数调用 有类型的数据对象
需要中间人 需要 broker 需知道服务端地址 无,端到端自组织
语义 “把这个包送过去” “调用你的方法” “这块数据的最新值是 X”
谁认识谁 生产者认识队列名 客户端认识服务端 两端都只认识 Topic 名
中间件理解数据吗 不理解 理解签名 理解每个字段

“以数据为中心” 的实际含义

中间件理解你传的是什么结构 —— 这就是为什么 DDS 强制你先写 IDL 文件描述数据。
因为理解数据,它才能做到这些只有理解数据才能做的事:

  • 按主键(key)把数据分成多个独立实例,各自维护历史
  • 只保留每个实例的最新 N 个值,旧的自动丢弃
  • 检测”某个发布者失联了”
  • 按内容过滤(只接收 temperature > 30 的样本)

消息队列做不到这些,因为对它来说消息就是一串字节。

一个视角转换

publisher.cpp:67 这一行:

1
writer.write(sample);

它的语义不是“发送一条消息”,而是”更新 Demo/SensorData 这个数据对象的当前值”。

这个区别现在看起来像咬文嚼字,但它会解释后面几乎所有的设计:为什么默认只保留最新
一帧、为什么丢帧是可接受的、为什么”发布端停止发送”不等于”数据归零”。

把 DDS 想象成一块所有进程都能看到的共享白板,而不是一根管道。发布端往白板上
写值,订阅端随时去看白板上现在写的是什么。管道关心”每个包都要送到”,白板只关心
“上面的值是不是最新的”。


3. 先把例程跑起来

例程只有 3 个源文件,各自职责:

1
2
3
idl/sensor.idl      数据类型定义(要传什么)
publisher.cpp 发布端,20Hz 发数据
subscriber.cpp 订阅端,回调收数据、main 循环 5Hz 消费

构建和运行:

1
2
3
./scripts/build.sh        # 构建(首次会自动跑 idlc 生成类型代码)
./scripts/run.sh both # 两端都起,跟随订阅端日志
./scripts/run.sh stop # 停止

日志在 logs/。你会看到这样的输出:

1
2
3
4
5
6
7
8
[SUB][main]     thread=140079814488320
[SUB] reader ready, waiting for data...
[SUB][main] no data yet (发布端起了吗?QoS 对得上吗?)
[SUB][main] no data yet (发布端起了吗?QoS 对得上吗?)
[SUB][callback] FIRST sample, thread=140079686215360 <-- 注意这个 id 和 main 不同
[SUB][main] frame=1 (+0) temp=25.2499 gyro=[0.05,0.999,0] status=0 | rx_total=1 latency=44867.6us
[SUB][main] frame=5 (+4) temp=26.237 gyro=[0.247,0.969,0] status=0 | rx_total=5 latency=43940us
[SUB][main] frame=9 (+4) temp=27.1748 gyro=[0.435,0.900,0] status=0 | rx_total=9 latency=43264.3us

这段输出里有 4 个值得注意的现象,本文后面会逐个解释:

  1. 两个 thread id 不一样 —— 回调不在 main 线程上(第 8 节)
  2. frame 每次 +4 —— 发布 20Hz、消费 5Hz,中间的帧被丢了(第 7 节)
  3. latency 是 44 毫秒,而且在逐渐变小 —— 这个数字测错了(第 10 节)
  4. 开头有几次 no data yet —— 连接建立需要时间(第 6 节)

4. 数据模型:IDL、Topic、类型名

为什么要写 IDL

DDS 要求你用 IDL(Interface Definition Language,OMG 的标准语言)描述数据结构。
idl/sensor.idl 是本例程的全部数据定义:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
module demo {
module msg {

struct Header {
long long timestamp; // 纳秒
long long frame_id; // 自增帧号,用来发现丢帧
};

// topic: "Demo/SensorData"
struct SensorData {
Header header;
float temperature; // ℃
float gyro[3]; // rad/s,定长数组
octet status; // 0=正常 1=告警
};

};
};

写 IDL 而不是直接写 C++ struct,有三个原因:

  1. 跨语言:同一份 IDL 可以生成 C++、C、Python、Java、Rust 的类型,不同语言写的
    进程能互通
  2. 跨平台:序列化格式(CDR)是标准的,x86 和 ARM、32 位和 64 位之间字节序和对齐
    都由中间件处理
  3. 中间件需要理解结构:第 2 节说的那些能力(按 key 分实例、内容过滤)都依赖这个

IDL 到 C++ 的映射

构建时 idlc 工具把 sensor.idl 编译成 build/sensor.hpp(663 行)和
build/sensor.cpp(89 行)。生成的是普通 C++ 类:

1
2
3
4
5
6
7
8
9
// build/sensor.hpp
class SensorData
{
private:
::demo::msg::Header header_;
float temperature_ = 0.0f;
std::array<float, 3> gyro_ = { }; // ← IDL 的 float[3]
uint8_t status_ = 0; // ← IDL 的 octet
};

映射规则:

IDL C++ 备注
module namespace 嵌套关系保持
long long int64_t
long int32_t 注意不是 C++ 的 long!
short int16_t
unsigned long uint32_t
octet uint8_t
boolean bool
double double
float[3] std::array<float,3> 定长,值语义,栈上
sequence<float> std::vector<float> 变长,会堆分配
string std::string 变长,会堆分配

这里有一条实时系统的重要规则:实时路径上的 IDL 应该只用定长类型。
sensor.idl:6 的注释点到了这一点 —— std::array 是值语义、栈上、零堆分配;
而 sequence / string 每次收发都可能触发 malloc,在 500Hz 的控制回路里,
malloc 的锁竞争和不确定延迟是真实风险。

Topic = 名字 + 类型

publisher.cpp:28 和 subscriber.cpp:91 是同一行代码:

1
dds::topic::Topic<demo::msg::SensorData> topic(participant, "Demo/SensorData");

Topic 的身份由两个维度构成:

  1. Topic 名:字符串 "Demo/SensorData"。斜杠只是命名习惯,DDS 里没有层级含义
  2. 类型名:demo::msg::SensorData

第二个维度是初学者最容易忽视的。它的实际来源在生成代码里,build/sensor.hpp:165:

1
2
3
4
template <> constexpr const char* TopicTraits<::demo::msg::SensorData>::getTypeName()
{
return "demo::msg::SensorData";
}

idlc 把 IDL 里的 module demo { module msg { struct SensorData 拼成了字符串
"demo::msg::SensorData"。这个字符串在连接建立阶段会随端点信息广播出去,
两端逐字节比较。

后果:把 IDL 里的 module demo 改成 module foo(只改一端),两端就再也连不上。
不是因为 C++ 类型不兼容,而是因为这个字符串变了。

这在多方协作时尤其要命:如果对接的是别人编译好的二进制(比如另一块板子上的底层
软件),它广播的类型名是固定的,你这边 IDL 里 module 名写错一个字母,就永远收不到
数据 —— 而且没有任何错误提示(第 11 节详述)。


5. 四层实体:那四行初始化各自在做什么

publisher.cpp:25-43 有四步初始化。它们不是仪式,每层有明确职责:

1
2
3
4
5
6
DomainParticipant(0)          域 = 隔离边界;一个进程一份(重量级)
├── Topic<T>("名字") 名字+类型的绑定声明
├── Publisher DataWriter 的容器
│ └── DataWriter<T> 真正干活的:write()
└── Subscriber
└── DataReader<T> 真正干活的:take() / read()

DomainParticipant:域是硬隔离

1
dds::domain::DomainParticipant participant(0);

参数 0 是 domain ID。规则很简单:不同 domain 的进程互相完全不可见,
哪怕在同一台机器、用同一个 topic 名、同一个类型。

这个隔离是物理层面的,不是逻辑过滤 —— domain ID 直接编码进 UDP 端口号
(下一节会看到)。所以它是最省事的隔离手段:同一台机器上跑两套互不干扰的系统
(比如仿真和实机、开发和测试),改 domain ID 就够了。

反过来说,如果整个项目所有进程都硬编码 DomainParticipant(0),就等于放弃了这个
隔离维度 —— 哪天要在同一台机器上跑双实例,这是第一个要改的地方。

Participant 是重量级对象。 我实测了它的开销 —— 一个只有 1 个 writer 的
publisher 进程:

1
2
$ ls /proc/<pid>/task    # 列出线程
publisher gc dq.builtins dq.user tev recv recvMC recvUC

一个 participant = 8 条线程(含 main)。加上 4 个 UDP socket。

所以”一个进程只创建一个 participant”是铁律。不要在函数里随手创建,更不要为每个
topic 创建一个。

Publisher / Subscriber:多数情况下是纯样板

1
2
dds::pub::Publisher publisher(participant);
dds::sub::Subscriber subscriber(participant);

说实话,这两层在 90% 的代码里没有实际作用。它们存在是为了承载跨多个 writer 的
组语义
—— 主要是 Presentation QoS 的 GROUP 级别,语义是”这几个 writer 的更新
要作为一个原子批次被订阅端看到”。

本例程没用这个特性,绝大多数项目也不用。所以看到这两层不用多想:

  • 它们不是性能相关的连接池
  • 它们不是线程的宿主
  • 它们只是 API 要求的中间层

DataWriter / DataReader:真正的主角

1
2
dds::pub::DataWriter<demo::msg::SensorData> writer(publisher, topic, qos);
dds::sub::DataReader<demo::msg::SensorData> reader(subscriber, topic, qos);

这一层是唯一有数据流经过的地方,也是绝大部分 QoS 的宿主。学 DDS 时,注意力应该
主要放在这一层。


6. 发现机制:没有 broker,两端怎么找到彼此

这是 DDS 最”魔法”的部分。两个进程,谁也不知道对方的 IP 和端口,谁也没配置过对方的
地址,启动后就能自动通信。

底层:RTPS over UDP

DDS 的线上协议叫 RTPS(Real-Time Publish-Subscribe),跑在 UDP 上。我实测了
publisher 进程打开的端口:

1
2
3
4
5
$ ss -unap | grep publisher
UNCONN 0.0.0.0:7400 ← SPDP 多播:发现 participant
UNCONN 0.0.0.0:7401 ← 用户数据多播
UNCONN 192.168.207.177:46087 ← 单播
UNCONN 0.0.0.0:36977 ← 单播

端口 7400 由公式 7400 + 250 × domain_id 得出 —— 这就是上一节说的
“domain 隔离是物理层面的”:不同 domain 用不同端口,包本身就到不了对方。

阶段 1:SPDP —— 找到对方进程

SPDP(Simple Participant Discovery Protocol):participant 创建后,周期性往
多播地址 239.255.0.1:7400 广播一条”我是谁”:

  • 自己的 GUID(全局唯一 ID)
  • 支持的协议版本
  • 自己的单播地址和端口
  • 租约时长(多久没心跳就认为我死了)

同时它也在这个端口上监听别人的广播。两个 participant 收到彼此的广播,就完成了
第一阶段握手。

关键点:这是多播(multicast)。 由此推出几个实践后果:

  • 同机通信也走网络协议栈(回环多播),不是共享内存(除非显式配置
    Iceoryx 之类的共享内存传输)
  • 多播被防火墙拦、被交换机禁、Docker 默认桥接网络不转发多播 → 发现失败
  • 跨网段需要路由器支持多播转发,否则要手写对端单播地址(CycloneDDS 通过
    CYCLONEDDS_URI 环境变量配置 <Peers>)

排查”两端都起了但收不到数据”,确认多播是否通畅永远是第一步。

阶段 2:SEDP —— 交换端点清单

SEDP(Simple Endpoint Discovery Protocol):SPDP 握手成功后,两个 participant
用单播互相发送内置 topic 的数据,内容是”我有哪些 writer 和 reader,各自的
topic 名、类型名、QoS 配置是什么”。

对端拿到清单后,逐条比对:

1
topic 名相同? → 类型名相同? → QoS 兼容? → 建立数据通路(match)

三个条件全部通过才 match。任何一条不过,静默丢弃 —— 无日志、无异常、无返回码。

这就是 DDS 排障困难的根本原因,第 11 节会专门讲。

一个设计细节:发现数据和用户数据分开处理

回看上一节实测的 8 条线程,有两条是投递队列:

1
2
dq.builtins   发现数据的投递队列
dq.user 用户数据的投递队列 ← 你的回调跑在这里

它们是分开的。这个设计有实际意义:你的回调把 dq.user 卡住时,发现机制仍然
正常工作,participant 不会被对端误判为失联。

实测:连接建立要多久

这个数字很重要,因为它是工程约束。我写了一个带计时的版本来测(用
publication_matched_status() 精确捕捉 match 时刻)。

场景 A:订阅端已经在运行,发布端后启动(跑 3 次)

1
2
3
[T] participant_created +2.19ms
[T] writer_created +2.82ms
[T] MATCHED +5.05ms (writer 创建之后仅 2.24ms)

场景 B:两端几乎同时启动(跑 3 次)

1
[T] MATCHED +5.56ms / +5.70ms / +5.68ms  (writer 创建之后约 2.3ms)

场景 C:发布端先启动,订阅端 1.5 秒后启动

1
[T] MATCHED +1509ms   ← 订阅端在 1500ms 才起,所以 match 本身只用了约 9ms

结论:在本机、多播通畅的情况下,match 只需 2-10 毫秒。 DDS 的发现机制比多数人
想象的快得多。

那开头的几次 no data yet 是怎么来的?

logs/subscriber.log 开头有 2-3 次 no data yet,按 200ms 一次算是 400-600ms。
既然 match 只要几毫秒,这几百毫秒是哪来的?

我用阶段计时测了订阅端(复刻 run.sh both 的启动方式):

1
2
3
4
5
[T] main_entered   +0.018ms
[T] participant +2.82ms
[T] reader_ready +3.60ms ← reader 在 3.6ms 就绪了
[T] FIRST_DATA +308.55ms frame=0 ← 首帧,而且是 frame=0
[T] MATCHED +308.78ms

答案清楚了:那 300ms 是在等发布端启动。 run.sh:66 里 publisher 比 subscriber
晚起 sleep 0.3,订阅端的 reader 在 3.6ms 就准备好了,剩下的时间纯粹在等对方。

而且注意 FIRST_DATA 收到的是 frame=0 —— 发布端的第一帧,一帧都没丢。

那为什么 main 循环打了 2-3 次 no data yet?因为 main 每 200ms 才看一次缓存
(subscriber.cpp:164),300ms 的等待期内它会看 1-2 次,都还没数据。这是 main
自己的采样周期决定的,与 DDS 无关。

这个澄清有实际价值:如果你以为”DDS 发现要几百毫秒”,就会在启动逻辑里留很长的
等待窗口,或者误以为连接慢是中间件的问题。实际上 match 是毫秒级的,几百毫秒的等待
几乎总是对端还没起来。


7. QoS:匹配是一份契约,不是配置项

QoS(Quality of Service,服务质量)是 DDS 里最重要也最容易出错的概念。DDS 标准定义了
22 种 QoS 策略。好消息是:日常只需要真正理解 3-4 个。

请求/提供模型(RxO)

初学者最常见的误解是”两边 QoS 写一样就行”。实际规则更精确:

Reader 请求(Requested)的服务等级,不能高于 Writer 提供(Offered)的。

这叫 RxO 契约(Requested ≤ Offered)。它是一个偏序关系,不是相等关系。

本例程的两处 QoS 配置是核心对照点。publisher.cpp:39-41:

1
2
3
4
// Writer:提供 BestEffort
dds::pub::qos::DataWriterQos qos;
qos << dds::core::policy::Reliability::BestEffort()
<< dds::core::policy::History::KeepLast(1);

subscriber.cpp:101-103:

1
2
3
4
// Reader:请求 BestEffort  ← 必须 ≤ Writer 提供的
dds::sub::qos::DataReaderQos qos;
qos << dds::core::policy::Reliability::BestEffort()
<< dds::core::policy::History::KeepLast(1);

Reliability:可靠性

两个取值,Reliable > BestEffort:

  • BestEffort:尽力而为。丢了就丢了,不重传
  • Reliable:可靠。丢包会检测并重传,保证按序无丢失到达

兼容矩阵:

Writer 提供 ↓ / Reader 请求 → BestEffort Reliable
BestEffort ✅ match ❌ 静默不匹配
Reliable ✅ match ✅ match

我实测验证了那个 ❌ 格子:把订阅端改成 Reliability::Reliable()(其他一切不变,
topic 名和类型名完全一致),运行 12 秒:

1
2
发布端:正常发送,[PUB] frame=0 / frame=20 / frame=40 ...
订阅端:no data yet ×63 次,收到数据帧数 = 0

63 次全部静默,零报错,零异常。 这就是 DDS 最经典的陷阱。

为什么高频实时流要选 BestEffort

publisher.cpp:32-38 的注释点到了要害。这个选择不是”图省事”,而是 Reliable 在这个
场景下有害:

  1. 重传的数据已经过期。 20Hz 的传感器数据,丢了第 100 帧,等重传到达时第 103 帧
    已经到了。用一帧 150ms 前的旧数据做控制,比直接跳过更糟
  2. 重传挤占带宽和 CPU,而且恰好发生在网络已经拥塞的时候 —— 正是最不该增加负载
    的时刻
  3. Reliable 会阻塞 writer。 队列满时 write() 会阻塞(受 ResourceLimits 和
    max_blocking_time 约束)。一个 20Hz 定周期循环被 write() 阻塞,节拍就崩了

选 QoS 的判断标准:传的是状态还是事件

这是一条实用的判断规则:

状态量 事件量
例子 温度、角速度、位姿、电池电量 急停按下、任务完成、模式切换
语义 “现在的值是 X” “发生了一件事”
漏掉一个 无所谓,下一帧就修正了 不可恢复
该选 BestEffort + KeepLast(1) Reliable + KeepAll

本例程传的是状态量(温度、陀螺仪),所以 BestEffort + KeepLast(1) 是正确的。

但要警惕:如果把事件量塞进同一条 BestEffort 的消息里(比如在传感器消息里加一个
“急停标志位”),那个标志位就会跟着变成”可以丢的”。这是真实项目里容易埋下的隐患。

History:保留多少历史

  • KeepLast(N):每个实例只保留最新 N 个样本,更新的挤掉旧的
  • KeepAll:全部保留(受 ResourceLimits 限制),队列满了 writer 会阻塞

本例程用 KeepLast(1) —— 只保留最新一帧。这直接产生了日志里那个 +4 现象。

生产快于消费:中间帧被丢弃

发布端 20Hz,订阅端 main 循环 5Hz(subscriber.cpp:164)。日志:

1
2
3
4
frame=1  (+0)
frame=5 (+4)
frame=9 (+4)
frame=13 (+4)

+4 是必然的:20 ÷ 5 = 4。中间的 3 帧被 KeepLast(1) 挤掉了。

我跑了 75 秒统计 gap 分布:

1
2
3
367 次 (+4)
7 次 (+3) ← 时钟漂移导致的,第 10 节解释
1 次 (+0) ← 首帧

这是故意的设计,不是缺陷。 高频状态流里,消费不过来的中间帧本来就没有价值 ——
你要的是”现在的值”,不是”每一个历史值”。

状态语义的必然后果:静默 ≠ 归零

publisher.cpp:48-51 记录的这一条,是这套范式最容易造成事故的地方:

上游停发不会自动归零 —— 下游会一直用最后一帧

因为语义是”这个数据对象的最新值”,而不是”消息流”。发布端进程崩了、网线断了,
订阅端缓存里还是崩溃前那一帧。如果那一帧是”前进 1.5 m/s”,机器人会一直往前跑。

DDS 其实提供了检测手段:

  • Liveliness QoS + on_liveliness_changed 回调:检测”发布者还活着吗”
  • Deadline QoS + on_requested_deadline_missed 回调:检测”数据是否按承诺的周期到达”

但本例程没有用 —— subscriber.cpp:109-110 只订阅了一个状态位:

1
2
dds::core::status::StatusMask mask = dds::core::status::StatusMask::none();
mask |= dds::core::status::StatusMask::data_available();

所以只能靠应用层约定:”刹车必须主动发零速,不能靠沉默”、”关机时锁定零速度”。

这是一个值得认真考虑的改进方向:应用层约定依赖每个开发者都记得写;而
Deadline QoS 是中间件强制的。给关键控制流加上 Deadline + 对应回调,能把
“上游失联”从”需要人记得处理”变成”中间件保证会告诉你”。代价是多处理一个回调,
收益是消除一整类静默失效。

其他 RxO 策略

同样遵循”Reader 不能比 Writer 要求更严”的策略:

策略 含义 偏序
Durability 晚加入的订阅端能否收到历史数据 Volatile < TransientLocal < Transient < Persistent
Deadline 承诺/要求的最大数据间隔 Reader 请求的周期不能短于 Writer 承诺的
Ownership 多个 writer 时谁说话算数 SHARED / EXCLUSIVE 必须相同
Liveliness 存活检测方式和租约时长 Reader 要求的租约不能短于 Writer 承诺的

记住”Reader 不能比 Writer 要求更严“这一条,其余都能推导出来。

其中 Durability 对启动顺序有实际影响:默认 Volatile 意味着”晚起的订阅端收不到
之前发过的数据”。如果希望订阅端一启动就能拿到当前值(而不是等下一帧),
需要 Writer 用 TransientLocal。


8. 线程模型:整份例程最重要的一节

如果本文只能记住一节,应该是这一节。DDS 的线程模型是实际项目中事故的主要来源。

实测:8 条线程各自的职责

1
2
3
4
5
6
7
8
publisher     main 线程 —— 你的代码
recv RTPS 报文接收主循环
recvMC 多播接收
recvUC 单播接收
dq.user ★ 用户数据投递队列 —— 你的回调跑在这里
dq.builtins 发现数据投递队列(与 dq.user 隔离)
tev 定时事件:心跳、租约续约、重传定时器
gc 垃圾回收:延迟释放已删除的实体

回调不在你的线程上

subscriber.cpp:86 和 :77 分别打印了两个 thread id:

1
2
[SUB][main]     thread=140079814488320
[SUB][callback] FIRST sample, thread=140079686215360 ← 不同

on_data_available 是你注册的函数,但调用它的是中间件的 dq.user 线程,不是
你的 main。这就是”DDS 回调线程”的含义。

这个事实推出两条硬约束。

约束一:回调与主循环共享数据,必须加锁

subscriber.cpp:25-33 定义了共享缓存:

1
2
3
4
5
6
7
8
9
struct SharedState
{
std::mutex mutex;
demo::msg::SensorData latest;
bool has_data = false;

// 计数用 atomic 就够了,不需要锁
std::atomic<int64_t> rx_count{0};
};

注意这里的粒度选择是经过思考的,不是随手写的:

  • latest 和 has_data 必须一起用 mutex 保护。否则可能读到
    has_data == true 但 latest 只写了一半的状态
  • rx_count 是独立计数器,与其他字段无一致性要求,用 atomic 就够,不必进临界区

写入侧(回调线程),subscriber.cpp:68-72:

1
2
3
4
5
{
std::lock_guard<std::mutex> lock(g_state.mutex);
g_state.latest = s.data();
g_state.has_data = true;
}

读取侧(main 线程),subscriber.cpp:130-137:

1
2
3
4
5
6
7
8
9
{
std::lock_guard<std::mutex> lock(g_state.mutex);
if (g_state.has_data)
{
snapshot = g_state.latest; // 拷一份出来,尽快放锁
valid = true;
}
}
// 锁已释放,后面的计算和打印都在锁外

在锁内只做一次拷贝,拿到快照后立刻放锁,所有计算和 I/O 都在锁外。

这是控制系统里的标准做法,理由很直接:锁的持有时间决定了它对回调线程的阻塞上限。
如果把 subscriber.cpp:154-161 那一大段 std::cout 放进锁里,每次打印都会阻塞
DDS 接收线程。

约束二:回调里绝对不能做重活

subscriber.cpp:51-58 的注释引用了一次真实事故:一次同步日志写入在回调里卡了
934 毫秒。

934ms 是什么后果?回调运行在 dq.user 上 —— 这条线程被卡住,该 participant 的
所有用户 topic 投递全部停摆
。对一个 40Hz 的控制回路,934ms = 丢掉 37 个周期。

这不是”延迟略微增大”,这是控制器彻底失联接近 1 秒。

正确的模式就是 subscriber.cpp:60-80 写的:

1
2
3
4
5
6
7
8
9
10
11
12
auto samples = reader.take();
for (const auto &s : samples)
{
if (!s.info().valid()) continue;
{
std::lock_guard<std::mutex> lock(g_state.mutex);
g_state.latest = s.data(); // 只做一次拷贝
g_state.has_data = true;
}
...
}
// 立刻返回

take → 拷贝 → 立刻返回。 真正的计算交给按自己节拍运行的工作线程。

回调里禁止做的事:

  • 同步文件 I/O、日志落盘
  • 网络请求
  • 加锁等待另一个可能被长期持有的锁
  • 大量内存分配
  • 复杂算法(矩阵运算、路径规划)
  • sleep / 阻塞等待

生产者与消费者节拍解耦

subscriber.cpp:117-124 的注释说明了这个结构的意义:发布端 20Hz、订阅端 main 5Hz,
两者完全解耦,各按自己的节拍跑。

1
2
3
发布端 20Hz  ──write()──>  [DDS]  ──回调──>  g_state.latest
↑
main 循环 5Hz 读快照

这个模式是整个实时系统的骨架。它的价值在于:

  • 上下游频率无需一致,也无需协商
  • 没有队列积压 —— 慢的一方只是看到较少的中间值,不会拖慢快的一方
  • 上游抖动不传导给下游

info().valid() 为什么必须判

subscriber.cpp:66:

1
if (!s.info().valid()) continue;

DDS 投递的样本分两类:

  • 带数据的样本:valid() == true,data() 内容有效
  • 纯元数据样本:valid() == false,表示”实例状态变了” —— 比如实例被
    dispose()、或者写它的所有 writer 都消失了(NOT_ALIVE_NO_WRITERS)。
    这种样本的 data() 内容是未定义的

不判就会把垃圾数据当真实值写进缓存。

有意思的是:这类样本恰好携带了”发布端消失了”这个信息 —— 第 7 节说的失联检测,
线索其实有一半就在这里,只是当前代码把它 continue 掉了。

take() vs read():一个容易踩的坑

两者都是”从 reader 缓存取数据”,区别只在取完之后样本还留不留:

  • take() —— 读取并从 reader 缓存移除。样本只被消费一次
  • read() —— 读取但保留在缓存

还有一个容易忽略的细节:read() 默认只返回 NOT_READ 状态的样本。 样本被读过
一次后会被标记为 READ,下次裸 read() 默认把它过滤掉。想拿到全部(含已读过的)
需要显式指定:

1
2
3
4
5
auto query = dds::sub::status::DataState(
dds::sub::status::SampleState::any(),
dds::sub::status::ViewState::any(),
dds::sub::status::InstanceState::any());
auto samples = reader.select().state(query).read();

实测:回调里用 read() 会怎样

同一份代码只切换 read() / take(),发布端 20Hz 跑 10 秒:

KeepLast(1) 下 —— 两者没有差别:

回调次数 累计取出样本 CPU
take() 198 198 0%
read() 198 198 0%

KeepAll 下 —— 差异才显现:

回调次数 累计取出样本 CPU
take() 398 398 0%
read() 398 34850 0%

两个结论:

  1. 回调次数完全相同(398 次),不会因为缓存非空而反复触发。 回调由”是否有新样本
    到达”驱动,与你是否清空缓存无关。CycloneDDS 不会空转。
  2. CPU 都是 0%,没有忙循环。 但 read() 累计处理了 34850 个样本 —— 同样的旧数据
    被反复处理了约 87 倍。

KeepLast(1) 下两者数字一样,是因为深度只有 1、且新帧覆盖旧帧,几乎不存在
“已读但仍在缓存”的样本可供重复返回。

真实危害是逻辑错误,不是性能爆炸

read() 在回调里的问题是每次回调都重复处理历史样本:

1
2
3
回调#1: 处理 [1]
回调#2: 处理 [1,2] ← 1 又处理了一遍
回调#3: 处理 [1,2,3] ← 1,2 又处理了一遍

主动连续 read() 的实测(KeepAll):

1
2
3
第1次 read() 返回 11 个: 9 10 11 ... 19
第2次 read() 返回 14 个: 9 10 11 ... 22 ← 9-19 重复
第3次 read() 返回 16 个: 9 10 11 ... 24 ← 又重复

所以危害是:如果回调里做的是”累加计数””追加写日志””触发一次动作”,同一帧会被算
好几次
。工作量随缓存深度线性放大(KeepAll 下实测 87 倍),但那是 O(n) 放大,
不是死循环。

结论

回调里仍然应该用 take(),理由是语义匹配而非性能:

  • take() 语义上就是”这批数据我消费了”,天然对应”每帧处理一次”
  • read() 需要你自己额外管理”处理到哪了”(记录 frame_id,或用 SampleState
    精确筛选),容易写出重复处理的 bug

read() 的正当用途是”只想看一眼当前值,不消费”—— 监控线程、调试打印。
如果需求是”多个消费者都要看到全部数据”,正确做法不是用 read(),而是
各建一个 DataReader —— 每个 DataReader 有自己独立的缓存,DDS 会把数据分别投递
给每一个匹配的 reader,互不影响。实测:两个独立 reader 各自 take,都收到全部 77 帧
(交集 77);而两个线程共用同一个 reader 时,A 拿 100 帧、B 拿 55 帧、交集 0 ——
先 take 的把样本移走了。详见第 12 节实验 7。

本例程用的是 take()(subscriber.cpp:60),正确。

进阶模式:独占线程

README.md:157-166 提到一种更彻底的做法:为每个订阅者再起一条独立线程,
on_data_available 里只做「置标志 + notify」就返回,真正的 take() 在那条线程里做。

这是把上面两条约束彻底化:连一次拷贝都不在中间件线程上做,dq.user 的占用降到
接近零。代价是多一条线程和一次上下文切换。

但这个模式有一个经典陷阱,README 里记录了真实踩坑:如果那条线程只在某个
Dispose() 方法里 join,而析构函数是空的,那么:

1
2
// std::thread 的析构语义:如果还 joinable,直接 std::terminate()
~std::thread() { if (joinable()) std::terminate(); }

不是泄漏、不是自动 detach,是立刻终止整个进程。 所以”忘记调用 Dispose()“的
惩罚不是资源泄漏,而是退出时崩溃。

正确做法是在析构函数里调 Stop(),让 RAII 完成清理。把”必须手动调 Dispose”
留成一条没有编译期保障的口头约定,等于给每个新增订阅者的人埋一个坑。


9. 实例与 Key:从消息管道到分布式表

这是本例程刻意没有使用、但 DDS 最有特色的能力。理解它,你才算真正理解
“以数据为中心”。

当前状态:keyless

build/sensor.hpp 里有这一行:

1
2
3
4
template <> constexpr bool TopicTraits<::demo::msg::SensorData>::isKeyless()
{
return true;
}

因为 sensor.idl 里没有任何 @key 标注。这意味着整个 Demo/SensorData topic
只有一个数据实例 —— 所有 write() 都在更新同一个对象。

用第 2 节的白板比喻:白板上只有一个格子,所有人往同一个格子里写。

加上 key 会怎样

如果 IDL 改成这样:

1
2
3
4
5
6
7
struct SensorData {
@key long sensor_id; // ← 加一个 key 字段
Header header;
float temperature;
float gyro[3];
octet status;
};

那么每个不同的 sensor_id 就是一个独立实例,各自维护独立的历史。
KeepLast(1) 变成”为每个实例各保留最新一帧“。

白板从一个格子变成一张表:sensor_id 是主键,每行独立。

实测对照

我做了这个实验:同一个 topic 上跑两个发布端,分别用 sensor_id=1 和
sensor_id=2,订阅端每秒 take() 一次,看能拿到几个样本。

有 key 的版本(@key long sensor_id):

1
2
3
4
5
[SUB] take() 返回 2 个样本: {id=2 frame=9   temp=40} {id=1 frame=9   temp=30}
[SUB] take() 返回 2 个样本: {id=1 frame=29 temp=30} {id=2 frame=29 temp=40}
[SUB] take() 返回 2 个样本: {id=2 frame=49 temp=40} {id=1 frame=49 temp=30}
[SUB] take() 返回 2 个样本: {id=1 frame=69 temp=30} {id=2 frame=69 temp=40}
...(10 次全部返回 2 个样本)

每次稳定拿到 2 个样本,两路数据都在,各自的 frame 都连续。

无 key 的对照版本(同样两个发布端,仅去掉 @key):

1
2
3
4
5
6
[SUB] take() 返回 1 个样本: {id=2 frame=9   temp=40}
[SUB] take() 返回 1 个样本: {id=2 frame=29 temp=40}
[SUB] take() 返回 1 个样本: {id=1 frame=49 temp=30} ← 变成 1 了
[SUB] take() 返回 1 个样本: {id=2 frame=69 temp=40} ← 又变回 2
[SUB] take() 返回 1 个样本: {id=1 frame=129 temp=30}
...

每次只有 1 个样本,id 在 1 和 2 之间随机跳 —— 两个发布端在互相覆盖同一个格子。

可以在生成代码里确认 key 是否生效,build/sensor.cpp:

1
keylist.add_key_endpoint(std::list<uint32_t>{0});   // 有 key 时才有这一行

这个能力什么时候用

任何”同一类数据有多个来源/多个目标”的场景:

场景 key 字段 不用 key 的代价
多个关节的状态 joint_id 开 N 个 topic,或自己拼数组
多路遥控器输入 device_id 互相覆盖,或开 N 个 topic
多个检测到的目标 track_id 无法表达”这个目标消失了”
多台机器人 robot_id 数据混在一起

用 key 比开 N 个 topic 干净得多:topic 数量不随实例数增长,订阅端一次 take()
就能拿到所有实例的最新值,而且中间件会告诉你”某个实例消失了”
(就是第 8 节说的那种 valid() == false 的样本)。

这就是”分布式表”的含义:DDS 维护的不是一根管道,而是一张所有进程都能看到的、
自动同步的表。key 是主键,QoS 决定每行保留多少版本。


10. 测量延迟:一个测错了的数字

这一节讲一个具体的测量错误,因为它揭示的原理在实时系统里到处都是。

现象

subscriber.cpp:146-148 计算了一个”端到端延迟”:

1
2
3
const int64_t now = ...;  // 当前时间
const double latency_us =
static_cast<double>(now - snapshot.header().timestamp()) / 1000.0;

日志里这个数字是 44000μs(44 毫秒)左右,而且在逐渐变小:

1
2
latency=47080  →  46414  →  45080  →  ...  →  2917  →  1783  →  【50730】  →  49480  →  ...
↑ 跳回

我跑了 75 秒,发现它不是”递减后稳定”,而是一路降到 1.8ms,然后瞬间跳回 50.7ms,
再次递减
—— 一个周期约 9 秒的锯齿波。

跳变点的现场:

1
2
frame=177 (+4) ... latency=1783.32us
frame=180 (+3) ... latency=50730.6us ← gap 从 +4 变成 +3,latency 跳 +49ms

75 秒内 gap 分布:367 次 (+4)、7 次 (+3)、1 次 (+0)。7 个 +3 恰好对应
7 个锯齿周期。

真实机制:拍频(相位漂移)

一个自然的猜测是”积压”:发布快、消费慢,所以读到的总是积压了一会儿的旧帧,
随着追赶延迟逐渐缩小。

但这个解释站不住,有两个反证:

  1. KeepLast(1) 的历史深度是 1,物理上无法积压
  2. 日志里 rx_total=169 恰好等于 frame=169 —— 一帧没丢、一帧没积压

真实机制是这样:

g_state.latest 被回调以 20Hz(每 50ms)覆写,main 以 5Hz(每 200ms)采样它。所以

1
2
latency = (main 采样时刻) − (最近一次覆写的那帧的发布时刻)
= 相位差 + 真实传输延迟

相位差的取值范围是 [0, 50ms) —— 取决于 main 的采样点恰好落在两次覆写之间的
哪个位置。这就是那个 44ms 的主要成分。

而两端的周期来自两个独立的、都不精确的时钟:sleep_for(50ms) 实际略大于 50ms,
sleep_for(200ms) 实际略大于 200ms,且两者不成精确的 4 倍关系。因此
T_sub 与 4 × T_pub 有一个微小差值 δ(实测约 1ms)。

于是每个订阅周期,相位就漂移 δ:

1
2
3
4
5
6
7
相位单调漂移  →  latency 单调递减
↓
减到接近 0 时,下一次采样落到了"上一帧还没被覆盖"的位置
↓
只前进 3 帧(gap +3),相位回卷一整个发布周期(+50ms)
↓
锯齿的垂直边

锯齿周期 = 50ms ÷ 1ms ≈ 50 个订阅周期 = 10 秒。实测 45 个采样点后跳变(9 秒),吻合。

正确的测法:在回调里打点

我在回调里加了打点(数据刚到达的时刻),和 main 侧对比:

回调侧(真实传输延迟),499 个样本:

1
min=130.9μs   median=282.5μs   p95=386.2μs   max=464.9μs

main 侧(同一次运行):

1
46637 → 45749 → 44890 → 44205 → ... (μs,锯齿递减)

中位 282 微秒 vs 44 毫秒 —— 差了 156 倍,两个数量级。

本机同进程间的 DDS 传输延迟实际是 0.3 毫秒量级(而且这 282μs 里还含一次
std::cout)。日志里那个 44ms 的绝大部分是订阅端自己的采样相位。

三层教训

1. 这个数值本身没有意义。 它在 1.8ms 到 50.7ms 之间锯齿摆动,测的主要是订阅端
自己的调度相位。真实传输延迟藏在锯齿的下界附近。

2. 想量传输延迟,必须在回调里打点。 那里才是”数据刚到达”的时刻。main 循环里
已经掺进了最多 200ms 的等待。

3. 一般性原理:任何”周期性采样另一个周期性信号”的测量,都会得到拍频,而不是被测
量的真值。

这在实时系统里到处都是:

1
2
500Hz 的硬件状态  被  40Hz 的控制器采样   → 拍频
100Hz 的控制指令 被 20Hz 的下发线程采样 → 拍频

每一处都有一个锯齿周期。所以延迟指标必须明确”从哪到哪“,否则量到的是自己的
调度周期。正确做法是把链路拆成”同一线程内的两个时刻”这样的可测段,分段打点。

顺带:跨机器时这个减法根本不成立

本例程两端在同一台机器、都用 system_clock,所以时间戳可以直接相减。

跨板子/跨机器时不行:两台机器的 system_clock 之间有 NTP/PTP 偏差,量级可能
远超你要测的延迟本身,甚至算出负数延迟。

跨机器测延迟必须先解决时钟同步(PTP,精度可到微秒级),或者改用”回环往返时间 ÷ 2”
的测法。


11. 静默失效:DDS 最反直觉的地方

到这里可以总结 DDS 最需要适应的一点:

绝大多数配置错误都不会报错,只会静默地什么都不发生。

原因在第 6 节讲过:match 的三个条件(topic 名、类型名、QoS)任何一条不满足,
中间件就是不建立通路。它不认为这是错误 —— 在一个动态的分布式系统里,”某个端点和我
不匹配”是完全正常的状态。

对比一下你熟悉的错误处理:

出错时
TCP connect 到错误端口 立刻 ECONNREFUSED
HTTP 请求错误路径 返回 404
调用不存在的 RPC 方法 返回 UNIMPLEMENTED
DDS QoS 不匹配 什么都不发生

怎么主动发现”没连上”

不要靠”等数据看有没有来”,DDS 提供了直接查询的接口:

1
2
3
4
5
// 发布端:有几个订阅者匹配上了?
writer.publication_matched_status().current_count()

// 订阅端:有几个发布者匹配上了?
reader.subscription_matched_status().current_count()

这是我在第 6 节测 match 耗时用的方法。建议在任何 DDS 程序的启动自检里加上它 ——
把”静默失效”变成一条明确的日志:

1
2
3
4
5
6
// 起来后等一小会儿,检查是否匹配上
std::this_thread::sleep_for(std::chrono::milliseconds(500));
if (reader.subscription_matched_status().current_count() == 0) {
std::cerr << "[ERROR] 500ms 内没有匹配到任何发布者!"
<< " 检查:topic 名 / 类型名 / QoS / domain ID / 多播\n";
}

第 6 节实测 match 只要 2-10ms,所以 500ms 的窗口足够宽松了。

还可以订阅对应的状态回调(本例程没用,只订了 data_available):

1
2
mask |= dds::core::status::StatusMask::subscription_matched();
// 然后覆写 on_subscription_matched(),match 变化时会被调用

观察发现过程的工具

CycloneDDS 自带诊断手段:

1
2
3
4
5
6
# 打开中间件日志(会打印发现、match、QoS 不兼容的细节)
export CYCLONEDDS_URI='<Tracing><Verbosity>finest</Verbosity><OutputFile>dds.log</OutputFile></Tracing>'
./build/subscriber

# 抓包看 RTPS 报文(Wireshark 有 RTPS 解析器)
sudo tcpdump -i any -n 'udp port 7400 or udp port 7401'

Verbosity 设成 finest 后,QoS 不兼容会在日志里留下痕迹 —— 这是排查静默失效最
直接的手段。

一个隐蔽的坑:多个发布者

scripts/run.sh:34-39 记录了一个真实踩坑,值得单独讲,因为它是 DDS 范式的直接后果:

孤儿 publisher 的后果很隐蔽:它和新起的 publisher 在同一个 topic 上同时是 writer,
订阅端会交替收到两者的数据,frame_id 在两个序列间来回跳。看起来像”DDS 乱序/串数据”,
实际是多了一个发布者。DDS 本身不限制一个 topic 有多少 writer,所以不会报错。

DDS 的 topic 是多对多的:N 个 writer、M 个 reader 都能挂在同一个 topic 上,
这是设计特性而非漏洞。所以”上一轮进程没清干净”会表现为数据看起来乱序 ——
一个极难定位的现象,因为你会先去怀疑中间件。

(第 9 节的 keyless 对照实验其实就是这个现象的受控复现:两个发布端互相覆盖,
id 在 1 和 2 之间随机跳。)

DDS 对此提供了 Ownership QoS:设成 EXCLUSIVE 后,同一实例同时只有
ownership_strength 最高的 writer 能生效,其余被忽略。对”同一路数据只应有一个权威
来源”的场景(几乎所有控制指令都是),这比靠脚本清进程可靠得多。本例程用的是默认的
SHARED。

版本演进的坑:可扩展性

build/sensor.cpp 里每个字段都标着 extensibility::ext_final:

1
2
props.push_back(entity_properties_t(0, 0, false, bb_unset, extensibility::ext_final));  //root
props.push_back(...get_bit_bound<float>(), extensibility::ext_final, false)); //::temperature

这是因为 IDL 里没写 @final / @appendable / @mutable,idlc 隐式取了 final
(CMakeLists.txt:30 那个 WARNINGS no-implicit-extensibility 就是在消这个告警)。

三种可扩展性的差别很大:

标注 语义 加字段后
@final 布局固定 旧版本收不到新版本的数据
@appendable 可在尾部追加字段 新旧版本可互通,旧版忽略新字段
@mutable 字段带 ID,可任意增删改序 最灵活,开销最大

当前是 final,意味着:给 SensorData 加一个字段,就必须同时重新编译并同时部署
所有收发两端。

对一个跨板子的系统(比如主控 + 底层板),这是很硬的运维约束 —— 如果底层软件是别人
编译的二进制,你没法要求它跟你同步升级。

如果 IDL 有演进需求,@appendable 是应该认真考虑的默认选择。 这是我在这份代码里
看到的、影响最深远但最不显眼的一个技术决策。


12. 速查表与动手实验

静默失效排查清单

按从最常见到最少见排列:

# 症状 原因 定位方法
1 完全收不到 QoS 不兼容(Reader 请求 > Writer 提供) 逐条对比 Reliability / Durability / Deadline
2 完全收不到 类型名不一致(module 名、struct 名) 对比 getTypeName() 返回的字符串
3 完全收不到 topic 名笔误、大小写不符 字符串逐字节比对
4 完全收不到 domain ID 不同 检查 DomainParticipant(N) 的 N
5 完全收不到 多播不通(防火墙 / Docker 桥接 / 跨网段) ss -unap 看 7400 端口;抓 RTPS 包
6 加字段后旧端收不到 @final 可扩展性 检查 IDL 有无 @appendable
7 数据看起来乱序 同 topic 有多个 writer(孤儿进程) pgrep -f;或用 Ownership::EXCLUSIVE
8 收到垃圾值 没判 info().valid() 检查回调里的过滤
9 该 topic 后续全部卡住 回调里做了重活/阻塞 在回调首尾打时间戳
10 晚起的订阅端收不到当前值 Durability 默认 Volatile Writer 改 TransientLocal
11 上游挂了下游还在动 状态语义 + 无 Liveliness/Deadline 加 Deadline QoS,或应用层强制发零
12 延迟数字奇怪/周期性摆动 跨线程周期采样的拍频 改在回调里打点
13 同一帧被处理多次(计数偏大、日志重复) 回调里用了 read() 而非 take() 改成 take();read() 需自己管理已处理位置
14 退出时崩溃 std::thread 析构时仍 joinable 析构函数里 join

核心概念一句话总结

概念 一句话
DCPS 传的不是消息,是”某块数据的当前值”
Domain 硬隔离边界,不同 domain 互相完全不可见
Topic 名字 + 类型名,两者都必须一致才能通信
Participant 重量级,8 条线程,一个进程只建一个
DataWriter/Reader 唯一有数据流经过的层,QoS 的宿主
发现 SPDP(多播找进程)→ SEDP(单播换清单)→ match
RxO Reader 请求的服务等级不能高于 Writer 提供的
BestEffort 状态量用它:丢了就丢,重传过期数据无意义
KeepLast(1) 只留最新值;生产快于消费时中间帧被丢弃
回调线程 不是你的线程;必须加锁、必须快速返回
take vs read take() 取走、read() 留下;回调里用 take(),否则同一帧会被重复处理
Key 加了它,topic 从”一个格子”变成”一张表”
静默失效 DDS 不匹配时什么都不发生,用 matched_status() 主动查

建议动手做的实验

按难度递增。前 4 个例程注释里已经提到,后 4 个是本文实测过的。

实验 1:看回调线程(1 分钟)

跑 ./scripts/run.sh both,对比两个 thread id。这是第 8 节全部内容的起点。

实验 2:QoS 不匹配(5 分钟)

把 subscriber.cpp:102 的 Reliability::BestEffort() 改成 Reliable(),重新构建运行。

预期:topic 名和类型完全没变,但订阅端全程静默、零报错。实测 12 秒内
no data yet ×63、收到 0 帧。

实验 3:类型名参与匹配(5 分钟)

把 idl/sensor.idl 里的 module demo 改成 module foo,只重新构建一端。

预期:topic 名一样但两端匹配不上。

实验 4:看丢帧(1 分钟)

观察 frame 后面的 (+N),稳定是 +4。把 subscriber.cpp:164 的
sleep_for(200ms) 改成 50ms,+N 会变成 +1。

实验 5:在回调里量真实延迟(15 分钟)★推荐

在 on_data_available 里加打点:

1
2
3
4
5
const int64_t now_cb = std::chrono::duration_cast<std::chrono::nanoseconds>(
std::chrono::system_clock::now().time_since_epoch()).count();
const double cb_us = static_cast<double>(now_cb - s.data().header().timestamp()) / 1000.0;
std::cout << "[CB] frame=" << s.data().header().frame_id()
<< " cb_latency=" << cb_us << "us\n";

预期:回调侧是稳定的几百微秒(实测中位 282μs),main 侧是 1.8-50ms 的锯齿。
这一个实验能把第 10 节讲的拍频彻底坐实。

实验 6:加 key 看实例(30 分钟)★最能体现 DDS 特色

给 sensor.idl 的 SensorData 加一个 @key long sensor_id;,让 publisher 从命令行
读 sensor_id,同时跑两个不同 id 的发布端,订阅端每秒 take() 一次打印样本数。

预期:take() 稳定返回 2 个样本。去掉 @key 重跑,变成每次 1 个、id 随机跳。

实验 7:两个消费者要不要各建 reader(20 分钟)★常见误区

分别试两种写法,各跑两个线程消费:

  • 两个独立 DataReader,各 take 自己的 → 实测两者都收到全部 77 帧,交集 77
  • 共用一个 DataReader,两线程都 take → 实测线程 A 拿 100 帧、B 拿 55 帧,交集 0

结论:要让多个消费者都看到全量数据,必须各建一个 DataReader。共用一个 reader 时
先 take 的会把样本移走。

实验 8:验证可扩展性(30 分钟)

给 SensorData 尾部加一个字段,只重新编译一端。@final 下会静默失配;
在 struct 前加 @appendable 后就能互通(旧端忽略新字段)。

这个实验对跨板子协作的场景最有参考价值。


附录一:本文实测数据汇总

环境:x86_64 / Ubuntu 24.04 / CycloneDDS 0.10.5 / 同机双进程

测量项 结果
Participant 线程数 8(main + gc + dq.builtins + dq.user + tev + recv + recvMC + recvUC)
Participant UDP socket 数 4(7400 多播、7401 多播、2 个单播)
DomainParticipant 构造耗时 约 2.2-2.8ms
DataWriter 构造耗时 约 0.6ms
match 耗时(订阅端已在) 2.16 / 2.24 / 2.24ms
match 耗时(两端同时起) 2.31 / 2.31 / 2.32ms
首帧 frame 号 frame=0(一帧未丢)
回调侧传输延迟(n=499) min 130.9μs / 中位 282.5μs / p95 386.2μs / max 464.9μs
main 侧”延迟”(同一次运行) 1.8-50.7ms 锯齿波,周期约 9 秒
gap 分布(75 秒) +4 ×367、+3 ×7、+0 ×1
QoS 失配(Reader 请求 Reliable) 12 秒内收到 0 帧,no data yet ×63,零报错
有 key + 2 个发布端 take() 稳定返回 2 个样本
无 key + 2 个发布端 take() 返回 1 个样本,id 随机跳
两个独立 DataReader 各自 take 各收到全部 77 帧,交集 77(互不影响)
共用一个 DataReader,两线程 take A 100 帧 / B 55 帧,交集 0(互相抢走)
回调 take() vs read()(KeepLast(1)) 回调各 198 次,取出样本各 198 个,CPU 均 0%
回调 take() vs read()(KeepAll) 回调各 398 次;取出样本 398 vs 34850;CPU 均 0%
write() 到无订阅端(Volatile) frame 0..4 未送达,晚到的订阅端只收到之后的帧
write() 到无订阅端(TransientLocal) 补发 1 帧(keyless 单实例只留最新)

附录二:完整源代码

例程全部源码,便于对照正文的行号引用。文件按阅读顺序排列。

idl/sensor.idl(28 行)

数据类型定义。全部通信内容的唯一来源,idlc 据此生成 C++ 类型。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
// 最小示例数据类型。
//
// 刻意模仿 sentigent_soc 的写法,方便对照:
// - module 两层嵌套(demo::msg),对应项目里的 hal::msg / neuro_control::msg
// - 带一个 Header 子结构(时间戳 + 自增帧号),项目里每条消息都有
// - 定长数组字段,DDS 里定长数组是值语义,不涉及堆分配
//
// 类型名(demo::msg::SensorData)是 DDS 匹配的依据之一:topic 名相同但
// 类型名不同,两端不会 match —— 这就是 rc_input.idl 里那段警告说的坑。

module demo {
module msg {

struct Header {
long long timestamp; // 纳秒
long long frame_id; // 自增帧号,用来发现丢帧
};

// topic: "Demo/SensorData"
struct SensorData {
Header header;
float temperature; // ℃
float gyro[3]; // rad/s,定长数组
octet status; // 0=正常 1=告警
};

};
};

publisher.cpp(79 行)

发布端。四层实体初始化 + 20Hz 定周期发布循环。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
// 发布端:20Hz 发布 Demo/SensorData。
//
// 对应项目里的 srt::common::Publisher(common/include/dds/publisher.hpp)。
// 那一层封装做的事就是这里的 participant → topic → publisher → writer 四步,
// 外加一个 Dispose()。

#include <dds/dds.hpp>
#include "sensor.hpp" // idlc 生成

#include <chrono>
#include <cmath>
#include <iostream>
#include <thread>

int64_t nowNs()
{
return std::chrono::duration_cast<std::chrono::nanoseconds>(
std::chrono::system_clock::now().time_since_epoch()).count();
}

int main()
{
// ── DDS 四层对象 ────────────────────────────────────────────────────
// domain 0 = 默认域。同域内的进程才能互相发现(项目里统一 Init(0))。
dds::domain::DomainParticipant participant(0);

// Topic = 名字 + 类型 的绑定。两端必须【名字和类型都一致】才 match。
dds::topic::Topic<demo::msg::SensorData> topic(participant, "Demo/SensorData");

dds::pub::Publisher publisher(participant);

// QoS:这里显式设成 BestEffort + KeepLast(1),即"只要最新一帧,丢了不重传"。
// 这正是项目里 RC 遥控器那一路的配置(senti_controller.cpp:108)。
// 高频实时流适合这个:重传一帧过期数据毫无意义,反而挤占带宽。
//
// ★ 订阅端必须设成【兼容】的 QoS,否则 topic 名和类型都对也不会 match。
// Reliability 的兼容规则:发布端 BestEffort 只能配 BestEffort 订阅端;
// 发布端 Reliable 则两种订阅端都能配。
dds::pub::qos::DataWriterQos qos;
qos << dds::core::policy::Reliability::BestEffort()
<< dds::core::policy::History::KeepLast(1);

dds::pub::DataWriter<demo::msg::SensorData> writer(publisher, topic, qos);

std::cout << "[PUB] writer ready on Demo/SensorData, 20Hz. Ctrl-C to stop.\n";

// ── 20Hz 发布循环 ───────────────────────────────────────────────────
// 注意这是"定周期发布",不是"变化时发布"。整个运控链路都是这个模式:
// 每一级按自己的节拍重放最新值。好处是没有队列积压,代价是
// 【上游停发不会自动归零】—— 下游会一直用最后一帧,所以每个环节
// 都需要显式的"必须发零"约定(见项目里 RC 刹车 / latchVelocityZero)。
int64_t frame = 0;
while (true)
{
demo::msg::SensorData sample;
sample.header().timestamp() = nowNs();
sample.header().frame_id() = frame;

// 造点会变化的数据,方便观察
const float t = static_cast<float>(frame) * 0.05f;
sample.temperature() = 25.0f + 5.0f * std::sin(t);
sample.gyro()[0] = std::sin(t);
sample.gyro()[1] = std::cos(t);
sample.gyro()[2] = 0.0f;
sample.status() = (sample.temperature() > 29.0f) ? 1 : 0;

writer.write(sample);

if (frame % 20 == 0) // ~1Hz 打印
{
std::cout << "[PUB] frame=" << frame
<< " temp=" << sample.temperature()
<< " status=" << static_cast<int>(sample.status()) << "\n";
}

++frame;
std::this_thread::sleep_for(std::chrono::milliseconds(50));
}
}

subscriber.cpp(166 行)

订阅端。回调只缓存、main 按 5Hz 自己的节拍消费。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
// 订阅端:演示 DDS 回调线程,以及项目里"回调只缓存、主循环按自己节拍消费"的模式。
//
// 运行后观察打印的 thread id:
// [SUB][main] main 线程
// [SUB][callback] on_data_available 所在线程 —— 【不是】main
// 这就是"DDS 回调线程"的含义:你注册的函数由中间件的线程调起。
//
// 对应项目里的 srt::common::Subscriber + DdsDataReaderListener
// (common/include/dds/subscriber.hpp、dds_data_reader_listener.hpp)。

#include <dds/dds.hpp>
#include "sensor.hpp" // idlc 生成

#include <atomic>
#include <chrono>
#include <iostream>
#include <mutex>
#include <sstream>
#include <thread>

// ── 共享缓存:回调线程写,main 线程读 ──────────────────────────────────
// 这是整个例程的重点。回调和主循环在不同线程,所以【必须加锁】。
// 项目里 ControllerSelector::action_mutex_ 保护三路 CommandIntent 缓存,
// 就是同一件事(controller_selector.h:95)。
struct SharedState
{
std::mutex mutex;
demo::msg::SensorData latest;
bool has_data = false;

// 计数用 atomic 就够了,不需要锁
std::atomic<int64_t> rx_count{0};
};

SharedState g_state;

std::string threadIdStr()
{
std::ostringstream oss;
oss << std::this_thread::get_id();
return oss.str();
}

// ── Listener:DDS 数据到达时被中间件调用 ───────────────────────────────
// 继承 NoOpDataReaderListener 只覆写 on_data_available,其余状态回调用默认空实现。
class SensorListener : public dds::sub::NoOpDataReaderListener<demo::msg::SensorData>
{
public:
void on_data_available(dds::sub::DataReader<demo::msg::SensorData> &reader) override
{
// ★ 这里运行在 DDS 的接收线程上,不是 main。
//
// 因此这个函数里【不能做重活、不能阻塞】:卡住它会拖住该 topic 后续
// 数据的接收。项目 neuro_control/src/main.cpp:33 记了一次真实事故 ——
// 一次同步日志写入卡了 934ms,直接影响控制环。
//
// 正确做法就是下面这样:take → 拷进缓存 → 立刻返回。
// 真正的计算留给按自己节拍跑的工作线程(项目里是 40Hz 的 NC_RL)。

auto samples = reader.take();
for (const auto &s : samples)
{
// info().valid() 必须判:DDS 会投递"实例已消失"这类元数据样本,
// 那种样本的 data() 内容无意义。项目里同样有这个判断
// (dds_data_reader_listener.hpp:111)。
if (!s.info().valid()) continue;

{
std::lock_guard<std::mutex> lock(g_state.mutex);
g_state.latest = s.data();
g_state.has_data = true;
}

const int64_t n = ++g_state.rx_count;
if (n == 1)
{
std::cout << "[SUB][callback] FIRST sample, thread=" << threadIdStr()
<< " <-- 注意这个 id 和 main 不同\n";
}
}
}
};

int main()
{
std::cout << "[SUB][main] thread=" << threadIdStr() << "\n";

dds::domain::DomainParticipant participant(0);

// topic 名和类型都必须与发布端完全一致
dds::topic::Topic<demo::msg::SensorData> topic(participant, "Demo/SensorData");

dds::sub::Subscriber subscriber(participant);

// ★ QoS 必须与发布端兼容。发布端是 BestEffort + KeepLast(1),
// 这里也设成 BestEffort。若这里改成 Reliable(),则 topic 名和类型
// 都对,但【收不到任何数据】—— 这是 DDS 最常见的排查陷阱,
// 项目 rc_input.idl 的注释专门警告过。
// 想自己验证:把下面 BestEffort() 换成 Reliable() 重编译,会看到
// 订阅端一直静默。
dds::sub::qos::DataReaderQos qos;
qos << dds::core::policy::Reliability::BestEffort()
<< dds::core::policy::History::KeepLast(1);

SensorListener listener;

// 只关心 data_available 这一个状态位。项目里也是这么设的
// (subscriber.hpp:80)。
dds::core::status::StatusMask mask = dds::core::status::StatusMask::none();
mask |= dds::core::status::StatusMask::data_available();

dds::sub::DataReader<demo::msg::SensorData> reader(subscriber, topic, qos);
reader.listener(&listener, mask);

std::cout << "[SUB] reader ready, waiting for data...\n";

// ── main 循环:按【自己的】5Hz 节拍消费缓存 ─────────────────────────
// 发布端 20Hz,这里 5Hz —— 两者解耦,各按自己节拍跑。
// 这正是项目的核心结构:
// HalRobotState 500Hz 到达 → NC 回调只缓存 → NC_RL 40Hz 取用
// 手柄 HID 事件驱动到达 → reader 线程写 state_ → 控制线程 100Hz 快照
//
// 消费频率低于生产频率时,中间的帧就被丢弃了(KeepLast(1) 只留最新)。
// 观察 frame_id 的跳跃就能看到:每次约 +4。
int64_t prev_frame = -1;
while (true)
{
demo::msg::SensorData snapshot;
bool valid = false;
{
std::lock_guard<std::mutex> lock(g_state.mutex);
if (g_state.has_data)
{
snapshot = g_state.latest; // 拷一份出来,尽快放锁
valid = true;
}
}

if (!valid)
{
std::cout << "[SUB][main] no data yet (发布端起了吗?QoS 对得上吗?)\n";
}
else
{
// 端到端延迟:发布端打的时间戳 vs 现在。同机通常几百 us。
const int64_t now = std::chrono::duration_cast<std::chrono::nanoseconds>(
std::chrono::system_clock::now().time_since_epoch()).count();
const double latency_us = static_cast<double>(now - snapshot.header().timestamp()) / 1000.0;

const int64_t f = snapshot.header().frame_id();
const int64_t gap = (prev_frame < 0) ? 0 : (f - prev_frame);
prev_frame = f;

std::cout << "[SUB][main] frame=" << f
<< " (+" << gap << ")"
<< " temp=" << snapshot.temperature()
<< " gyro=[" << snapshot.gyro()[0] << "," << snapshot.gyro()[1]
<< "," << snapshot.gyro()[2] << "]"
<< " status=" << static_cast<int>(snapshot.status())
<< " | rx_total=" << g_state.rx_count.load()
<< " latency=" << latency_us << "us\n";
}

std::this_thread::sleep_for(std::chrono::milliseconds(200)); // 5Hz
}
}

CMakeLists.txt(113 行)

构建配置。两条依赖解析路径:标准安装(路 A)/ 预编译库兜底(路 B)。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
cmake_minimum_required(VERSION 3.16)
project(dds_test CXX)

set(CMAKE_CXX_STANDARD 17)
set(CMAKE_CXX_STANDARD_REQUIRED ON)
if(NOT CMAKE_BUILD_TYPE)
set(CMAKE_BUILD_TYPE Debug)
endif()

# ─────────────────────────────────────────────────────────────────────────────
# 依赖解析:两条路,优先用标准安装
#
# 路 A(推荐):系统里装好了 CycloneDDS + CycloneDDS-CXX 的 CMake 包。
# 直接用官方的 idlcxx_generate(),最干净。
# 路 B(兜底):只有 sentigent_soc/third_party/dds 里的预编译库 +
# 一个能用的 idlc/idlcxx。手动调 idlc 生成。
#
# README.md 里写了怎么装。
# ─────────────────────────────────────────────────────────────────────────────
find_package(CycloneDDS-CXX QUIET)

if(CycloneDDS-CXX_FOUND)
message(STATUS "[dds_test] 使用系统安装的 CycloneDDS-CXX (路 A)")

# 官方提供的 CMake 函数:把 .idl 编译成 C++ 类型,产出一个可链接的 target。
# WARNINGS no-implicit-extensibility 只是消掉一个与本例无关的告警。
idlcxx_generate(
TARGET sensor_msg
FILES ${CMAKE_CURRENT_SOURCE_DIR}/idl/sensor.idl
WARNINGS no-implicit-extensibility
)
set(DDS_LIBS CycloneDDS-CXX::ddscxx sensor_msg)
set(DDS_INCLUDES "")

else()
message(STATUS "[dds_test] 未找到 CycloneDDS-CXX 的 CMake 包,回退到预编译库 (路 B)")

# 复用 sentigent_soc 自带的预编译库。可用 -DSOC_ROOT=... 覆盖。
set(SOC_ROOT "/home/hucw/workspace/sentigent_soc" CACHE PATH "sentigent_soc 根目录")
set(DDS_PREBUILT "${SOC_ROOT}/third_party/dds")

if(NOT EXISTS "${DDS_PREBUILT}/include/ddscxx/dds/dds.hpp")
message(FATAL_ERROR
"找不到预编译 DDS:${DDS_PREBUILT}\n"
"请装 CycloneDDS-CXX(见 README.md),或用 -DSOC_ROOT=<路径> 指定。")
endif()

# 选架构目录
if(CMAKE_SYSTEM_PROCESSOR MATCHES "aarch64|arm64")
set(DDS_LIB_DIR "${DDS_PREBUILT}/lib/aarch64")
else()
set(DDS_LIB_DIR "${DDS_PREBUILT}/lib/x86_64")
endif()

# 找 idlc。装了 cyclonedds-tools 就在 PATH 里;源码装的可能在 ~/.local/bin。
find_program(IDLC_EXECUTABLE idlc
HINTS ENV HOME PATH_SUFFIXES .local/bin
PATHS /usr/bin /usr/local/bin $ENV{HOME}/.local/bin)
if(NOT IDLC_EXECUTABLE)
message(FATAL_ERROR "找不到 idlc。见 README.md「安装 idlc」一节。")
endif()

# 找 idlcxx 后端插件(生成 C++ 而非 C 的关键;apt 的 cyclonedds-tools 不含它)。
find_library(IDLCXX_PLUGIN NAMES cycloneddsidlcxx
PATHS /usr/lib /usr/local/lib $ENV{HOME}/.local/lib
PATH_SUFFIXES x86_64-linux-gnu aarch64-linux-gnu)
if(NOT IDLCXX_PLUGIN)
message(FATAL_ERROR
"找到了 idlc 但缺 C++ 后端 libcycloneddsidlcxx。\n"
"apt 的 cyclonedds-tools 只有 C 后端,必须装 cyclonedds-cxx。见 README.md。")
endif()
message(STATUS "[dds_test] idlc=${IDLC_EXECUTABLE} plugin=${IDLCXX_PLUGIN}")

# 手动跑 idlc:sensor.idl → sensor.hpp + sensor.cpp
set(GEN_DIR "${CMAKE_CURRENT_BINARY_DIR}/generated")
file(MAKE_DIRECTORY ${GEN_DIR})
add_custom_command(
OUTPUT ${GEN_DIR}/sensor.hpp ${GEN_DIR}/sensor.cpp
COMMAND ${IDLC_EXECUTABLE} -l ${IDLCXX_PLUGIN} -o ${GEN_DIR}
${CMAKE_CURRENT_SOURCE_DIR}/idl/sensor.idl
DEPENDS ${CMAKE_CURRENT_SOURCE_DIR}/idl/sensor.idl
COMMENT "[dds_test] idlc: sensor.idl -> sensor.hpp/.cpp"
VERBATIM)

add_library(sensor_msg STATIC ${GEN_DIR}/sensor.cpp)
target_include_directories(sensor_msg PUBLIC
${GEN_DIR}
${DDS_PREBUILT}/include/ddscxx
${DDS_PREBUILT}/include)
target_link_directories(sensor_msg PUBLIC ${DDS_LIB_DIR})
target_link_libraries(sensor_msg PUBLIC ddscxx ddsc)

# rpath:让可执行文件运行时能找到预编译的 .so,免得每次手动设
# LD_LIBRARY_PATH。
set(CMAKE_BUILD_RPATH "${DDS_LIB_DIR}")

set(DDS_LIBS sensor_msg)
set(DDS_INCLUDES ${GEN_DIR})
endif()

# ─────────────────────────────────────────────────────────────────────────────
# 两个可执行文件
# ─────────────────────────────────────────────────────────────────────────────
add_executable(publisher publisher.cpp)
add_executable(subscriber subscriber.cpp)

foreach(tgt publisher subscriber)
target_link_libraries(${tgt} PRIVATE ${DDS_LIBS})
if(DDS_INCLUDES)
target_include_directories(${tgt} PRIVATE ${DDS_INCLUDES})
endif()
target_compile_options(${tgt} PRIVATE -Wall -Wextra)
endforeach()

scripts/build.sh(25 行)

构建脚本。首次会自动跑 idlc 生成类型代码。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
#!/usr/bin/env bash
# 构建 DDS 最小例程。
# ./scripts/build.sh 增量构建
# ./scripts/build.sh clean 先删 build 再全量构建
set -euo pipefail

ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
BUILD_DIR="$ROOT/build"

if [[ "${1:-}" == "clean" ]]; then
echo "[build] 清理 $BUILD_DIR"
rm -rf "$BUILD_DIR"
fi

mkdir -p "$ROOT/logs"

cmake -S "$ROOT" -B "$BUILD_DIR" -DCMAKE_BUILD_TYPE=Debug
cmake --build "$BUILD_DIR" -j"$(nproc)"

echo
echo "[build] 完成。可执行文件:"
echo " $BUILD_DIR/publisher"
echo " $BUILD_DIR/subscriber"
echo
echo "下一步: ./scripts/run.sh both"

scripts/run.sh(83 行)

运行脚本。含孤儿进程清理与行缓冲处理。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
#!/usr/bin/env bash
# 运行 DDS 最小例程。日志统一进 logs/。
#
# ./scripts/run.sh pub 只跑发布端(前台)
# ./scripts/run.sh sub 只跑订阅端(前台)
# ./scripts/run.sh both 两个都后台跑,实时跟随订阅端日志(推荐先用这个)
# ./scripts/run.sh stop 停掉后台进程
set -euo pipefail

ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
BUILD_DIR="$ROOT/build"
LOG_DIR="$ROOT/logs"
mkdir -p "$LOG_DIR"

if [[ ! -x "$BUILD_DIR/publisher" || ! -x "$BUILD_DIR/subscriber" ]]; then
echo "[run] 还没构建,先跑 ./scripts/build.sh" >&2
exit 1
fi

stop_all() {
for name in publisher subscriber; do
pidfile="$LOG_DIR/$name.pid"
if [[ -f "$pidfile" ]]; then
pid="$(cat "$pidfile")"
if kill -0 "$pid" 2>/dev/null; then
kill "$pid" 2>/dev/null || true
echo "[run] 已停止 $name (pid=$pid)"
fi
rm -f "$pidfile"
fi
done

# 兜底:按可执行文件全路径扫一遍漏网的。
#
# 这一步不是多余的。pid 文件记错、SIGTERM 没生效、或上一轮脚本被 Ctrl-C
# 打断,都会留下孤儿进程 —— 而孤儿 publisher 的后果很隐蔽:它和新起的
# publisher 在【同一个 topic 上同时是 writer】,订阅端会交替收到两者的
# 数据,frame_id 在两个序列间来回跳。看起来像"DDS 乱序/串数据",实际是
# 多了一个发布者。DDS 本身不限制一个 topic 有多少 writer,所以不会报错。
local leftover
leftover="$(pgrep -f "^$BUILD_DIR/(publisher|subscriber)$" 2>/dev/null || true)"
if [[ -n "$leftover" ]]; then
echo "[run] 发现孤儿进程,强制清理: $leftover"
# shellcheck disable=SC2086
kill -9 $leftover 2>/dev/null || true
fi
}

case "${1:-both}" in
pub)
exec stdbuf -oL "$BUILD_DIR/publisher" 2>&1 | tee "$LOG_DIR/publisher.log"
;;
sub)
exec stdbuf -oL "$BUILD_DIR/subscriber" 2>&1 | tee "$LOG_DIR/subscriber.log"
;;
both)
stop_all
# 先起订阅端:DDS 的发现是双向的,谁先起都行,但订阅端先起能
# 观察到第一帧,便于看清 match 过程。
#
# stdbuf -oL 是必须的:stdout 重定向到文件时 libc 默认块缓冲(4KB),
# 进程被 kill 时缓冲区内容直接丢失 —— 日志会是空的。行缓冲后每行
# 立即落盘,tail -f 也才能实时跟随。
stdbuf -oL "$BUILD_DIR/subscriber" > "$LOG_DIR/subscriber.log" 2>&1 &
echo $! > "$LOG_DIR/subscriber.pid"
sleep 0.3
stdbuf -oL "$BUILD_DIR/publisher" > "$LOG_DIR/publisher.log" 2>&1 &
echo $! > "$LOG_DIR/publisher.pid"

echo "[run] 两端已启动。日志: $LOG_DIR/"
echo "[run] 下面跟随订阅端输出,Ctrl-C 退出跟随(进程仍在后台)。"
echo "[run] 停止: ./scripts/run.sh stop"
echo
tail -f "$LOG_DIR/subscriber.log"
;;
stop)
stop_all
;;
*)
echo "用法: $0 {pub|sub|both|stop}" >&2
exit 1
;;
esac