【问题标题】:How to add multiple consumers to consume from kafka stream?如何从kafka流中添加多个消费者来消费?
【发布时间】:2019-11-18 20:29:45
【问题描述】:

我最近开始使用 .NET 核心开发 Kafka 流应用程序。我已按照教程进行操作:https://medium.com/@srigumm/building-realtime-streaming-applications-using-net-core-and-kafka-ad45ed081b31

我已经构建了一个基本的生产者-消费者应用程序,生产者在其中获取输入数据并将其推送到 kafka-topic 中。消费者可以订阅主题并使用其中的数据。我还可以通过使用新的生产者将这些数据推送到另一个主题。但我无法做的是初始化多个消费者从同一个主题消费。

在 appsettings.json 中:

  "consumer": {
    "bootstrapservers": "localhost:9092", //specify your kafka broker address
    "groupid": "csharp-consumer",
    "enableautocommit": true,
    "statisticsintervalms": 5000,
    "sessiontimeoutms": 6000,
    "autooffsetreset": 0,
    "enablepartitioneof": true,
    "SaslMechanism": 0, //0 for GSSAPI
    //"SaslKerberosKeytab":"filename.keytab", //specify your keytab file here
    "SaslKerberosPrincipal": "youralias@DOMAIN.COM", //specify your alias here
    "SaslKerberosServiceName": "kafka"
    //"SaslKerberosKinitCmd":"kinit -k -t %{sasl.kerberos.keytab} %{sasl.kerberos.principal}"
  },

processOrderServices.cs:

namespace Api.Services
{

    public class ProcessOrdersService : BackgroundService
    {
        private readonly ConsumerConfig consumerConfig;
        private readonly ProducerConfig producerConfig;
        //----------------------
        //private readonly ConsumerConfig consumerConfig2;
        public ProcessOrdersService(ConsumerConfig consumerConfig, ProducerConfig producerConfig)
        {
            this.producerConfig = producerConfig;
            this.consumerConfig = consumerConfig;
            //----------------------
            //this.consumerConfig2 = consumerConfig;
        }
        protected override async Task ExecuteAsync(CancellationToken stoppingToken)
        {
            Console.WriteLine("OrderProcessing Service Started\n\n");

            while (!stoppingToken.IsCancellationRequested)
            {
                var consumerHelper = new ConsumerWrapper(consumerConfig, "orderrequests");
                //var consumerHelper2 = new ConsumerWrapper(consumerConfig, "orderrequests");
                string orderRequest = consumerHelper.readMessage();
                //consumerHelper2.DisplayMessage();

                //Deserilaize 
                OrderRequest order = JsonConvert.DeserializeObject<OrderRequest>(orderRequest);
                //TODO:: Process Order
                Console.WriteLine($"Info: OrderHandler => Processing the order for {order.productname}\n\n");
                order.status = OrderStatus.COMPLETED;

                //Write to ReadyToShip Queue

                var producerWrapper = new ProducerWrapper(producerConfig,"readytoship");
                await producerWrapper.writeMessage(JsonConvert.SerializeObject(order));
                //--------------------------
               // var consumerHelper2 = new ConsumerWrapper(consumerConfig2, "orderrequests");
                //string processedOrder = consumerHelper2.readMessage();
                //OrderRequest order2 = JsonConvert.DeserializeObject<OrderRequest>(processedOrder);
                //Console.WriteLine($"Info: OrderHandler => Delivered the order for {order2.productname}\n\n");
                //order2.status = OrderStatus.DELIVERED;
                //----------------------------
            }
        }
    }
} 

ConsumerWrapper.cs:

namespace Api
{

    public class ConsumerWrapper
    {
        private string _topicName;
        private ConsumerConfig _consumerConfig;
        private Consumer<string,string> _consumer;
        private static readonly Random rand = new Random();
        public ConsumerWrapper(ConsumerConfig config,string topicName)
        {
            this._topicName = topicName;
            this._consumerConfig = config;
            this._consumer = new Consumer<string,string>(this._consumerConfig);
            this._consumer.Subscribe(topicName);
        }
        public string readMessage(){
            var consumeResult = this._consumer.Consume();
            return consumeResult.Value;
        }
        public void DisplayMessage()
        {
            var consumeResult = this._consumer.Consume();
            Console.WriteLine(consumeResult.Value);
            Console.WriteLine($"Info: OrderHandler => Delivered the order for {consumeResult.Value}\n\n");
            return;
        }
    }
} 

我希望能够多次调用 Consumer 类并能够阅读同一主题。我知道需要创建多个分区/组 ID 才能做到这一点。但我无法弄清楚在哪里以及如何做到这一点。

【问题讨论】:

  • 我对 .NET 不太熟悉,但是要让多个消费者从同一个主题消费,您只需为消费者指定不同的组 id,这是您在创建消费者时设置的配置.请查看kafka.apache.org/documentation/#consumerconfigs 了解更多信息。
  • 能不能看一下我上面加的appsettings.json代码。我是否需要添加另一个具有不同组 ID 的消费者对象?这是你的意思吗?。
  • 据我从上面的代码中了解到,您需要创建不同的 appsettings.json 并为单独的消费者使用不同的 groupid。

标签: c# asp.net .net apache-kafka stream


【解决方案1】:

您可以通过在Kafka中使用组ID概念来做到这一点,只需为多个消费者使用相同的组ID,以避免重复消费来自同一主题的数据。

【讨论】:

  • 你能详细说明一下吗?我想启用重复使用数据。那么,我该怎么办?另外,您能否修改代码或点我需要进行更改的地方
  • 我不熟悉.net,但我可以在java Properties props = new Properties(); 中为您提供帮助props.put("bootstrap.servers", "localhost:9092"); props.put("group.id", "consumer"); props.put("key.deserializer", StringDeserializer.class.getName()); props.put("value.deserializer", StringDeserializer.class.getName()); KafkaConsumer consumer = new KafkaConsumer(props);
  • 我们应该编写具有不同组 id 的多个消费者,因为多个消费者将消费相同的数据。 props.put("group.id", "consumer1l");
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-04-07
  • 1970-01-01
  • 1970-01-01
  • 2013-08-30
  • 2017-09-23
相关资源
最近更新 更多