写了一个全功能的 RocketMQ .NET 客户端:EventHorizon.RocketMQ

发布时间:2026/8/4 8:58:31

写了一个全功能的 RocketMQ .NET 客户端:EventHorizon.RocketMQ
目录前言功能概览它不是官方客户端的替代品快速接入先启动本地测试环境gRPC通过 Proxy 连接 RocketMQ 5Remoting通过 NameServer 发现 Broker多集群或多实例发送消息消费消息gRPC 的消费模型PushConsumerclassic Remoting 的消费模型PushConsumerOpenTelemetry把消息链路接进现有观测体系测试不只验证“能发送一条消息”使用前需要知道的边界总结项目地址eventhorizon-cli/EventHorizon.RocketMQNuGetEventHorizon.RocketMQ.Grpc | EventHorizon.RocketMQ.Remoting前言#我一直很喜欢 RocketMQ。它并不只是吞吐高在电商这类高并发场景中从可靠投递、顺序消费到延迟与事务消息RocketMQ 提供了一套相当完整的消息能力。很久以前我就想写一个 .NET 的 RocketMQ 客户端。那时 RocketMQ 5 还没有发布客户端主要使用传统的 Remoting 协议。Remoting 协议经过多年演进除了消息收发外还涉及 NameServer 路由、Broker 通信、心跳、负载均衡、消费进度、重试和事务消息等大量细节实现完整客户端的工作量很大。当时项目只写了一个开头后来受限于时间和精力便搁置了。后来 RocketMQ 5 引入了 gRPC 协议并发布了官方 C# 客户端。官方客户端已能满足通过 RocketMQ 5 Proxy 接入的消息收发需求新项目接入 RocketMQ 时应优先评估官方客户端。不过官方 C# 客户端主要面向 gRPC 协议而实际环境中仍有大量采用 NameServer 发现 Broker 的 Remoting 部署。除此之外.NET 应用通常还需要和 DI、Generic Host、Options、日志以及 OpenTelemetry 更自然地集成。我希望这个项目不仅能完成消息的收发也能提供更好的使用体验应用可以通过 DI 注册和注入 Producer、Consumer 与 Admin API让客户端跟随 Generic Host 一起启动和停止需要连接多个集群时可以使用 Keyed Service 区分不同客户端同时通过 OpenTelemetry 将消息发送、接收和消费处理纳入现有的 Trace 和 Metrics 体系。借助最近的 Vibe Coding我重新把这个项目捡了起来。在实现 RocketMQ 5 gRPC 客户端的同时也补充了 classic Remoting 协议的支持并尝试以更符合 .NET 应用习惯的方式组织两种协议下常见的 Producer、Consumer 和 Admin API。功能概览#✅表示客户端已实现对应 API。实际是否可用还取决于 RocketMQ 服务端版本、Broker 配置和 Proxy 能力。功能RocketMQ 5 gRPCclassic Remoting说明连接目标✅✅gRPC 连接 RocketMQ 5 ProxyRemoting 连接 NameServer 和 Broker。普通消息发送✅✅两种协议都支持。FIFO/顺序消息发送✅✅两种协议都通过MessageGroup选择稳定队列最终的顺序语义还取决于所选 Consumer 模型和服务端能力。定时/延迟消息发送✅✅服务端需要支持定时投递。撤回定时/延迟消息✅✅必须在消息投递前执行。优先级消息✅✅服务端需要支持优先级消息。Lite 消息发送✅✅依赖 Lite Topic 及服务端配置。事务消息✅✅gRPC 需要配置事务检查器两种协议都要求业务保存本地事务结果。批量发送—✅Remoting Producer 支持一批消息必须属于同一个 Topic。单向发送—✅Remoting Producer 支持调用完成不代表 Broker 已存储消息。指定队列发送—✅Remoting Producer 支持。队列选择器—✅Remoting Producer 支持。请求-响应消息—✅Remoting Producer 支持。SimpleConsumer✅—显式拉取、确认、修改不可见时间、转入死信队列。gRPC PushConsumer 并发消费✅—客户端主动查询分配并通过长轮询获取消息。gRPC PushConsumer FIFO 消费✅—同一MessageGroup内保持顺序。gRPC FIFO 消费加速✅—不同MessageGroup并行同组仍保持顺序。LitePushConsumer✅—需要支持SyncLiteSubscription的 Proxy。PullConsumer—✅显式指定队列和 Offset 拉取消息。LitePullConsumer—✅支持客户端轮询、分配、Seek、暂停/恢复和提交当前仅支持集群消费。POPConsumer—✅支持 Receipt 确认与不可见时间续期仅适用于普通 Topic。Remoting PushConsumer 并发消费—✅支持批量回调与部分前缀确认。Remoting PushConsumer FIFO 消费—✅单消息顺序分发。集群消费✅✅具体消费模型取决于所选 Consumer。广播消费—✅Remoting PushConsumer 支持。Tag 过滤✅✅两种协议都支持。SQL 过滤✅✅Broker 需要启用相应能力。运行时订阅变更✅✅Consumer 可在运行时调整订阅。重试与死信处理✅✅支持范围随 Consumer 模型变化。只读 Admin API—✅查询 Topic、队列、消息和消费位点等信息。DI 注册✅✅分别使用AddRocketMQGrpc和AddRocketMQRemoting。Generic Host 生命周期✅✅Producer、Consumer 可随 Host 启动和停止。Options、日志✅✅使用 .NET 常见配置与日志模型。默认客户端注册✅✅可直接通过构造函数注入。Keyed Service 多客户端注册✅✅可区分多个集群或同一角色的多个实例。OpenTelemetry Trace✅✅覆盖发送、接收、消费处理与确认等操作。OpenTelemetry Metrics✅✅应用自行选择 exporter 和观测后端。TLS、ACL、Namespace✅✅具体配置方式随协议不同。上表列出了项目当前实现的主要能力。为了节约篇幅本文不会逐一介绍每个 API也不会把所有消息类型和消费者模型都展开。下面只选择几个最常用的场景协议选择、客户端注册、发送消息、消费消息和 OpenTelemetry 接入。完整的功能说明、协议差异和可运行示例可以参考项目仓库中的 README 与 samples。它不是官方客户端的替代品#RocketMQ 5 的 gRPC 协议是面向多语言客户端设计的新协议。它需要 RocketMQ 5 服务端和 Proxy使用的是重新设计过的 API。对于能够使用这条链路的新项目应优先评估官方 C# 客户端。EventHorizon.RocketMQ不尝试在 API 层面兼容或替代官方客户端。它提供的是另一种 .NET 客户端实现并同时保留 classic Remoting 的接入能力。两条协议的边界如下维度RocketMQ 5 gRPCclassic Remoting连接路径应用 - Proxy - NameServer/Broker应用 - NameServer - Broker客户端 PackageEventHorizon.RocketMQ.GrpcEventHorizon.RocketMQ.Remoting典型场景已部署 RocketMQ 5 Proxy希望使用新协议 API需要对接已有经典集群或使用 Pull、POP、Admin、请求-响应等经典能力公开模型EventHorizon.RocketMQ.Grpc.*EventHorizon.RocketMQ.Remoting.*能否直接混用模型不可以不可以同一应用可以同时引用两个 Package但它们的Message、发送回执、异常和 Consumer 模型都是协议专用的类型。需要同时使用两种协议时应该显式为同名类型设置别名而不是把其中一个协议的模型传递给另一个协议。using GrpcMessage EventHorizon.RocketMQ.Grpc.Producer.Message; using RemotingMessage EventHorizon.RocketMQ.Remoting.Producer.Message;快速接入#按应用实际使用的协议安装对应的 Package只有同时使用两种协议时才需要同时安装两个dotnet add package EventHorizon.RocketMQ.Grpc dotnet add package EventHorizon.RocketMQ.Remoting两个 Package 的公开 API 是各自独立的不会在客户端 API 层面相互依赖。先启动本地测试环境#下面的示例使用仓库提供的 本地 Apache RocketMQ 5.5.0 环境。在仓库根目录执行docker compose -f test-environments/rocketmq/compose.yaml up -d --wait环境会启动 NameServer、Broker、cluster-mode Proxy 和 Dashboard宿主机上的 gRPC Proxy 地址是127.0.0.1:8081Remoting 的 NameServer 地址是127.0.0.1:9876。一次性resource-init服务会在Compose 命令返回前自动创建普通 Topiceventhorizon-test-topic因此接下来的发送和订阅示例可以直接使用它。后文会使用 OpenTelemetry 查看 Metrics 和 Trace可一并启动仓库提供的 Grafana OTEL LGTM 环境docker compose -f test-environments/otel-lgtm/compose.yaml up -d --waitGrafana 地址为 http://127.0.0.1:3000默认账号和密码均为admin。该环境同时提供 OTLP Collector后面的示例会把客户端遥测数据导出到这里。如果需要验证其他场景也可以自行创建 Topic打开本地 RocketMQ Dashboard在 Topic 管理页面创建也可以在环境启动后执行命令docker compose -f test-environments/rocketmq/compose.yaml exec broker sh mqadmin updateTopic \ -n nameserver:9876 -c DefaultCluster -t my-test-topic本地环境为了便于验证暴露了上述端口生产环境仍应由运维显式准备 Topic、Consumer Group、ACL、TLS 和 Broker 对外可达地址。gRPC通过 Proxy 连接 RocketMQ 5#下面的示例注册一个 gRPC Producer。Endpoint可以配置一个或多个以分号分隔的 Proxy 地址当应用运行在 Generic Host 中时客户端会随 Host 一起启动和停止。using EventHorizon.RocketMQ.Grpc; using EventHorizon.RocketMQ.Grpc.Producer; var builder WebApplication.CreateBuilder(args); var rocketMQ builder.Services.AddRocketMQGrpc(options { // 仓库测试环境暴露的 gRPC Proxy多个地址可用分号分隔。 options.Endpoint 127.0.0.1:8081; // 控制单次请求、路由缓存和客户端心跳的时序。 options.RequestTimeout TimeSpan.FromSeconds(3); options.RouteCacheDuration TimeSpan.FromMinutes(5); options.HeartbeatInterval TimeSpan.FromSeconds(30); // 测试环境未启用 TLS/ACL。生产环境按需启用 TLS并从安全配置读取 AK/SK 和可选 SecurityToken。 options.UseTLS false; // options.AccessKey builder.Configuration[RocketMQ:AccessKey]; // options.AccessSecret builder.Configuration[RocketMQ:AccessSecret]; // options.SecurityToken builder.Configuration[RocketMQ:SecurityToken]; // 多租户环境可按需设置 Namespace为 Topic 和 Consumer Group 添加前缀。 // options.Namespace tenant-a; }); rocketMQ.AddGrpcProducer();Remoting通过 NameServer 发现 Broker#classic Remoting 客户端首先连接 NameServer再根据路由信息直接连接对应的 Brokerusing EventHorizon.RocketMQ.Remoting; using EventHorizon.RocketMQ.Remoting.Producer; var builder WebApplication.CreateBuilder(args); var rocketMQ builder.Services.AddRocketMQRemoting(options { // 仓库测试环境的 NameServer多个地址同样可以用分号分隔。 options.NamesrvAddr 127.0.0.1:9876; }); rocketMQ.AddRemotingProducer(options { options.GroupName eventhorizon-blog-remoting-producer; });这意味着应用进程必须能访问 NameServer 返回的每个 Broker 地址。对于容器、Kubernetes 或跨网络部署Broker 对外注册的地址是否可达和客户端代码本身同样重要。多集群或多实例#一个 Host 连接多个集群或者注册同一角色的多个实例时可以使用 .NET Keyed Service。注册名同时也是解析服务时使用的 keyusing EventHorizon.RocketMQ.Grpc; using EventHorizon.RocketMQ.Grpc.Producer; using Microsoft.Extensions.DependencyInjection; builder.Services .AddRocketMQGrpc(local-rocketmq, options options.Endpoint 127.0.0.1:8081) .AddGrpcProducer(); public sealed class OrderPublisher( [FromKeyedServices(local-rocketmq)] IGrpcProducer producer) { }默认注册仍然使用普通构造函数注入。只有在不使用 Generic Host、直接从独立ServiceProvider解析客户端时才需要自行调用客户端的StartAsync和StopAsync。发送消息#注册完成后直接注入协议专用的 Producer 即可。下面是最常见的 gRPC 发送路径using System.Text; using EventHorizon.RocketMQ.Grpc.Producer; public sealed class OrderPublisher(IGrpcProducer producer) { public async Taskstring PublishAsync( string orderId, CancellationToken cancellationToken) { var message new Message(eventhorizon-test-topic, Encoding.UTF8.GetBytes(orderId)) { Tag created, MessageGroup orderId }; message.Keys.Add(orderId); message.Properties[source] checkout; var receipt await producer.SendAsync(message, cancellationToken); return receipt.MessageId; } }设置MessageGroup会将消息标记为 FIFO并按该 Group 选择稳定队列同组由 FIFO Consumer 串行处理不同组可在可用并发度内并行。普通消息不需要设置它。gRPC 的 FIFO、定时/延迟、优先级和 Lite 是互斥的消息类型不能在同一条消息上同时设置事务消息也不能与它们组合。分别使用方式如下设置DeliveryTimestamp发送定时或延迟消息在投递前可使用发送回执中的RecallHandle撤回。设置非负Priority发送优先级消息。设置LiteTopic发送 Lite 消息。使用SendTransactionAsync发送事务消息并在 Producer 启动前为事务 Topic 配置检查器。事务消息不会天然保证业务一致性。事务检查器必须能根据持久化的本地业务状态给出Commit、Rollback或Unknown不能只依赖进程内内存状态。Remoting Producer 也提供普通、FIFO、延迟、优先级、Lite 和事务消息除此之外还支持批量发送、单向发送、指定队列、队列选择器和请求-响应。发送结果是RemotingSendResult部分 Broker 返回的非成功发送状态不会直接表现为异常因此业务应检查返回状态。消费消息#gRPC 的消费模型#gRPC 的 Consumer 模型较少但把“应用自己控制接收”与“客户端自动分发”分得很清楚模型适用情况IGrpcSimpleConsumer业务自己决定何时ReceiveAsync、每次取多少消息、何时确认、修改不可见时间或转发死信。IGrpcPushConsumer希望客户端自动管理分配、长轮询、分发与处理结果业务只实现 Handler。IGrpcLitePushConsumer使用 LiteTopic并希望沿用 Push 的自动分发模型服务端还需要准备 LITE parent topic 与 Consumer Group bind topic。gRPC 的 Push 和 LitePush 都是客户端主动发起分配查询和ReceiveMessage长轮询不是 Broker 主动向应用建立推送连接。SimpleConsumer则把接收和确认时机明确交给业务代码。PushConsumer#业务不需要自己写长轮询循环时可以使用IGrpcPushConsumer。它由客户端自动查询分配结果、通过长轮询接收消息并分发给 Handlerusing EventHorizon.RocketMQ.Grpc.Consumer; using EventHorizon.RocketMQ.Grpc.Consumer.Push; using Microsoft.Extensions.DependencyInjection; rocketMQ.AddGrpcPushConsumerOrderCreatedHandler(ServiceLifetime.Scoped, options { options.GroupName eventhorizon-blog-grpc-push-consumer; options.MaxConcurrency 8; options.MaxDeliveryAttempts 16; options.InvisibleDuration TimeSpan.FromSeconds(30); options.ConsumeTimeout TimeSpan.FromMinutes(15); options.Subscribe(eventhorizon-test-topic, new FilterExpression(created)); }); public sealed class OrderCreatedHandler : IGrpcPushMessageHandler { public ValueTaskConsumeResult HandleAsync( GrpcMessageView message, CancellationToken cancellationToken) { // 处理 message.Body。 return ValueTask.FromResult(ConsumeResult.Success); } }Handler 返回Success时确认消息返回Retry时请求重新投递返回DeadLetter时将消息转入死信队列。非 FIFO 消费中处理超过ConsumeTimeout后客户端会停止延长消息的不可见时间使其后续重新投递。若业务代码忽略取消令牌而继续执行客户端无法强制终止它因此 Handler 必须响应取消请求并保证处理幂等。需要自行控制接收批次、确认时机和不可见时间时应选择上表中的IGrpcSimpleConsumer使用 LiteTopic 时则选择IGrpcLitePushConsumer并完成相应的服务端资源准备。classic Remoting 的消费模型#Remoting 提供的消费模型更多应根据业务需要选择模型适用情况IRemotingPullConsumer业务明确管理队列和 Offset并决定何时发起拉取。IRemotingLitePullConsumer希望由客户端管理分配、轮询、提交、暂停/恢复和 Seek。IRemotingPopConsumer需要显式选择队列并使用 Receipt 和不可见时间完成确认。IRemotingPushConsumer希望自动分配与回调分发支持集群、广播、并发批量和 FIFO 消费。PushConsumer#Remoting PushConsumer 同样由客户端完成队列分配、拉取和回调分发但 Handler 的入参是一批消息using EventHorizon.RocketMQ.Remoting.Consumer; using EventHorizon.RocketMQ.Remoting.Consumer.Push; using Microsoft.Extensions.DependencyInjection; rocketMQ.AddRemotingPushConsumerRemotingOrderCreatedHandler( ServiceLifetime.Scoped, options { options.GroupName eventhorizon-blog-remoting-push-consumer; options.MaxConcurrency 8; // 最多同时执行的 Handler 数。 options.BatchSize 32; // 单次从 Broker 拉取的最大消息数。 options.ConsumeMessageBatchSize 16; // 单次交给 Handler 的最大消息数。 options.MaxDeliveryAttempts 16; options.Subscribe(eventhorizon-test-topic, new FilterExpression(created)); }); public sealed class RemotingOrderCreatedHandler : IRemotingPushMessageHandler { public ValueTaskConsumeResult HandleAsync( IReadOnlyListRemotingMessageView messages, RemotingPushConsumeContext context, CancellationToken cancellationToken) { foreach (var message in messages) { // 根据 message.Body 执行业务处理。 } return ValueTask.FromResult(ConsumeResult.Success); } }BatchSize控制单次长轮询向 Broker 请求的消息上限ConsumeMessageBatchSize控制单次回调交给 Handler 的消息上限。仅在并发、非 FIFO 批次中Handler 返回Success前可将context.AckIndex设为最后一条已成功消息的从零开始索引以确认连续前缀并重试余下消息Retry和DeadLetter始终作用于整个批次。带MessageGroup的 FIFO 消息和ConsumeOrderly都按单条串行处理广播模式没有 Broker 的重试或死信流程。无论使用哪种 ConsumerRocketMQ 的重试语义都要求业务具备幂等性。网络超时、Consumer 崩溃或确认结果丢失时一条已处理消息仍可能再次到达。OpenTelemetry把消息链路接进现有观测体系#两个协议 Package 都内置了 OpenTelemetry Activity 和 Metrics不需要额外安装独立的 instrumentation Package。应用负责选择 resource、采样器、processor 和 exporterbuilder.Services.AddOpenTelemetry() .WithTracing(tracing tracing .AddAspNetCoreInstrumentation() // 应用自己创建的业务 span 也需要显式注册 source。 .AddSource(GrpcConsumerMessageHandler.ActivitySourceName) .AddRocketMQGrpcInstrumentation() .AddRocketMQRemotingInstrumentation() .AddOtlpExporter()) .WithMetrics(metrics metrics .AddAspNetCoreInstrumentation() .AddRocketMQGrpcInstrumentation() .AddRocketMQRemotingInstrumentation() .AddOtlpExporter());客户端会为 Producer 发送、非空接收、Push/LitePush 处理和 Consumer 确认等操作产生遥测数据。Producer 的sendActivity 会将当前 Trace Context 写入消息属性消息经 Broker 到达 Consumer 后PushConsumer 会读取这份上下文并通过Activity Link建立关联。因此发送端和消费端的遥测数据可以关联查看。这里需要区分“上下文已经传到 Consumer”和“瀑布图一定是一棵严格的父子树”。PushConsumer 先记录本地的receive再以它为父节点创建processProducer 的send则作为Activity Link关联到消费侧。这样既保留了跨 Broker 的因果关系也能正确处理长轮询和批量接收避免把多个不同的生产端硬接成一个父节点。Handler 执行时当前 Activity 已经是 PushConsumer 创建的process。业务代码只需创建自己的 Activity不需要再解析或传递消息属性它会自然挂在process下面。下面以库存预占为例模拟处理耗时并创建业务 Activityusing System.Diagnostics; using EventHorizon.RocketMQ.Grpc.Consumer; using EventHorizon.RocketMQ.Grpc.Consumer.Push; internal sealed class GrpcConsumerMessageHandler : IGrpcPushMessageHandler { internal const string ActivitySourceName EventHorizon.RocketMQ.Sample.Consumer; private static readonly ActivitySource ActivitySource new(ActivitySourceName); public async ValueTaskConsumeResult HandleAsync( GrpcMessageView message, CancellationToken cancellationToken) { // 模拟库存预占或订单状态落库等业务处理。 using var activity ActivitySource.StartActivity(inventory.reserve, ActivityKind.Internal); await Task.Delay(TimeSpan.FromMilliseconds(650), cancellationToken); return ConsumeResult.Success; } }消费端本地的父子树仍然反映实际执行顺序receive - process - inventory.reserve。Producer 的send不会被强行设为process的父节点而是作为process的Activity Link关联起来在 Grafana 中展开process的References即可看到对应的View Linked Span。Consumer Trace receive -- process |-- inventory.reserve - - Activity Link / View Linked Span - - - Producer Trace: send这条 Link 说明 Trace Context 已随消息跨越 Producer、Broker 到达 Consumer同时本地 Trace 树仍保留消息消费的异步与批量处理语义。仓库提供了一个端到端的 OpenTelemetry 示例。其中ProducerWebApi 和 ConsumerWebApi 分别运行在独立的 ASP.NET Core Minimal API Host 中。每个 Host 都注册 gRPC 与 Remoting 客户端Producer 通过协议专用路由发送消息Consumer 则使用彼此独立的 Consumer Group避免两种协议共享消费位点。示例会向 OTLP/gRPC Collector 导出数据并提供 Grafana Dashboard用于查看发送量、发送时延、处理时延和 Producer/Consumer Trace。下面是本地示例分别发送 gRPC 和 Remoting 消息后的 Producer dashboardConsumer dashboard 则展示接收、处理、确认和提交等消费侧操作在 Grafana Drilldown 中打开 Consumer 的 Trace并展开process的References后可以同时看到业务inventory.reserve子 span以及跳转到 Producersend的关联本地运行示例时先按前面的步骤启动 RocketMQ 测试环境再启动 Grafana OTEL LGTMdocker compose -f test-environments/otel-lgtm/compose.yaml up -d --wait ASPNETCORE_URLShttp://127.0.0.1:5242 \ dotnet run --no-launch-profile --project samples/opentelemetry/ConsumerWebApi ASPNETCORE_URLShttp://127.0.0.1:5241 \ dotnet run --no-launch-profile --project samples/opentelemetry/ProducerWebApi启动后访问http://127.0.0.1:5241/swagger调用/messages/grpc和/messages/remoting即可分别发送消息Consumer 状态位于http://127.0.0.1:5242/consumers。Grafana 默认地址为http://127.0.0.1:3000可查看预置的 Producer 和 Consumer dashboard。测试不只验证“能发送一条消息”#实现 RocketMQ 客户端时很多问题不会出现在简单的发送示例中。路由刷新、协议编码、并发回调、消费位点、异常响应、Broker 重连和服务端配置差异都需要有针对性的验证。项目将测试拆成两层单元测试覆盖 gRPC、Remoting 和兼容性相关的客户端行为不依赖网络或 Docker适合快速、确定性地验证纯客户端逻辑。集成测试使用 Testcontainers 启动真实的 Apache RocketMQ 5.5.0 环境。gRPC 测试覆盖 Proxy、NameServer、Broker 的真实链路Remoting 测试覆盖 NameServer 与 Broker 的直接通信路径。集成测试既有单 Broker 基线也有三个 Broker 的路由场景。gRPC 路径覆盖发送、消费、事务、Lite、死信和多 Broker 路由Remoting 路径覆盖 Admin、Pull、发送、请求-响应、事务、撤回和多 Broker 路由。每个 fixture 使用隔离 Docker 网络和动态端口不依赖开发机上的手工 Compose 环境。CI 会采集单元测试和 Docker 集成测试的覆盖率并上传到 Codecov。高覆盖率不能替代生产环境验证但它至少能让每次改动同时经过纯客户端行为和真实协议链路的检查。使用前需要知道的边界#EventHorizon.RocketMQ明确区分客户端能力、协议差异和服务端前提gRPC Package 只连接 RocketMQ Proxy不直接连接 NameServer 或 Broker它要求 RocketMQ 5 的 gRPC 服务端能力。Remoting Package 需要应用能够访问 NameServer 发现到的 Broker 地址容器网络与 Broker 对外注册地址必须正确。两个 Package 可以并存但公开模型、异常、发送结果和 Consumer 语义不能混用。Tag 过滤由客户端 API 支持SQL 过滤、POP、Lite、优先级、延迟消息撤回等功能还依赖 Broker 或 Proxy 的配置与版本。PushConsumer不等于服务端主动推送gRPC 的PushConsumer/LitePushConsumer与 classic Remoting 的PushConsumer都由客户端发起轮询或拉取。消费确认和重试提供的是消息系统语义不会自动解决重复处理。事务检查、重试 Handler 和消费业务都需要幂等设计。OpenTelemetry 包负责产生 Activity 与 MetricsCollector、Exporter、采样、存储和告警策略仍由应用和平台负责。总结#EventHorizon.RocketMQ是面向 .NET 8 及更高版本的非官方 Apache RocketMQ 客户端同时支持 RocketMQ 5 gRPC 与 classic Remoting。它以更贴合 .NET 应用习惯的方式提供 Producer、Consumer、Admin、DI、Generic Host、日志和 OpenTelemetry 支持。选择路径可以简单理解为新项目已部署 RocketMQ 5 Proxy 时应先评估官方 C# gRPC 客户端若选择本项目的 gRPC 实现则使用EventHorizon.RocketMQ.Grpc。需要对接 NameServer 和 Broker或依赖 Pull、LitePull、POP、Admin、请求-响应等经典能力时使用EventHorizon.RocketMQ.Remoting。同一应用有两类协议需求时可以同时注册两个 Package但应保持类型和客户端角色的协议边界。项目地址eventhorizon-cli/EventHorizon.RocketMQ完整功能说明仓库 README可运行示例samples测试说明本地与集成测试

