You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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 &quot;maprfs&quot;
    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 &quot;maprfs&quot;
    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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.27 00:07:32