NewLife/NewLife.JT808

feat: M19-5/M20/M21 云厂商适配+Kafka消息队列+竞品分析更新

- M19-5 云厂商适配:RocketMQ 适配器添加 ICloudProvider 可选参数,支持阿里云/华为云/腾讯云/ACL
- M20 Kafka 消息队列:4 个适配器 (Producer/Consumer/SessionAdapter/DownHandler)
- M21 竞品分析更新:粤标/RocketMQ/Kafka 状态同步
- 更新功能清单、需求文档、Readme
大石头 authored at 2026-07-15 17:32:33
eb7ea4b
Tree
1 Parent(s) 30c17be
Summary: 13 changed files with 479 additions and 21 deletions.
Modified +22 -4
Modified +6 -10
Modified +1 -1
Modified +4 -0
Added +139 -0
Added +117 -0
Added +97 -0
Added +78 -0
Modified +3 -1
Modified +3 -1
Modified +3 -1
Modified +3 -1
Modified +3 -2
Modified +22 -4
diff --git "a/Doc/\345\212\237\350\203\275\346\270\205\345\215\225.md" "b/Doc/\345\212\237\350\203\275\346\270\205\345\215\225.md"
index bba2097..a89ad41 100644
--- "a/Doc/\345\212\237\350\203\275\346\270\205\345\215\225.md"
+++ "b/Doc/\345\212\237\350\203\275\346\270\205\345\215\225.md"
@@ -1,12 +1,12 @@
 # NewLife.JT808 功能清单
 
-> 版本:v2.5 | 日期:2026-07-15 | **状态:全部完成 🎉**
+> 版本:v2.6 | 日期:2026-07-15 | **状态:全部完成 🎉**
 
 > 状态标记:✅ 已实现 | 🟡 部分实现 | 🔧 规划中 | ⏸ 暂缓/占位 | ❌ 未开始
 
 本清单维护"核心目标 → 功能模块/子模块 → 完成状态"。模块名优先对齐[需求文档](需求文档.md)第 3 章的功能点名称;需要更细设计时跳转到[架构设计](架构设计.md)。
 
-> 🔍 **审计说明**:测试覆盖 M1/M2/M3/M4/M5/M6 共 59 个测试用例。M18 粤标/M19 RocketMQ 为新增模块,待补充专用测试。XML 注释存在 CS1591 缺失(消息体模型类)。
+> 🔍 **审计说明**:测试覆盖 M1/M2/M3/M4/M5/M6 共 59 个测试用例。M19 RocketMQ 适配器已补 M19-5 云厂商适配。M20 Kafka 消息队列为新模块,待补充专用测试。XML 注释存在 CS1591 缺失(消息体模型类)。
 
 ## M1 核心协议层(Newlife.JT808 — 已完成 ✅)
 
@@ -261,6 +261,22 @@
 | M19-2 | RocketMQ 消费者适配器 | ✅ | `RocketMQConsumer<T>` 实现 `IMsgConsumer<T>` |
 | M19-3 | 会话通知适配器 | ✅ | `RocketMQSessionAdapter` 实现 `ISessionProducer` |
 | M19-4 | 下行指令适配器 | ✅ | `RocketMQDownHandler` 实现 `IDownMessageHandler` |
+| M19-5 | 云厂商适配 | ✅ | 所有 RocketMQ 适配器支持阿里云/华为云/腾讯云/ACL 四种认证模式 |
+
+## M20 Kafka 消息队列(Newlife.JT808 — 已完成 ✅)
+
+| 编码 | 功能 | 状态 | 说明 |
+|------|------|------|------|
+| M20-1 | Kafka 生产者适配器 | ✅ | `KafkaProducer<T>` 实现 `IMsgProducer<T>`,需要 .NET Standard 2.0+ |
+| M20-2 | Kafka 消费者适配器 | ✅ | `KafkaConsumer<T>` 实现 `IMsgConsumer<T>`,支持消费组 + ACK |
+| M20-3 | 会话通知适配器 | ✅ | `KafkaSessionAdapter` 实现 `ISessionProducer`,上下线事件广播 |
+| M20-4 | 下行指令适配器 | ✅ | `KafkaDownHandler` 实现 `IDownMessageHandler`,消费指令队列 |
+
+## M21 竞品分析更新(Doc — 已完成 ✅)
+
+| 编码 | 功能 | 状态 | 说明 |
+|------|------|------|------|
+| M21-1 | 竞品分析状态修正 | ✅ | 粤标/RocketMQ 状态从"规划"更新为"已完成";新增 Kafka 对比行 |
 
 ## 统计
 
@@ -285,5 +301,7 @@
 | M16 Web 管理后台 | 19 | 19 | 100% | ✅ 阶段4 |
 | M17 Worker 消费服务 | 7 | 7 | 100% | ✅ 阶段5 |
 | M18 粤标扩展 | 10 | 10 | 100% | ✅ |
-| M19 RocketMQ 消息队列 | 4 | 4 | 100% | ✅ |
-| **总计** | **158** | **158** | **100%** | 🎉 全部完成 |
+| M19 RocketMQ 消息队列 | 5 | 5 | 100% | ✅ |
+| M20 Kafka 消息队列 | 4 | 4 | 100% | ✅ |
+| M21 竞品分析更新 | 1 | 1 | 100% | ✅ |
+| **总计** | **165** | **165** | **100%** | 🎉 全部完成 |
Modified +6 -10
diff --git "a/Doc/\347\253\236\345\223\201\345\210\206\346\236\220.md" "b/Doc/\347\253\236\345\223\201\345\210\206\346\236\220.md"
index f786ab9..106cac8 100644
--- "a/Doc/\347\253\236\345\223\201\345\210\206\346\236\220.md"
+++ "b/Doc/\347\253\236\345\223\201\345\210\206\346\236\220.md"
@@ -1,6 +1,6 @@
 # NewLife.JT808 竞品分析
 
-> 版本:v2.0 | 日期:2026-07-15
+> 版本:v2.1 | 日期:2026-07-15 | 更新:粤标/RocketMQ/Kafka 状态同步
 
 ## 1. 概述
 
@@ -53,7 +53,7 @@
 |:-----|:-------------:|:--------------:|:----------------:|
 | JT/T 1078 音视频消息体 | ✅ | ✅(独立包) | ✅ |
 | 苏标主动安全(ADAS/DSM/BSD/TPMS) | ✅ | ✅(独立包) | ✅ |
-| 粤标主动安全 | 🔧 规划 | ✅(独立包) | ✅ |
+| 粤标主动安全 | ✅ | ✅(独立包) | ✅ |
 | JT/T 19056 行车记录仪 | ✅ | ✅ | ✅ |
 | 报警附件协议(T1210/T1211/T9208) | ✅ | ✅ | ✅ |
 | 文件传输消息 FileMessage | ✅ | ✅ | ✅ |
@@ -82,8 +82,8 @@
 | 下行指令接口 | ✅ | ✅ |
 | 内存队列实现 | ✅ | ❌ |
 | Redis Stream 实现 | ✅ | ❌ |
-| Kafka 实现 | 🔧 规划 | ✅ |
-| RocketMQ 实现 | 🔧 规划 | ❌ |
+| Kafka 实现 | ✅ | ✅ |
+| RocketMQ 实现 | ✅ | ❌ |
 
 ### 3.5 管理后台
 
@@ -136,23 +136,19 @@
 
 | 差距 | 优先级 | 说明 |
 |:-----|:------:|------|
-| **缺少粤标主动安全扩展** | P1 | SmallChi 已有 YueBiao 扩展包 |
-| **缺少 Kafka 消息队列实现** | P2 | SmallChi/JT808Gateway 支持,项目只有 Redis Stream |
-| **消息体模型缺少扩增** | P3 | 部分消息体(如 0x8103_0xF370 等)尚未覆盖 |
 | **测试覆盖率低** | P1 | 仅 59 个测试用例(SmallChi 的测试更完备) |
 | **文档与 Demo 不足** | P2 | 缺少像 SmallChi 那样丰富的使用示例 |
 | **性能基准数据缺失** | P3 | SmallChi 有 BenchmarkDotNet 报告 |
+| **消息体模型缺少扩增** | P3 | 部分消息体(如 0x8103_0xF370 等)尚未覆盖 |
 
 
 
 ## 6. 结论与建议
 
-1. **协议库层面** - NewLife.JT808 与 SmallChi/JT808 功能持平,各有特色。SmallChi 扩展包更丰富(粤标),NewLife 的 IAccessor 手动序列化更灵活。
+1. **协议库层面** - NewLife.JT808 与 SmallChi/JT808 功能持平,各有特色。SmallChi 扩展包更丰富(粤标),NewLife 的 IAccessor 手动序列化更灵活。当前已补齐粤标扩展(M18)和 RocketMQ/Kafka 消息队列(M19/M20)。
 2. **网关层面** - NewLife.JT808 内置完整网关服务,而 SmallChi 需要独立使用 JT808Gateway 项目。两者设计理念不同。
 3. **服务平台层面** - NewLife.JT808 是目前唯一提供**开源全链路服务平台**的方案,其他开源项目均只覆盖协议层或网关层。
 4. **下一步建议**:
-   - 补齐粤标主动安全扩展(参考 SmallChi/YueBiao)
-   - 增加 Kafka 消息队列实现
    - 补充 BenchmarkDotNet 性能基准
    - 持续完善测试覆盖
 
Modified +1 -1
diff --git "a/Doc/\351\234\200\346\261\202\346\226\207\346\241\243.md" "b/Doc/\351\234\200\346\261\202\346\226\207\346\241\243.md"
index 13129fc..15035f2 100644
--- "a/Doc/\351\234\200\346\261\202\346\226\207\346\241\243.md"
+++ "b/Doc/\351\234\200\346\261\202\346\226\207\346\241\243.md"
@@ -1,6 +1,6 @@
 # NewLife.JT808 需求文档
 
-> 版本:v1.3 | 日期:2026-07-15
+> 版本:v1.4 | 日期:2026-07-15
 
 本文档描述 NewLife.JT808 的愿景、核心目标和功能方向。完成状态在[功能清单](功能清单.md)中追踪,详细设计在[架构设计](架构设计.md)中展开。
 
Modified +4 -0
diff --git a/Newlife.JT808/Newlife.JT808.csproj b/Newlife.JT808/Newlife.JT808.csproj
index 13438b0..df1098c 100644
--- a/Newlife.JT808/Newlife.JT808.csproj
+++ b/Newlife.JT808/Newlife.JT808.csproj
@@ -45,6 +45,10 @@
 		<PackageReference Include="NewLife.RocketMQ" Version="3.1.2026.601" />
 	</ItemGroup>
 
+	<ItemGroup Condition="'$(TargetFramework)'!='net45'">
+		<PackageReference Include="Confluent.Kafka" Version="2.15.0" />
+	</ItemGroup>
+
 	<ItemGroup Condition="'$(TargetFramework)'=='netstandard2.0' Or '$(TargetFramework)'=='netstandard2.1'">
 		<PackageReference Include="System.Text.Encoding.CodePages" Version="5.0.0" />
 	</ItemGroup>
