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

时间: 2024-10-10 02:05:20

NET中解决KafKa多线程发送多主题的问题的相关文章

30_Django中关于使用ajax发送请求中`csrf_token`的问题和解决

目录 Django中关于csrf_token的问题和解决 1. 在form表单中加上{% csrf_token %} 2. 使用装饰器的方法,可以不用在form表单中添加上{% csrf_token %} 了 3. 使用中间件,为每个请求都获取到token (使用这个方法,所有的表单都不需要在form标签中添加上{% csrf_token %} 了) Django中关于csrf_token的问题和解决 在需要发送POST请求的html文件中引入下面的js代码, 我为其命名为get_csrf_to

在Asp.Net Core中集成Kafka(中)

在上一篇中我们主要介绍如何在Asp.Net Core中同步Kafka消息,通过上一篇的操作我们发现上面一篇中介绍的只能够进行简单的首发kafka消息并不能够消息重发.重复消费.乐观锁冲突等问题,这些问题在实际的生产环境中是非常要命的,如果在消息的消费方没有做好必须的幂等性操作,那么消费者重复消费的问题会比较严重的,另外对于消息的生产者来说,记录日志的方式也不是足够友好,很多时候在后台监控程序中我们需要知道记录更多的关于消息的分区.偏移等更多的消息.而在消费者这边我们更多的需要去解决发送方发送重复

iOS开发中GCD在多线程方面的理解

GCD为Grand Central Dispatch的缩写. Grand Central Dispatch (GCD)是Apple开发的一个多核编程的较新的解决方法.在Mac OS X 10.6雪豹中首次推出,并在最近引入到了iOS4.0. GCD是一个替代诸如NSThread等技术的很高效和强大的技术.GCD完全可以处理诸如数据锁定和资源泄漏等复杂的异步编程问题. GCD可以完成很多事情,但是这里仅关注在iOS应用中实现多线程所需的一些基础知识. 在开始之前,需要理解是要提供给GCD队列的是代

深入解析PHP中的(伪)多线程与多进程

本篇文章是对PHP中的(伪)多线程与多进程进行了详细的分析介绍,需要的朋友参考下 (伪)多线程:借助外力利用WEB服务器本身的多线程来处理,从WEB服务器多次调用我们需要实现多线程的程序.QUOTE:我们知道PHP本身是不支持多线程的, 但是我们的WEB服务器是支持多线程的.也就是说可以同时让多人一起访问. 这也是我在PHP中实现多线程的基础.假设我们现在运行的是a.php这个文件. 但是我在程序中又请求WEB服务器运行另一个b.php那么这两个文件将是同时执行的.(PS: 一个链接请求发送之后

如何在优雅地Spring 中实现消息的发送和消费

本文将对rocktmq-spring-boot的设计实现做一个简单的介绍,读者可以通过本文了解将RocketMQ Client端集成为spring-boot-starter框架的开发细节,然后通过一个简单的示例来一步一步的讲解如何使用这个spring-boot-starter工具包来配置,发送和消费RocketMQ消息. 作者简介:辽天,阿里巴巴技术专家,Apache RocketMQ 内核控,拥有多年分布式系统研发经验,对Microservice.Messaging和Storage等领域有深刻

源码分析 Kafka 消息发送流程(文末附流程图)

温馨提示:本文基于 Kafka 2.2.1 版本.本文主要是以源码的手段一步一步探究消息发送流程,如果对源码不感兴趣,可以直接跳到文末查看消息发送流程图与消息发送本地缓存存储结构. 从上文 初识 Kafka Producer 生产者,可以通过 KafkaProducer 的 send 方法发送消息,send 方法的声明如下: Future<RecordMetadata> send(ProducerRecord<K, V> record) Future<RecordMetada

C#中异步和多线程的区别

C#中异步和多线程的区别是什么呢?异步和多线程两者都可以达到避免调用线程阻塞的目的,从而提高软件的可响应性.甚至有些时候我们就认为异步和多线程是等同的概念.但是,异步和多线程还是有一些区别的.而这些区别造成了使用异步和多线程的时机的区别. 异步和多线程的区别之异步操作的本质 所有的程序最终都会由计算机硬件来执行,所以为了更好的理解异步操作的本质,我们有必要了解一下它的硬件基础. 熟悉电脑硬件的朋友肯定对DMA这个词不陌生,硬盘.光驱的技术规格中都有明确DMA的模式指标,其实网卡.声卡.显卡也是有

python中urllib2与多线程使用

问题提出 几天前,我在上一篇博客中写了如何使用urllib2模块来批量下载wallheaven上的图片资源,但是在我几次运行下来之后发现了一个非常严重的问题,如果下载图片数量非常多的话,程序需要运行很长时间.所以显然这样不是一个很好的解决方法,所以后来我在程序中加入了多线程,程序性能提升了何止数倍,下面是具体的解决过程. 问题解决 从我上一边的博客中不难看出,第一次的下载程序每次只能下载一张图片,这样完全浪费了计算机的内存资源和网络资源.所以之后我加入了多线程,每次可以根据不同的需求开启更多的下

java中使HttpDelete可以发送body信息

java中使HttpDelete可以发送body信息RESTful api中用到了DELETE方法,android开发的同事遇到了问题,使用HttpDelete执行DELETE操作的时候,不能携带body信息,研究了很久之后找到了解决方法. 我们查看httpclient-4.2.3的源码可以发现,methods包下面包含HttpGet, HttpPost, HttpPut, HttpDelete等类来实现http的常用操作. 其中,HttpPost继承自HttpEntityEnclosingRe