ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

fhevm listener 的 broker 消息中间件抽象层:用一套 Rust API 同时驾驭 Redis Streams 与 RabbitMQ

fhevm listener 的 broker 消息中间件抽象层:用一套 Rust API 同时驾驭 Redis Streams 与 RabbitMQ fhevm listener 的 broker 消息中间件抽象层用一套 Rust API 同时驾驭 Redis Streams 与 RabbitMQ【免费下载链接】fhevmFHEVM, a full-stack framework for integrating Fully Homomorphic Encryption (FHE) with blockchain applications项目地址: https://gitcode.com/GitHub_Trending/fh/fhevm导读listener是 fhEVM 全栈框架中负责监听链上事件、向下游工作节点分发消息的关键服务而broker正是支撑其消息流转的底层抽象组件。本文以 listener/crates/shared/broker/README.md 为核心骨架结合 broker 源码 与端到端测试系统讲解如何用同一份 Rust 代码在 Redis Streams 与 RabbitMQ 之间无缝切换、如何通过 Topic/Namespace 建模路由、如何利用错误分类 熔断器构建健壮的消费管线。读完本文你将掌握broker的全部核心 API、配置项语义与底层实现原理可直接在 listener 或其他 fhEVM 组件中落地一套可复用的消息基础设施。组件定位listener 中的统一消息抽象broker是 fhEVM listener 工作区中的一个共享 crate在 listener/Cargo.toml 中作为crates/shared/broker引入其定位在 lib.rs 的模块注释中写得很清楚为 Redis Streams 和 RabbitMQ 提供统一接口支持直连路由Direct Routing、扇出Fanout、竞争消费Competing Consumers、熔断器Circuit Breaker以及瞬态/永久错误分类。它解决的问题很典型listener 既要消费链上事件又要向 relayer、coprocessor 等下游节点分发任务而底层消息基础设施在不同部署环境下可能是 Redis Streams 或 RabbitMQ。broker将「业务逻辑」与「基础设施选择」解耦——应用代码只面对Broker、Topic、Consumer、Publisher这些抽象后端差异被收敛到构造阶段。在 listener_core 的入口 中可以看到真实用法程序根据配置中的BrokerType选择Broker::redis_with_ensure_publish(...)或Broker::amqp(...).with_ensure_publish(...).build()之后的发布/消费代码完全不感知后端差异。库文件还强制要求至少启用redis或amqp其中一个 featurelib.rs保证编译期类型安全。快速上手发布与消费的最小闭环先看最小可运行示例与 README 的 Quick Start 保持一致use broker::{Broker, Topic, routing, AsyncHandlerPayloadOnly}; #[tokio::main] async fn main() - Result(), Boxdyn std::error::Error { // 连接 broker二选一业务代码无需关心差异 let broker Broker::redis(redis://localhost:6379).await?; // 或: Broker::amqp(amqp://localhost:5672).build().await?; // 发布端以 namespace 为作用域 let publisher broker.publisher(ethereum).await?; publisher.publish(blocks, serde_json::json!({number: 12345})).await?; // 消费端Topic 消费组 handler let topic Topic::new(routing::BLOCKS).with_namespace(ethereum); let handler AsyncHandlerPayloadOnly::new(|block: serde_json::Value| async move { println!(Processing block: {:?}, block); Ok::(), std::convert::Infallible(()) }); broker.consumer(topic) .group(indexer) .consumer_name(pod-1) .prefetch(100) .run(handler).await?; Ok(()) }这段代码有几点值得注意Broker::redis(url)是异步构造内部通过RedisConnectionManager::new_with_retry(url)建立带重试的连接lib.rsBroker::amqp(url)返回一个 builder需要调用.build().await才得到Brokerlib.rs如果拿不准 URL 协议可以使用Broker::from_url(url)它会根据redis:///amqp://前缀自动分派lib.rs。RabbitMQ 默认拓扑main / retry / dlx 三个共享交换机Broker::amqp(...)默认使用全局共享交换机拓扑在 lib.rs 中定义了三个常量常量交换机名职责AMQP_DEFAULT_MAIN_EXCHANGEmain主交换机正常消息路由AMQP_DEFAULT_RETRY_EXCHANGEretry重试交换机带 TTL 的延迟重试AMQP_DEFAULT_DLX_EXCHANGEdlx死信交换机消息最终失败去处broker 会在发布/消费建立时自动声明这套拓扑应用代码无需任何显式 AMQP 拓扑调用。这一默认行为在单元测试 lib.rs 的测试模块 中做了断言。如果需要与默认不同的命名builder 提供两个扩展点lib.rs// 前缀模式自动派生 listener、listener.retry、listener.dlx let broker Broker::amqp(amqp://localhost:5672) .with_exchange_prefix(listener) .build().await?; // 完全自定义拓扑 let broker Broker::amqp(amqp://localhost:5672) .with_topology(ExchangeTopology::new(my-main, my-retry, my-dlx)) .build().await?;发布可靠性ensure_publishBroker的两种后端都支持「确保发布」语义lib.rsRedisBroker::redis_with_ensure_publish(url, true)会让发布者在每次XADD后执行WAIT 1 500确认复制到至少一个副本WAIT最多等待 500msAMQP.with_ensure_publish(true)会开启confirm_select发布确认模式发布者等待 broker 对每次发布的Confirmation::Acklib.rs。listener 生产代码正是通过这个开关控制复制级可靠性见 main.rs。Topic后端无关的路由标识符Topic是整个抽象层的路由核心。它由可选的namespace链/服务维度如ethereum、polygon和必填的routing消息类型维度如blocks、forks组成topic.rs。let topic Topic::new(blocks).with_namespace(ethereum); assert_eq!(topic.key(), ethereum.blocks); // 全限定 key assert_eq!(topic.dead_key(), ethereum.blocks:dead); // 死信 key assert_eq!(topic.routing_segment(), blocks); // 未加 namespace 的段Topic的派生规则topic.rskey()有 namespace 时为{namespace}.{routing}否则就是{routing}dead_key()恒为{key}:deadrouting_segment()去掉 namespace 后的原始路由段空/空白 namespace 会被当作None处理topic.rs便捷构造函数Topic::namespaced(ethereum, blocks)等价于Topic::new(blocks).with_namespace(ethereum)。运行时 broker 会按后端做不同映射Rediskey 映射为 stream / dead-stream 名称AMQPkey 成为 broker 管理的共享交换机上的 routing key。所以应用代码永远只操作Topic基础设施接线差异由后端自行消化。对应实现可分别参考 Redis 侧的 StreamTopology 与 AMQP 侧的 consumer.rs。预定义 routing keyREADME 中的routing::BLOCKS/routing::FORKS实际定义在更底层的 primitives cratelistener 内部用到的完整路由常量包括fetch-new-blocks, fetch-final-block, backtrack-reorg, control.watch, control.unwatch, clean-blocks, clean-final-blocks, new-event, final-event, catchup, range-catchup, catchup-event, final-catchup, range-final-catchup, final-catchup-event以及一组面向具体消费者的路由构造函数例如consumer_new_event_routing(consumer_id)会生成{consumer_id}.new-event。这些常量通过pub mod routing在 primitives/src/lib.rs 导出broker 在文档与示例中将其复导出为routing::*使用。路由模型同一拓扑的两种后端表达以Topic::new(blocks).with_namespace(ethereum)搭配.group(indexer)为例两种后端的路由拓扑分别如下。Redis Streams 模型发布 消费Redis 侧的核心要点XADD写入主 stream消费者通过XREADGROUP从消费组读消息处理失败永久错误、重试耗尽的消息会被写入{key}:dead死信 stream。RabbitMQ 模型发布 消费AMQP 侧的队列命名规则在 consumer.rs 的resolve_amqp_queue_name中实现有 namespace 时队列名为{namespace}.{group}如ethereum.indexer若 group 本身已带{namespace}.前缀则原样保留避免重复限定无 namespace 时直接使用 group 名。对应测试可在 consumer.rs 测试模块 中找到。关键概念对照表术语RabbitMQRedis全局拓扑交换机main/retry/dlx 派生队列主 stream 死信 stream交换机固定共享交换机mainN/ANamespaceRouting key 前缀Stream 前缀RoutingRouting key 后缀Stream 后缀Topic完整 routing keynamespace.routing完整 stream 名namespace.routingGroup逻辑组名队列派生为namespace.group消费者组队列名由group派生为namespace.groupN/A消费者名Consumer tagConsumer name后端无关的通用用法发布与消费的代码路径在两种后端上完全一致只有 broker 构造不同。README 给出了一个同时运行blocksforks两个消费者的多主题示例这里完整复现并做要点标注use broker::{routing, Broker, Topic, AsyncHandlerPayloadOnly}; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] struct BlockEvent { number: u64 } #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] struct ForkEvent { at_block: u64 } async fn build_broker( broker_type: str, broker_url: str, ) - ResultBroker, Boxdyn std::error::Error { let broker match broker_type { redis Broker::redis(broker_url).await?, amqp Broker::amqp(broker_url).build().await?, other return Err(format!(unsupported broker type: {other}).into()), }; Ok(broker) } async fn run_pipeline(broker: Broker) - Result(), Boxdyn std::error::Error { let blocks_topic Topic::new(routing::BLOCKS).with_namespace(ethereum); let forks_topic Topic::new(routing::FORKS).with_namespace(ethereum); // Publisher 以 namespace 为作用域 let publisher broker.publisher(ethereum).await?; publisher .publish_to_topic(blocks_topic, BlockEvent { number: 12345 }) .await?; publisher .publish_to_topic(forks_topic, ForkEvent { at_block: 12340 }) .await?; // 消费端配置在后端间完全一致只有 group 与 topic 不同 let block_handler AsyncHandlerPayloadOnly::new(|evt: BlockEvent| async move { println!([blocks] processing block {}, evt.number); Ok::(), std::convert::Infallible(()) }); let fork_handler AsyncHandlerPayloadOnly::new(|evt: ForkEvent| async move { println!([forks] processing fork at {}, evt.at_block); Ok::(), std::convert::Infallible(()) }); let broker_blocks broker.clone(); let broker_forks broker.clone(); tokio::select! { result broker_blocks .consumer(blocks_topic) .group(indexer) .consumer_name(pod-1) .prefetch(100) .run(block_handler) { eprintln!(blocks consumer exited: {:?}, result); } result broker_forks .consumer(forks_topic) .group(fork-handler) .consumer_name(pod-1) .prefetch(50) .run(fork_handler) { eprintln!(forks consumer exited: {:?}, result); } } Ok(()) }值得深入说明的细节Broker是Clone的lib.rs内部共享Arc连接管理器克隆成本极低可放心在tokio::spawn中分发publish_to_topic有 namespace 一致性校验发布器的 namespace 必须与 topic 的 namespace 匹配包括都为 unscoped否则返回BrokerError::NamespaceMismatchpublisher.rs发布器还提供publish_batch(routing, [T])批量发布publisher.rsAsyncHandlerPayloadOnly::new内部会自动做serde_json反序列化handler 收到的T直接来自msg.payload的 JSON 解析handler.rs。错误分类与 Handler 变体broker将错误分为瞬态错误基础设施故障与永久错误消息 payload 本身有问题该分类与后端无关并驱动两类行为重试预算瞬态错误获得无限重试消息没问题是基础设施坏了永久错误计入max_retries耗尽后进入死信熔断器连续瞬态错误超过阈值时消费者完全暂停消费避免故障期间污染死信队列冷却期过后处理一条探测消息成功则熔断器闭合、恢复消费。三种 Handler 变体对照Handler闭包签名错误分类行为AsyncHandlerPayloadOnlyFn(T) - impl FutureOutput Result(), E所有闭包错误统一包装为HandlerError::Execution永久。适合只需要成功/失败二分的场景。AsyncHandlerPayloadClassifiedFn(T) - impl FutureOutput Result(), HandlerError保留闭包返回的错误分类。适合需要区分瞬态/永久失败的场景。AsyncHandlerNoArgsFn() - impl FutureOutput Result(), E所有错误包装为永久。适合纯副作用 handler心跳、定时器。此外还有第四种AsyncHandlerWithContextFn(T, MessageMetadata) - impl FutureOutput Result(), E当 handler 需要消息 ID、主题、投递次数等投递元数据时使用handler.rs。MessageMetadata是后端无关的结构体包含id、topic、delivery_count字段。底层错误模型HandlerError 与 AckDecision错误分类的根基是 handler.rs 中的HandlerError枚举Deserializationserde_json反序列化失败自动转换Executionhandler 执行失败永久Transient基础设施失败瞬态触发熔断器。两个便捷构造器handler.rs// 基础设施类错误 → 无限重试 触发熔断器 db.save(block).await.map_err(HandlerError::transient)?; // 业务逻辑类错误 → 计入 max_retries重置熔断器瞬态计数 verify_block(block).map_err(HandlerError::permanent)?;同时handler 的成功路径还可以通过返回AckDecision显式控制消息去向handler.rs变体AMQP (RabbitMQ)Redis StreamsAckbasic_ackXACKNackbasic_nack(requeue: true)→ 回主队列留在 PEL由 ClaimSweeper 负责重试Dead直接发布到 DLQ 再basic_ackXADD到死信 stream 再XACKDelay(d)发布到重试交换机带每条消息的expirationTTL再basic_ackRedis Streams 无逐消息 TTL等价于NackPEL claim_min_idle控制时机如何选择 HandlerAsyncHandlerPayloadOnly所有错误等价时使用。消费者重试最多max_retries次然后死信最简方案AsyncHandlerPayloadClassifiedhandler 需要调用外部基础设施数据库、RPC、API需要瞬态失败触发熔断器、永久失败坏 payload直接走重试/DLQ 路径时使用。README 给出了完整示例use broker::{AsyncHandlerPayloadClassified, HandlerError}; let handler AsyncHandlerPayloadClassified::new(|block: BlockEvent| async move { // 瞬态基础设施坏了不是消息的问题 // → 无限重试触发熔断器 db.save(block).await.map_err(HandlerError::transient)?; // 永久消息本身非法 // → 计入 max_retries重置熔断器瞬态计数 verify_block(block).map_err(HandlerError::permanent)?; Ok(()) });关键区别AsyncHandlerPayloadOnly把所有闭包错误都包装成HandlerError::Execution永久因此瞬态基础设施故障永远不会触发熔断器AsyncHandlerPayloadClassified让闭包直接返回HandlerError::Transient或HandlerError::Execution保留消费循环所依赖的错误分类。handler 的成功/失败路径最终会统一归一化为HandlerOutcomeAck/Nack/Dead/Delay/Transient/Permanent见 handler.rs供熔断器与 ACK/NACK 逻辑使用。熔断器状态机与完整示例状态机┌──────────┐ threshold consecutive ┌──────────┐ cooldown ┌───────────┐ │ Closed │ transient failures ──► │ Open │ expires ─► │ Half-Open │ │ (normal) │ │ (paused) │ │ (probe) │ └──────────┘ └──────────┘ └───────────┘ ▲ │ │ probe succeeds │ └───────────────────────────────────────────────────────────────┘ probe fails → back to Open状态行为Closed正常消费。瞬态失败递增计数成功/永久失败重置计数。Open不消费任何消息。等待冷却期到期。Half-Open处理一条探测消息。成功 → Closed失败 → Open重新开始冷却。源码级实现细节熔断器的核心实现在 traits/circuit_breaker.rs几个值得注意的工程细节默认参数L11-L39failure_threshold 5、cooldown_duration 30s、half_open_timeout 5 × cooldownOpen 状态下的成功是 no-op在同一批XREADGROUP消息里如果消息 N 触发了熔断、消息 N1 却成功了这个「偶然成功」不得绕过冷却期直接闭合电路——唯一的合法Open → Closed路径是「冷却到期 → Half-Open 探测 → 成功」L184-L211永久失败重置瞬态计数L254-L260消息问题不该累积到熔断器上Half-Open 探测互锁probe_in_flight标志保证同一 Half-Open 窗口内只有一个探测任务在途若探测任务异常退出panic/发送失败clear_probe_in_flight让熔断器留在 Half-Open 等待下一轮重新探测L289-L311Half-Open 超时兜底如果 Half-Open 期间消费者的 PEL 被清空例如 ClaimSweeper 把条目认领到了别的消费者探测消息迟迟不来half_open_expired会让熔断器重新回到 Open 而非盲目闭合L313-L341。这些行为都有对应的单元测试覆盖例如success_while_open_does_not_close_circuitL454-L474验证了 Open 状态下的偶然成功不会闭合电路。Prometheus 指标熔断器在挂上标签with_labels(backend, topic)后会发射三组指标L100-L155broker_circuit_breaker_state{backend,topic}0closed1open2half-openbroker_circuit_breaker_consecutive_failures{backend,topic}连续瞬态失败计数broker_circuit_breaker_trips_total{backend,topic}熔断次数每次跳闸时累加。完整示例含熔断器use std::time::Duration; use broker::{routing, Broker, Topic, AsyncHandlerPayloadClassified, HandlerError}; async fn run(broker: Broker) - Result(), Boxdyn std::error::Error { let topic Topic::new(routing::BLOCKS).with_namespace(ethereum); let handler AsyncHandlerPayloadClassified::new(|block: BlockEvent| async move { db.save(block).await.map_err(HandlerError::transient)?; verify_block(block).map_err(HandlerError::permanent)?; Ok(()) }); broker.consumer(topic) .group(indexer) .prefetch(100) .max_retries(5) // 永久错误5 次尝试后进 DLQ .circuit_breaker(3, Duration::from_secs(30)) // 3 个瞬态错误 → 暂停 30s .run(handler) .await?; Ok(()) }熔断器在 Redis Streams 与 RabbitMQ 上行为完全一致且状态在两种后端的重连之间都能保持。e2e 测试如 redis_e2e.rs专门验证了一个微妙场景Half-Open 的探测消息必须来自 PELXREADGROUP 0而不是新消息以保证消息处理顺序不被打乱——这是用prefetch_count1使测试确定性的。队列检查Queue InspectionBroker::exists()检查某 topic 对应的队列或 stream 是否已创建不需要消费组存在lib.rsuse broker::{Broker, Topic, routing}; let topic Topic::new(routing::BLOCKS).with_namespace(ethereum); if !broker.exists(topic).await? { // 队列/stream 尚未创建 —— 跳过深度检查 }后端机制返回true的条件Redis对 key 执行TYPE命令key 存在且类型为streamAMQP被动queue_declarebroker 确认队列存在404 →false除exists()外Broker还提供了三组深度/空闲检查方法lib.rsqueue_depths(topic, group)返回主队列、重试队列、死信队列的消息数传入 group 时Redis 还会通过XINFO GROUPS填充pendingPEL 计数与lag未投递条目可用has_pending_work()判断消费者是否即将收到消息——这是「是否需要发布种子消息」的主要判断依据is_empty(topic, group)单次后端往返的快速检查判断消费者是否会收到零条消息用于决定是否发布种子消息启动处理循环is_empty_or_pending(topic, group)判断消费组是否已追平完全空闲或至多一条消息在途pending 1 lag 0专为去重守卫设计——prefetch 为 1 时单条 pending 意味着已有消费者在处理调用方可跳过重复工作。消费模式直连、扇出与竞争消费直连路由Direct Routing发布者发到特定 routing key只有匹配的队列/组能收到// 只接收 blocks 消息 broker.consumer(Topic::new(blocks).with_namespace(ethereum)) .group(block-processor) .run(handler).await?;扇出Fanout多组广播多个 group 绑定到同一 routing key各自独立收到全部消息副本let blocks_topic Topic::new(blocks).with_namespace(ethereum); // 两个 group 各自独立收到所有 block 消息 broker.consumer(blocks_topic).group(indexer).run(handler1).await?; broker.consumer(blocks_topic).group(analytics).run(handler2).await?;竞争消费Competing Consumers组内负载均衡同一 group、不同 consumer_name消息在组内负载均衡分发// Pod 1 broker.consumer(topic) .group(indexer) .consumer_name(pod-1) .run(handler).await?; // Pod 2 —— 与 pod-1 分摊负载 broker.consumer(topic) .group(indexer) .consumer_name(pod-2) .run(handler).await?;这三种模式对应 README 开头声明的三大特性Direct Routing / Fanout / Competing Consumers。在 Redis 侧竞争消费依赖XREADGROUP GROUP {group} {consumer}的组内多消费者语义在 AMQP 侧同一队列被多个 consumer tag 消费天然具备负载均衡。未指定consumer_name时builder 会用{group}-{uuid_v7}自动生成唯一实例名consumer.rs。配置选项全解通用选项两种后端均适用broker.consumer(topic) .group(worker) // 必填组名 .consumer_name(pod-1) // 可选实例名 .prefetch(100) // 缓冲消息数默认10 .max_retries(5) // 死信前的最大重试次数默认3 .circuit_breaker(3, Duration::from_secs(30)) .run(handler).await?;默认值定义在 config.rs 的ConsumerConfigprefetch 10、max_retries 3、consumer_name None自动生成、circuit_breaker None默认关闭。group是唯一必填项build()时会校验缺失则返回BrokerError::MissingGroupconsumer.rs。各通用选项的后端落地方式consumer.rs选项RedisAMQPprefetchXREADGROUP的read_countQoS 预取计数consumer_name组内消费者名consumer tagmax_retries死信前重试上限同上circuit_breaker(threshold, cooldown)连续瞬态错误阈值 冷却时长同上Redis 专属选项broker.consumer(topic) .group(worker) .redis_block_ms(10000) // 空 stream 阻塞等待 10s .redis_claim_min_idle(60) // 60s 后认领卡住的消息默认30s .redis_claim_interval(10) // 每 10s 巡检一次默认10s .run(handler).await?;对应的默认值在 config.rs 的RedisOptionsblock_ms 5000、claim_min_idle_secs 30、claim_interval_secs 10。claim_*选项驱动后台 ClaimSweeper 任务——消费者崩溃后留在 PEL 中的消息会被其他消费者按claim_min_idle认领保证至少一次投递语义。这些选项在 AMQP 后端上会被静默忽略consumer.rs。AMQP 专属选项broker.consumer(topic) .group(worker) .amqp_retry_delay(10) // 重试间隔 10s默认5s .amqp_routing_pattern(ethereum.blocks.#) // 通配符订阅 .run(handler).await?;默认值在 config.rs 的AmqpOptionsretry_delay_secs 5routing_key_pattern未设置时使用 namespace 限定的默认 key{namespace}.{routing}。通配符订阅利用了 RabbitMQ topic exchange 的#/*语义。这些选项在 Redis 后端上会被静默忽略consumer.rs。构建与运行消费者两种时机方式一立即构建并运行broker.consumer(topic) .group(worker) .run(handler).await?;方式二先构建后运行适合需要并行编排多个消费者的场景// 先构建消费者 let blocks_consumer broker.consumer(Topic::new(blocks).with_namespace(ethereum)) .group(block-processor) .prefetch(100) .build()?; let forks_consumer broker.consumer(Topic::new(forks).with_namespace(ethereum)) .group(fork-handler) .build()?; // 用 tokio::spawn 运行 let h1 tokio::spawn(blocks_consumer.run(block_handler)); let h2 tokio::spawn(forks_consumer.run(fork_handler)); futures::future::join_all(vec![h1, h2]).await;两种方式本质相同run(handler)内部就是build()?之后调用consumer.run(handler)consumer.rs。此外build()得到的Consumer还有一个重要方法ensure_topology()consumer.rsAMQP声明交换机、队列与绑定。必须在检查队列深度is_empty或发布种子消息之前调用否则没有绑定队列时 AMQP 会静默丢弃消息Redisno-opstream 在首次XADD时自动创建。多 Namespace 部署一进程多链listener 经常需要同时服务多条链ethereum / polygon / arbitrumbroker的 namespace 机制让这变得非常简单let broker Broker::redis(redis://localhost:6379).await?; for ns in [ethereum, polygon, arbitrum] { let topic Topic::new(blocks).with_namespace(ns); let broker broker.clone(); tokio::spawn(async move { broker.consumer(topic) .group(indexer) // AMQP 队列变为 {ns}.indexer .consumer_name(pod-1) .run(handler).await }); }在 Redis 侧各链的 stream 天然隔离为ethereum.blocks、polygon.blocks、arbitrum.blocks在 AMQP 侧队列派生为ethereum.indexer、polygon.indexer、arbitrum.indexer互不干扰。优雅停机CancellationToken使用tokio_util::sync::CancellationToken可以让消费者干净地停止。broker复导出该类型lib.rs用法如下use broker::{Broker, Topic, routing, CancellationToken}; let token CancellationToken::new(); let consumer broker.consumer(topic) .group(indexer) .with_cancellation(token.clone()) .build()?; // 在任务中运行 let handle tokio::spawn(consumer.run(handler)); // 之后触发优雅停机 token.cancel(); handle.await??;取消语义token 被取消后消费者完成当前消息批次并排空在途结果后停止consumer.rs不会粗暴中断正在处理的消息配合basic_ack/XACK可以做到尽量不丢消息、不重复处理。自定义 Routing KeyTopic::new接受任意字符串作为 routing segment因此可以自由扩展业务路由let topic Topic::new(erc20-transfers).with_namespace(ethereum);listener 中大量使用这种自定义路由fetch-new-blocks、backtrack-reorg、control.watch、catchup-event等见 primitives/src/routing.rs。当业务需要比预定义常量更细的粒度时例如按合约地址、按消费者 ID 路由直接构造即可。测试与验证从单元测试到端到端broker的可靠性建立在三层测试之上可作为你使用该库的参照单元测试Topic的 key/dead_key 派生topic.rs、ConsumerConfig默认值与 builderconfig.rs、AMQP 拓扑默认值与队列名派生consumer.rs、熔断器完整状态机circuit_breaker.rs端到端测试需要 Docker 的 Redis 与 AMQP 集成测试位于 tests/ 目录。其中 redis_e2e.rs 覆盖了熔断器探测顺序、PEL 认领、死信路由等关键行为README 头注释建议通过make test-e2e-redis运行它会清理命名容器生产用法验证listener_core/src/main.rs 展示了真实服务中如何根据配置选择后端并注入ensure_publish。小结broker的价值在于用一套类型安全的 Rust API 屏蔽了 Redis Streams 与 RabbitMQ 的全部差异让业务代码只关注Topic建模、handler 逻辑与错误分类。核心设计可以概括为四点Topic/Namespace 路由模型{namespace}.{routing}在两种后端上映射为 stream 名或 routing key天然支持多链/多服务隔离错误分类驱动重试与熔断瞬态错误无限重试 触发熔断永久错误计重试上限 进死信二者边界清晰、行为可预测三种消费模式开箱即用直连、扇出、竞争消费配合默认参数即可支撑从单机到多 Pod 的水平扩展可观测性与生命周期完备Prometheus 熔断指标、队列深度检查、CancellationToken优雅停机以及ensure_topology/ensure_publish等生产级细节。对于需要同时对接 Redis 与 RabbitMQ 的 fhEVM 组件listener、relayer、coprocessor 均可受益broker提供了现成的答案——你只需在构造阶段决定后端其余代码原样复用。【免费下载链接】fhevmFHEVM, a full-stack framework for integrating Fully Homomorphic Encryption (FHE) with blockchain applications项目地址: https://gitcode.com/GitHub_Trending/fh/fhevm创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表