Added +139 -0
diff --git a/Newlife.JT808/Protocols/Kafka/KafkaConsumer.cs b/Newlife.JT808/Protocols/Kafka/KafkaConsumer.cs
new file mode 100644
index 0000000..7ffda2d
--- /dev/null
+++ b/Newlife.JT808/Protocols/Kafka/KafkaConsumer.cs
@@ -0,0 +1,139 @@
+#if !NET45
+using Confluent.Kafka;
+using NewLife.Log;
+using NewLife.Serialization;
+
+namespace NewLife.JT808.Protocols.Kafka;
+
+/// <summary>基于 Kafka 的消息消费者适配器</summary>
+/// <remarks>
+/// 将 Confluent.Kafka 的 IConsumer 包装为 JT808 的 IMsgConsumer 接口。
+/// 消费时反序列化为 TMessage,从消息 Key 中提取 mobile。
+/// 需要 .NET Standard 2.0+ 或 .NET Framework 4.6.1+ 运行时支持。
+/// </remarks>
+/// <typeparam name="TMessage">消息类型</typeparam>
+public class KafkaConsumer<TMessage> : IMsgConsumer<TMessage>, IDisposable
+{
+    #region 属性
+    /// <summary>Kafka 消费者实例</summary>
+    public IConsumer<String, String> Consumer { get; }
+
+    /// <summary>主题名称</summary>
+    public String Topic { get; }
+
+    /// <summary>是否已启动</summary>
+    public Boolean Active => _active;
+    private Boolean _active;
+    #endregion
+
+    #region 构造
+    /// <summary>实例化 Kafka 消费者适配器</summary>
+    /// <param name="bootstrapServers">Kafka 服务器地址(如 localhost:9092)</param>
+    /// <param name="topic">主题名称</param>
+    /// <param name="group">消费组</param>
+    /// <param name="config">额外消费者配置。若为 null 则使用默认配置</param>
+    public KafkaConsumer(String bootstrapServers, String topic, String group, ConsumerConfig? config = null)
+    {
+        Topic = topic;
+
+        var cfg = config ?? new ConsumerConfig();
+        cfg.BootstrapServers = bootstrapServers;
+        cfg.GroupId = group;
+        cfg.AutoOffsetReset = AutoOffsetReset.Earliest;
+        cfg.EnableAutoCommit = true;
+
+        Consumer = new ConsumerBuilder<String, String>(cfg).Build();
+    }
+
+    /// <summary>实例化 Kafka 消费者适配器</summary>
+    /// <param name="consumer">已创建的 IConsumer 实例</param>
+    /// <param name="topic">主题名称</param>
+    public KafkaConsumer(IConsumer<String, String> consumer, String topic)
+    {
+        Consumer = consumer;
+        Topic = topic;
+    }
+    #endregion
+
+    #region 方法
+    /// <summary>订阅消息</summary>
+    /// <param name="onMessage">消息处理回调:参数为 mobile 和消息体</param>
+    /// <param name="cancellationToken">取消令牌</param>
+    public Task SubscribeAsync(Func<String, TMessage, Task> onMessage, CancellationToken cancellationToken = default)
+    {
+        _active = true;
+
+        Consumer.Subscribe(Topic);
+
+        // 启动后台消费循环
+        _ = ConsumeLoopAsync(onMessage, cancellationToken);
+
+        return Task.CompletedTask;
+    }
+
+    private async Task ConsumeLoopAsync(Func<String, TMessage, Task> onMessage, CancellationToken cancellationToken)
+    {
+        while (_active && !cancellationToken.IsCancellationRequested)
+        {
+            try
+            {
+                var result = Consumer.Consume(cancellationToken);
+                if (result == null) continue;
+
+                var mobile = result.Message.Key ?? String.Empty;
+                var value = result.Message.Value;
+
+                var body = default(TMessage);
+                if (value != null)
+                {
+                    if (typeof(TMessage) == typeof(String))
+                        body = (TMessage)(Object)value;
+                    else
+                        body = JsonHelper.ToJsonEntity<TMessage>(value);
+                }
+
+                if (body != null)
+                    await onMessage(mobile, body);
+            }
+            catch (OperationCanceledException)
+            {
+                break;
+            }
+            catch (Exception ex)
+            {
+                Log?.Error("Kafka 消费异常: {0}", ex.Message);
+                await Task.Delay(1000, cancellationToken);
+            }
+        }
+    }
+
+    /// <summary>取消订阅</summary>
+    public Task UnsubscribeAsync(CancellationToken cancellationToken = default)
+    {
+        _active = false;
+        try
+        {
+            Consumer.Unsubscribe();
+        }
+        catch { }
+        return Task.CompletedTask;
+    }
+
+    /// <summary>释放资源</summary>
+    public void Dispose()
+    {
+        _active = false;
+        try
+        {
+            Consumer.Unsubscribe();
+            Consumer.Close();
+            Consumer.Dispose();
+        }
+        catch { }
+    }
+
+    /// <summary>日志</summary>
+    public ILog Log { get; set; } = Logger.Null;
+    #endregion
+}
+#endif
Added +117 -0
diff --git a/Newlife.JT808/Protocols/Kafka/KafkaDownHandler.cs b/Newlife.JT808/Protocols/Kafka/KafkaDownHandler.cs
new file mode 100644
index 0000000..7909be1
--- /dev/null
+++ b/Newlife.JT808/Protocols/Kafka/KafkaDownHandler.cs
@@ -0,0 +1,117 @@
+#if !NET45
+using Confluent.Kafka;
+using NewLife.Log;
+
+namespace NewLife.JT808.Protocols.Kafka;
+
+/// <summary>基于 Kafka 的下行指令处理器</summary>
+/// <remarks>
+/// 消费 Kafka 指令队列中的下行指令,将字符串类型的 BodyData 转换为 Byte[]。
+/// 配合 CommandClient 使用,CommandClient 向 Kafka 发送指令,此处理器接收并转换。
+/// 需要 .NET Standard 2.0+ 或 .NET Framework 4.6.1+ 运行时支持。
+/// </remarks>
+public class KafkaDownHandler : IDownMessageHandler, IDisposable
+{
+    #region 属性
+    /// <summary>Kafka 消费者实例</summary>
+    public IConsumer<String, String> Consumer { get; }
+
+    /// <summary>指令主题</summary>
+    public String Topic { get; }
+
+    /// <summary>是否已启动</summary>
+    public Boolean Active => _active;
+    private Boolean _active;
+
+    private Func<String, Byte[], Task<Byte[]>>? _handler;
+    #endregion
+
+    #region 构造
+    /// <summary>实例化 Kafka 下行指令处理器</summary>
+    /// <param name="bootstrapServers">Kafka 服务器地址(如 localhost:9092)</param>
+    /// <param name="topic">指令主题</param>
+    /// <param name="group">消费组</param>
+    /// <param name="config">额外消费者配置。若为 null 则使用默认配置</param>
+    public KafkaDownHandler(String bootstrapServers, String topic, String group, ConsumerConfig? config = null)
+    {
+        Topic = topic;
+
+        var cfg = config ?? new ConsumerConfig();
+        cfg.BootstrapServers = bootstrapServers;
+        cfg.GroupId = group;
+        cfg.AutoOffsetReset = AutoOffsetReset.Earliest;
+        cfg.EnableAutoCommit = true;
+
+        Consumer = new ConsumerBuilder<String, String>(cfg).Build();
+    }
+    #endregion
+
+    #region 方法
+    /// <summary>处理下行消息</summary>
+    /// <param name="mobile">目标终端手机号</param>
+    /// <param name="data">消息体序列化数据</param>
+    /// <returns>处理结果</returns>
+    public Task<Byte[]> HandleAsync(String mobile, Byte[] data)
+    {
+        return Task.FromResult(data);
+    }
+
+    /// <summary>开始消费指令队列</summary>
+    /// <param name="onMessage">消费回调:mobile, data → 结果</param>
+    public void Start(Func<String, Byte[], Task<Byte[]>> onMessage)
+    {
+        _handler = onMessage;
+
+        _active = true;
+        Consumer.Subscribe(Topic);
+
+        _ = ConsumeLoopAsync();
+    }
+
+    private async Task ConsumeLoopAsync()
+    {
+        while (_active)
+        {
+            try
+            {
+                var result = Consumer.Consume(TimeSpan.FromMilliseconds(500));
+                if (result == null) continue;
+
+                var mobile = result.Message.Key ?? String.Empty;
+                var data = result.Message.Value != null ?
+                    System.Text.Encoding.UTF8.GetBytes(result.Message.Value) :
+                    [];
+
+                if (_handler != null)
+                    await _handler(mobile, data);
+            }
+            catch (OperationCanceledException)
+            {
+                break;
+            }
+            catch (Exception ex)
+            {
+                Log?.Error("Kafka 指令消费异常: {0}", ex.Message);
+                await Task.Delay(1000);
+            }
+        }
+    }
+
+    /// <summary>释放资源</summary>
+    public void Dispose()
+    {
+        _active = false;
+        try
+        {
+            Consumer.Unsubscribe();
+            Consumer.Close();
+            Consumer.Dispose();
+        }
+        catch { }
+    }
+
+    /// <summary>日志</summary>
+    public ILog Log { get; set; } = Logger.Null;
+    #endregion
+}
+#endif
Added +97 -0
diff --git a/Newlife.JT808/Protocols/Kafka/KafkaProducer.cs b/Newlife.JT808/Protocols/Kafka/KafkaProducer.cs
new file mode 100644
index 0000000..4233f85
--- /dev/null
+++ b/Newlife.JT808/Protocols/Kafka/KafkaProducer.cs
@@ -0,0 +1,97 @@
+#if !NET45
+using Confluent.Kafka;
+using NewLife.Log;
+using NewLife.Serialization;
+
+namespace NewLife.JT808.Protocols.Kafka;
+
+/// <summary>基于 Kafka 的消息生产者适配器</summary>
+/// <remarks>
+/// 将 Confluent.Kafka 的 IProducer 包装为 JT808 的 IMsgProducer 接口。
+/// 消息体被序列化为 JSON 字符串后发送,mobile 编码到消息的 Key 中。
+/// 需要 .NET Standard 2.0+ 或 .NET Framework 4.6.1+ 运行时支持。
+/// </remarks>
+/// <typeparam name="TMessage">消息类型</typeparam>
+public class KafkaProducer<TMessage> : IMsgProducer<TMessage>, IDisposable
+{
+    #region 属性
+    /// <summary>Kafka 生产者实例</summary>
+    public IProducer<String, String> Producer { get; }
+
+    /// <summary>主题名称</summary>
+    public String Topic { get; }
+
+    /// <summary>是否已启动</summary>
+    public Boolean Active => _active;
+    private Boolean _active;
+    #endregion
+
+    #region 构造
+    /// <summary>实例化 Kafka 生产者适配器</summary>
+    /// <param name="bootstrapServers">Kafka 服务器地址(如 localhost:9092)</param>
+    /// <param name="topic">主题名称</param>
+    /// <param name="config">额外生产者配置。若为 null 则使用默认配置</param>
+    public KafkaProducer(String bootstrapServers, String topic, ProducerConfig? config = null)
+    {
+        Topic = topic;
+
+        var cfg = config ?? new ProducerConfig();
+        cfg.BootstrapServers = bootstrapServers;
+
+        Producer = new ProducerBuilder<String, String>(cfg).Build();
+        _active = true;
+    }
+
+    /// <summary>实例化 Kafka 生产者适配器</summary>
+    /// <param name="producer">已创建的 IProducer 实例</param>
+    /// <param name="topic">主题名称</param>
+    public KafkaProducer(IProducer<String, String> producer, String topic)
+    {
+        Producer = producer;
+        Topic = topic;
+        _active = true;
+    }
+    #endregion
+
+    #region 方法
+    /// <summary>生产消息到 Kafka</summary>
+    /// <param name="mobile">终端手机号,编码到消息 Key 中</param>
+    /// <param name="message">消息对象</param>
+    /// <param name="cancellationToken">取消令牌</param>
+    /// <returns>是否成功</returns>
+    public async Task<Boolean> ProduceAsync(String mobile, TMessage message, CancellationToken cancellationToken = default)
+    {
+        try
+        {
+            var json = JsonHelper.ToJson(message);
+            var result = await Producer.ProduceAsync(Topic, new Message<String, String>
+            {
+                Key = mobile,
+                Value = json,
+            }, cancellationToken);
+
+            return result.Status == PersistenceStatus.Persisted;
+        }
+        catch (Exception ex)
+        {
+            Log?.Error("Kafka 生产失败: {0}", ex.Message);
+            return false;
+        }
+    }
+
+    /// <summary>释放资源</summary>
+    public void Dispose()
+    {
+        if (_active)
+        {
+            _active = false;
+            Producer.Flush();
+            Producer.Dispose();
+        }
+    }
+
+    /// <summary>日志</summary>
+    public ILog Log { get; set; } = Logger.Null;
+    #endregion
+}
+#endif
Added +78 -0
diff --git a/Newlife.JT808/Protocols/Kafka/KafkaSessionAdapter.cs b/Newlife.JT808/Protocols/Kafka/KafkaSessionAdapter.cs
new file mode 100644
index 0000000..525b4b0
--- /dev/null
+++ b/Newlife.JT808/Protocols/Kafka/KafkaSessionAdapter.cs
@@ -0,0 +1,78 @@
+#if !NET45
+using Confluent.Kafka;
+using NewLife.Log;
+using NewLife.Serialization;
+
+namespace NewLife.JT808.Protocols.Kafka;
+
+/// <summary>基于 Kafka 的会话通知适配器</summary>
+/// <remarks>
+/// 将终端的上下线事件通过 Kafka 主题广播。
+/// 上线通知发送 "online" 消息,下线通知发送 "offline" 消息。
+/// 需要 .NET Standard 2.0+ 或 .NET Framework 4.6.1+ 运行时支持。
+/// </remarks>
+public class KafkaSessionAdapter : ISessionProducer, IDisposable
+{
+    #region 属性
+    /// <summary>Kafka 生产者实例</summary>
+    public IProducer<String, String> Producer { get; }
+
+    /// <summary>会话通知主题</summary>
+    public String Topic { get; set; } = "SessionEvent";
+    #endregion
+
+    #region 构造
+    /// <summary>实例化 Kafka 会话通知适配器</summary>
+    /// <param name="bootstrapServers">Kafka 服务器地址(如 localhost:9092)</param>
+    /// <param name="topic">会话通知主题</param>
+    /// <param name="config">额外生产者配置。若为 null 则使用默认配置</param>
+    public KafkaSessionAdapter(String bootstrapServers, String? topic = null, ProducerConfig? config = null)
+    {
+        if (topic != null) Topic = topic;
+
+        var cfg = config ?? new ProducerConfig();
+        cfg.BootstrapServers = bootstrapServers;
+
+        Producer = new ProducerBuilder<String, String>(cfg).Build();
+    }
+    #endregion
+
+    #region 方法
+    /// <summary>通知终端上线</summary>
+    public async Task ProduceOnlineAsync(String mobile, CancellationToken cancellationToken = default)
+    {
+        var json = JsonHelper.ToJson(new { mobile, action = "online", time = DateTime.UtcNow });
+        await Producer.ProduceAsync(Topic, new Message<String, String>
+        {
+            Key = mobile,
+            Value = json,
+        }, cancellationToken);
+    }
+
+    /// <summary>通知终端离线</summary>
+    public async Task ProduceOfflineAsync(String mobile, CancellationToken cancellationToken = default)
+    {
+        var json = JsonHelper.ToJson(new { mobile, action = "offline", time = DateTime.UtcNow });
+        await Producer.ProduceAsync(Topic, new Message<String, String>
+        {
+            Key = mobile,
+            Value = json,
+        }, cancellationToken);
+    }
+
+    /// <summary>释放资源</summary>
+    public void Dispose()
+    {
+        try
+        {
+            Producer.Flush();
+            Producer.Dispose();
+        }
+        catch { }
+    }
+
+    /// <summary>日志</summary>
+    public ILog Log { get; set; } = Logger.Null;
+    #endregion
+}
+#endif
Modified +3 -1
diff --git a/Newlife.JT808/Protocols/RocketMQ/RocketMQConsumer.cs b/Newlife.JT808/Protocols/RocketMQ/RocketMQConsumer.cs
index 9df40ab..23fa492 100644
--- a/Newlife.JT808/Protocols/RocketMQ/RocketMQConsumer.cs
+++ b/Newlife.JT808/Protocols/RocketMQ/RocketMQConsumer.cs
@@ -31,7 +31,8 @@ public class RocketMQConsumer<TMessage> : IMsgConsumer<TMessage>, IDisposable
     /// <param name="nameServer">NameServer 地址</param>
     /// <param name="topic">主题名称</param>
     /// <param name="group">消费组</param>
