如何使用Java+Spark读取指定文件夹全部CSV文件并合并为单个CSV
这个需求完全可以实现,现有代码运行异常主要是版本兼容性、路径配置、依赖缺失三类问题导致,具体修复方案如下:
常见异常原因排查
- 版本兼容性问题:Spark 2.0+版本已经将CSV读取能力内置到官方源码中,不需要再引入第三方的
com.databricks.spark.csv依赖,旧第三方依赖和当前Spark版本不兼容是触发直接崩溃的最常见原因。 - 路径规则问题:如果在集群模式运行,你填写的本地路径
/Users/input/*.csv只有Driver节点可以访问,Executor节点无法读取对应路径的文件会直接导致任务终止。 - 配置缺失问题:没有明确指定CSV分隔符、编码、列类型规则,可能出现不同文件列解析异常触发任务失败。
标准实现代码(Spark 2.x及以上版本,无需额外依赖)
import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SparkSession; public class MergeCsvDemo { public static void main(String[] args) { // 初始化SparkSession,集群运行时删除.master("local[*]")配置即可 SparkSession spark = SparkSession.builder() .appName("MergeMultiCsv") .master("local[*]") .getOrCreate(); // 直接读取input文件夹,Spark会自动扫描路径下所有符合规则的CSV文件 Dataset<Row> mergedDf = spark.read() .format("csv") // 确认第一行是列名 .option("header", "true") // 明确指定CSV分隔符,默认是逗号,按需调整即可 .option("sep", ",") // 自动推断列类型,也可以自定义schema避免类型推断误差 .option("inferSchema", "true") // 只读取后缀为.csv的文件,过滤目录下其他格式的文件 .option("pathGlobFilter", "*.csv") // 直接填文件夹路径即可,适配文件夹内CSV数量动态变化的场景 .load("/Users/input/"); // 测试验证合并后的数据 mergedDf.show(); // 如果需要输出为单个CSV文件,合并分区后输出即可 mergedDf.coalesce(1) .write() .option("header", "true") .format("csv") .save("/Users/output/merged_result"); spark.stop(); } }
附加注意事项
- 如果你使用的是Spark 1.6及以下的低版本,才需要引入databricks的CSV依赖,需要严格保证依赖版本和Spark版本匹配,Maven依赖参考如下:
<dependency> <groupId>com.databricks</groupId> <artifactId>spark-csv_2.11</artifactId> <version>1.5.0</version> </dependency>
- 集群运行时需要将input文件夹上传到HDFS、S3等分布式存储中,填写分布式存储路径,保证所有Executor都可以访问到对应文件。
- 不需要手动指定通配符
*.csv,Spark会自动识别路径下所有符合规则的CSV文件,新增文件时直接放到input目录即可,不需要修改代码。
内容的提问来源于stack exchange,提问作者JoeyOC
相关产品推荐
相关产品推荐

