资讯详情

资讯详情

在 CAP 中使用 Redis Streams 作为消息传输器:配置、原理与最佳实践

后端消息队列微服务【免费下载链接】CAP基于最终一致性的微服务分布式事务解决方案也是一种采用 Outbox 模式的事件总线。项目地址https://gitcode.com/dotnetcore/CAP点击查看免费下载导读Redis Streams 是 Redis 5.0 引入的仅追加append-only日志式数据结构天然适合充当消息中间件。本文以 CAP 官方文档为核心结合DotNetCore.CAP.RedisStreams包在仓库中的真实源码实现系统讲解如何将 Redis Streams 接入 CAP 事件总线、全部配置项的含义与默认值、消息的发布与消费流程以及流清理等生产环境注意事项。读完本文你将能够独立完成 Redis Streams 传输器的安装、配置与调优并理解其底层消费组与消息确认机制。一、为什么选择 Redis Streams 作为 CAP 的传输器CAP 是一个基于最终一致性理念、采用 Outbox 模式的分布式事务解决方案与事件总线。其消息传输器Transporter是可插拔的官方仓库中已内置 Kafka、RabbitMQ、Azure Service Bus、NATS、Pulsar、AWS SQS、Redis Streams 等多种实现。Redis 本身是开源BSD 许可、内存型的数据结构存储系统可用作数据库、缓存和消息中间件。Redis Stream 是 Redis 5.0 引入的一种新数据类型它以仅附加的数据结构抽象地模拟了日志数据结构——消息按顺序追加、按 ID 索引、可被多个消费者组独立消费这与事件总线「一次写入、多方订阅」的语义高度吻合。因此在只需要轻量级消息传递、且基础设施中已存在 Redis 的场景下Redis Streams 是 CAP 传输器的理想选择。二、安装与基础配置要使用 Redis Streams 传输器首先需要从 NuGet 安装以下包PM Install-Package DotNetCore.CAP.RedisStreams安装完成后在Startup.cs或 .NET 6 的Program.cs的ConfigureServices方法中添加基于 Redis Stream 的配置public void ConfigureServices(IServiceCollection services) { services.AddCap(capOptions { capOptions.UseRedis(redisOptions { //redisOptions }); }); }从源码看UseRedis扩展方法定义在 src/DotNetCore.CAP.RedisStreams/CapOptions.Redis.Extensions.cs共提供三个重载UseRedis()完全使用默认配置UseRedis(string connection)直接传入逗号分隔的连接字符串内部会调用ConfigurationOptions.Parse(connection)解析为ConfigurationOptionsUseRedis(ActionCapRedisOptions configure)通过委托精细配置各项参数这是最常用的形式。三者最终都会通过options.RegisterExtension(new RedisOptionsExtension(configure))注册扩展。在 src/DotNetCore.CAP.RedisStreams/ICapOptionsExtension.Redis.cs 中可以看到注册时会向容器依次注入IRedisStreamManager流管理器、IConsumerClientFactory消费者客户端工厂、ITransport传输器即RedisTransport以及IRedisConnectionPool连接池实现完整的传输器闭环。三、Redis Streams 配置参数详解CAP 直接对外提供的 Redis Stream 配置参数如下NAMEDESCRIPTIONTYPEDEFAULTConfigurationRedis 连接配置StackExchange.RedisConfigurationOptionsConfigurationOptionsStreamEntriesCount读取时从 stream 返回的条目数uint10ConnectionPoolSize连接池数uint10OnConsumeError消费消息发生错误时调用的回调函数FuncConsumeErrorContext, Tasknull这些属性定义在 src/DotNetCore.CAP.RedisStreams/CapOptions.Redis.cs 的CapRedisOptions类中。逐一说明其作用与底层影响3.1 Configuration连接配置类型为ConfigurationOptions对应 StackExchange.Redis 的原生连接配置对象。这是整个传输器的连接基石内部属性Endpoint Configuration?.ToString() ?? string.Empty会将其序列化为 Broker 地址字符串同时RedisTransport与RedisConsumerClient中暴露的BrokerAddress都直接读取该值。注意一个容易踩坑的默认行为在CapRedisOptionsPostConfigure见 src/DotNetCore.CAP.RedisStreams/ICapOptionsExtension.Redis.cs中若Configuration为 null 则自动创建空实例若未显式添加任何 EndPoint会自动补上IPAddress.Loopback, 0并调用SetDefaultPorts()默认 6379。这意味着你不配置连接地址时CAP 会默认尝试连接本机 Redis——建议生产中务必显式指定。3.2 StreamEntriesCount单次读取条目数类型为uint默认值 10。它决定每次调用StreamReadGroupAsync时从 stream 拉取的条目数量直接控制消费的吞吐与批次粒度。默认值由CapRedisOptionsPostConfigure中的if (options.StreamEntriesCount default) options.StreamEntriesCount 10;兜底赋值。在 src/DotNetCore.CAP.RedisStreams/IRedisStream.Manager.Default.cs 的TryReadConsumerGroupAsync中可以看到它的真实用途CAP 会先按 key 的 HashSlot 对 stream 分组再对每组执行StreamReadGroupAsync(..., (int)_options.StreamEntriesCount)并Task.WhenAll并发读取最后合并结果。它还是判断「待处理消息是否已消费完毕」的关键PollStreamsPendingMessagesAsync中若某轮读取到的条目数全部小于该值则认为历史积压消息已排空从而退出恢复循环转入新消息监听。3.3 ConnectionPoolSize连接池大小类型为uint默认值 10。它控制RedisConnectionPool中维护的IConnectionMultiplexer懒连接数量见 src/DotNetCore.CAP.RedisStreams/IConnectionPool.Default.cs。连接池在构造函数中一次性创建指定数量的AsyncLazyRedisConnectionConnectAsync时会优先返回尚未创建空闲的连接并按ConnectionCapacity排序选取负载最轻的连接复用从而在 Redis 与 CAP 之间分摊连接压力。默认值同样由PostConfigure兜底为 10。3.4 OnConsumeError消费错误回调类型为FuncConsumeErrorContext, Task默认 null。当消费消息发生错误时被调用其上下文ConsumeErrorContext是一个 record携带两个字段Exception Exception异常对象与StreamEntry? Entry出错的流条目。在 src/DotNetCore.CAP.RedisStreams/IConsumerClient.Redis.cs 的ConsumeAsync中消费异常会先记录日志、构造LogMessageEventArgsMqLogType.RedisConsumeError随后尝试调用_options.Value.OnConsumeError?.Invoke(...)若回调自身抛异常则单独记录日志而不会中断主流程最后通过OnLogCallback上报。注意这里的回调只是让你感知与处理消费异常消息的确认/重试仍由 CAP 的存储与重试机制负责——失败消息不会在这里被 ACK而是等待 CAP 后续的重试处理。四、更细粒度的原生 Redis 配置如果你需要更多 StackExchange.Redis 原生配置选项超时、密码、SSL、连接复用等可以在Configuration选项中设置services.AddCap(capOptions { capOptions.UseRedis(redisOptions { // redis options. redisOptions.Configuration.EndPoints.Add(IPAddress.Loopback, 0); }); });Configuration是 StackExchange.Redis 的ConfigurationOptions官方文档中提供了完整的可用选项列表。除逐个属性赋值外UseRedis的字符串重载也支持直接传入 StackExchange.Redis 的逗号分隔配置语法例如仓库示例 samples/Samples.Redis.SqlServer/Program.cs 中的写法redis.Configuration ConfigurationOptions.Parse(redis-node-0:6379,passwordcap);五、消息在 Redis Streams 中的流转原理结合源码一条消息从发布到消费的完整链路如下1. 发布阶段。发布方通过ICapPublisher发布消息后RedisTransport.SendAsync见 src/DotNetCore.CAP.RedisStreams/ITransport.Redis.cs调用_redis.PublishAsync(message.GetName(), message.AsStreamEntries())。其中消息名message.GetName()作为stream 的 key即一个 topic 对应一个 streammessage.AsStreamEntries()由 src/DotNetCore.CAP.RedisStreams/TransportMessage.Redis.cs 的RedisMessage静态类实现将 CAP 消息拆成两个字段——headers头信息 JSON和body正文二进制 JSON通过StreamAddAsync追加到 stream发送失败时捕获异常并包装为PublisherSentFailedException返回OperateResult.Failed。2. 消费组创建。消费者订阅 topic 时RedisConsumerClient.SubscribeAsync会对每个 topic 调用CreateStreamWithConsumerGroupAsync。其底层见 src/DotNetCore.CAP.RedisStreams/IRedisStream.Manager.Extensions.cs实现了幂等创建若 stream 已存在则检查StreamGroupInfoAsync中是否已有同名消费者组存在则跳过不存在或 stream 不存在时调用StreamCreateConsumerGroupAsync(stream, consumerGroup, StreamPosition.NewMessages)创建并捕获GroupAlreadyExists错误保证并发安全。3. 消费阶段先补历史再听新消息。RedisConsumerClient.ListeningForMessagesAsync见 src/DotNetCore.CAP.RedisStreams/IConsumerClient.Redis.cs的执行策略非常关键先调用PollStreamsPendingMessagesAsync从StreamPosition.Beginning开始消费积压的未确认消息用于崩溃恢复——重启后先处理完 pending 消息避免消息丢失待 pending 排空后判断依据即上文提到的StreamEntriesCount比较逻辑再启动PollStreamsLatestMessagesAsync从StreamPosition.NewMessages开始持续监听新消息每次轮询间隔为传入的timeout每个 stream entry 通过RedisMessage.Create(entry, _groupId)还原为TransportMessage并把消费组 ID 写入Headers.Group反序列化失败缺 headers/body、JSON 非法会抛出对应的RedisConsumeMissingHeadersException等专用异常消费成功后由 CAP 调用CommitAsync内部执行StreamAcknowledgeAsync(stream, group, id)完成 ACK见RedisStreamManager.Ack并释放并发信号量RejectAsync则只释放信号量、不做 ACK把重试交给 CAP 的存储机制。4. 并发控制。RedisConsumerClient通过SemaphoreSlim(_groupConcurrent)控制同一组内的并发消费数量当_groupConcurrent 0时每条消息以Task.Run异步并发消费并受信号量限流为 0 时则串行await消费。该值来自 CAP 订阅者组配置可从源码推断其语义为「同一消费组内最大并发数」。六、完整示例Redis Streams SQL Server仓库中的 samples/Samples.Redis.SqlServer 演示了 Redis Streams 传输器与 SQL Server 存储、Dashboard 的完整组合。启动配置如下builder.Services.AddCap(options { options.UseRedis(redis { redis.Configuration ConfigurationOptions.Parse(redis-node-0:6379,passwordcap); redis.OnConsumeError context { throw new InvalidOperationException(); }; }); options.UseSqlServer(Serverdb;Databasemaster;Usersa;PasswordPssw0rd;EncryptFalse); options.UseDashboard(); });发布与订阅见 samples/Samples.Redis.SqlServer/Controllers/HomeController.cs[HttpGet] public async Task Publish([FromQuery] string message test-message) { await _publisher.PublishAsync(message, new Person() { Age 11, Name James }); } [CapSubscribe(test-message)] [CapSubscribe(test-message-1)] [CapSubscribe(test-message-2)] [CapSubscribe(test-message-3)] [NonAction] public void Subscribe(Person p, [FromCap] CapHeader header) { _logger.LogInformation(${header[Headers.MessageName]} subscribed with value -- p); }该示例展示了 CAP 的核心用法一个方法可通过多个[CapSubscribe]同时订阅多个 topic每个 topic 即一个 Redis stream发布的消息名与订阅名一一对应消息体自动反序列化为Person强类型对象。七、流清理注意事项生产必读由于 Redis Streams没有自动删除所有已被全部消费组确认ACK的消息的特性见 Redis 官方 issue #5774stream 中的数据会随消息不断累积、持续占用内存。因此你需要考虑编写脚本定期清理已被所有消费者组确认的旧消息例如通过XINFO GROUPS查询各组的lag/last-delivered-id或结合XACK与XDEL/XTRIM按策略裁剪评估消息留存需求与内存容量合理设置清理周期如每天凌晨执行对于超长留存的历史消息也可考虑用XTRIM按长度或时间MAXLEN/MINID裁剪 stream但需确认不影响尚未消费的 pending 消息。八、小结Redis Streams 让 CAP 在不引入额外消息中间件的情况下获得可靠的事件总线能力。本篇文章覆盖了从 NuGet 安装、UseRedis三种重载、四个核心配置项Configuration、StreamEntriesCount、ConnectionPoolSize、OnConsumeError到消息发布/消费/ACK 的完整源码链路并通过仓库示例给出了可落地的完整配置。最后再次提醒务必显式配置 Redis 连接地址默认回落到本机并在生产环境规划好 stream 的定期清理策略。赞分享后端消息队列微服务【免费下载链接】CAP基于最终一致性的微服务分布式事务解决方案也是一种采用 Outbox 模式的事件总线。项目地址https://gitcode.com/dotnetcore/CAP点击查看免费下载相关推荐使用 Redis Streams 作为 CAP 消息传输器配置、消费原理与运维要点使用 Redis Streams 作为 CAP 消息传输器配置、消费原理与运维要点 导读 本文聚焦 CAP 的 Redis Streams 传输模块围绕 d后端消息队列微服务消息路由CAP 使用 AWS SNS SQS 作为消息传输器原理、配置与源码解析CAP 使用 AWS SNS SQS 作为消息传输器原理、配置与源码解析 导读 本文基于 CAP .NET Core CAP https://link.后端消息队列微服务消息路由CAP 集成 Redis Streams 消息传输器配置详解、参数说明与源码级原理剖析CAP 集成 Redis Streams 消息传输器配置详解、参数说明与源码级原理剖析 导读 本文以 CAP 的 Redis Streams 传输器为主题系后端消息队列微服务消息路由上一篇S905L3 电视盒刷机教程把 UNT400G 改造为 Armbian Linux 服务器下一篇解决PT助手Plus高内存占用问题6个实用优化技巧创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →