【问题标题】:Create multiple consumer group for multiple topic为多个主题创建多个消费者组
【发布时间】:2016-07-18 13:24:31
【问题描述】:

就我而言,我需要多个主题,每个主题都与多个消费者相关联。我想为每个主题设置一个消费者组。我在kafka .net 客户端中没有找到任何方法,因此我可以动态创建消费者组并将主题与该消费者组链接。我正在使用kafka 0.9.0 版本,请告诉我是否需要更改为kafka server 设置或Zookeeper

【问题讨论】:

    标签: apache-kafka


    【解决方案1】:

    我不确定您所说的“为每个主题设置一个消费者组”是什么意思。如果您启动一个新的消费者组(即,使用相应的 group-ID 启动一个消费者),消费者组将决定它想要消费的主题(只需订阅主题)。它没有特殊的配置,因此,总是在主题订阅时动态创建一个组。

    更新(.net 客户端):

    我不熟悉 .net 客户端。但是,根据 Github 页面 (github.com/Jroland/kafka-net) 看来,尚不支持消费组。

    但是,您似乎可以使用白名单来仅读取某些分区。因此,您可以手动分配负载:

    来自https://github.com/Jroland/kafka-net#consumer-1

    如果未提供白名单,则所有分区都将被消耗,为每个分区领导者创建一个 KafkaConnection

    【讨论】:

    【解决方案2】:

    我使用 Microsoft .NET kafka 构建了一个快速原型,链接如下。不确定它是否解决了您的问题。

    但是,我强烈推荐你使用这个库,因为它包含比 kafka-net 更多的功能(例如,支持 zookeeper 来维护偏移量、主题组等)

    https://github.com/Microsoft/CSharpClient-for-Kafka

    示例代码

    这将向 kafka 发送 10 条消息,并在消费者收到消息时将消息输出到控制台。

    static void Main(string[] args)
        {
            Task.Factory.StartNew(() =>
            {
                ConsumerConfiguration consumerConfig = new ConsumerConfiguration
                {
                    AutoCommit = true, 
                    AutoCommitInterval = 1000, 
                    GroupId = "group1",
                    ConsumerId = "1",
                    AutoOffsetReset = OffsetRequest.SmallestTime,
                    NumberOfTries = 20,
                    ZooKeeper = new ZooKeeperConfiguration("localhost:2181", 30000, 30000, 2000)           
                };
                var consumer = new ZookeeperConsumerConnector(consumerConfig, true);
                var dictionaryMapping = new Dictionary<string, int>();
                dictionaryMapping.Add("topic1", 1);
    
                var streams = consumer.CreateMessageStreams(dictionaryMapping, new DefaultDecoder());
    
                var messageStream = streams["topic1"][0];
    
                foreach (var message in messageStream.GetCancellable(new CancellationToken()))
                {
                    Console.WriteLine("Response: P{0},O{1} : {2}", message.PartitionId, message.Offset, Encoding.UTF8.GetString(message.Payload));
    
                    //If you set AutoCommit to false, you can commit by yourself from this command.
                    //consumer.CommitOffsets()     
                }
            });
    
    
            var brokerConfig = new BrokerConfiguration()
            {
                BrokerId = 1,
                Host = "localhost",
                Port = 9092
            };
            var config = new ProducerConfiguration(new List<BrokerConfiguration> { brokerConfig });
            config.CompressionCodec = CompressionCodecs.DefaultCompressionCodec;
            config.ProducerRetries = 3;
            config.RequiredAcks = -1;            
            var kafkaProducer = new Producer(config);
    
            byte[] payloadData = Encoding.UTF8.GetBytes("Test Message");
            var inputMessage = new Message(payloadData);
            var data = new ProducerData<string, Message>("topic1", inputMessage);
    
            for (int i = 0; i < 10; i++)
            {
                kafkaProducer.Send(data);
            }
    
            Console.ReadLine();
        }
    

    希望对您有所帮助。

    【讨论】:

      猜你喜欢
      • 2020-10-12
      • 2020-06-03
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-08-31
      • 1970-01-01
      • 2017-01-26
      • 1970-01-01
      相关资源
      最近更新 更多