Java自学者论坛

 找回密码
 立即注册

手机号码,快捷登录

恭喜Java自学者论坛(https://www.javazxz.com)已经为数万Java学习者服务超过8年了!积累会员资料超过10000G+
成为本站VIP会员,下载本站10000G+会员资源,会员资料板块,购买链接:点击进入购买VIP会员

JAVA高级面试进阶训练营视频教程

Java架构师系统进阶VIP课程

分布式高可用全栈开发微服务教程Go语言视频零基础入门到精通Java架构师3期(课件+源码)
Java开发全终端实战租房项目视频教程SpringBoot2.X入门到高级使用教程大数据培训第六期全套视频教程深度学习(CNN RNN GAN)算法原理Java亿级流量电商系统视频教程
互联网架构师视频教程年薪50万Spark2.0从入门到精通年薪50万!人工智能学习路线教程年薪50万大数据入门到精通学习路线年薪50万机器学习入门到精通教程
仿小米商城类app和小程序视频教程深度学习数据分析基础到实战最新黑马javaEE2.1就业课程从 0到JVM实战高手教程MySQL入门到精通教程
查看: 341|回复: 0

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

[复制链接]
  • TA的每日心情
    奋斗
    11 小时前
  • 签到天数: 757 天

    [LV.10]以坛为家III

    2034

    主题

    2092

    帖子

    70万

    积分

    管理员

    Rank: 9Rank: 9Rank: 9

    积分
    707886
    发表于 2021-7-22 11:06:04 | 显示全部楼层 |阅读模式

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

     

    哎...今天够累的,签到来了1...
    回复

    使用道具 举报

    您需要登录后才可以回帖 登录 | 立即注册

    本版积分规则

    QQ|手机版|小黑屋|Java自学者论坛 ( 声明:本站文章及资料整理自互联网,用于Java自学者交流学习使用,对资料版权不负任何法律责任,若有侵权请及时联系客服屏蔽删除 )

    GMT+8, 2024-7-2 23:36 , Processed in 0.059933 second(s), 27 queries .

    Powered by Discuz! X3.4

    Copyright © 2001-2021, Tencent Cloud.

    快速回复 返回顶部 返回列表