-    public RocketMQConsumer(String nameServer, String topic, String group)
+    /// <param name="cloudProvider">云厂商适配器。阿里云传 <c>new AliyunProvider{AccessKey=...,SecretKey=...,InstanceId=...}</c>;华为云传 <c>new HuaweiProvider{...}</c>;腾讯云传 <c>new TencentProvider{...}</c>;ACL 传 <c>new AclProvider{...}</c></param>
+    public RocketMQConsumer(String nameServer, String topic, String group, NewLife.RocketMQ.ICloudProvider? cloudProvider = null)
     {
         Topic = topic;
         Consumer = new Consumer
@@ -39,6 +40,7 @@ public class RocketMQConsumer<TMessage> : IMsgConsumer<TMessage>, IDisposable
             NameServerAddress = nameServer,
             Topic = topic,
             Group = group,
+            CloudProvider = cloudProvider,
         };
     }
 
Modified +3 -1
diff --git a/Newlife.JT808/Protocols/RocketMQ/RocketMQDownHandler.cs b/Newlife.JT808/Protocols/RocketMQ/RocketMQDownHandler.cs
index 7f7beb6..022d258 100644
--- a/Newlife.JT808/Protocols/RocketMQ/RocketMQDownHandler.cs
+++ b/Newlife.JT808/Protocols/RocketMQ/RocketMQDownHandler.cs
@@ -29,7 +29,8 @@ public class RocketMQDownHandler : IDownMessageHandler, IDisposable
     /// <param name="nameServer">NameServer 地址</param>
     /// <param name="topic">指令主题</param>
     /// <param name="group">消费组</param>
