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

Spark从HDFS加载CSV文件求助:本地内存不足需迁移至HDFS运行

从HDFS加载CSV文件替代本地加载的Spark解决方案

嘿,刚好能帮你搞定这个问题!本地内存不够跑Spark程序,迁移到HDFS环境是个很合理的选择,我来给你讲清楚怎么改代码,还有需要注意的细节:

核心修改点

其实替换起来超简单,只需要把本地文件路径换成HDFS的路径格式就行。HDFS的路径格式一般有两种:

  • 完整路径:hdfs://<namenode主机名>:<端口号>/你的文件在HDFS的路径.csv,比如hdfs://localhost:9000/user/griplur/video_games_sales.csv
  • 相对路径:如果你的Spark程序是在Hadoop集群节点上运行,并且已经配置了Hadoop环境变量,也可以直接写/user/griplur/video_games_sales.csv(Spark会自动识别这是HDFS路径)

修改后的完整代码

我把你的代码补全并修改了加载路径的部分,你可以直接参考:

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SaveMode;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.types.StructType;

public class VideoGamesSale {
    public static void main(String[] args) {
        // 初始化SparkSession
        SparkSession spark = SparkSession.builder()
                .appName("VideoGamesSaleAnalysis")
                // 如果是提交到集群运行,记得去掉local[*],由集群管理器分配资源
                .master("local[*]")
                .getOrCreate();

        // 可选:如果你的CSV有固定Schema,可以在这里自定义(比自动推断更高效)
        // StructType customSchema = new StructType()
        //         .add("Rank", "integer")
        //         .add("Name", "string")
        //         .add("Platform", "string")
        //         .add("Year", "integer")
        //         .add("Genre", "string")
        //         .add("Publisher", "string")
        //         .add("NA_Sales", "double")
        //         .add("EU_Sales", "double")
        //         .add("JP_Sales", "double")
        //         .add("Other_Sales", "double")
        //         .add("Global_Sales", "double");

        // 从HDFS加载CSV文件
        Dataset<Row> df = spark.read()
                .option("header", "true") // 如果CSV文件带表头,一定要打开这个选项
                .option("inferSchema", "true") // 自动推断字段类型,也可以替换成上面的customSchema
                // 这里替换成你实际的HDFS文件路径
                .csv("hdfs://localhost:9000/user/griplur/video_games_sales.csv");

        // 这里可以写你的业务逻辑,比如数据清洗、统计分析等
        df.show(5); // 示例:打印前5条数据验证加载成功

        // 如果需要把处理结果写回HDFS,也用类似的路径格式就行
        // df.write()
        //         .mode(SaveMode.Overwrite)
        //         .csv("hdfs://localhost:9000/user/griplur/processed_video_games.csv");

        // 记得关闭SparkSession
        spark.stop();
    }
}

额外注意事项

  • 环境配置:如果是在本地开发机连接远程HDFS,需要设置HADOOP_CONF_DIR环境变量,指向Hadoop集群的配置目录(里面有core-site.xml、hdfs-site.xml这些文件),这样Spark才能找到HDFS的namenode地址。
  • 权限问题:确保运行Spark的用户拥有HDFS目标文件的读取权限,不然会抛出权限拒绝的错误。
  • 集群运行:提交到Spark集群运行时,要把master("local[*]")去掉,或者改成集群的master地址(比如yarn或者spark://master:7077)。

内容的提问来源于stack exchange,提问作者griplur

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 07:04:02