Spark读取CSV推断Schema时如何限制FileScan仅扫描少量行?
你遇到的这个问题确实挺常见的——Spark默认的inferSchema=true机制,哪怕设置了samplingRatio,底层的FileScan还是会读取全部文件内容,只是在推断Schema时用采样的子集。要让推断Schema时只扫描少量行,有几个实用的方案:
方案1:先读取小批量数据生成Schema,再复用该Schema
这个思路很直接:先只读取少量行来推断出Schema,然后用这个预生成的Schema去读取全量数据。这样第一次读取时只会扫描你指定的行数,避免全量扫描。
示例代码:
// 先读取100行数据来推断Schema val tempDF = spark.read .option("header", "true") .option("inferSchema", "true") .csv("/path/to/csvs") .limit(100) // 获取推断出的Schema val inferredSchema = tempDF.schema // 用这个Schema读取全量数据 val fullDF = spark.read .option("header", "true") .schema(inferredSchema) .csv("/path/to/csvs")
⚠️ 注意:如果你的数据中某些列的类型在前面100行没有体现(比如前100行都是整数,但后面有字符串),这种方法可能会导致类型推断不准确。如果有这种情况,你可能需要调整读取的行数,或者结合对数据的了解做微调。
方案2:使用Spark配置限制Schema推断时的读取行数
从Spark 2.3版本开始,提供了一个专门的配置项spark.sql.csv.inferSchema.maxRows,可以直接限制推断Schema时最多读取的行数。设置这个参数后,Spark在推断Schema时只会读取指定数量的行,不会扫描全量文件。
你可以在SparkSession初始化时设置,或者在读取数据前临时设置:
// 方式1:初始化SparkSession时设置 val spark = SparkSession.builder() .appName("SchemaInferDemo") .config("spark.sql.csv.inferSchema.maxRows", 100) .getOrCreate() // 方式2:读取数据前临时设置 spark.conf.set("spark.sql.csv.inferSchema.maxRows", 100) val df = spark.read .option("header", "true") .option("inferSchema", "true") .csv("/path/to/csvs")
这个方案是最直接的,它从底层限制了Schema推断阶段的读取量,完美解决你说的FileScan全量扫描的问题。
为什么samplingRatio达不到预期效果?
你提到设置samplingRatio < 1.0没用,这是因为这个参数的作用是在已经读取的全量数据中采样一部分来推断Schema,而不是限制读取的行数。也就是说,Spark还是会先扫描所有文件的内容,然后从中抽取指定比例的数据来做类型推断,所以FileScan的输入大小还是全量数据的大小,这就是为什么你看不到扫描量减少的原因。
如果以上方案还不能解决你的问题,或者你有特殊的数据源场景,欢迎补充细节~
内容的提问来源于stack exchange,提问作者Jedi