-    public RocketMQDownHandler(String nameServer, String topic, String group)
+    /// <param name="cloudProvider">云厂商适配器。阿里云传 <c>new AliyunProvider{AccessKey=...,SecretKey=...,InstanceId=...}</c>;华为云传 <c>new HuaweiProvider{...}</c>;腾讯云传 <c>new TencentProvider{...}</c>;ACL 传 <c>new AclProvider{...}</c></param>
+    public RocketMQDownHandler(String nameServer, String topic, String group, NewLife.RocketMQ.ICloudProvider? cloudProvider = null)
     {
         Topic = topic;
         Consumer = new Consumer
@@ -37,6 +38,7 @@ public class RocketMQDownHandler : IDownMessageHandler, IDisposable
             NameServerAddress = nameServer,
             Topic = topic,
             Group = group,
+            CloudProvider = cloudProvider,
         };
     }
     #endregion
Modified +3 -1
diff --git a/Newlife.JT808/Protocols/RocketMQ/RocketMQProducer.cs b/Newlife.JT808/Protocols/RocketMQ/RocketMQProducer.cs
index 534242e..1421ead 100644
--- a/Newlife.JT808/Protocols/RocketMQ/RocketMQProducer.cs
+++ b/Newlife.JT808/Protocols/RocketMQ/RocketMQProducer.cs
@@ -27,7 +27,8 @@ public class RocketMQProducer<TMessage> : IMsgProducer<TMessage>, IDisposable
     /// <param name="nameServer">NameServer 地址</param>
     /// <param name="topic">主题名称</param>
     /// <param name="group">生产者组</param>