相关新闻

引流DFM工具 VS 专业商用PCBA DFM软件,企业该如何取舍?

引流DFM工具 VS 专业商用PCBA DFM软件,企业该如何取舍?

2026/8/4 8:58:31

电子产品持续向高密度、小型化迭代,DFM可制造性前置审查已成为研发NPI、工艺量产必备环节。不少硬件工程师、企业采购都会陷入疑问:PCB厂、SMT平台推出的引流版DFM零成本,是否能替代付费商用PCBA DFM软件? 当前全球成熟商用PCBA D…

《碧蓝幻想Relink》终盘开荒指南:资源管理与副本策略详解

《碧蓝幻想Relink》终盘开荒指南:资源管理与副本策略详解

2026/8/4 8:48:31

1. 先搞清楚“开荒压力减半”到底指什么《碧蓝幻想:Relink》进入终盘,特别是面对“无尽黄昏”这类高难内容时,很多玩家的压力点其实很集中:不是打不过,而是“刷不动”。这里的“开荒”更多是指从零开始,为角…

基于NLP的新闻视频内容分析:从语音转写到情感分析的技术实践

基于NLP的新闻视频内容分析:从语音转写到情感分析的技术实践

2026/8/4 8:48:31

这次我们来看一个名为“CNN NEWS OUTFRONT 20260725-0700”的视频内容分析项目。这并非一个传统的软件开发工具或AI模型,而是一期具体的CNN新闻节目片段。对于技术博客读者而言,其核心价值在于提供了一个分析媒体内容、信息传播以及如何通过技术手段&…

