客户端协调与请求所有权
Rust 客户端将面向应用的生产、消费 API 与路由发现、Broker 连接、重试和消费协调连接起来。共享机制属于应用提供的 ClientRuntime;创建生产者不会隐式建立备用运行时。
所有权与共享状态
应用创建 RuntimeOwner,将子服务上下文和遥测句柄传给 ClientRuntime::try_new,再把返回的 Arc<ClientRuntime> 共享给各客户端门面。公开类型从 rocketmq_client_rust 或其 prelude 导入。内部工厂和处理器模块用于理解实现,不是应用应依赖的导入路径。
MQClientInstance 汇集生产者/消费者注册、主题路由缓存、Broker 地址状态、心跳协调、传输回调及定时刷新/再平衡工作。共享运行时所有权不代表每个门面拥有相同的组、订阅或客户端身份。身份和组配置仍决定 Broker 如何观察客户端。
从主题到请求
- 解析配置的 NameServer 端点,按明确选择的发现实现进行可选发现。
- 获取或刷新主题路由,生成发布/订阅队列视图。
- 按本次操作的路由策略选择队列,解析其 Broker 端点。
- 携带剩余截止时间,通过具有准入约束的受管理传输提交请求。
- 同时解释传输结果和 Broker 结果,更新路由/故障信息,选择允许的重试。
路由缓存减少发现流量,但 Broker 迁移或角色变化后可能过期。路由更新协调器及定时刷新维护本地视图;具体操作仍需处理路由缺失及通告地址不可达。NameServer 连通性和 Broker 连通性是两项不同依赖。
普通生产者重试循环携带共享截止时间,在工作前检查,并将派生出的单次尝试截止时间传入发送核心。它不会为每次尝试重新获得完整超时预算。重试策略结合路由错误、发送状态、剩余次数和通信模式,可以刷新路由或选择其他 Broker。选择器及显式指定队列的 API 有自身语义,不能假设所有发送重载都采用相同故障转移行为。
完成具有多个层次
| 观察到的结果 | 应用已知事实 | 仍不确定的内容 |
|---|---|---|
| 发送前本地校验/准入拒绝 | 本次尝试未通过该路径到达 Broker | 之前某次尝试是否成功 |
| 本地写入完成 | 传输完成本地写操作 | Broker 处理和持久化 |
| 有效发送响应 | Broker 返回了特定发送状态 | 业务消费,以及超出该状态的持久性 |
| 分发后响应超时或断连 | 未观察到确定响应 | Broker 是否已接纳或完成请求 |
| 客户端工作取消 | 本地等待/工作按所有者规则停止 | 远端操作是否仍会完成 |
回调、单向发送与等待响应发送的完成契约不同。外层 Ok 不总是 SendOk 确认;应检查所选 API 的返回或回调路径,详见发送消息。
应用重试需要结合幂等性和整体截止时间。在客户端有界重试之外包一层无限循环,会破坏该边界并可能产生重复效果。
消费协调
| API | 协调与完成 |
|---|---|
| Push 并发消费 | 客户端拉取/接收并调度回调;根据回调结果确认或重试 |
| Push 顺序消费 | 队列所有权/锁及顺序回调约束单队列处理 |
| Lite Pull |