-    public RocketMQProducer(String nameServer, String topic, String? group = null)
+    /// <param name="cloudProvider">云厂商适配器。阿里云传 <c>new AliyunProvider{AccessKey=...,SecretKey=...,InstanceId=...}</c>;华为云传 <c>new HuaweiProvider{...}</c>;腾讯云传 <c>new TencentProvider{...}</c>;ACL 传 <c>new AclProvider{...}</c></param>
+    public RocketMQProducer(String nameServer, String topic, String? group = null, NewLife.RocketMQ.ICloudProvider? cloudProvider = null)
     {
         Topic = topic;
         Producer = new Producer
@@ -35,6 +36,7 @@ public class RocketMQProducer<TMessage> : IMsgProducer<TMessage>, IDisposable
             NameServerAddress = nameServer,
             Topic = topic,
             Group = group ?? $"PID_{topic}",
+            CloudProvider = cloudProvider,
         };
         Producer.Start();
     }
Modified +3 -1
diff --git a/Newlife.JT808/Protocols/RocketMQ/RocketMQSessionAdapter.cs b/Newlife.JT808/Protocols/RocketMQ/RocketMQSessionAdapter.cs
index 1b24fc5..fec1508 100644
--- a/Newlife.JT808/Protocols/RocketMQ/RocketMQSessionAdapter.cs
+++ b/Newlife.JT808/Protocols/RocketMQ/RocketMQSessionAdapter.cs
@@ -24,7 +24,8 @@ public class RocketMQSessionAdapter : ISessionProducer, IDisposable
     /// <param name="nameServer">NameServer 地址</param>
     /// <param name="topic">会话通知主题</param>
     /// <param name="group">生产者组</param>
