NET中解决KafKa多线程发送多主题的问题

简介:   一般在KafKa消费程序中消费可以设置多个主题,那在同一程序中需要向KafKa发送不同主题的消息,如异常需要发到异常主题,正常的发送到正常的主题,这时候就需要实例化多个主题,然后逐个发送。   在NET中用RdKafka组件来做消息处理,在Nuget中引用。

  一般在KafKa消费程序中消费可以设置多个主题,那在同一程序中需要向KafKa发送不同主题的消息,如异常需要发到异常主题,正常的发送到正常的主题,这时候就需要实例化多个主题,然后逐个发送。

  在NET中用RdKafka组件来做消息处理,在Nuget中引用。

  在程序中初始化Producer,并创建多个Topic

        private string comtopic = "topic1";
        private string errtopic = "topic2";
        private string kfkip = "192.168.80.32:9092";
        Topic topic = null;
        Topic errTopic = null;

        public ExcuteFlow()
        {
            try
            {
                Producer producer = new Producer(kfkip);
                topic = producer.Topic(comtopic);
                errTopic = producer.Topic(errtopic);
            }
            catch (RdKafkaException ex)
            {
                LogHelper.Error("KafKa初始化KafKa异常 ", ex);
            }
            catch (Exception ex)
            {
                LogHelper.Error("KafKa初始化异常", ex);
            }

        }

  在程序中发送其中一个主题:

          try
            {

                if (topic != null)
                {
                    byte[] datas = Encoding.UTF8.GetBytes(JsonHelper.ToJson(flowCommond));
                    Task<DeliveryReport> deliveryReport = topic.Produce(datas);
                    var unused = deliveryReport.ContinueWith(task =>
                    {
                        LogHelper.Info("内容:{flowCommond.ID} 发送到分区:{task.Result.Partition}, Offset 为: {task.Result.Offset}");
                    });
                }
                else
                {
                    throw new Exception("发送消息到KafKa topic 为空");
                }
            }
            catch (RdKafkaException ex)
            {
                LogHelper.Error("发送消息到KafKa  KafKa异常", ex);
            }
            catch (Exception ex)
            {
                LogHelper.Error("发送消息到KafKa异常", ex);
            }

  flowCommond为要发送的对象内容,格式化为Json字符串再发送。

  另一个主题一样处理。

   这里实现一个线程里面发送多个主题,那下面实现多个线程中如何发送多个主题。

  多线程中如果每个线程都new Producer(kfkip) 一次,那KafKa的连接很快会被占满。

  那这里就用单例模式来解决这个问题,每次要用到Producer时检查一下是否已经存在Producer实例,若存在则直接用不用再生成。

    /// <summary>
    /// 单例模式的实现
    /// </summary>
    public class SingleProduct : Producer
    {
        // 定义一个静态变量来保存类的实例
        private static SingleProduct uniqueInstance;
        // 定义一个标识确保线程同步
        private static readonly object locker = new object();
        // 定义私有构造函数,使外界不能创建该类实例
        private SingleProduct(string brokerList) : base(brokerList)
        {
        }

        /// <summary>
        /// 定义公有方法提供一个全局访问点,同时你也可以定义公有属性来提供全局访问点
        /// </summary>
        /// <returns></returns>
        public static SingleProduct GetInstance()
        {
            // 当第一个线程运行到这里时,此时会对locker对象 "加锁",
            // 当第二个线程运行该方法时,首先检测到locker对象为"加锁"状态,该线程就会挂起等待第一个线程解锁
            // lock语句运行完之后(即线程运行完之后)会对该对象"解锁"
            if (uniqueInstance == null)
            {
                lock (locker)
                {
                    // 如果类的实例不存在则创建,否则直接返回
                    if (uniqueInstance == null)
                    {
                        string kfkip = System.Configuration.ConfigurationManager.AppSettings["KfkIP"];

                        try
                        {
                            uniqueInstance = new SingleProduct(kfkip);
                            LogHelper.Error("单例模式 实例化 SingleProduct");
                        }
                        catch (RdKafkaException ex)
                        {
                            LogHelper.Error("单例模式 KafKa初始化KafKa异常 ", ex);
                        }
                        catch (Exception ex)
                        {
                            LogHelper.Error("单例模式 KafKa初始化异常", ex);
                        }
                    }
                }
            }

            return uniqueInstance;
        }
    }

   然后在初始化的代码中替换Producer producer = new Producer(kfkip);为 Producer producer = SingleProduct.GetInstance();

  OK!以上就完成了多线程多主题的消息发送。

 

目录
相关文章
|
消息中间件 负载均衡 Kafka
Kafka学习---2、kafka生产者、异步和同步发送API、分区、生产经验(一)
Kafka学习---2、kafka生产者、异步和同步发送API、分区、生产经验(一)
|
消息中间件 网络协议 安全
【Kafka从入门到成神系列 八】Kafka 多线程消费者及TCP连接
【Kafka从入门到成神系列 八】Kafka 多线程消费者及TCP连接
【Kafka从入门到成神系列 八】Kafka 多线程消费者及TCP连接
|
消息中间件 存储 Docker
.Net Core对于`RabbitMQ`封装分布式事件总线
.Net Core对于`RabbitMQ`封装分布式事件总线
197 1
.Net Core对于`RabbitMQ`封装分布式事件总线
|
消息中间件 缓存 Kafka
Kafka学习---2、kafka生产者、异步和同步发送API、分区、生产经验(二)
Kafka学习---2、kafka生产者、异步和同步发送API、分区、生产经验(二)
|
消息中间件 Java Kafka
Java 最常见的面试题:kafka 同时设置了 7 天和 10G 清除数据,到第五天的时候消息达到了 10G,这个时候 kafka 将如何处理?
Java 最常见的面试题:kafka 同时设置了 7 天和 10G 清除数据,到第五天的时候消息达到了 10G,这个时候 kafka 将如何处理?
|
消息中间件 监控
十五、.net core(.NET 6)搭建RabbitMQ消息队列生产者和消费者的简单方法
搭建RabbitMQ简单通用的直连方法 如果还没有MQ环境,可以参考上一篇的博客: https://www.cnblogs.com/weskynet/p/14877932.html
691 0
十五、.net core(.NET 6)搭建RabbitMQ消息队列生产者和消费者的简单方法
|
消息中间件 缓存 网络协议
【Kafka从入门到成神系列 四】Kafka 消息丢失及 TCP 管理
【Kafka从入门到成神系列 四】Kafka 消息丢失及 TCP 管理
【Kafka从入门到成神系列 四】Kafka 消息丢失及 TCP 管理
|
消息中间件 存储 安全
Broker消息设计--Kafka从入门到精通(十三)
Broker消息设计--Kafka从入门到精通(十三)
|
消息中间件 安全 Java
Kafka源码解析之SocketServer(下)
Kafka源码解析之SocketServer
151 0
Kafka源码解析之SocketServer(下)
|
消息中间件 缓存 运维
Kafka源码解析之SocketServer(上)
Kafka源码解析之SocketServer
139 0
Kafka源码解析之SocketServer(上)