本地Spark处理100GB CSV数据计数耗时久,求优化方案
问题
我是Hadoop与Spark的初学者,在配备4核CPU、32GB内存的本地系统中处理近100GB数据时,仅执行数据计数操作就耗时1.5小时,怀疑操作存在问题,特附上代码寻求帮助。
代码
public class ReadCSV { public static void main(String[] args) { long startTime = System.currentTimeMillis(); SparkSession sparkSession = SparkSession.builder().appName("CsvReader") .master("local[4]") .config("spark.sql.debug.maxToStringFields", 100) .getOrCreate(); SparkContext sparkContext = sparkSession.sparkContext(); sparkContext.setLogLevel("ERROR"); try { String filePath = "/mnt/vol2/OpenLR/openlr/processedData/Monday/*/*/*.csv"; Dataset<Row> dataset = sparkSession.read() .option("header", "true") .option("inferSchema", "true") .csv(filePath); System.out.println("Total: " + dataset.count()); System.out.println("Time taken to complete: " + (System.currentTimeMillis() - startTime)); } catch (Exception e) { e.printStackTrace(); } sparkContext.stop(); sparkSession.close(); } }
优化方案
关闭Schema自动推断:
inferSchema参数会让Spark全量扫描数据来推断字段类型,额外消耗大量IO时间。建议手动定义Schema,示例如下:// 替换为你的实际字段名和类型 StructType customSchema = new StructType() .add("field1", StringType) .add("field2", LongType) .add("field3", DoubleType); Dataset<Row> dataset = sparkSession.read() .option("header", "true") .schema(customSchema) .csv(filePath);调整Spark内存与任务参数:本地模式默认内存分配不足,可通过配置提升资源利用率:
SparkSession sparkSession = SparkSession.builder().appName("CsvReader") .master("local[4]") .config("spark.sql.debug.maxToStringFields", 100) .config("spark.driver.memory", "8g") // 驱动内存分配8GB .config("spark.executor.memory", "16g") // 执行器内存分配16GB .config("spark.sql.shuffle.partitions", "8") // 分区数设为CPU核心数的2倍,减少 shuffle 开销 .getOrCreate();合并小文件:如果目标路径下是大量小CSV文件,Spark需要创建大量任务处理,调度开销会显著增加。可以先将小文件合并为1GB左右的大文件再处理,比如用HDFS的
getmerge命令或Spark自身的重分区操作。验证存储介质性能:数据存储在
/mnt/vol2,确认该挂载点是SSD还是HDD,HDD的随机读写性能远低于SSD,这也是影响速度的关键因素之一。
内容的提问来源于stack exchange,提问作者Ajit Sharma
相关产品推荐
相关产品推荐

