Apache Spark Parquet 数据源完全指南:读写、分区发现、Schema 合并与列级加密 Apache Spark Parquet 数据源完全指南读写、分区发现、Schema 合并与列级加密【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/sparkParquet 是一种被众多大数据处理系统广泛支持的列式存储格式本指南以 Apache Spark当前仓库官方文档为主体系统讲解 Spark SQL 如何读写 Parquet 文件、自动保留 Schema、发现分区信息、合并异构 Schema、转换 Hive metastore Parquet 表以及 Spark 3.2 起支持的列级加密与 KMS 集成。读完本文你将掌握从 DataFrame API、SQL 到底层配置项的全链路 Parquet 实战技能并能基于仓库源码理解各选项的真实生效机制。Parquet 数据源概述Parquet 是面向分析型负载设计的列式存储格式其特点是按列组织数据、自带 Schema 描述self-describing并内建统计信息能够为谓词下推、列裁剪等优化提供基础。Spark SQL 对 Parquet 提供了开箱即用的读写支持并且会自动保留原始数据的 Schema。需要特别注意的是在读取 Parquet 文件时所有列都会被自动转换为可空nullable列这是为了兼容性考虑而做的统一处理。在 examples/src/main/resources 下仓库提供了users.parquet、people.json等示例数据以及dir1/含file1.parquet、file2.parquet等目录可用于验证本文中的全部示例。以编程方式加载与保存 ParquetParquet 文件可以像其他数据源一样通过DataFrameReader/DataFrameWriter或 SQL 语句进行读写。下面以examples/src/main/resources/people.json为输入演示完整的「读取 JSON → 写入 Parquet → 回读 Parquet → 注册临时视图 → SQL 查询」链路。完整示例代码位于 examples/src/main/python/sql/datasource.pyparquet_example函数、examples/src/main/scala/org/apache/spark/examples/sql/SQLDataSourceExample.scala、examples/src/main/java/org/apache/spark/examples/sql/JavaSQLDataSourceExample.java 与 examples/src/main/r/RSparkSQLExample.R。PythonpeopleDF spark.read.json(examples/src/main/resources/people.json) # DataFrames 可以保存为 Parquet 文件并保留 Schema 信息 peopleDF.write.parquet(people.parquet) # Parquet 文件是自描述的因此回读时 Schema 被完整保留 parquetFile spark.read.parquet(people.parquet) # Parquet 文件还可以注册为临时视图再通过 SQL 语句使用 parquetFile.createOrReplaceTempView(parquetFile) teenagers spark.sql(SELECT name FROM parquetFile WHERE age 13 AND age 19) teenagers.show() # ------ # | name| # ------ # |Justin| # ------Scalaimport spark.implicits._ val peopleDF spark.read.json(examples/src/main/resources/people.json) // DataFrames 可以保存为 Parquet 文件并保留 Schema 信息 peopleDF.write.parquet(people.parquet) val parquetFileDF spark.read.parquet(people.parquet) parquetFileDF.createOrReplaceTempView(parquetFile) val namesDF spark.sql(SELECT name FROM parquetFile WHERE age BETWEEN 13 AND 19) namesDF.map(attributes Name: attributes(0)).show() // ------------ // | value| // ------------ // |Name: Justin| // ------------JavaDatasetRow peopleDF spark.read().json(examples/src/main/resources/people.json); // DataFrames 可以保存为 Parquet 文件并保留 Schema 信息 peopleDF.write().parquet(people.parquet); DatasetRow parquetFileDF spark.read().parquet(people.parquet); parquetFileDF.createOrReplaceTempView(parquetFile); DatasetRow namesDF spark.sql(SELECT name FROM parquetFile WHERE age BETWEEN 13 AND 19); DatasetString namesDS namesDF.map( (MapFunctionRow, String) row - Name: row.getString(0), Encoders.STRING()); namesDS.show();Rdf - read.df(examples/src/main/resources/people.json, json) write.parquet(df, people.parquet) parquetFile - read.parquet(people.parquet) createOrReplaceTempView(parquetFile, parquetFile) teenagers - sql(SELECT name FROM parquetFile WHERE age 13 AND age 19)SQLCREATE TEMPORARY VIEW parquetTable USING org.apache.spark.sql.parquet OPTIONS ( path examples/src/main/resources/people.parquet ) SELECT * FROM parquetTable此外也可以直接使用parquet.前缀的路径式 SQL 语法读取例如 datasource.py 中的df spark.sql(SELECT * FROM parquet.examples/src/main/resources/users.parquet)分区发现Partition Discovery表分区是 Hive 等系统中常见的优化手段分区表的数据通常存放在不同目录下分区列的值被编码在每个分区目录的路径中。所有内置文件数据源包括 Text/CSV/JSON/ORC/Parquet都能自动发现并推断分区信息。例如将人口数据按两个额外分区列gender和country组织成如下目录结构path └── to └── table ├── gendermale │ ├── ... │ ├── countryUS │ │ └── data.parquet │ ├── countryCN │ │ └── data.parquet │ └── ... └── genderfemale ├── ... ├── countryUS │ └── data.parquet ├── countryCN │ └── data.parquet └── ...只要把path/to/table传给SparkSession.read.parquet或SparkSession.read.loadSpark SQL 就会自动从路径中提取分区信息。此时返回的 DataFrame Schema 变为root |-- name: string (nullable true) |-- age: long (nullable true) |-- gender: string (nullable true) |-- country: string (nullable true)分区列类型推断与关闭方式注意分区列的数据类型是自动推断出来的。目前支持数值类型、date、timestamp和string类型。如果不想自动推断分区列类型可以通过配置spark.sql.sources.partitionColumnTypeInference.enabled关闭该配置默认值为true当类型推断被关闭时分区列统一按string类型处理。分区发现的基准路径行为从 Spark 1.6.0 开始分区发现默认只在给定的路径下寻找分区。以上面的示例来说如果用户把path/to/table/gendermale传给SparkSession.read.parquet或SparkSession.read.load那么gender不会被当作分区列。如果确实需要指定分区发现的基准路径可以在数据源选项中设置basePath。例如当数据路径为path/to/table/gendermale且basePath设置为path/to/table/时gender就会被识别为分区列。Schema 合并Schema Merging与 Protocol Buffer、Avro、Thrift 类似Parquet 也支持 Schema 演进schema evolution用户可以先从一个简单的 Schema 开始再按需逐步添加新列最终得到一批 Schema 不同但彼此兼容的 Parquet 文件。Parquet 数据源能够自动检测这种场景并合并所有这些文件的 Schema。由于 Schema 合并是一个相对昂贵的操作且在大多数场景下并非必需Spark 从 1.5.0 起默认关闭了该功能。可以通过以下两种方式开启读取 Parquet 文件时设置数据源选项mergeSchema为true见下方示例设置全局 SQL 配置spark.sql.parquet.mergeSchema为true。以 datasource.py 中的parquet_schema_merging_example为例先创建两个 Schema 不同的 DataFrame 并写入key1、key2两个分区目录再以mergeSchematrue回读整个分区表Pythonsc spark.sparkContext # 创建简单 DataFrame写入分区目录 squaresDF spark.createDataFrame(sc.parallelize(range(1, 6)) .map(lambda i: Row(singlei, doublei ** 2))) squaresDF.write.parquet(data/test_table/key1) # 在另一个分区目录中写入新 DataFrame新增 triple 列、丢弃 double 列 cubesDF spark.createDataFrame(sc.parallelize(range(6, 11)) .map(lambda i: Row(singlei, triplei ** 3))) cubesDF.write.parquet(data/test_table/key2) # 读取分区表 mergedDF spark.read.option(mergeSchema, true).parquet(data/test_table) mergedDF.printSchema() # 最终 Schema 由各 Parquet 文件中的所有列加上分区列共同组成 # root # |-- double: long (nullable true) # |-- single: long (nullable true) # |-- triple: long (nullable true) # |-- key: integer (nullable true)Scalaimport spark.implicits._ // 创建简单 DataFrame写入分区目录 val squaresDF spark.sparkContext.makeRDD(1 to 5).map(i (i, i * i)).toDF(value, square) squaresDF.write.parquet(data/test_table/key1) // 在另一个分区目录中写入新 DataFrame新增 cube 列、丢弃 square 列 val cubesDF spark.sparkContext.makeRDD(6 to 10).map(i (i, i * i * i)).toDF(value, cube) cubesDF.write.parquet(data/test_table/key2) // 读取分区表 val mergedDF spark.read.option(mergeSchema, true).parquet(data/test_table) mergedDF.printSchema() // root // |-- value: int (nullable true) // |-- square: int (nullable true) // |-- cube: int (nullable true) // |-- key: int (nullable true)Java版本的runParquetSchemaMergingExample位于 JavaSQLDataSourceExample.java逻辑相同先以 POJOSquare、Cube构造两个 DataFrame 写入key1/key2再以mergeSchematrue读取并打印合并后的 Schema。R版本位于 RSparkSQLExample.R使用read.df(data/test_table, parquet, mergeSchema true)完成同样的操作。与respectSummaryFiles的关联需要说明的是spark.sql.parquet.mergeSchema为true时Parquet 数据源会合并所有数据文件的 Schema否则 Schema 取自 summary 文件若没有 summary 文件则随机取一个数据文件。另一个专家级选项spark.sql.parquet.respectSummaryFiles默认false与 Schema 合并行为相关详见下文配置表。Hive metastore Parquet 表转换当读取 Hive metastore 中的 Parquet 表、或向非分区 Hive metastore Parquet 表写入时Spark SQL 会优先使用自己的 Parquet 支持而非 Hive SerDe以获得更好的性能。该行为由配置spark.sql.hive.convertMetastoreParquet控制默认开启true。设置为false后Spark SQL 将回退使用 Hive SerDe 处理 Parquet 表。Hive/Parquet Schema 协调从表 Schema 处理的角度看Hive 与 Parquet 之间存在两个关键差异Hive 大小写不敏感而 Parquet 大小写敏感Hive 认为所有列都可空而 Parquet 中的可空性nullability是有意义的。因此在把 Hive metastore Parquet 表转换为 Spark SQL Parquet 表时必须对两边的 Schema 进行协调reconcile。协调规则如下两个 Schema 中同名字段必须具有相同的数据类型无论可空性如何协调后的字段采用 Parquet 一侧的数据类型从而尊重 Parquet 的可空性语义。协调后的 Schema 恰好包含 Hive metastore Schema 中定义的那些字段只出现在 Parquet Schema 中的字段会被丢弃只出现在 Hive metastore Schema 中的字段会以可空字段的形式补充进协调后的 Schema。元数据刷新为了提升性能Spark SQL 会缓存 Parquet 元数据当 Hive metastore Parquet 表转换开启时这些被转换表的元数据同样会被缓存。如果这些表被 Hive 或其他外部工具更新过需要手动刷新元数据以保证一致性四种 API 与 SQL 的刷新方式如下Python / Scala# spark 是已存在的 SparkSession spark.catalog.refreshTable(my_table)// spark 是已存在的 SparkSession spark.catalog.refreshTable(my_table)Javaspark.catalog().refreshTable(my_table);RrefreshTable(my_table)SQLREFRESH TABLE my_table;列级加密Columnar Encryption自 Spark 3.2 起配合 Apache Parquet 1.12Spark 支持对 Parquet 表进行列级加密。加密模型信封加密Parquet 采用信封加密envelope encryption实践文件各部分使用「数据加密密钥」DEKData Encryption Keys加密而 DEK 再由「主加密密钥」MEKMaster Encryption Keys加密。其中 DEK 由 Parquet 为每个加密文件/列随机生成MEK 则在用户选定的密钥管理服务KMSKey Management Service中生成、存储与管理。Parquet Maven 仓库提供了包含 mock KMS 实现的 jarparquet-hadoop-tests.jar允许用户只使用 spark-shell 就完成列加密与解密的演示无需实际部署 KMS 服务器下载该 jar 并放入 Spark 的jars目录即可。使用 mock KMS 进行加密读写演示Python通过 Spark 任务的--conf传入 Hadoop 配置属性# 设置 hadoop 配置属性例如 # --conf spark.hadoop.parquet.encryption.kms.client.class\ # org.apache.parquet.crypto.keytools.mocks.InMemoryKMS\ # --conf spark.hadoop.parquet.encryption.key.list\ # keyA:AAECAwQFBgcICQoLDA0ODw , keyB:AAECAAECAAECAAECAAECAA\ # --conf spark.hadoop.parquet.crypto.factory.class\ # org.apache.parquet.crypto.keytools.PropertiesDrivenCryptoFactory # 写入加密的 DataFrame 文件。 # 列 square 将使用主密钥 keyA 保护。 # Parquet 文件 footer 将使用主密钥 keyB 保护。 squaresDF.write\ .option(parquet.encryption.column.keys, keyA:square)\ .option(parquet.encryption.footer.key, keyB)\ .parquet(/path/to/table.parquet.encrypted) # 读取加密的 DataFrame 文件 df2 spark.read.parquet(/path/to/table.parquet.encrypted)Scalasc.hadoopConfiguration.set(parquet.encryption.kms.client.class, org.apache.parquet.crypto.keytools.mocks.InMemoryKMS) // 显式指定主密钥base64 编码——仅 mock InMemoryKMS 需要 sc.hadoopConfiguration.set(parquet.encryption.key.list, keyA:AAECAwQFBgcICQoLDA0ODw , keyB:AAECAAECAAECAAECAAECAA) // 激活由 Hadoop 属性驱动的 Parquet 加密 sc.hadoopConfiguration.set(parquet.crypto.factory.class, org.apache.parquet.crypto.keytools.PropertiesDrivenCryptoFactory) // 写入加密的 DataFrame 文件 squaresDF.write. option(parquet.encryption.column.keys, keyA:square). option(parquet.encryption.footer.key, keyB). parquet(/path/to/table.parquet.encrypted) // 读取加密的 DataFrame 文件 val df2 spark.read.parquet(/path/to/table.parquet.encrypted)Java版本对应代码位于 JavaSQLDataSourceExample.java通过sc.hadoopConfiguration().set(...)设置同样三个属性随后以.option(parquet.encryption.column.keys, keyA:square)与.option(parquet.encryption.footer.key, keyB)写入加密文件并以spark.read().parquet(...)读取。KMS Client 接口InMemoryKMS类仅用于演示 Parquet 加密功能不应在生产环境使用。主加密密钥必须由用户所在组织部署的生产级 KMS 系统保管和管理。要让 Spark 在生产环境落地 Parquet 加密需要为 KMS 服务器实现一个客户端类。Parquet 提供了用于开发此类客户端的插件接口KmsClientpublic interface KmsClient { // 包装密钥——使用主密钥加密 public String wrapKey(byte[] keyBytes, String masterKeyIdentifier); // 解密解包密钥——使用主密钥解密 public byte[] unwrapKey(String wrappedKey, String masterKeyIdentifier); // 初始化参数可选 public void initialize(Configuration configuration, String kmsInstanceID, String kmsInstanceURL, String accessToken); }在 parquet-java 仓库中提供了面向开源 KMSVault的客户端示例VaultClient。生产级 KMS 客户端应与组织的安全管理员协同设计并由具有访问控制管理经验的开发者构建。客户端类创建完成后可通过parquet.encryption.kms.client.class参数传给应用普通 Spark 用户即可像上面的加密读写示例一样直接使用。双重信封加密默认情况下Parquet 实现的是「双重信封加密」double envelope encryption模式该模式能最大程度减少 Spark 执行器与 KMS 服务器之间的交互DEK 使用「密钥加密密钥」KEKKey Encryption Keys由 Parquet 随机生成加密KEK 再在 KMS 中由 MEK 加密加密结果与 KEK 本身会缓存在 Spark 执行器内存中。如果希望使用常规的信封加密可以将parquet.encryption.double.wrapping参数设置为false。数据源选项Data Source OptionParquet 的数据源选项可以通过以下方式设置DataFrameReader/DataFrameWriter/DataStreamReader/DataStreamWriter的.option/.options方法CREATE TABLE USING DATA_SOURCE语句的OPTIONS子句参见 CREATE TABLE USING DATA_SOURCE 语法文档。属性名默认值含义作用域datetimeRebaseModespark.sql.parquet.datetimeRebaseModeInRead配置的值指定DATE、TIMESTAMP_MILLIS、TIMESTAMP_MICROS逻辑类型的值从儒略历Julian重定基rebase到公历Proleptic Gregorian的模式。支持EXCEPTION读取到两种历法下存在歧义的古代日期/时间戳时报错、CORRECTED不重定基直接加载、LEGACY将古代日期/时间戳从儒略历重定基到公历readint96RebaseModespark.sql.parquet.int96RebaseModeInRead配置的值指定 INT96 时间戳从儒略历重定基到公历的模式。支持EXCEPTION、CORRECTED、LEGACY含义同上readmergeSchemaspark.sql.parquet.mergeSchema配置的值是否合并从所有 Parquet part-file 收集到的 Schema。该选项会覆盖spark.sql.parquet.mergeSchemareadcompressionsnappy保存文件时使用的压缩编解码器可以是下列大小写不敏感的短名称之一none、uncompressed、snappy、gzip、lzo、brotli、lz4、lz4_raw、zstd。该选项会覆盖spark.sql.parquet.compression.codecwrite其他通用文件源选项请参见 Generic Files Source Options。从源码实现看这些选项的解析集中在 sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/parquet/ParquetOptions.scalacompressionCodecClassName的解析遵循严格的优先级compression数据源选项→parquet.compression即ParquetOutputFormat.COMPRESSION→spark.sql.parquet.compression.codec全局配置并且会对短名称做大小写不敏感的归一化映射mergeSchema、datetimeRebaseModeInRead、int96RebaseModeInRead等选项均支持「选项优先、全局配置兜底」的取值策略数据源选项会覆盖同名 SQL 配置常量定义MERGE_SCHEMA、COMPRESSION、DATETIME_REBASE_MODE、INT96_REBASE_MODE与文档中的选项一一对应。ConfigurationParquet 相关配置全表Parquet 的配置既可以通过spark.conf.set设置也可以通过 SQL 执行SET keyvalue命令设置。下表完整列出了 Spark SQL 中与 Parquet 相关的全部配置项。属性名默认值含义起始版本spark.sql.parquet.binaryAsStringfalse一些其他 Parquet 生成系统尤其是 Impala、Hive 以及旧版 Spark SQL在写出 Parquet Schema 时不区分二进制数据与字符串。该开关让 Spark SQL 把二进制数据解释为字符串以兼容这些系统1.1.1spark.sql.parquet.int96AsTimestamptrue一些 Parquet 生成系统尤其是 Impala 和 Hive将 Timestamp 存储为 INT96。该开关让 Spark SQL 把 INT96 数据解释为时间戳以兼容这些系统1.3.0spark.sql.parquet.int96TimestampConversionfalse控制将 INT96 数据转换为时间戳时是否应用时间戳调整用于兼容 Impala 写出的数据。之所以必要是因为 Impala 存储 INT96 数据时使用的时区偏移与 Hive 和 Spark 不同2.3.0spark.sql.parquet.outputTimestampTypeINT96设置 Spark 向 Parquet 文件写数据时使用的 Parquet 时间戳类型。INT96是 Parquet 中非标准但常用的时间戳类型TIMESTAMP_MICROS是 Parquet 标准类型存储自 Unix 纪元起的微秒数TIMESTAMP_MILLIS同样为标准类型但精度为毫秒意味着 Spark 需要截断时间戳值的微秒部分2.3.0spark.sql.parquet.compression.codecsnappy写 Parquet 文件时使用的压缩编解码器。若表级选项/属性中同时指定了compression或parquet.compression优先级为compressionparquet.compressionspark.sql.parquet.compression.codec。可接受值none、uncompressed、snappy、gzip、lzo、brotli、lz4、lz4_raw、zstd。注意brotli需要安装BrotliCodec1.1.1spark.sql.parquet.filterPushdowntrue设为true时启用 Parquet 过滤器下推优化1.2.0spark.sql.parquet.aggregatePushdownfalse设为true时聚合会下推到 Parquet 进行优化。支持MIN、MAX和COUNT聚合表达式MIN/MAX支持布尔、整数、浮点与日期类型COUNT支持所有数据类型。如果某个 Parquet 文件的 footer 缺少统计信息会抛出异常3.3.0spark.sql.hive.convertMetastoreParquettrue设为false时Spark SQL 对 Parquet 表使用 Hive SerDe 而非内置支持1.1.1spark.sql.parquet.mergeSchemafalse为true时Parquet 数据源合并所有数据文件的 Schema否则 Schema 取自 summary 文件若无 summary 文件则随机取一个数据文件1.5.0spark.sql.parquet.respectSummaryFilesfalse为true时假定所有 Parquet part-file 与 summary 文件一致合并 Schema 时忽略 part-file为false默认时合并所有 part-file。该选项属于专家级选项只有在完全理解其含义后才应开启1.5.0spark.sql.parquet.writeLegacyFormatfalse为true时以 Spark 1.4 及更早版本的方式写数据。例如 decimal 值将以 Apache Parquet 的定长字节数组格式写出Hive、Impala 等系统使用该格式为false时使用 Parquet 新格式例如 decimal 以 int 为基础的格式。若 Parquet 输出要供不支持新格式的系统使用应设为true1.6.0spark.sql.parquet.enableVectorizedReadertrue启用向量化 Parquet 解码2.0.0spark.sql.parquet.enableNestedColumnVectorizedReadertrue为嵌套列如 struct、list、map启用向量化 Parquet 解码。要求spark.sql.parquet.enableVectorizedReader已启用3.3.0spark.sql.parquet.recordLevelFilter.enabledfalse为true时使用下推过滤器启用 Parquet 原生的记录级过滤。该配置仅在spark.sql.parquet.filterPushdown启用且未使用向量化读取器时生效可通过将spark.sql.parquet.enableVectorizedReader设为false来确保不使用向量化读取器2.3.0spark.sql.parquet.columnarReaderBatchSize4096Parquet 向量化读取器每个批次包含的行数。该数值需要仔细权衡以最小化开销并避免读取数据时发生 OOM2.4.0spark.sql.parquet.fieldId.write.enabledtrueField ID 是 Parquet Schema 规范的原生字段。启用后Parquet 写入器会将 Spark Schema 中携带的 field Id 元数据写入 Parquet Schema3.3.0spark.sql.parquet.fieldId.read.enabledfalse启用后Parquet 读取器使用请求的 Spark Schema 中的 field ID若存在查找 Parquet 字段而非使用列名3.3.0spark.sql.parquet.fieldId.read.ignoreMissingfalse当 Parquet 文件没有任何 field ID但 Spark 读取 Schema 使用 field ID 读取时启用该开关则静默返回 null否则报错3.3.0spark.sql.parquet.inferTimestampNTZ.enabledtrue启用时Schema 推断阶段将带有isAdjustedToUTC false注解的 Parquet 时间戳列推断为TIMESTAMP_NTZ类型否则所有 Parquet 时间戳列都被推断为TIMESTAMP_LTZ。注意 Spark 写入文件时会把输出 Schema 写入 Parquet footer 元数据并在读取时加以利用因此该配置只影响非 Spark 写出的 Parquet 文件的 Schema 推断3.4.0spark.sql.parquet.timeType.allowIsAdjustedToUtcReadfalse启用时Schema 推断阶段将带有isAdjustedToUTC true注解的 Parquet TIME 列推断为TIME类型以兼容 Apache Arrow 等写入方否则 Schema 推断会拒绝此类列并报错。该配置仅影响 Schema 推断使用显式用户指定的TIMESchema 读取始终成功因为 Spark 无时区的TIME类型无论哪种方式都解码为同一天内的时间4.4.0spark.sql.parquet.datetimeRebaseModeInReadEXCEPTIONDATE、TIMESTAMP_MILLIS、TIMESTAMP_MICROS逻辑类型的值从儒略历重定基到公历的模式EXCEPTION读到两种历法下歧义的古代日期/时间戳时报错、CORRECTED不重定基按原样读取、LEGACY读取 Parquet 时从旧版混合历法重定基到公历。仅当 Parquet 文件的写入方信息如 Spark、Hive未知时该配置才生效3.0.0spark.sql.parquet.datetimeRebaseModeInWriteEXCEPTIONDATE、TIMESTAMP_MILLIS、TIMESTAMP_MICROS逻辑类型的值从公历重定基到儒略历的模式EXCEPTION写到歧义古代日期/时间戳时报错、CORRECTED不重定基按原样写入、LEGACY写 Parquet 时从公历重定基到旧版混合历法3.0.0spark.sql.parquet.int96RebaseModeInReadEXCEPTIONINT96时间戳类型从儒略历重定基到公历的模式EXCEPTION/CORRECTED/LEGACY语义同上。仅当写入方信息未知时生效3.1.0spark.sql.parquet.int96RebaseModeInWriteEXCEPTIONINT96时间戳类型从公历重定基到儒略历的模式EXCEPTION/CORRECTED/LEGACY语义同上3.1.0关键配置的实战提示历法重定基Rebase三选项EXCEPTION模式下遇到两种历法下存在歧义的「古代」日期/时间戳会直接报错适合对数据正确性要求严格的场景CORRECTED不做任何重定基适合确认数据本身一致的新写入场景LEGACY则用于兼容 Spark 3.0 之前的旧数据。注意InRead系列配置仅在 Parquet 文件写入方信息未知时才生效——由 Spark 写出的文件会把自己的写入信息记录在文件中。压缩编解码器优先级表级compression选项优先于表级parquet.compression再优先于全局spark.sql.parquet.compression.codec。仓库中的 ParquetOptions.scala 明确注释了这一优先级顺序。向量化读取spark.sql.parquet.enableVectorizedReader与spark.sql.parquet.enableNestedColumnVectorizedReader共同控制向量化解码若要使用 Parquet 原生记录级过滤recordLevelFilter.enabled则需要显式关闭向量化读取器。Field ID 读写写侧默认开启true读侧默认关闭false。开启读侧 Field ID 匹配后若文件缺少 Field ID 且未开启ignoreMissing读取会报错而非静默返回 null。性能优化路径与底层实现佐证Parquet 数据源之所以高效依赖的是多个层面协同工作谓词下推与列裁剪由spark.sql.parquet.filterPushdown默认开启驱动配合 Parquet footer 中的统计信息在读取阶段跳过无关行块聚合下推spark.sql.parquet.aggregatePushdown默认关闭可将MIN/MAX/COUNT下推到 Parquet其实现依赖每个文件 footer 的统计信息缺失统计信息时会抛出异常向量化解码spark.sql.parquet.enableVectorizedReader将批量解码默认批次 4096 行由spark.sql.parquet.columnarReaderBatchSize控制与行式解码区分开显著降低解码开销。这些行为的实现载体是 sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/parquet/ParquetFileFormat.scala文件格式的读写主逻辑与 ParquetOptions.scala选项解析后者同时定义了数据源选项与 SQL 配置的映射关系及优先级规则。仓库中sql/core模块下的 Parquet 测试用例sql/core/src/test中相关*.scala对分区发现、Schema 合并、历法重定基等行为均有覆盖可作为深入研读实现细节的入口。总结本指南完整覆盖了 Spark SQL 使用 Parquet 数据源的五个核心主题编程式读写与 Schema 保留、分区发现与basePath/类型推断控制、Schema 合并mergeSchema、Hive metastore Parquet 表转换与元数据刷新以及 Spark 3.2 的列级加密与 KMS 集成。配合末尾的数据源选项表与配置全表含默认值、优先级与版本信息读者可以据此在真实集群上快速落地 Parquet 读写、加密与调优方案需要深入原理时可进一步阅读 ParquetFileFormat.scala 与 ParquetOptions.scala 确认每个开关的实际生效路径。【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考