【问题标题】:How can I write messages in producer at once and read 1 message per minute with consumer?如何一次在生产者中编写消息并与消费者每分钟读取 1 条消息?
【发布时间】:2022-01-10 23:13:45
【问题描述】:

如何一次在生产者中写入消息并与消费者每分钟读取 1 条消息?

我可以使用的配置属性

note : please "max.poll.records" note that I cannot use the method

我的消费阶层:

var settings = ConfigurationManager.KafkaSettings.Topics[Topics.FaturaKaydetViaTp];

LogManager.Logger.Debug("Consumer initiating for {topic}", settings.TopicName);

using (var consumer = new ConsumerBuilder<Ignore, MailMessage>(consumerConfig).SetValueDeserializer(new ObjectDeserializer<MailMessage>()).Build())
{
    LogManager.Logger.Debug("Consumer initiated");

    LogManager.Logger.Debug("Subscribing for {topic}", settings.TopicName);
    
    consumer.Subscribe(settings.TopicName);

    try
    {
        while (true)
        {
            try
            {
                
                var cr = consumer.Consume();

                LogManager.Logger.Debug("Message received for '{topic}' at: '{topicPartitionOffset}'.", settings.TopicName, cr.TopicPartitionOffset);
                if (HandleOnMessage(cr.Value))
                    if (ConfigurationManager.KafkaSettings.AutoCommit == false)
                        consumer.Commit(cr);
            }
            catch (ConsumeException e)
            {
                LogManager.Logger.Fatal(e, "ConsumeException");
            }
        }
    }
    catch (OperationCanceledException)
    {
        // Ensure the consumer leaves the group cleanly and final offsets are committed.
        consumer.Close();
    }
}

【问题讨论】:

  • “每分钟一个”。为什么?如果您尝试进行批处理,Kafka 是错误的工具。如果您在消费调用之间等待太久,这最终会导致不必要的消费组重新平衡

标签: c# apache-kafka kafka-consumer-api confluent-kafka-dotnet


【解决方案1】:

一般策略需要您暂停消费者,让运行循环的线程休眠,然后重新订阅/恢复消费者。

这需要首先在ConsumerBuilder 中设置PartitionsAssignedHandler,然后保存返回的分区分配,因为Consumer.Pause()Consumer.Resume() 方法需要这些分配。

IEnumerable<TopicPartition> partitions; // TODO: assign this from handler
using (var consumer = ... ) { // set PartitionsAssignedHandler in ConsumerBuilder and set partitions above. 

        consumer.Subscribe(settings.TopicName);
        while (true)
        {
            try
            {
                // TODO: Might need to check if currently paused, somehow
                List<TopicPartitionError> resumeErrors = consumer.Resume(partitons);
                // TODO: handle resume errors

                var cr = consumer.Consume(); // This already gets one record from the assigned, resumed partitions

                LogManager.Logger.Debug("Message received for '{topic}' at: '{topicPartitionOffset}'.", cr.Topic, cr.TopicPartitionOffset);

                // Pause and sleep
                consumer.Pause(partitons);
                Thread.Sleep(1000 * 60); // 1 minute
            }
            catch (ConsumeException e)
            {
                LogManager.Logger.Fatal(e, "ConsumeException");
            }
        }

请记住,如果您运行多个实例,在暂停和恢复时总是会重新平衡,这会导致处理严重延迟,可能导致每分钟消耗超过一个。

【讨论】:

    【解决方案2】:

    首先,您需要将max.poll.records 设置为1,否则您将获取每个consumer.poll(Duration) 上的所有可用记录。

    同样Duration 传递给consumer.poll(…) 不会强制等待,它的工作方式不同。来自文档:

    如果有可用记录,此方法会立即返回。否则,它将等待通过的超时。如果超时,将返回一个空记录集

    要每分钟获取一个,您需要在短时间内轮询,而不是等待 60 秒(不准确),或者使用其他一些工具以 60 秒的间隔进行轮询。

    【讨论】:

    • 请“max.poll.records”注意我不能使用github.com/confluentinc/confluent-kafka-dotnet/issues/1451的方法
    • @GülsenKeskin - 因此您可以忽略这部分,因为在 C# 客户端中您无需任何额外的努力即可一一收到消息 :) - 您只需要整理间隔和轮询持续时间。
    【解决方案3】:

    根据您使用的 api,consumer.poll() 方法接受持续时间(以毫秒为单位)的参数。

    Java 中的示例代码:

    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(60000));
    

    文档中的更多信息:https://kafka.apache.org/26/javadoc/index.html?org/apache/kafka/clients/consumer/KafkaConsumer.html

    编辑:添加代码后。

    consumer.consume 方法也接受 duration.ms 的参数

    支持文档:https://docs.confluent.io/5.0.0/clients/confluent-kafka-dotnet/api/Confluent.Kafka.Consumer.html#Confluent_Kafka_Consumer_Consume_Confluent_Kafka_Message__System_Int32_

    【讨论】:

    • 我正在使用 confluent kafka,据我所知消费者没有 poll 方法
    • confluent kafka 只是 kafka 的一个发行版。你用什么连接到kafka的实例? python/java 还是只是一个 cli 消费者?
    • 我正在使用 c# ..
    • C# consumer.poll 方法也接受持续时间作为参数。公共无效投票(时间跨度超时)docs.confluent.io/5.0.0/clients/confluent-kafka-dotnet/api/…
    • 谢谢。那么你能告诉我一个如何使用它的例子吗?
    猜你喜欢
    • 1970-01-01
    • 2015-11-10
    • 1970-01-01
    • 1970-01-01
    • 2023-01-21
    • 2020-08-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多