Spark Worker节点无法访问添加的本地文件问题求助
问题原因
你遇到的问题本质是:Docker容器内的Spark Worker无法直接访问宿主机的本地文件路径,而sparkContext.addFile()传递的是宿主机的本地文件路径,Worker容器内不存在该路径,且Spark的文件分发机制未正常生效(Worker没有从驱动的HTTP服务器下载文件)。
以下是无需HDFS的解决方案:
方案1:Docker共享卷挂载(最直接)
在docker-compose.yml的spark-worker服务中添加卷挂载,把宿主机存放addresses.csv的目录映射到Worker容器内部:
services: spark: # 保持原有配置不变 spark-worker: # 保持原有配置不变 volumes: - /你本地存放csv的文件夹路径:/opt/spark/data
然后修改Java代码,直接读取容器内的路径:
Dataset<Row> csv = spark.read().csv("/opt/spark/data/addresses.csv");
注意:替换
/你本地存放csv的文件夹路径为实际路径,比如Windows下可以用/c/Users/patry/xxx,确保Docker有权限访问该目录。
方案2:修复Spark文件分发机制
当使用addFile时,Spark驱动会启动HTTP服务器让Worker下载文件,你需要确保Worker能访问到这个服务器:
- 在SparkConf中添加固定的文件服务器端口,并开放该端口:
SparkConf conf = new SparkConf() .setAppName("Simple Application") .set("spark.driver.port", "65040") .set("spark.driver.host", "192.168.18.15") .set("spark.fileserver.port", "65041") // 固定文件服务器端口 .set("spark.files.overwrite", "true") .setMaster("spark://localhost:7077");
- 修改
docker-compose.yml的spark-worker端口映射,新增文件服务器端口:
spark-worker: # 原有配置不变 ports: - '8081:8081' - '65040:65040' - '65041:65041' # 映射文件服务器端口
- 先把classpath里的文件复制到本地临时目录再添加(避免classpath路径的问题):
// 从classpath提取文件到本地临时文件 InputStream inputStream = getClass().getResourceAsStream("/addresses.csv"); File tempFile = File.createTempFile("addresses", ".csv"); Files.copy(inputStream, tempFile.toPath(), StandardCopyOption.REPLACE_EXISTING); sparkContext.addFile(tempFile.getAbsolutePath()); Dataset<Row> csv = spark.read().csv(SparkFiles.get("addresses.csv"));
- 测试Worker容器能否ping通宿主机的
192.168.18.15,确保网络连通。
方案3:广播变量传递小文件(适合小体积csv)
如果csv文件不大,可以把文件内容通过广播变量发送给Worker,在Worker端生成临时文件再读取:
// 读取classpath中的文件内容 String fileContent = new String(Files.readAllBytes(Paths.get(getClass().getResource("/addresses.csv").getPath()))); Broadcast<String> fileBroadcast = sparkContext.broadcast(fileContent); // 在Worker端生成临时文件并读取 Dataset<Row> csv = spark.read().csv( sparkContext.parallelize(Collections.singletonList(fileBroadcast.value())) .map(content -> { File tempFile = File.createTempFile("addresses", ".csv"); Files.write(tempFile.toPath(), content.getBytes()); return tempFile.getAbsolutePath(); }, Encoders.STRING()) .first() );
内容的提问来源于stack exchange,提问作者P. Kobędza
相关产品推荐
相关产品推荐

