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

简介:

一般在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!以上就完成了多线程多主题的消息发送。

本文转自欢醉博客园博客,原文链接http://www.cnblogs.com/zhangs1986/p/7285525.html如需转载请自行联系原作者


欢醉

相关文章
|
1月前
|
并行计算 安全 Java
C# .NET面试系列四:多线程
<h2>多线程 #### 1. 根据线程安全的相关知识,分析以下代码,当调用 test 方法时 i > 10 时是否会引起死锁? 并简要说明理由。 ```c# public void test(int i) { lock(this) { if (i > 10) { i--; test(i); } } } ``` 在给定的代码中,不会发生死锁。死锁通常是由于两个或多个线程互相等待对方释放锁而无法继续执行的情况。在这个代码中,只有一个线程持有锁,且没有其他线程参与,因此不
105 3
|
消息中间件 Java Kafka
Spring Boot集成Kafka动态创建消费者与动态删除主题(实现多消费者的发布订阅模型)
Spring Boot集成Kafka动态创建消费者与动态删除主题(实现多消费者的发布订阅模型)
17069 1
Spring Boot集成Kafka动态创建消费者与动态删除主题(实现多消费者的发布订阅模型)
|
消息中间件 监控 Java
图解Kafka线程模型及其设计缺陷
图解Kafka线程模型及其设计缺陷
图解Kafka线程模型及其设计缺陷
|
7月前
|
消息中间件 存储 Kafka
Kafka主题,分区,副本介绍
今天分享一下kafka的主题(topic),分区(partition)和副本(replication),主题是Kafka中很重要的部分,消息的生产和消费都要以主题为基础,一个主题可以对应多个分区,一个分区属于某个主题,一个分区又可以对应多个副本,副本分为leader和follower。
94 0
|
7月前
|
消息中间件 存储 缓存
十分钟,了解Kafka的Sender线程
十分钟,了解Kafka的Sender线程
59 0
|
消息中间件 网络协议 安全
【Kafka从入门到成神系列 八】Kafka 多线程消费者及TCP连接
【Kafka从入门到成神系列 八】Kafka 多线程消费者及TCP连接
【Kafka从入门到成神系列 八】Kafka 多线程消费者及TCP连接
|
11月前
|
API Apache
Apache Kafka-通过API获取主题所有分区的积压消息数量
Apache Kafka-通过API获取主题所有分区的积压消息数量
100 0
|
消息中间件 数据采集 监控
Python 基于Python结合pykafka实现kafka生产及消费速率&主题分区偏移实时监控
Python 基于Python结合pykafka实现kafka生产及消费速率&主题分区偏移实时监控
176 0
|
消息中间件 存储 监控
KafKa主题、分区、副本、消息代理
Kafka将主题拆分为多个分区,不同的分区存在不同的服务器上,这样就使kafka具有拓展性,可以通过调整分区的数量和节点的数量,来线性对Kafka进行拓展,分区是一个线性增长的不可变日志,当消息存储到分区中之后,消息就不可变更,kafka为每条消息设置一个偏移量也就是offset,offset可以记录每条消息的位置,kafka可以通过偏移量对消息进行提取,但是没法对消息的内容进行检索和查询,偏移量在每个分区中是唯一的不可重复,并且它是递增的,不同分区间偏移量可以重复。
133 0
|
数据采集 Java C++
【.NET 6】多线程的几种打开方式和代码演示
多线程无处不在,平常的开发过程中,应该算是最常用的基础技术之一了。以下通过Thread、ThreadPool、再到Task、Parallel、线程锁、线程取消等方面,一步步进行演示多线程的一些基础操作。欢迎大家围观。如果大佬们有其他关于多线程的拓展,也欢迎在评论区进行留言,大佬们的知识互助,是.net生态发展的重要一环,欢迎大佬们进行留言,帮助更多的人。
196 0
【.NET 6】多线程的几种打开方式和代码演示

热门文章

最新文章