如何在Spark的Worker/Task节点上用单线程读取Parquet文件?
让Spark单线程读取Parquet文件的方法
完全可行,你可以通过下面几种方式实现:
1. 调整Spark配置强制单分区读取
在执行查询前,先设置几个关键配置,让Spark不对目标文件做拆分,只生成一个任务(对应单线程):
// 先配置参数 spark.conf.set("spark.sql.files.maxPartitionBytes", "1073741824") // 设成比你的Parquet文件更大的值,比如1GB spark.conf.set("spark.default.parallelism", "1") spark.conf.set("spark.sql.shuffle.partitions", "1") // 再执行原查询 spark.sql("select * from parquet.`/Users/MyUser/TEST/testcompression/part-00009-asdfasdf-e829-421d-b14f-asdfasdf.c000.snappy.parquet`") .show(5,false)
核心是spark.sql.files.maxPartitionBytes:Spark默认会按这个大小拆分文件,只要它比你的目标文件大,就不会拆分,自然只会启动一个Task,用单线程执行。另外两个参数是兜底,防止后续操作触发多线程。
2. 用Read API显式指定单分区
也可以换个写法,用spark.read读取后强制重分区为1,再用SQL查询:
// 读取文件并强制重分区为1 val parquetDF = spark.read .parquet("/Users/MyUser/TEST/testcompression/part-00009-asdfasdf-e829-421d-b14f-asdfasdf.c000.snappy.parquet") .repartition(1) // 注册成临时表再查询 parquetDF.createOrReplaceTempView("single_part_parquet") spark.sql("select * from single_part_parquet").show(5,false)
repartition(1)会把所有数据合并到一个分区,后续查询就只会用一个线程处理。
注意点
- 单线程只适合小文件或者调试场景,要是你的Parquet文件特别大,单线程读取会很慢,生产环境别这么用。
- 这么做能避免多线程带来的日志杂乱、资源争抢问题,适合排查问题的时候用。
内容的提问来源于stack exchange,提问作者sojim2
相关产品推荐
相关产品推荐

