【问题标题】:Kafka 0.10 Java Client TimeoutException: Batch containing 1 record(s) expiredKafka 0.10 Java Client TimeoutException:包含 1 条记录的批次已过期
【发布时间】:2016-11-10 15:14:32
【问题描述】:

我有一个单节点、多 (3) 个代理 Zookeeper / Kafka 设置。我正在使用 Kafka 0.10 Java 客户端。

我写了以下简单的远程(在与 Kafka 不同的服务器上)生产者(在代码中我用 MYIP 替换了我的公共 IP 地址):

Properties config = new Properties();
try {
    config.put(ProducerConfig.CLIENT_ID_CONFIG, InetAddress.getLocalHost().getHostName());
    config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "MYIP:9092, MYIP:9093, MYIP:9094");
    config.put(ProducerConfig.ACKS_CONFIG, "all");
    config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
    config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer");
    producer = new KafkaProducer<String, byte[]>(config);
    Schema.Parser parser = new Schema.Parser();
    schema = parser.parse(GATEWAY_SCHEMA);
    recordInjection = GenericAvroCodecs.toBinary(schema);
    GenericData.Record avroRecord = new GenericData.Record(schema);
    //Filling in avroRecord (code not here)
    byte[] bytes = recordInjection.apply(avroRecord);

    Future<RecordMetadata> future = producer.send(new ProducerRecord<String, byte[]>(datasetId+"", "testKey", bytes));
    RecordMetadata data = future.get();
} catch (Exception e) {
    e.printStackTrace();
}

我的 3 个代理的服务器属性如下所示(在 3 个不同的服务器属性文件中,broker.id 为 0、1、2,侦听器为 PLAINTEXT://:9092、PLAINTEXT://:9093、PLAINTEXT:/ /:9094 和 host.name 是 10.2.0.4、10.2.0.5、10.2.0.6)。 这是第一个服务器属性文件:

broker.id=0
listeners=PLAINTEXT://:9092
num.network.threads=3
num.io.threads=8
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600
log.dirs=/tmp/kafka1-logs
num.partitions=1
num.recovery.threads.per.data.dir=1
log.retention.hours=168
log.segment.bytes=1073741824
log.retention.check.interval.ms=300000
zookeeper.connect=localhost:2181
zookeeper.connection.timeout.ms=6000

当我执行代码时,我得到以下异常:

java.util.concurrent.ExecutionException: org.apache.kafka.common.errors.TimeoutException: Batch containing 1 record(s) expired due to timeout while requesting metadata from brokers for 100101-0
    at org.apache.kafka.clients.producer.internals.FutureRecordMetadata.valueOrError(FutureRecordMetadata.java:65)
    at org.apache.kafka.clients.producer.internals.FutureRecordMetadata.get(FutureRecordMetadata.java:52)
    at org.apache.kafka.clients.producer.internals.FutureRecordMetadata.get(FutureRecordMetadata.java:25)
    at com.nr.roles.gateway.GatewayManager.addTransaction(GatewayManager.java:212)
    at com.nr.roles.gateway.gw.service(gw.java:126)
    at javax.servlet.http.HttpServlet.service(HttpServlet.java:790)
    at org.eclipse.jetty.servlet.ServletHolder.handle(ServletHolder.java:821)
    at org.eclipse.jetty.servlet.ServletHandler.doHandle(ServletHandler.java:583)
    at org.eclipse.jetty.server.handler.ContextHandler.doHandle(ContextHandler.java:1158)
    at org.eclipse.jetty.servlet.ServletHandler.doScope(ServletHandler.java:511)
    at org.eclipse.jetty.server.handler.ContextHandler.doScope(ContextHandler.java:1090)
    at org.eclipse.jetty.server.handler.ScopedHandler.handle(ScopedHandler.java:141)
    at org.eclipse.jetty.server.handler.HandlerCollection.handle(HandlerCollection.java:109)
    at org.eclipse.jetty.server.handler.HandlerWrapper.handle(HandlerWrapper.java:119)
    at org.eclipse.jetty.server.Server.handle(Server.java:517)
    at org.eclipse.jetty.server.HttpChannel.handle(HttpChannel.java:308)
    at org.eclipse.jetty.server.HttpConnection.onFillable(HttpConnection.java:242)
    at org.eclipse.jetty.io.AbstractConnection$ReadCallback.succeeded(AbstractConnection.java:261)
    at org.eclipse.jetty.io.FillInterest.fillable(FillInterest.java:95)
    at org.eclipse.jetty.io.SelectChannelEndPoint$2.run(SelectChannelEndPoint.java:75)
    at org.eclipse.jetty.util.thread.strategy.ExecuteProduceConsume.produceAndRun(ExecuteProduceConsume.java:213)
    at org.eclipse.jetty.util.thread.strategy.ExecuteProduceConsume.run(ExecuteProduceConsume.java:147)
    at org.eclipse.jetty.util.thread.QueuedThreadPool.runJob(QueuedThreadPool.java:654)
    at org.eclipse.jetty.util.thread.QueuedThreadPool$3.run(QueuedThreadPool.java:572)
    at java.lang.Thread.run(Thread.java:745)
 Caused by: org.apache.kafka.common.errors.TimeoutException: Batch containing 1 record(s) expired due to timeout while requesting metadata from brokers for 100101-0

有人知道我错过了什么吗?任何帮助,将不胜感激。非常感谢

【问题讨论】:

  • 我也尝试过与上述相同的操作,但只有一个代理(在端口 9092 上)。我仍然得到完全相同的异常。我确保远程机器上的 broker 和 zookeeper 端口是打开的,我可以从 Producer 机器上远程登录它们。

标签: java apache-kafka kafka-producer-api


【解决方案1】:

我遇到了同样的问题。

您应该更改您的 kafka server.properties 以指定 IP 地址。 例如:

PLAINTEXT://<b>YOUIP</b>:9093

如果没有,kafka 将使用主机名,如果生产者无法获取主机,即使可以 telnet 也无法向 kafka 发送消息。

【讨论】:

  • 我的属性在我的 server.properties 中设置为 listeners=PLAINTEXT://domain_name:9092,但仍然收到此异常 org.apache.kafka.common.errors.TimeoutException: Batch Expired java.util.concurrent.ExecutionException,从外部服务器连接时,我也将 request.timeout.ms 增加到更高的值。
【解决方案2】:

BOOTSTRAP_SERVERS_CONFIG 配置中的端口信息不正确(MYIP:9092)。

正如您在 server.properties 中提到的“PLAINTEXT://:9093, PLAINTEXT://:9093, PLAINTEXT://:9094”。

【讨论】:

  • 对不起,我在这里打错了,在 server.properties 中是“PLAINTEXT://:9092, PLAINTEXT://:9093, PLAINTEXT://:9094”。所以BOOTSTRAP_SERVERS_CONFIG 端口是正确的。
【解决方案3】:

This 回答分享了一些见解。您可以增加 request.timeout.ms 生产者配置,这将允许客户端在批次过期之前将其排队更长时间。

您可能还想查看batch.sizelinger.ms 配置并找到适合您情况的最佳配置。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-03-06
    • 2018-04-28
    • 2017-05-24
    • 1970-01-01
    • 1970-01-01
    • 2012-10-10
    • 2018-09-30
    • 2015-01-01
    相关资源
    最近更新 更多