修正命名错误
xiyunfei authored at 2023-04-09 21:10:58
3.09 KiB
NewLife.JT808
#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