【问题标题】:java.nio.channels.ClosedChannelException while Consuming message from storm spoutjava.nio.channels.ClosedChannelException 同时消费来自风暴喷口的消息
【发布时间】:2018-10-26 13:35:45
【问题描述】:

我已经编写了风暴拓扑,它使用 kafka spout 从 kafka 获取数据,它在我的本地环境中运行良好,但在集群中

我收到以下错误:

2018-05-16 18:25:59.358 o.a.s.k.ZkCoordinator Thread-25-kafkaSpout-executor[20 20] [INFO] 任务 [1/1] 刷新分区管理器连接 2018-05-16 18:25:59.359 oaskDynamicBrokersReader Thread-25-kafkaSpout-executor[20 20] [INFO] 从 Zookeeper 读取分区信息:GlobalPartitionInformation{topic=data-ops, partitionMap={0=uat-datalake-node2 .org:6667}} 2018-05-16 18:25:59.359 oaskKafkaUtils Thread-25-kafkaSpout-executor[20 20] [INFO] 任务 [1/1] 分配 [Partition{host=uat-datalake-node2.org:6667, topic=数据操作,分区=0}] 2018-05-16 18:25:59.360 o.a.s.k.ZkCoordinator Thread-25-kafkaSpout-executor[20 20] [INFO] 任务 [1/1] 已删除分区管理器:[] 2018-05-16 18:25:59.360 o.a.s.k.ZkCoordinator Thread-25-kafkaSpout-executor[20 20] [INFO] 任务 [1/1] 新分区管理器:[] 2018-05-16 18:25:59.360 o.a.s.k.ZkCoordinator Thread-25-kafkaSpout-executor[20 20] [INFO] Task [1/1] 完成刷新 2018-05-16 18:25:59.361 k.c.SimpleConsumer Thread-25-kafkaSpout-executor[20 20] [INFO] 由于错误重新连接: java.nio.channels.ClosedChannelException 在 kafka.network.BlockingChannel.send(BlockingChannel.scala:110) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer.liftedTree1$1(SimpleConsumer.scala:85) [kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer.kafka$consumer$SimpleConsumer$$sendRequest(SimpleConsumer.scala:83) [kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1$$anonfun$apply$mcV$sp$1.apply$mcV$sp(SimpleConsumer.scala:132) [kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1$$anonfun$apply$mcV$sp$1.apply(SimpleConsumer.scala:132) [kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1$$anonfun$apply$mcV$sp$1.apply(SimpleConsumer.scala:132) [kafka_2.10-0.10.2.1.jar:?] 在 kafka.metrics.KafkaTimer.time(KafkaTimer.scala:33) [kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1.apply$mcV$sp(SimpleConsumer.scala:131) [kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1.apply(SimpleConsumer.scala:131) [kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1.apply(SimpleConsumer.scala:131) [kafka_2.10-0.10.2.1.jar:?] 在 kafka.metrics.KafkaTimer.time(KafkaTimer.scala:33) [kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer.fetch(SimpleConsumer.scala:130) [kafka_2.10-0.10.2.1.jar:?] 在 kafka.javaapi.consumer.SimpleConsumer.fetch(SimpleConsumer.scala:47) [kafka_2.10-0.10.2.1.jar:?] 在 org.apache.storm.kafka.KafkaUtils.fetchMessages(KafkaUtils.java:191) [storm-kafka-1.0.1.jar:1.0.1] 在 org.apache.storm.kafka.PartitionManager.fill(PartitionManager.java:189) [storm-kafka-1.0.1.jar:1.0.1] 在 org.apache.storm.kafka.PartitionManager.next(PartitionManager.java:138) [storm-kafka-1.0.1.jar:1.0.1] 在 org.apache.storm.kafka.KafkaSpout.nextTuple(KafkaSpout.java:135) [storm-kafka-1.0.1.jar:1.0.1] 在 org.apache.storm.daemon.executor$fn__6505$fn__6520$fn__6551.invoke(executor.clj:651) [storm-core-1.0.1.2.5.3.0-37.jar:1.0.1.2.5.3.0- 37] 在 org.apache.storm.util$async_loop$fn__554.invoke(util.clj:484) [storm-core-1.0.1.2.5.3.0-37.jar:1.0.1.2.5.3.0-37] 在 clojure.lang.AFn.run(AFn.java:22) [clojure-1.7.0.jar:?] 在 java.lang.Thread.run(Thread.java:748) [?:1.8.0_144] 2018-05-16 18:26:09.372 o.a.s.k.KafkaUtils Thread-25-kafkaSpout-executor[20 20] [WARN] 获取消息时出现网络错误: java.net.SocketTimeoutException 在 sun.nio.ch.SocketAdaptor$SocketInputStream.read(SocketAdaptor.java:211) ~[?:1.8.0_144] 在 sun.nio.ch.ChannelInputStream.read(ChannelInputStream.java:103) ~[?:1.8.0_144] 在 java.nio.channels.Channels$ReadableByteChannelImpl.read(Channels.java:385) ~[?:1.8.0_144] 在 org.apache.kafka.common.network.NetworkReceive.readFromReadableChannel(NetworkReceive.java:81) ~[kafka-clients-0.10.2.1.jar:?] 在 kafka.network.BlockingChannel.readCompletely(BlockingChannel.scala:129) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.network.BlockingChannel.receive(BlockingChannel.scala:120) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer.liftedTree1$1(SimpleConsumer.scala:99) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer.kafka$consumer$SimpleConsumer$$sendRequest(SimpleConsumer.scala:83) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1$$anonfun$apply$mcV$sp$1.apply$mcV$sp(SimpleConsumer.scala:132) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1$$anonfun$apply$mcV$sp$1.apply(SimpleConsumer.scala:132) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1$$anonfun$apply$mcV$sp$1.apply(SimpleConsumer.scala:132) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.metrics.KafkaTimer.time(KafkaTimer.scala:33) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1.apply$mcV$sp(SimpleConsumer.scala:131) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1.apply(SimpleConsumer.scala:131) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1.apply(SimpleConsumer.scala:131) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.metrics.KafkaTimer.time(KafkaTimer.scala:33) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer.fetch(SimpleConsumer.scala:130) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.javaapi.consumer.SimpleConsumer.fetch(SimpleConsumer.scala:47) ~[kafka_2.10-0.10.2.1.jar:?] 在 org.apache.storm.kafka.KafkaUtils.fetchMessages(KafkaUtils.java:191) [storm-kafka-1.0.1.jar:1.0.1] 在 org.apache.storm.kafka.PartitionManager.fill(PartitionManager.java:189) [storm-kafka-1.0.1.jar:1.0.1] 在 org.apache.storm.kafka.PartitionManager.next(PartitionManager.java:138) [storm-kafka-1.0.1.jar:1.0.1] 在 org.apache.storm.kafka.KafkaSpout.nextTuple(KafkaSpout.java:135) [storm-kafka-1.0.1.jar:1.0.1] 在 org.apache.storm.daemon.executor$fn__6505$fn__6520$fn__6551.invoke(executor.clj:651) [storm-core-1.0.1.2.5.3.0-37.jar:1.0.1.2.5.3.0- 37] 在 org.apache.storm.util$async_loop$fn__554.invoke(util.clj:484) [storm-core-1.0.1.2.5.3.0-37.jar:1.0.1.2.5.3.0-37] 在 clojure.lang.AFn.run(AFn.java:22) [clojure-1.7.0.jar:?] 在 java.lang.Thread.run(Thread.java:748) [?:1.8.0_144] 2018-05-16 18:26:09.373 o.a.s.k.KafkaSpout Thread-25-kafkaSpout-executor[20 20] [WARN] 获取失败 org.apache.storm.kafka.FailedFetchException:java.net.SocketTimeoutException 在 org.apache.storm.kafka.KafkaUtils.fetchMessages(KafkaUtils.java:199) ~[storm-kafka-1.0.1.jar:1.0.1] 在 org.apache.storm.kafka.PartitionManager.fill(PartitionManager.java:189) ~[storm-kafka-1.0.1.jar:1.0.1] 在 org.apache.storm.kafka.PartitionManager.next(PartitionManager.java:138) ~[storm-kafka-1.0.1.jar:1.0.1] 在 org.apache.storm.kafka.KafkaSpout.nextTuple(KafkaSpout.java:135) [storm-kafka-1.0.1.jar:1.0.1] 在 org.apache.storm.daemon.executor$fn__6505$fn__6520$fn__6551.invoke(executor.clj:651) [storm-core-1.0.1.2.5.3.0-37.jar:1.0.1.2.5.3.0- 37] 在 org.apache.storm.util$async_loop$fn__554.invoke(util.clj:484) [storm-core-1.0.1.2.5.3.0-37.jar:1.0.1.2.5.3.0-37] 在 clojure.lang.AFn.run(AFn.java:22) [clojure-1.7.0.jar:?] 在 java.lang.Thread.run(Thread.java:748) [?:1.8.0_144] 引起:java.net.SocketTimeoutException 在 sun.nio.ch.SocketAdaptor$SocketInputStream.read(SocketAdaptor.java:211) ~[?:1.8.0_144] 在 sun.nio.ch.ChannelInputStream.read(ChannelInputStream.java:103) ~[?:1.8.0_144] 在 java.nio.channels.Channels$ReadableByteChannelImpl.read(Channels.java:385) ~[?:1.8.0_144] 在 org.apache.kafka.common.network.NetworkReceive.readFromReadableChannel(NetworkReceive.java:81) ~[kafka-clients-0.10.2.1.jar:?] 在 kafka.network.BlockingChannel.readCompletely(BlockingChannel.scala:129) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.network.BlockingChannel.receive(BlockingChannel.scala:120) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer.liftedTree1$1(SimpleConsumer.scala:99) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer.kafka$consumer$SimpleConsumer$$sendRequest(SimpleConsumer.scala:83) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1$$anonfun$apply$mcV$sp$1.apply$mcV$sp(SimpleConsumer.scala:132) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1$$anonfun$apply$mcV$sp$1.apply(SimpleConsumer.scala:132) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1$$anonfun$apply$mcV$sp$1.apply(SimpleConsumer.scala:132) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.metrics.KafkaTimer.time(KafkaTimer.scala:33) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1.apply$mcV$sp(SimpleConsumer.scala:131) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1.apply(SimpleConsumer.scala:131) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer$$anonfun$fetch$1.apply(SimpleConsumer.scala:131) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.metrics.KafkaTimer.time(KafkaTimer.scala:33) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.consumer.SimpleConsumer.fetch(SimpleConsumer.scala:130) ~[kafka_2.10-0.10.2.1.jar:?] 在 kafka.javaapi.consumer.SimpleConsumer.fetch(SimpleConsumer.scala:47) ~[kafka_2.10-0.10.2.1.jar:?] 在 org.apache.storm.kafka.KafkaUtils.fetchMessages(KafkaUtils.java:191) ~[storm-kafka-1.0.1.jar:1.0.1] ... 7 更多

【问题讨论】:

    标签: java apache-kafka apache-storm


    【解决方案1】:

    当 Storm 工作人员尝试从 Kafka 代理读取数据时,您似乎遇到了超时。也许它们之间的连接不稳定或缓慢?

    也就是说,堆栈跟踪似乎表明消费者已经重新连接,所以如果这种情况很少发生,您可能只是在工作人员和 Kafka 之间的连接中遇到了问题。

    如果这种情况经常发生,并且您确定连接稳定,我会尝试在 Kafka 邮件列表https://kafka.apache.org/contact 上询问。如果您发布您的问题以及您使用的是哪个 Kafka 版本,他们可能会告诉您是否存在可能导致使用者套接字超时的问题。

    【讨论】:

    • 我有另一个单独的拓扑,它在 kafka 主题中发布数据,它的工作正常问题在消费者端。这里我使用的是 kafka 版本 0.10.0,关于连接它可以用于插入,但我不确定连接稳定性。请告知可能的分辨率提示。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2014-08-16
    • 2014-10-18
    • 2012-05-30
    • 2012-01-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多