【问题标题】:Kafka Producer error Expiring 10 record(s) for TOPIC:XXXXXX: 6686 ms has passed since batch creation plus linger timeKafka Producer 错误正在过期 TOPIC:XXXXXX 的 10 条记录:自批处理创建以来已过去 6686 毫秒加上逗留时间
【发布时间】:2018-03-26 19:12:30
【问题描述】:

卡夫卡版本:0.10.2.1,

Kafka Producer error Expiring 10 record(s) for TOPIC:XXXXXX: 6686 ms has passed since batch creation plus linger time
org.apache.kafka.common.errors.TimeoutException: Expiring 10 record(s) for TOPIC:XXXXXX: 6686 ms has passed since batch creation plus linger time

【问题讨论】:

  • 当生产者无法根据其拥有的元数据将数据发送到它认为负责消息的代理时,您会收到此错误。卡夫卡经纪人死了还是你的制片人当时有连接问题?
  • 我全天也间歇性地收到此错误。寻找答案
  • 当我更改我的 kafka 生产者“max.request.size”:“4713360”,“acks”:“all”,“timeout.ms”:“18000”,“batch. size": "100000", -- 这是字节大小.. "linger.ms":"100", "retries": "5", "min.insync.replicas":"2", "buffer.memory ":"66554432", "request.timeout.ms":"90000","block.on.buffer.full","true" 基本上是 linger.ms 和 batch.size 和 block.on.buffer.full 起主要作用这里

标签: apache-kafka


【解决方案1】:

发生此异常是因为您以比发送记录快得多的速度排队记录。

当您调用 send 方法时,ProducerRecord 将存储在内部缓冲区中,以发送给代理。该方法在 ProducerRecord 被缓冲后立即返回,无论它是否已发送。

记录被分组发送到代理,以减少每条消息的传输偷听并提高吞吐量。

将记录添加到批次后,发送该批次有一个时间限制,以确保它已在指定的时间内发送。这由 Producer 配置参数 request.timeout.ms 控制,默认为 30 秒。

如果批处理的排队时间超过了超时限制,则会抛出异常。该批次中的记录将从发送队列中删除。

生产者配置 block.on.buffer.full、metadata.fetch.timeout.ms 和 timeout.ms 已被删除。它们最初在 Kafka 0.9.0.0 中被弃用。

因此尝试增加 request.timeout.ms

不过,如果您对吞吐量有任何问题,也可以参考以下blog

【讨论】:

  • 不幸的是,这个链接对我来说已经死了。
  • @CristianoFontes 是的,不幸的是,它已经关闭了。但是上面已经涵盖了很多答案。
  • 可能发生这种情况的一种情况是,当您开始工作并以比平时更快的速度处理积压工作时。 Kafka 可能难以处理流量激增的情况。
  • @arglee this "request.timeout.ms" 表示“Sender thread”发送缓冲消息所花费的网络时间?
【解决方案2】:

我收到了同样的消息,我修复了它从 zookeeper 清理 kafka 数据。之后它就可以工作了。

【讨论】:

  • 您是刚刚从 ZK 中清除了某些内容,还是删除了 ZK 的所有 Kafka 托管数据?
  • 就我而言,我清理了所有 ZK 数据。
  • 谁能指导我如何做到这一点?我正面临着确切的问题,除了清理数据外,都试过了!我愿意清理所有 zk 数据,但不明白怎么做!
【解决方案3】:

我在 aks 集群中遇到了同样的问题,只是重新启动 kafka 和 zookeeper 服务器解决了这个问题。

【讨论】:

    【解决方案4】:

    假设一个主题有 100 个分区 (0-99)。 Kafka 允许您通过指定特定分区为主题生成记录。面临我试图生成分区 > 99 的问题,因为经纪人拒绝这些记录。

    【讨论】:

      【解决方案5】:

      我们尝试了一切,但没有成功。

      1. 生产者批量减少,request.timeout.ms 增加。
      2. 重启了目标 kafka 集群,还是不行。
      3. 检查了目标 kafka 集群上的复制,也正常工作。
      4. 在 prodcuer 属性中添加了 retries、retries.backout.ms。
      5. 在 kafka prodcuer 属性中添加了 linger.time。

      最后,我们的案例是 kafka 集群本身存在问题,我们无法从 2 个服务器中获取其间的元数据。

      当我们将目标 kafka 集群更改为我们的开发箱时,它工作正常。

      【讨论】:

        【解决方案6】:

        当 wither brokers/topics/partitions 无法与生产者联系或生产者在队列之前超时时,会出现此问题。

        我发现即使是现场经纪人也可能遇到这个问题。就我而言,主题分区领导者指向非活动代理 ID。要解决此问题,您必须将这些领导者迁移到活跃的代理。

        对受影响的主题使用主题重新分配工具。 话题迁移:https://kafka.apache.org/21/documentation.html#basic_ops_automigrate

        【讨论】:

          【解决方案7】:

          卡夫卡码头箱

          花了很多时间找出发生了什么,包括更改 server.propertiesproducer.properties 和我的代码 (Eclipse)。这对我不起作用(我从笔记本电脑向 Linux 服务器上的 Kafka Docker 发送消息)

          我清理了 Kafka 和 Zookeeper,并通过 docker-compose.yml 重新安装了它们(我是新手)。请查看我的docker-compose.yml 文件并按照我如何将这些 IP 更改为我的 Linux 服务器的 IP

          bitnami/kafka

          bitnami/kafka

          到...

          bitnami-changed

          而 10.5.1.30 是我的 Linux 服务器的 IP 地址

          wurstmeister 卡夫卡

          wurstmeister

          之后,我运行了我的代码,结果如下:

          result

          完整代码:

          import java.util.Properties;
          import java.util.concurrent.Future;
          
          import org.apache.kafka.clients.producer.KafkaProducer;
          import org.apache.kafka.clients.producer.Producer;
          import org.apache.kafka.clients.producer.ProducerRecord;
          import org.apache.kafka.clients.producer.RecordMetadata;
          
          public class SimpleProducer {
              public static void main(String[] args) throws Exception {
                  try {
                      String topicName = "demo";
                      Properties props = new Properties();
                      props.put("bootstrap.servers", "10.5.1.30:9092");
                      props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
                      props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
                      Producer<String, String> producer = new KafkaProducer<String, String>(props);
                      Future<RecordMetadata> f = producer.send(new ProducerRecord<String, String>(topicName, "Eclipse3"));
                      System.out.println("Message sent successfully, total of message is: " + f.get().toString());
                      producer.close();
                  } catch (Exception e) {
                      System.out.println(e.getMessage());
                  }
                  System.out.println("Successful");
          
              }
          }
          

          希望对您有所帮助。和平!!!

          【讨论】:

            猜你喜欢
            • 2018-04-28
            • 2017-06-30
            • 2016-11-10
            • 2020-10-07
            • 2021-11-06
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 2013-11-30
            相关资源
            最近更新 更多