状态、订阅与实时控制
快照、订阅和 Servo 解决不同问题:快照读取一次数据,订阅按周期拉取数据,Servo 向设备连续发送控制目标。
读取方式对比
| 方式 | 数据方向 | 适用场景 |
|---|---|---|
state.snapshot() | 设备到应用,一次 | 预检、命令后确认 |
state.subscribe() | 设备到应用,连续 | 监控、记录、闭环算法输入 |
touch.subscribe() | 设备到应用,连续 | 触觉监控和算法输入 |
health.subscribe() | 设备到应用,连续 | 运行状态和故障监控 |
open_servo() | 应用到设备,连续 | 实时位置或阻抗目标 |
订阅生命周期
subscription = hand.state.subscribe(period=0.02)
try:
while True:
state = await subscription.next()
process_state(state)
finally:
subscription.close()SDK 不保存订阅历史。处理速度低于读取速度时,应用应使用有界队列并定义丢弃或降采样策略,不能无限堆积数据。
有界队列与丢帧策略
队列容量必须根据订阅周期、消费者最坏处理时间和允许的最大缓存延迟确定。例如,State 周期为 20 ms、最多允许缓存 640 ms 时,容量上限可取 ceil(0.64 / 0.02) = 32。32 只是计算示例,不是所有设备和任务的固定推荐值。
| 数据用途 | 建议隔离方式 | 队列满时的默认策略 | 要求 |
|---|---|---|---|
| State 控制输入或 UI | 独立小容量队列 | 丢最旧帧,保留最新状态 | 记录丢帧数,通过时间戳检测间隔 |
| Touch 实时算法 | 与 State 分开的有界队列 | 丢最旧帧或按明确比例降采样 | 使用 sequence 检测缺帧,不混入 Health 队列 |
| Health 与故障 | 独立高优先级队列或事件通道 | 故障出现、恢复和安全状态变化不得静默丢弃 | 队列溢出时报警并进入安全处理流程 |
| 原始数据记录 | 独立写入任务和持久化缓冲 | 默认不允许静默丢帧 | 无法施加背压时应报警、停止记录或明确标记数据缺口 |
不要用同一个队列承载所有遥测
高频 State 或 Touch 数据可能持续占满队列。Health、故障和安全状态必须使用独立通道,不能被普通遥测帧挤出。
以下 Python 与 C++ 示例使用相同的“丢最旧帧”策略,并维护 received、processed、dropped 和 max_depth 指标:
import asyncio
from dataclasses import dataclass
@dataclass
class QueueMetrics:
received: int = 0
processed: int = 0
dropped: int = 0
max_depth: int = 0
def offer_latest(
queue: asyncio.Queue,
sample,
metrics: QueueMetrics,
) -> None:
metrics.received += 1
if queue.full():
queue.get_nowait() # Drop the oldest State sample.
queue.task_done()
metrics.dropped += 1
queue.put_nowait(sample)
metrics.max_depth = max(metrics.max_depth, queue.qsize())
async def forward_states(hand) -> None:
queue = asyncio.Queue(maxsize=32)
metrics = QueueMetrics()
subscription = hand.state.subscribe(period=0.02)
async def read_subscription() -> None:
while True:
state = await subscription.next()
offer_latest(queue, state, metrics)
async def process_states() -> None:
while True:
state = await queue.get()
try:
# Replace this with the application's non-blocking handoff.
print(state.timestamp, state.positions_deg)
metrics.processed += 1
finally:
queue.task_done()
tasks = [
asyncio.create_task(read_subscription()),
asyncio.create_task(process_states()),
]
try:
await asyncio.gather(*tasks)
finally:
subscription.close()
for task in tasks:
task.cancel()
await asyncio.gather(*tasks, return_exceptions=True)
print(metrics)#include <revo3/revo3.hpp>
#include <algorithm>
#include <chrono>
#include <condition_variable>
#include <cstddef>
#include <cstdio>
#include <deque>
#include <exception>
#include <mutex>
#include <optional>
#include <stdexcept>
#include <thread>
#include <utility>
struct QueueMetrics {
std::size_t received = 0;
std::size_t processed = 0;
std::size_t dropped = 0;
std::size_t max_depth = 0;
};
template <typename T>
class BoundedLatestQueue {
public:
explicit BoundedLatestQueue(std::size_t capacity) : capacity_(capacity) {
if (capacity_ == 0) {
throw std::invalid_argument("queue capacity must be positive");
}
}
bool push_latest(T value) {
std::lock_guard<std::mutex> lock(mutex_);
if (closed_) {
return false;
}
++metrics_.received;
if (queue_.size() == capacity_) {
queue_.pop_front();
++metrics_.dropped;
}
queue_.push_back(std::move(value));
metrics_.max_depth = std::max(metrics_.max_depth, queue_.size());
ready_.notify_one();
return true;
}
std::optional<T> pop() {
std::unique_lock<std::mutex> lock(mutex_);
ready_.wait(lock, [this] { return closed_ || !queue_.empty(); });
if (queue_.empty()) {
return std::nullopt;
}
T value = std::move(queue_.front());
queue_.pop_front();
return value;
}
void mark_processed() {
std::lock_guard<std::mutex> lock(mutex_);
++metrics_.processed;
}
void close() {
{
std::lock_guard<std::mutex> lock(mutex_);
closed_ = true;
}
ready_.notify_all();
}
QueueMetrics metrics() const {
std::lock_guard<std::mutex> lock(mutex_);
return metrics_;
}
private:
const std::size_t capacity_;
mutable std::mutex mutex_;
std::condition_variable ready_;
std::deque<T> queue_;
QueueMetrics metrics_;
bool closed_ = false;
};
int main() {
revo3::Manager manager;
auto hand = manager.connect_auto();
auto subscription = hand.state().subscribe(std::chrono::milliseconds(20));
using State = decltype(subscription.next());
BoundedLatestQueue<State> queue(32);
std::exception_ptr reader_error;
std::thread reader([&] {
try {
for (std::size_t index = 0; index < 500; ++index) {
if (!queue.push_latest(subscription.next())) {
break;
}
}
} catch (...) {
reader_error = std::current_exception();
}
queue.close();
});
std::thread worker([&] {
while (auto state = queue.pop()) {
std::printf("timestamp=%lld.%09lld J0=%.2f degree\n",
static_cast<long long>(state->timestamp.sec),
static_cast<long long>(state->timestamp.nsec),
state->motors.positions_deg[0]);
queue.mark_processed();
}
});
reader.join();
worker.join();
subscription.close();
const auto metrics = queue.metrics();
std::printf("received=%zu processed=%zu dropped=%zu max_depth=%zu\n",
metrics.received, metrics.processed,
metrics.dropped, metrics.max_depth);
if (reader_error) {
std::rethrow_exception(reader_error);
}
return 0;
}Python:offer_latest() 适用于只关心最新状态的 State/UI/闭环输入,不适用于必须完整留存的日志。耗时或阻塞的文件、数据库和网络写入应放到工作线程或独立进程;如果记录链路无法跟上采集速度,应显式报警或停止记录,不能悄悄套用“丢最旧帧”。
**C++:**示例使用固定数量的读取让程序可以自然结束。实际长时间服务应在停止信号到达后先关闭订阅和队列,再 join() 生产者与消费者线程。跨线程传递回调数据时,还必须按 SDK 声明的生命周期复制所需字段。close() 会唤醒等待中的消费者,生产者异常会在两个线程结束后重新抛出。
Health 与故障不能复用这个 BoundedLatestQueue:应使用独立队列或事件通道,保留故障出现、恢复和安全状态变化;发生溢出时进入报警或安全停止流程,而不是执行 pop_front()。
无论使用哪种语言,丢帧策略都应写入应用配置,并通过日志或监控暴露;不能只依赖队列实现中的隐含行为。
Servo 会话
ServoSession 具有独立控制权和生命周期。应用必须:
- 在打开会话前完成布局和健康预检。
- 使用明确的发送周期和命令超时。
- 监控 State 与 Health,但不要假设订阅频率等于控制频率。
- 在异常和正常退出路径关闭会话。
- 断线后丢弃旧会话并重新连接。
任何固定频率都依赖主机、适配器、总线负载和固件;文档不承诺无条件固定控制频率。接口见Motion API和遥测 API。