Apache Kafka 中间件
包名: Snet.Kafka | 类: KafkaOperate | 基类: MqAbstract<KafkaOperate, KafkaData.Basics>
通过 Apache Kafka 提供高吞吐量的分布式发布/订阅消息。支持生产者和消费者操作,可配置安全策略和 SASL 认证。
概述
KafkaOperate 连接到 Apache Kafka 集群,通过 IMq(IProducer + IConsumer)提供消息生产和消费功能。它继承自 MqAbstract<KafkaOperate, KafkaData.Basics>。
安装
dotnet add package Snet.Kafka
快速开始
using Snet.Kafka;
var kafka = await KafkaOperate.InstanceAsync(new KafkaData.Basics
{
BootstrapServers = "localhost:9092"
});
await kafka.OnAsync();
// 绑定事件再消费
kafka.OnDataEventAsync += async (sender, e) =>
{
if (e.Status)
Console.WriteLine($"接收到数据: {e.ResultData}");
else
Console.WriteLine($"消费失败: {e.Message}");
};
await kafka.ConsumeAsync("sensor-data");
// 生产消息(可选)
await kafka.ProduceAsync("sensor-data", "hello kafka", System.Text.Encoding.UTF8);
// 后续:取消消费
// await kafka.UnConsumeAsync("sensor-data");
// await kafka.DisposeAsync();
生产与消费
发布与订阅是 IMq 的一等操作——先绑定数据事件,再消费;发布可用字符串或原始字节:
// 1) 先绑定事件再消费——收到的消息在此到达
kafka.OnDataEventAsync += async (sender, e) =>
{
if (e.Status)
Console.WriteLine($"收到: {e.ResultData}");
else
Console.WriteLine($"消费失败: {e.Message}");
};
// 2) 消费:订阅主题
await kafka.ConsumeAsync("sensor-data");
// 3) 生产:发布字符串消息(UTF-8)
await kafka.ProduceAsync("sensor-data", "hello kafka", System.Text.Encoding.UTF8);
// 4) 生产:发布原始字节
var bytes = System.Text.Encoding.UTF8.GetBytes("payload");
await kafka.ProduceAsync("sensor-data", bytes);
// 后续:停止消费
// await kafka.UnConsumeAsync("sensor-data");
ProduceAsync(topic, string, Encoding?)与ProduceAsync(topic, byte[])均可发布;ConsumeAsync/UnConsumeAsync管理主题订阅。事件见:事件。
配置
| 参数 | 类型 | 默认值 | 描述 |
|---|---|---|---|
SN |
string | — | 序列号 / 设备标识符 |
BootstrapServers |
string | — | 逗号分隔的 host:port 列表,用于初始集群连接 |
SecurityProtocol |
enum | Plaintext | 安全协议:Plaintext, Ssl, SaslPlaintext, SaslSsl |
SaslMechanism |
enum | Gssapi | SASL 机制:Gssapi, Plain, ScramSha256, ScramSha512, OAuthBearer |
SaslKerberosServiceName |
string | "snet" |
SASL GSSAPI 认证的 Kerberos 服务名称 |
GroupId |
string | "snet" |
固定消费组——同组实例分摊分区消息 |
AutoOffsetReset |
enum | Latest | 偏移量重置行为:Latest, Earliest, Error |
ResponseType |
enum | Content | 响应类型 |
说明
仅 Kerberos (GSSAPI) 实际可用(Plain/SCRAM 在现有 Basics 中无凭据注入路径——没有 SASL 用户名/密码字段)。
消费组说明: GroupId(默认 "snet")为固定消费组——同组实例分摊消息,已提交偏移量跨重启延续(从上次位置继续消费,而非每次从 Latest 开始);不同组各自全量消费。AutoOffsetReset 仅在组内无已提交偏移量时生效。偏移量为手动提交(EnableAutoCommit=false),每条消息在事件派发后提交。
支持的操作
| 操作 | 方法 | 描述 |
|---|---|---|
| 连接 | OnAsync() |
初始化生产者/消费者/管理员客户端(不做真实连接探测——Confluent 客户端惰性连接) |
| 断开 | OffAsync() |
断开与集群的连接 |
| 发布 | ProduceAsync(topic, string, Encoding?) |
向主题发布字符串消息 |
| 发布 | ProduceAsync(topic, byte[]) |
向主题发布原始字节 |
| 发布(带 Key) | Produce(topic, key, content) / ProduceAsync(topic, key, content, token) |
携带分区键发布消息(内容为字符串或 byte[]) |
| 创建主题 | CreateTopics(List<string>) / CreateTopicsAsync(List<string>, token) |
在集群上创建主题 |
| 订阅 | ConsumeAsync(topic) |
作为消费者订阅主题 |
| 取消订阅 | UnConsumeAsync(topic) |
取消订阅主题 |
主题分区
Kafka 主题被划分为多个分区以实现可扩展性。主题内的消息按分区排序:
- 生产者:指定键来控制分区路由;相同键的消息发送到同一分区
- 消费者:同一组中的消费者在它们之间分配分区,实现并行处理
事件
| 事件 | 签名 | 描述 |
|---|---|---|
OnDataEvent |
EventHandler<EventDataResult> |
收到消费消息时触发 |
OnDataEventAsync |
EventHandlerAsync<EventDataResult> |
OnDataEvent 的异步变体 |
OnInfoEvent |
EventHandler<EventInfoResult> |
收到信息和状态消息时触发 |
OnInfoEventAsync |
EventHandlerAsync<EventInfoResult> |
OnInfoEvent 的异步变体 |
OnLanguageEvent |
EventHandler<EventLanguageResult> |
语言变更时触发 |
OnLanguageEventAsync |
EventHandlerAsync<EventLanguageResult> |
OnLanguageEvent 的异步变体 |
