【问题标题】:Expiring 1 record(s) for xxxxx: 30030 ms has passed since batch creation plus linger timexxxxx 的 1 条记录到期:自批次创建以来已过去 30030 毫秒加上延迟时间
【发布时间】:2018-04-28 04:39:08
【问题描述】:

我的用例: 使用 Postman,我调用了 Spring boot 肥皂端点。端点创建一个 KafkaProducer 并向特定主题发送消息。我还有一个 TaskScheduler 来使用该主题。

问题: 调用soap将消息推送到主题时,出现此错误:

2017-11-14 21:29:31.463 错误 6389 --- [广告 |生产者-3] DomainEntityProducer:即将到期的 1 条记录 DomainEntityCommandStream-0:自批处理创建以来已过去 30030 毫秒 加上逗留时间 2017-11-14 21:29:31.464 错误 6389 --- [nio-8080-exec-6] DomainEntityProducer: org.apache.kafka.common.errors.TimeoutException: Expiring 1 record(s) 对于 DomainEntityCommandStream-0:自批处理以来已过去 30030 毫秒 创造加上逗留时间

这是我用来推送主题的方法:

public DomainEntity push(DomainEntity pDomainEntity) throws Exception {
    logger.log(Level.INFO, "streaming...");
    wKafkaProperties.put("bootstrap.servers", "localhost:9092");
    wKafkaProperties.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
    wKafkaProperties.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
    KafkaProducer wKafkaProducer = new KafkaProducer(wKafkaProperties);
    ProducerRecord wProducerRecord = new ProducerRecord("DomainEntityCommandStream", getJSON(pDomainEntity));
    wKafkaProducer.send(wProducerRecord, (RecordMetadata r, Exception e) -> {
        if (e != null) {
            logger.log(Level.SEVERE, e.getMessage());
        }
    }).get();
    return pDomainEntity;
}

使用命令外壳脚本

./kafka-console-producer.sh --broker-list 10.0.1.15:9092 --topic 域实体命令流

./kafka-console-consumer.sh --boostrap-server 10.0.1.15:9092 --topic DomainEntityCommandStream --从头开始

效果很好。

通过Stackoverflow上的一些相关问题,我试图清除主题:

./kafka-topics.sh --zookeeper 10.0.1.15:9092 --alter --topic DomainEntityCommandStream --config retention.ms=1000

查看 kafka 日志,我发现保留时间已更改。

但是,不走运,我得到了同样的错误。

payload 小得离谱,为什么要更改 batch.size?

<soapenv:Envelope xmlns:soapenv="http://schemas.xmlsoap.org/soap/envelope/"
                  xmlns:gs="http://soap.problem.com">
   <soapenv:Header/>
   <soapenv:Body>
      <gs:streamDomainEntityRequest>
         <gs:domainEntity>
                <gs:name>12345</gs:name>
                <gs:value>Quebec</gs:value>
                <gs:version>666</gs:version>
            </gs:domainEntity>
      </gs:streamDomainEntityRequest>
   </soapenv:Body>
</soapenv:Envelope>

【问题讨论】:

  • 顺便说一句,我在 docker 容器中使用 zookeeper 和 kafka
  • 我发现了我的错误。在 Docker 容器中使用 kafka 时,需要在 yml 文件中指定 KAFKA_ADVERTISED_HOST_NAME。这使得容器外的生产者和消费者能够与 kafka 进行交互。
  • 你使用了哪个值?

标签: kafka-producer-api


【解决方案1】:

使用 Docker 和 Kafka 0.11.0.1 镜像需要在容器中添加以下环境参数:

KAFKA_ZOOKEEPER_CONNECT = X.X.X.X:XXXX(您的 zookeeper IP 或域:PORT 默认 2181)

KAFKA_ADVERTISED_HOST_NAME = X.X.X.X(您的 kafka IP 或域)

KAFKA_ADVERTISED_PORT = XXXX(你的 kafka 端口号默认为 9092)

可选:

KAFKA_BROKER_ID = 999(某个值)

KAFKA_CREATE_TOPICS=test:1:1(在开始时创建的一些主题名称)

如果它不起作用并且您仍然收到相同的消息(“Expiring X record(s) for xxxxx: XXXXX ms has been given since batch creation plus linger time”),您可以尝试从 zookeeper 中清理 kafka 数据。

【讨论】:

    猜你喜欢
    • 2018-03-26
    • 2017-06-30
    • 1970-01-01
    • 2016-11-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多