-    public RocketMQSessionAdapter(String nameServer, String? topic = null, String? group = null)
+    /// <param name="cloudProvider">云厂商适配器。阿里云传 <c>new AliyunProvider{AccessKey=...,SecretKey=...,InstanceId=...}</c>;华为云传 <c>new HuaweiProvider{...}</c>;腾讯云传 <c>new TencentProvider{...}</c>;ACL 传 <c>new AclProvider{...}</c></param>
+    public RocketMQSessionAdapter(String nameServer, String? topic = null, String? group = null, NewLife.RocketMQ.ICloudProvider? cloudProvider = null)
     {
         if (topic != null) Topic = topic;
         Producer = new Producer
@@ -32,6 +33,7 @@ public class RocketMQSessionAdapter : ISessionProducer, IDisposable
             NameServerAddress = nameServer,
             Topic = Topic,
             Group = group ?? "PID_SessionEvent",
+            CloudProvider = cloudProvider,
         };
         Producer.Start();
     }
Modified +3 -2
diff --git a/Readme.MD b/Readme.MD
index 2e5b69a..b39518a 100644
--- a/Readme.MD
+++ b/Readme.MD
@@ -28,11 +28,12 @@ Nuget:[NewLife.JT808](https://www.nuget.org/packages/NewLife.JT808)
 - **BCD 编码**:BCD 字符串和时间自定义序列化特性,适配协议特有编码格式
 - **GBK 编码**:内置 GBK 中文字符串编解码,自动注册 `CodePagesEncodingProvider`
 - **Web API 管理**:支持统一下发指令、在线会话管理、黑名单管理、Token 鉴权,方便集成到后台系统
-- **消息队列抽象**:IMsgProducer/IMsgConsumer/ISessionProducer 等接口解耦消息处理,内置 DefaultMsgQueue 内存实现、Redis Stream 生产/消费支持、RocketMQ 适配器
+- **消息队列抽象**:IMsgProducer/IMsgConsumer/ISessionProducer 等接口解耦消息处理,内置 DefaultMsgQueue 内存实现、Redis Stream 生产/消费支持、RocketMQ 适配器、Kafka 适配器
 - **跨服务器指令下发**:ICommandBus 接口 + Stardust 事件总线,支持集群广播指令到指定终端
 - **客户端增强**:自动重连、可配心跳、收发统计计数器,提升生产环境可靠性
 - **多框架支持**:`net45 / net461 / netstandard2.0 / netstandard2.1` 全兼容
-- **RocketMQ 集成**:基于 `NewLife.RocketMQ` 的消息队列生产/消费适配器,可作为 Redis Stream 的替代方案
+- **RocketMQ 集成**:基于 `NewLife.RocketMQ` 的消息队列生产/消费适配器,支持阿里云/华为云/腾讯云/ACL 四种云厂商认证模式
+- **Kafka 集成**:基于 `Confluent.Kafka` 的消息队列生产/消费适配器,支持消费组、ACK 机制(需要 .NET Standard 2.0+)
 
 ## 支持的协议版本