5MB超轻量中文字体终极方案:WenQuanYi Micro Hei深度实战指南

5MB超轻量中文字体终极方案:WenQuanYi Micro Hei深度实战指南

2026/8/4 10:18:35

5MB超轻量中文字体终极方案:WenQuanYi Micro Hei深度实战指南 【免费下载链接】fonts-wqy-microhei Debian package for WenQuanYi Micro Hei (mirror of https://anonscm.debian.org/git/pkg-fonts/fonts-wqy-microhei.git) 项目地址: https://gitcode.com/gh_mi…

低成本共享ChatGPT Team与Codex API:代理网关方案与前端集成实践

低成本共享ChatGPT Team与Codex API:代理网关方案与前端集成实践

2026/8/4 10:18:35

1. 先搞清楚 GPT Team 会员、Codex 和前端对比到底在说什么 看到这个标题,很多人第一反应可能是“GPT-5.6”和“Fable5”这两个前端框架的对比评测。但仔细看,核心其实是两件事:一是如何低成本、稳定地使用 ChatGPT Team 会员 以及 Codex…

Cesium地形高度获取:原理、方法与应用场景

Cesium地形高度获取:原理、方法与应用场景

2026/8/4 10:18:35

1. 为什么需要获取地形高度? 在三维地理信息系统中,地形高度是最基础的空间数据之一。Cesium作为领先的WebGL三维地球引擎,其地形处理能力直接影响着场景的真实感和空间分析的准确性。获取地形高度的需求主要来自以下几个典型场景&#xff1a…

超元力MR大空间CS对战:文旅商业竞技业态落地实操分析

超元力MR大空间CS对战:文旅商业竞技业态落地实操分析

2026/8/4 10:18:35

一、行业现状与运营痛点走访近百家商场、FEC、景区项目后,共性经营矛盾集中在三点。其一,平峰客流分化严重,工作日日间场地空置率超60%,固定租金、人工成本持续消耗营收;其二,传统娱乐内容单次体验属性强&a…

方程解应用题专项训练方法与技巧

方程解应用题专项训练方法与技巧

2026/8/4 10:18:34

1. 方程解应用题专项训练的必要性 数学应用题一直是学生学习的难点和痛点。面对复杂的文字描述,很多同学往往无从下手,要么读不懂题意,要么列不出正确的方程。这个专项训练就是要帮助大家突破这个瓶颈,掌握用方程解决复杂应用题的…

MATLAB复小波变换与时频脊线提取技术详解

MATLAB复小波变换与时频脊线提取技术详解

2026/8/4 10:08:34

1. 项目概述:复小波变换与时频脊线提取 在信号处理领域,时频分析一直是研究非平稳信号特性的重要手段。传统傅里叶变换只能提供信号的全局频率信息,而无法反映频率成分随时间的变化情况。复小波变换作为一种时频分析工具,能够同时…

ncmdumpGUI:一键解锁网易云音乐ncm文件的终极解决方案

ncmdumpGUI:一键解锁网易云音乐ncm文件的终极解决方案

2026/8/3 4:49:52

ncmdumpGUI:一键解锁网易云音乐ncm文件的终极解决方案 【免费下载链接】ncmdumpGUI C#版本网易云音乐ncm文件格式转换,Windows图形界面版本 项目地址: https://gitcode.com/gh_mirrors/nc/ncmdumpGUI 你是否曾经从网易云音乐下载了心爱的歌曲&am…

分布式配置中心选型实战:Nacos与Consul在创业场景下的对比

分布式配置中心选型实战:Nacos与Consul在创业场景下的对比

2026/8/3 19:24:18

分布式配置中心选型实战:Nacos与Consul在创业场景下的对比工程导读:本文深入讨论 分布式配置中心选型实战:Nacos与Consul在创业场景下的对比 在生产工程实践中的核心落地方案。基于 分布式架构与微服务设计 视角,剖析实际痛点、架…

MoneyPrinterPlus实战指南:AI视频批量生成与自动化发布完整解决方案

MoneyPrinterPlus实战指南:AI视频批量生成与自动化发布完整解决方案

2026/8/3 20:38:37

MoneyPrinterPlus实战指南:AI视频批量生成与自动化发布完整解决方案 【免费下载链接】MoneyPrinterPlus AI一键批量生成各类短视频,自动批量混剪短视频,自动把视频发布到抖音,快手,小红书,视频号上,赚钱从来没有这么容易过! 支持本地语音模型chatTTS,fasterwhisper,…

3步解决Windows DLL缺失问题:VisualCppRedist AIO终极运行库修复方案

3步解决Windows DLL缺失问题:VisualCppRedist AIO终极运行库修复方案

2026/8/4 0:07:58

3步解决Windows DLL缺失问题:VisualCppRedist AIO终极运行库修复方案 【免费下载链接】vcredist AIO Repack for latest Microsoft Visual C Redistributable Runtimes 项目地址: https://gitcode.com/gh_mirrors/vc/vcredist 你是否曾经在打开游戏或软件时遇…

SingleFile终极指南:一键保存完整网页的5大核心功能

SingleFile终极指南:一键保存完整网页的5大核心功能

2026/8/4 0:07:58

SingleFile终极指南:一键保存完整网页的5大核心功能 【免费下载链接】SingleFile Web Extension for saving a faithful copy of a complete web page in a single HTML file 项目地址: https://gitcode.com/gh_mirrors/si/SingleFile 你是否曾经遇到过这样的…

国家中小学智慧教育平台电子课本下载终极方案:三步免费获取PDF教材

国家中小学智慧教育平台电子课本下载终极方案:三步免费获取PDF教材

2026/8/4 0:07:58

国家中小学智慧教育平台电子课本下载终极方案:三步免费获取PDF教材 【免费下载链接】tchMaterial-parser 国家中小学智慧教育平台 电子课本下载工具,帮助您从智慧教育平台中获取电子课本的 PDF 文件网址并进行下载,让您更方便地获取课本内容。…

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

2026/8/2 17:06:42

一天写完毕业论文在2026年已不再是天方夜谭。2026年最炸裂、实测能大幅提速的AI论文写作工具,覆盖选题构思、文献整理、内容生成、格式排版等核心场景,真正帮你高效搞定论文难题。 一、全流程王者:一站式搞定论文全链路(一天定稿首…

导师推荐!2026最新AI论文工具测评与实用推荐

导师推荐!2026最新AI论文工具测评与实用推荐

2026/8/3 7:25:44

2026年真正好用的AI论文工具,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。 一、…

告别游戏崩溃:XCOM 2模组管理器的智能革命

告别游戏崩溃:XCOM 2模组管理器的智能革命

2026/8/3 2:41:27

告别游戏崩溃:XCOM 2模组管理器的智能革命 【免费下载链接】xcom2-launcher The Alternative Mod Launcher (AML) is a replacement for the default game launchers from XCOM 2 and XCOM Chimera Squad. 项目地址: https://gitcode.com/gh_mirrors/xc/xcom2-lau…