
后端消息队列微服务消息路由【免费下载链接】CAPDistributed transaction solution in micro-service base on eventually consistency, also an eventbus with Outbox pattern项目地址https://gitcode.com/gh_mirrors/ca/CAP点击查看免费下载Apache Kafka 是 CAP 官方支持的核心消息传输器Transporter之一本指南基于仓库中 Kafka 传输器官方文档 展开结合 DotNetCore.CAP.Kafka 源码与仓库内 Kafka PostgreSQL 示例项目系统讲解从 NuGet 安装、UseKafka配置到MainConfig原生参数、CustomHeadersBuilder自定义消息头的完整用法并深入发送/消费两条链路的底层实现。读完你将掌握如何在 CAP 中启用 Kafka 作为消息队列、如何按需调优 Kafka 专属配置、如何通过自定义消息头对接异构系统以及 CAP 内部如何管理 Kafka 连接池与自动建 Topic。Kafka 在 CAP 中的角色Apache Kafka 是由 LinkedIn 发起并捐赠给 Apache 软件基金会的开源事件流平台使用 Scala 和 Java 编写。在 CAP 中Kafka 可以作为一个消息传输器Transporter使用CAP 负责 Outbox 模式的本地消息持久化与重试Kafka 负责消息的实际投递与消费。二者的边界非常清晰——CAP 通过ITransport抽象屏蔽具体消息队列差异KafkaCapOptionsExtension 在 DI 容器中注册了三个核心服务services.AddSingletonITransport, KafkaTransport(); services.AddSingletonIConsumerClientFactory, KafkaConsumerClientFactory(); services.AddSingletonIConnectionPool, ConnectionPool();即发送走KafkaTransport、消费由KafkaConsumerClientFactory创建客户端、连接复用依赖ConnectionPool理解这三者的职责有助于后面读懂各配置项的作用。安装与最小配置使用 Kafka 作为传输器需要先从 NuGet 安装以下包PM Install-Package DotNetCore.CAP.Kafka然后在Startup.cs或使用 minimal hosting 的Program.cs的ConfigureServices方法中添加配置public void ConfigureServices(IServiceCollection services) { // ... services.AddCap(x { x.UseKafka(opt { //KafkaOptions }); // x.UseXXX ... }); }UseKafka提供了两个重载见 CAP.Options.Extensions.csUseKafka(string bootstrapServers)直接传入 broker 地址以及UseKafka(ActionKafkaOptions configure)进行完整编程式配置。仓库示例 Sample.Kafka.PostgreSql/Program.cs 展示了最常见的组合——同时配置存储与传输builder.Services.AddCap(x { x.UsePostgreSql(AppConstants.DbConnectionString); // 存储 x.UseKafka(127.0.0.1:9092); // 传输 x.UseDashboard(); });注意 CAP 的最小配置要求是至少一个传输 一个存储参见 配置文档因此实际项目中UseKafka必须与UseSqlServer/UseMySql/UsePostgreSql/UseMongoDB/UseInMemoryStorage等存储配置成对出现。Kafka Options 参数详解CAP 直接提供的 Kafka 配置参数定义在 KafkaOptions完整参数如下名称说明类型默认值ServersBroker 服务器地址string必填MainConfiglibrdkafka 配置参数Dictionarystring, string见下方说明ConnectionPoolSize连接池大小int10CustomHeadersBuilder自定义订阅消息头FuncN/ARetriableErrorCodesConsumeException 时可重试的错误码IListErrorCode见代码TopicOptions新建 Topic 的分区数与副本因子配置KafkaTopicOptions-1逐个拆解如下Servers即bootstrap.servers以 CSV 形式列出初始 broker 列表host 或 host:port是MainConfig中该项的强类型入口。生产端与消费端都会读取它来建立连接。ConnectionPoolSize生产者连接池大小默认10。由 ConnectionPool 实现RentProducer()从并发队列取出空闲生产者不足时新建Return()在池未满时回收超出则直接Dispose()丢弃从而避免无上限地创建连接。RetriableErrorCodes消费过程中遇到ConsumeException时的可重试错误码集合。源码构造器中默认加入了 10 个错误码见 CAP.KafkaOptions.cs包括GroupLoadInProgress、Local_Retry、Local_TimedOut、RequestTimedOut、LeaderNotAvailable、NotLeaderForPartition、RebalanceInProgress、NotCoordinatorForGroup、NetworkException、GroupCoordinatorNotAvailable。消费循环里命中这些错误码时见 KafkaConsumerClient.cs会记日志并跳过本轮而不是直接失败。TopicOptionsKafkaTopicOptions类型包含NumPartitions新建 Topic 的分区数与ReplicationFactor副本因子默认均为-1即交给 Kafka 服务端按集群默认策略决定-1 通常表示采用 broker 端default.replication.factor/num.partitions。MainConfig注入 librdkafka 原生配置如果需要更多 Kafka 原生配置项可以在MainConfig配置字典中设置services.AddCap(capOptions { capOptions.UseKafka(kafkaOption { // kafka options. // kafkaOptions.MainConfig.Add(, ); }); });MainConfig是一个Dictionarystring, string支持的配置项完整清单可查阅 librdkafka 官方CONFIGURATION.md该文档随 librdkafka 项目维护仓库注释 CAP.KafkaOptions.cs 亦指向它。这些键值会被原样传递给ProducerConfig与ConsumerConfig源码见 IConnectionPool.Default.cs 与 KafkaConsumerClient.cs因此所有 librdkafka/Confluent.Kafka 支持的原生配置都能通过它透传。关闭 Topic 自动创建CAP 默认会在启动时自动创建所需的 Topic前提是 broker 允许。若想防止 CAP 自动创建 Topic改为由运维预先创建可以这样配置services.AddCap(capOptions { capOptions.UseKafka(kafkaOption { kafkaOption.MainConfig.Add(allow.auto.create.topics, false); }); });从源码看KafkaConsumerClient.FetchTopicsAsync 会先读取MainConfig中allow.auto.create.topics的值若为true默认行为则通过AdminClient调用CreateTopicsAsync按TopicOptions.NumPartitions与ReplicationFactor创建 TopicTopic 已存在时静默忽略already exists异常设为false后则完全跳过建 Topic 流程Topic 必须预先在 Kafka 集群中创建好。值得注意的是订阅 Topic 名称支持通配符FetchTopicsAsync内部会对订阅名调用Helper.WildcardToRegex转换为正则后再交给 Kafka这解释了为何自动建 Topic 使用的是转换后的正则名。生产者默认参数源码补充当MainConfig未显式指定时CAP 会为生产者补齐三个实用默认值见 IConnectionPool.Default.cs配置键默认值含义QueueBufferingMaxMessages10生产者本地缓冲队列最大消息数MessageTimeoutMs5000消息送达确认超时毫秒RequestTimeoutMs3000broker 请求超时毫秒消费端同样有默认补齐项见 KafkaConsumerClient.csAutoOffsetReset Earliest无提交 offset 时从头消费、EnableAutoCommit false关闭自动提交由 CAP 在消息成功处理后显式Commit这是 CAP 可靠消费的关键、GroupId默认取 CAP 的消费者组名。CustomHeadersBuilder自定义消息头当消息来自异构系统非 CAP 客户端发送时CAP 需要额外的头信息才能正确识别与路由消息通过CustomHeadersBuilder参数即可为订阅者补全这些自定义头。关于异构系统集成的详细说明见 消息Messaging文档。此外如果你希望把 broker 附带的额外上下文信息如 offset、partition写进消息也可以使用这个选项。例如x.UseKafka(opt { //... opt.CustomHeadersBuilder (kafkaResult, sp) new ListKeyValuePairstring, string { new KeyValuePairstring, string(my.kafka.offset, kafkaResult.Offset.ToString()), new KeyValuePairstring, string(my.kafka.partition, kafkaResult.Partition.ToString()) }; });随后即可在订阅方法中通过[FromCap] CapHeader读取这些自定义头[CapSubscribe(sample.kafka.postgrsql)] public void HeadersTest(DateTime value, [FromCap]CapHeader header) { var offset header[my.kafka.offset]; var partition header[my.kafka.partition]; }从源码看KafkaConsumerClient.ConsumeAsync 在把 broker 消息转成 CAP 的TransportMessage时先解析 Kafka 原生 Headers再依次注入Headers.Group消费者组名随后调用CustomHeadersBuilder(consumerResult, _serviceProvider)并把返回的键值对全部合并进消息头。注意CustomHeadersBuilder的入参是ConsumeResultstring, byte[]这意味着 offset、partition、timestamp、key 等ConsumeResult暴露的任何信息都可以被提取为自定义头。异构系统集成中的关键头从 CAP 3.0 起消息被拆分为 Header Body 传输Body 即用户Publish的原始内容不做包装Header 中携带 CAP 运行所需的关键信息。异构系统向 Kafka 发送消息时需要写入以下头详见 messaging.md键数据类型说明cap-msg-idlong消息 Id由雪花算法生成cap-msg-namestring消息名称cap-msg-typestring消息类型typeof(T).FullName非必需cap-senttimestring发送时间非必需cap-kafka-keystring按 Kafka Key 分区其中cap-kafka-key对应源码常量 KafkaHeaders.KafkaKey。KafkaTransport.SendAsync 在发送时会优先读取该头作为 Kafka 消息的Key用于分区路由未设置时回退为message.GetId()——即默认按消息 Id 分区保证同一消息的投递顺序性。发布时也可显式指定该头控制分区var headers new Dictionarystring, string?() { { cap-kafka-key, request.OrderId } }; _publisher.PublishOrderRequest(OrderRequest, request, headers);源码视角发送与消费链路发送链路KafkaTransport.SendAsync 的流程是从连接池RentProducer()→ 把 CAP 消息头逐条转成 KafkaHeaderUTF-8 编码→ 按上文规则确定Key→ProduceAsync投递到message.GetName()对应的 Topic。只有当PersistenceStatus为Persisted或PossiblyPersisted时才返回成功否则抛出PublisherSentFailedException包装为OperateResult.Failed——失败的发送会进入 CAP 的发送重试机制默认 3 次即时重试后按分钟递增最多 50 次见 messaging.md。连接池在finally中归还 producer。消费链路KafkaConsumerClient 的关键行为并行度控制构造函数按groupConcurrent创建SemaphoreSlimListeningAsync中每消费一条消息先Wait信号量再Task.Run异步处理CommitAsync提交 offset 后Release信号量从而限制同组并行消费数。可靠性EnableAutoCommit false仅当 CAP 的订阅执行器处理成功后OnMessageCallback正常返回才调用Commit手动提交处理失败走 CAP 消费重试流程。错误处理可重试错误码命中时记ConsumeRetries日志并继续连接级错误通过SetErrorHandler回调记录ServerConnError日志见 ConsumerClient_OnConsumeError。与 CAP 消息机制的配合Group 与 GroupConcurrentKafka 中Group对应Consumer Group参见 configuration.md相同Name且相同Group的订阅者共享消费只有一人收到相同Name但不同Group的订阅者都会收到消息。GroupConcurrent用于设置订阅者并行度若未指定GroupCAP 会自动以Name创建 Group。默认组名格式为cap.queue.{程序集名}也可通过CapOptions.DefaultGroupName自定义。事务消息与示例项目仓库的 Sample.Kafka.PostgreSql 是 Kafka PostgreSQL 组合的完整可运行示例其 ValuesController.cs 演示了 Ado.NET 事务、EF 事务、延迟消息与普通订阅的完整用法例如在本地数据库事务中原子地发布 Kafka 消息using (var transaction connection.BeginTransaction(producer, autoCommit: false)) { //your business code connection.Execute(INSERT INTO ..., transaction: (IDbTransaction)transaction.DbTransaction); producer.Publish(sample.kafka.postgrsql, DateTime.Now); transaction.Commit(); }该示例的 Program.cs 还附带了使用 KRaft 模式Kafka 3.7.0单节点无需 ZooKeeper启动本地 Kafka 的 Docker 命令可直接复制运行docker run -d --name kafka -p 9092:9092 -e KAFKA_NODE_ID1 -e KAFKA_PROCESS_ROLESbroker,controller -e KAFKA_LISTENERSPLAINTEXT://0.0.0.0:9092,CONTROLLER://:9093 -e KAFKA_ADVERTISED_LISTENERSPLAINTEXT://127.0.0.1:9092 -e KAFKA_CONTROLLER_LISTENER_NAMESCONTROLLER -e KAFKA_LISTENER_SECURITY_PROTOCOL_MAPCONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT -e KAFKA_CONTROLLER_QUORUM_VOTERS1localhost:9093 -e KAFKA_LOG_DIRS/var/lib/kafka/data -e KAFKA_AUTO_CREATE_TOPICS_ENABLEtrue -e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR1 -e KAFKA_OFFSETS_TOPIC_MIN_ISR1 -e KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR1 -e KAFKA_TRANSACTION_STATE_LOG_MIN_ISR1 apache/kafka:3.7.0启动后把UseKafka指向127.0.0.1:9092并访问~/without/transaction、~/adonet/transaction等端点即可观察消息的发布与Test2订阅者的消费输出。小结CAP Kafka 的组合要点可归纳为通过UseKafka注册传输器并与任一存储搭配用Servers指定 broker、用MainConfig透传 librdkafka 原生参数含关闭自动建 Topic用ConnectionPoolSize与RetriableErrorCodes调优连接与容错用CustomHeadersBuilder补齐异构系统消息头或注入 offset/partition 上下文订阅端用[FromCap] CapHeader读取这些头并可通过cap-kafka-key控制分区路由。配合 传输器总览、消息机制 与 存储选型 文档即可在生产环境中落地一套基于 Kafka 的可靠事件总线。赞分享后端消息队列微服务消息路由【免费下载链接】CAPDistributed transaction solution in micro-service base on eventually consistency, also an eventbus with Outbox pattern项目地址https://gitcode.com/gh_mirrors/ca/CAP点击查看免费下载相关推荐CAP 中使用 Apache Kafka 作为消息传输器配置、源码级原理与异构系统集成CAP 中使用 Apache Kafka 作为消息传输器配置、源码级原理与异构系统集成 导读 本文围绕开源分布式事务解决方案 The NCC / CAP ht后端消息队列微服务CAP 集成 Apache Pulsar以 Pulsar 作为消息传输器的配置与实现原理CAP 集成 Apache Pulsar以 Pulsar 作为消息传输器的配置与实现原理 Apache Pulsar 是诞生于 Yahoo!、现为 Apach后端消息队列微服务如何利用Apache Arrow与Kafka构建高效数据处理管道完整指南如何利用Apache Arrow与Kafka构建高效数据处理管道完整指南 Apache Arrow是一个多语言工具集专为加速数据交换和内存处理而设计。当与K数据工程大数据序列化数据分析上一篇Handsontable 编辑与校验实战级联下拉与行级校验错误汇总配方解析下一篇GenUI开源贡献者手册如何参与下一代UI框架开发创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考