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 = new KafkaOperate(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();
配置
| 参数 | 类型 | 默认值 | 描述 |
|---|---|---|---|
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 服务名称 |
AutoOffsetReset |
enum | Latest | 偏移量重置行为:Earliest, Latest |
ResponseType |
enum | Content | 响应类型 |
支持的操作
| 操作 | 方法 | 描述 |
|---|---|---|
| 连接 | OnAsync() |
连接到 Kafka 集群 |
| 断开 | OffAsync() |
断开与集群的连接 |
| 发布 | ProduceAsync(topic, string, Encoding?) |
向主题发布字符串消息 |
| 发布 | ProduceAsync(topic, byte[]) |
向主题发布原始字节 |
| 订阅 | 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 的异步变体 |
