- 后端
- 消息队列
- 微服务
- 消息路由
【免费下载链接】CAP
Distributed transaction solution in micro-service base on eventually consistency, also an eventbus with Outbox pattern
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.AddSingleton<ITransport, KafkaTransport>(); services.AddSingleton<IConsumerClientFactory, KafkaConsumerClientFactory>(); services.AddSingleton<IConnectionPool, 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.cs):UseKafka(string bootstrapServers)直接传入 broker 地址,以及UseKafka(Action<KafkaOptions> 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,完整参数如下:
| 名称 | 说明 | 类型 | 默认值 |
|---|---|---|---|
| Servers | Broker 服务器地址 | string | (必填) |
| MainConfig | librdkafka 配置参数 | Dictionary<string, string> | 见下方说明 |
| ConnectionPoolSize | 连接池大小 | int | 10 |
| CustomHeadersBuilder | 自定义订阅消息头 | Func<> | N/A |
| RetriableErrorCodes | ConsumeException 时可重试的错误码 | IList<ErrorCode> | 见代码 |
| 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)会记日志并跳过本轮,而不是直接失败。 - TopicOptions:
KafkaTopicOptions类型,包含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是一个Dictionary<string, 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创建 Topic,Topic 已存在时静默忽略already exists异常;设为false后则完全跳过建 Topic 流程,Topic 必须预先在 Kafka 集群中创建好。
值得注意的是,订阅 Topic 名称支持通配符:FetchTopicsAsync内部会对订阅名调用Helper.WildcardToRegex转换为正则后再交给 Kafka(这解释了为何自动建 Topic 使用的是转换后的正则名)。
生产者默认参数(源码补充)
当MainConfig未显式指定时,CAP 会为生产者补齐三个实用默认值(见 IConnectionPool.Default.cs):
| 配置键 | 默认值 | 含义 |
|---|---|---|
| QueueBufferingMaxMessages | 10 | 生产者本地缓冲队列最大消息数 |
| MessageTimeoutMs | 5000 | 消息送达确认超时(毫秒) |
| RequestTimeoutMs | 3000 | broker 请求超时(毫秒) |
消费端同样有默认补齐项(见 KafkaConsumerClient.cs):AutoOffsetReset = 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 List<KeyValuePair<string, string>> { new KeyValuePair<string, string>("my.kafka.offset", kafkaResult.Offset.ToString()), new KeyValuePair<string, 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的入参是ConsumeResult<string, byte[]>,这意味着 offset、partition、timestamp、key 等ConsumeResult暴露的任何信息都可以被提取为自定义头。
异构系统集成中的关键头
从 CAP 3.0 起,消息被拆分为 Header + Body 传输,Body 即用户Publish的原始内容不做包装,Header 中携带 CAP 运行所需的关键信息。异构系统向 Kafka 发送消息时,需要写入以下头(详见 messaging.md):
| 键 | 数据类型 | 说明 |
|---|---|---|
| cap-msg-id | long | 消息 Id,由雪花算法生成 |
| cap-msg-name | string | 消息名称 |
| cap-msg-type | string | 消息类型typeof(T).FullName(非必需) |
| cap-senttime | string | 发送时间(非必需) |
| cap-kafka-key | string | 按 Kafka Key 分区 |
其中cap-kafka-key对应源码常量 KafkaHeaders.KafkaKey。KafkaTransport.SendAsync 在发送时会优先读取该头作为 Kafka 消息的Key(用于分区路由),未设置时回退为message.GetId()——即默认按消息 Id 分区,保证同一消息的投递顺序性。发布时也可显式指定该头控制分区:
var headers = new Dictionary<string, string?>() { { "cap-kafka-key", request.OrderId } }; _publisher.Publish<OrderRequest>("OrderRequest", request, headers);源码视角:发送与消费链路
发送链路
KafkaTransport.SendAsync 的流程是:从连接池RentProducer()→ 把 CAP 消息头逐条转成 KafkaHeader(UTF-8 编码)→ 按上文规则确定Key→ProduceAsync投递到message.GetName()对应的 Topic。只有当PersistenceStatus为Persisted或PossiblyPersisted时才返回成功,否则抛出PublisherSentFailedException包装为OperateResult.Failed——失败的发送会进入 CAP 的发送重试机制(默认 3 次即时重试后按分钟递增,最多 50 次,见 messaging.md)。连接池在finally中归还 producer。
消费链路
KafkaConsumerClient 的关键行为:
- 并行度控制:构造函数按
groupConcurrent创建SemaphoreSlim,ListeningAsync中每消费一条消息先Wait信号量再Task.Run异步处理,CommitAsync提交 offset 后Release信号量,从而限制同组并行消费数。 - 可靠性:
EnableAutoCommit = false,仅当 CAP 的订阅执行器处理成功后(OnMessageCallback正常返回)才调用Commit手动提交;处理失败走 CAP 消费重试流程。 - 错误处理:可重试错误码命中时记
ConsumeRetries日志并继续;连接级错误通过SetErrorHandler回调记录ServerConnError日志(见 ConsumerClient_OnConsumeError)。
与 CAP 消息机制的配合
Group 与 GroupConcurrent
Kafka 中Group对应Consumer Group(参见 configuration.md):相同Name且相同Group的订阅者共享消费(只有一人收到);相同Name但不同Group的订阅者都会收到消息。GroupConcurrent用于设置订阅者并行度,若未指定Group,CAP 会自动以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_ID=1 ` -e KAFKA_PROCESS_ROLES=broker,controller ` -e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092,CONTROLLER://:9093 ` -e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://127.0.0.1:9092 ` -e KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER ` -e KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT ` -e KAFKA_CONTROLLER_QUORUM_VOTERS=1@localhost:9093 ` -e KAFKA_LOG_DIRS=/var/lib/kafka/data ` -e KAFKA_AUTO_CREATE_TOPICS_ENABLE=true ` -e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 ` -e KAFKA_OFFSETS_TOPIC_MIN_ISR=1 ` -e KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR=1 ` -e KAFKA_TRANSACTION_STATE_LOG_MIN_ISR=1 ` 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 的可靠事件总线。
- 后端
- 消息队列
- 微服务
- 消息路由
【免费下载链接】CAP
Distributed transaction solution in micro-service base on eventually consistency, also an eventbus with Outbox pattern
相关推荐
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
数据工程大数据序列化数据分析
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考