Kafka - Snet Docs

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 的异步变体

另请参阅