使用Storm客户端写入FS时遭遇maprfs及Kafka方法错误求助
问题解决:Storm写入文件系统时的两类错误
错误1:org.apache.hadoop.fs.UnsupportedFileSystemException: No FileSystem for scheme "maprfs"
原因
Storm的HdfsBolt无法识别maprfs协议,核心是缺少MapR文件系统的实现依赖,或未在配置中注册对应的文件系统类。
解决步骤
- 添加MapR依赖:如果用Maven构建,在
pom.xml中添加与集群MapR版本匹配的依赖:<dependency> <groupId>com.mapr.hadoop</groupId> <artifactId>maprfs</artifactId> <version>你的MapR集群版本号</version> </dependency> - 注册文件系统实现:在Storm作业配置或集群
storm.yaml中添加配置,指定maprfs对应的实现类:fs.maprfs.impl: com.mapr.fs.MapRFileSystem - 确保依赖在Classpath中:集群部署时,将MapR相关jar放入Storm节点的
lib目录;或作业打包时将依赖一并打入(注意避免版本冲突)。
错误2:java.lang.NoSuchMethodError: 'org.apache.kafka.clients.consumer.ConsumerRecords org.apache.kafka.clients.consumer.Consumer.poll(java.time.Duration)'
原因
Storm Kafka Spout与Kafka客户端版本不兼容。poll(Duration)是Kafka 2.0+客户端新增的方法,若项目混入低版本Kafka客户端(如0.11.x),或Storm依赖的Kafka版本与实际使用版本不一致,就会触发该错误。当前使用的Storm 2.6.1需搭配兼容的Kafka客户端版本。
解决步骤
- 对齐Kafka版本:Storm 2.6.1推荐搭配Kafka 2.8.x系列客户端。在
pom.xml中明确指定版本,并排除冲突依赖:<dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>2.8.1</version> <scope>compile</scope> </dependency> <dependency> <groupId>org.apache.storm</groupId> <artifactId>storm-kafka-client</artifactId> <version>2.6.1</version> <exclusions> <exclusion> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> </exclusion> </exclusions> </dependency> - 排查依赖冲突:执行
mvn dependency:tree查看依赖树,清理所有非指定版本的kafka-clients引用。 - 重新打包作业:清理旧编译产物,重新打包作业jar,确保仅包含兼容版本的Kafka客户端类。
详细错误日志
consumer=org.apache.kafka.clients.consumer.KafkaConsumer@305c2b21, topic-partitions=[]] 2024-03-28 10:57:12.259 o.a.s.k.s.KafkaSpout Thread-15-kafka-spout-executor[3, 3] [INFO] Partitions reassignment. [task-ID=3, consumer-group=topic03-group-id, consumer=org.apache.kafka.clients.consumer.KafkaConsumer@305c2b21, topic-partitions=[topic03-0]] 2024-03-28 10:57:12.264 o.a.k.c.c.i.AbstractCoordinator Thread-15-kafka-spout-executor[3, 3] [INFO] [Consumer clientId=consumer-1, groupId=topic03-group-id] Discovered group coordinator node6.ezm.tst:9092 (id: 2147483647 rack: null) 2024-03-28 10:57:12.271 o.a.k.c.c.i.Fetcher Thread-15-kafka-spout-executor[3, 3] [INFO] [Consumer clientId=consumer-1, groupId=topic03-group-id] Resetting offset for partition topic03-0 to offset 0. 2024-03-28 10:57:12.277 o.a.s.k.s.KafkaSpout Thread-15-kafka-spout-executor[3, 3] [INFO] Initialization complete 2024-03-28 10:57:12.278 o.a.s.k.s.KafkaSpout Thread-15-kafka-spout-executor[3, 3] [INFO] Partitions assignments has changed, updating metrics. 2024-03-28 10:57:12.278 o.a.s.k.s.m.KafkaOffsetMetricManager Thread-15-kafka-spout-executor[3, 3] [INFO] Registering metric for topicPartition: topic03-0 2024-03-28 10:57:12.279 o.a.s.k.s.m.KafkaOffsetTopicMetrics Thread-15-kafka-spout-executor[3, 3] [INFO] Create KafkaOffsetTopicMetrics for topic: topic03 2024-03-28 10:57:12.281 o.a.s.k.s.m.KafkaOffsetPartitionMetrics Thread-15-kafka-spout-executor[3, 3] [INFO] Running KafkaOffsetMetricSet 2024-03-28 10:57:12.285 o.a.s.u.Utils Thread-15-kafka-spout-executor[3, 3] [ERROR] Async loop died! java.lang.NoSuchMethodError: 'org.apache.kafka.clients.consumer.ConsumerRecords org.apache.kafka.clients.consumer.Consumer.poll(java.time.Duration)' at org.apache.storm.kafka.spout.KafkaSpout.pollKafkaBroker(KafkaSpout.java:360) ~[stormjar.jar:?] at org.apache.storm.kafka.spout.KafkaSpout.nextTuple(KafkaSpout.java:289) ~[stormjar.jar:?] at org.apache.storm.executor.spout.SpoutExecutor$2.call(SpoutExecutor.java:187) ~[storm-client-2.6.1.jar:2.6.1] at org.apache.storm.executor.spout.SpoutExecutor$2.call(SpoutExecutor.java:153) ~[storm-client-2.6.1.jar:2.6.1] at org.apache.storm.utils.Utils$1.run(Utils.java:398) [storm-client-2.6.1.jar:2.6.1] at java.base/java.lang.Thread.run(Thread.java:834) [?:?] 2024-03-28 10:57:12.275 o.a.s.u.Utils Thread-16-hdfs-bolt-executor[2, 2] [ERROR] Async loop died! java.lang.RuntimeException: Error preparing HdfsBolt: No FileSystem for scheme "maprfs" at org.apache.storm.hdfs.bolt.AbstractHdfsBolt.prepare(AbstractHdfsBolt.java:120) ~[stormjar.jar:?] at org.apache.storm.executor.bolt.BoltExecutor.init(BoltExecutor.java:128) ~[storm-client-2.6.1.jar:2.6.1] at org.apache.storm.executor.bolt.BoltExecutor.call(BoltExecutor.java:138) ~[storm-client-2.6.1.jar:2.6.1] at org.apache.storm.executor.bolt.BoltExecutor.call(BoltExecutor.java:54) ~[storm-client-2.6.1.jar:2.6.1] at org.apache.storm.utils.Utils$1.run(Utils.java:393) [storm-client-2.6.1.jar:2.6.1] at java.base/java.lang.Thread.run(Thread.java:834) [?:?] Caused by: org.apache.hadoop.fs.UnsupportedFileSystemException: No FileSystem for scheme "maprfs" at org.apache.hadoop.fs.FileSystem.getFileSystemClass(FileSystem.java:3585) ~[stormjar.jar:?] at org.apache.hadoop.fs.FileSystem.createFileSystem(FileSystem.java:3608) ~[stormjar.jar:?] at org.apache.hadoop.fs.FileSystem.access$300(FileSystem.java:174) ~[stormjar.jar:?] at org.apache.hadoop.fs.FileSystem$Cache.getInternal(FileSystem.java:3712) ~[stormjar.jar:?] at org.apache.hadoop.fs.FileSystem$Cache.get(FileSystem.java:3663) ~[stormjar.jar:?] at org.apache.hadoop.fs.FileSystem.get(FileSystem.java:557) ~[stormjar.jar:?] at org.apache.storm.hdfs.bolt.HdfsBolt.doPrepare(HdfsBolt.java:99) ~[stormjar.jar:?] at org.apache.storm.hdfs.bolt.AbstractHdfsBolt.prepare(AbstractHdfsBolt.java:118) ~[stormjar.jar:?] ... 5 more 2024-03-28 10:57:12.285 o.a.s.e.e.ReportError Thread-15-kafka-spout-executor[3, 3] [ERROR] Error java.lang.RuntimeException: java.lang.NoSuchMethodError: 'org.apache.kafka.clients.consumer.ConsumerRecords org.apache.kafka.clients.consumer.Consumer.poll(java.time.Duration)' at org.apache.storm.utils.Utils$1.run(Utils.java:413) ~[storm-client-2.6.1.jar:2.6.1] at java.base/java.lang.Thread.run(Thread.java:834) [?:?] Caused by: java.lang.NoSuchMethodError: 'org.apache.kafka.clients.consumer.ConsumerRecords org.apache.kafka.clients.consumer.Consumer.poll(java.time.Duration)' at org.apache.storm.kafka.spout.KafkaSpout.pollKafkaBroker(KafkaSpout.java:360) ~[stormjar.jar:?] at org.apache.storm.kafka.spout.KafkaSpout.nextTuple(KafkaSpout.java:289) ~[stormjar.jar:?] at org.apache.storm.executor.spout.SpoutExecutor$2.call(SpoutExecutor.java:187) ~[storm-client-2.6.1.jar:2.6.1] at org.apache.storm.executor.spout.SpoutExecutor$2.call(SpoutExecutor.java:153) ~[storm-client-2.6.1.jar:2.6.1] at org.apache.storm.utils.Utils$1.run(Utils.java:398) ~[storm-client-2.6.1.jar:2.6.1] ... 1 more
内容的提问来源于stack exchange,提问作者Saravana R
相关产品推荐
相关产品推荐

