HoodieDeltaStreamer报Filesystem closed异常求助
问题描述
环境版本:
- Hudi 0.10.1
- Spark 3.2.4
- Hadoop 3.3.5
运行的spark-submit命令:
spark-submit --class org.apache.hudi.utilities.deltastreamer.HoodieDeltaStreamer "file://$HUDI_HOME/docker/hoodie/hadoop/hive_base/target/hoodie-utilities.jar" --continuous --table-type COPY_ON_WRITE --source-class org.apache.hudi.utilities.sources.AvroKafkaSource --source-ordering-field submit_date --target-base-path "hdfs://172.16.0.132:9000/data-lake/raw-zone/tables/temp_cow" --target-table temp_cow --props "file://$HUDI_HOME/hudi-utilities/src/test/resources/delta-streamer-config/kafka-source-Table_May_live_v3.properties" --schemaprovider-class org.apache.hudi.utilities.schema.SchemaRegistryProvider --source-limit 50000 > /data/data-lake/logs/pull_from_kafka_topic1.log 2>&1 &
运行时出现的异常堆栈:
2023-06-27 12:23:37,238 INFO sources.AvroKafkaSource: About to read 0 from Kafka for topic :Table_May_live_v3 2023-06-27 12:23:37,238 INFO deltastreamer.DeltaSync: No new data, source checkpoint has not changed. Nothing to commit. Old checkpoint=(Option{val=Table_May_live_v3,0:200000}). New Checkpoint=(TempLMSTable_May_live_v3,0:200000) 2023-06-27 12:23:37,240 ERROR deltastreamer.HoodieDeltaStreamer: Shutting down delta-sync due to exception java.io.IOException: Filesystem closed at org.apache.hadoop.hdfs.DFSClient.checkOpen(DFSClient.java:494) at org.apache.hadoop.hdfs.DFSClient.getFileInfo(DFSClient.java:1729) at org.apache.hadoop.hdfs.DistributedFileSystem$29.doCall(DistributedFileSystem.java:1752) at org.apache.hadoop.hdfs.DistributedFileSystem$29.doCall(DistributedFileSystem.java:1749) at org.apache.hadoop.fs.FileSystemLinkResolver.resolve(FileSystemLinkResolver.java:81) at org.apache.hdfs.DistributedFileSystem.getFileStatus(DistributedFileSystem.java:1764) at org.apache.hadoop.fs.FileSystem.exists(FileSystem.java:1760) at org.apache.hudi.utilities.deltastreamer.DeltaSync.refreshTimeline(DeltaSync.java:244) at org.apache.hudi.utilities.deltastreamer.DeltaSync.syncOnce(DeltaSync.java:288) at org.apache.hudi.utilities.deltastreamer.HoodieDeltaStreamer$DeltaSyncService.lambda$startService$0(HoodieDeltaStreamer.java:640) at java.util.concurrent.CompletableFuture$AsyncSupply.run(CompletableFuture.java:1590) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) at java.lang.Thread.run(Thread.java:748)
排查与解决方案
1. 修复HDFS客户端缓存问题
Hudi 0.10.x在continuous模式下存在HDFS客户端连接复用的bug,导致已关闭的FileSystem被重复调用。
- 在spark-submit命令中添加以下配置:
--conf spark.hadoop.hdfs.client.cache.enabled=false \ --conf spark.hadoop.fs.hdfs.impl.disable.cache=true - 或者在Hudi的props配置文件中添加:
hoodie.hadoop.hdfs.client.cache.enabled=false fs.hdfs.impl.disable.cache=true
2. 移除冲突参数
--source-limit 50000参数用于单次同步测试,与continuous持续运行模式冲突,会触发资源提前释放:
- 修改spark-submit命令,删除
--source-limit 50000参数。
3. 调整HDFS连接超时配置
避免因连接超时导致FileSystem被关闭,在Hudi props配置文件中添加:
dfs.client.socket-timeout=300000 dfs.connection.timeout=300000 dfs.datanode.socket.write.timeout=300000
4. 升级Hudi版本
Hudi 0.11.0及以上版本修复了continuous模式下的HDFS连接泄漏问题,若上述临时方案无效,建议升级至0.11.0或更高稳定版本。
5. 检查HDFS集群状态
确认HDFS集群所有节点正常运行,网络连通性良好,无防火墙/安全组限制Spark节点与HDFS节点的通信。
内容的提问来源于stack exchange,提问作者Ankit Bansal
相关产品推荐
相关产品推荐

