Spark SQL Parquet 兼容性测试指南:Avro/Thrift Schema 与代码生成机制解析
2026/9/20 17:57:02 网站建设 项目流程
  • 大数据
  • 数据分析
  • 批处理
  • 流处理
  • 机器学习
  • 图计算

【免费下载链接】spark

Apache Spark - A unified analytics engine for large-scale data processing

项目地址:https://gitcode.com/gh_mirrors/sp/spark
点击查看免费下载

Apache Spark 的 SQL 模块在读写 Parquet 文件时,需要与 parquet-avro、parquet-thrift 等生态组件保持双向兼容。本文以 sql/core/src/test/README.md 为主线,系统讲解 Spark 中 Parquet 兼容性测试的组织结构、Schema 定义方式、代码生成流程与前置工具链。读完本文,你将掌握这套测试资产的目录含义、gen-avro.sh/gen-thrift.sh的更新流程,并能看懂 Avro IDL 与 Thrift schema 中每个测试字段的设计意图,从而具备为 Spark 贡献新的兼容性测试用例的实际能力。

测试资产的整体结构与设计意图

Parquet 兼容性测试的核心目标,是验证 Spark SQL 与 parquet-avro、parquet-thrift 写入/读取的 Parquet 文件在 Schema、类型映射和数据语义上保持兼容。为此,测试资产被刻意组织为"Schema 源文件 + 生成代码 + 生成脚本"三部分,目录结构如下:

sql/core/src/test/ ├── README.md # 本说明文件 ├── avro │ ├── parquet-compat.avdl # 测试用 Avro IDL(Schema 源) │ └── parquet-compat.avpr # !! NO TOUCH !! 由 Avro IDL 生成的协议文件 ├── gen-java # !! NO TOUCH !! 生成的 Java 代码 ├── scripts │ ├── gen-avro.sh # 用于生成 Avro Java 代码的脚本 │ └── gen-thrift.sh # 用于生成 Thrift Java 代码的脚本 └── thrift └── parquet-compat.thrift # 测试用 Thrift schema(Schema 源)

这里有两个关键约定值得注意:

  1. avro/*.avprgen-java/是生成产物,禁止手工修改。README 中以!! NO TOUCH !!明确标注,它们的唯一合法来源是*.avdl/*.thrift通过生成脚本产出,任何手工改动都会在下次生成时被覆盖,并可能导致测试基准与 Schema 源不一致。
  2. 生成代码直接入库(check in)。README 明确指出:"To avoid code generation during build time, Java code generated from testing Thrift schema and Avro IDL are also checked in."—— 即为了避免在构建阶段引入代码生成步骤(以及由此带来的工具链依赖、构建不确定性),由 Thrift schema 和 Avro IDL 生成的 Java 代码被预先提交到仓库中,构建时直接编译使用。这是测试基建中一种"一次性生成、构建期零依赖"的成熟做法。

Schema 更新流程:何时重新生成

测试 Schema 不是静态的。当需要增加新的类型映射用例、调整字段定义或补充边界条件时,就需要修改 Schema 源并同步更新生成代码。README 给出的规则非常明确:

When updating the testing Thrift schema and Avro IDL, please rungen-avro.shandgen-thrift.shaccordingly to update generated Java code.

即每次修改thrift/*.thriftavro/*.avdl后,必须依次运行对应脚本,使gen-java/下的 Java 代码与 Schema 源保持一致。提交 PR 时应同时包含 Schema 源文件、重新生成的.avprgen-java/变更,否则测试可能使用过期代码。

gen-avro.sh:Avro IDL 到 Java 的两步生成

gen-avro.sh 的核心逻辑为:

cd $(dirname $0)/.. BASEDIR=`pwd` cd - rm -rf $BASEDIR/gen-java mkdir -p $BASEDIR/gen-java for input in `ls $BASEDIR/avro/*.avdl`; do filename=$(basename "$input") filename="${filename%.*}" avro-tools idl $input > $BASEDIR/avro/${filename}.avpr avro-tools compile -string protocol $BASEDIR/avro/${filename}.avpr $BASEDIR/gen-java done

脚本先清空并重建gen-java/目录,然后对avro/下每个.avdl执行两步操作:

  1. avro-tools idl <input>:将 Avro IDL 文本编译为 JSON 格式的协议文件.avpr
  2. avro-tools compile -string protocol <avpr> <outputDir>:将.avpr协议编译为 Java 类,输出到gen-java/

其中-string选项强制以 JavaString表示 Avro 的string类型(而非Utf8),这直接影响后续生成的 POJO 字段类型,进而影响 parquet-avro 写入 Parquet 时的类型映射。

gen-thrift.sh:Thrift schema 到 Java 的生成

gen-thrift.sh 的逻辑与 Avro 侧对称:

cd $(dirname $0)/.. BASEDIR=`pwd` cd - rm -rf $BASEDIR/gen-java mkdir -p $BASEDIR/gen-java for input in `ls $BASEDIR/thrift/*.thrift`; do thrift --gen java -out $BASEDIR/gen-java $input done

同样先重建gen-java/,再通过thrift --gen java -out <outputDir> <schema>为每个.thrift文件生成 Java 代码。注意两个脚本共用同一个gen-java/输出目录,因此任意一侧的 Schema 变更都需要重新运行两个脚本中的对应者,且二者的产物共存于同一包结构下——这要求两侧的 namespace 设计相互隔离(见下文)。

前置条件:工具链安装

代码生成依赖两个外部 CLI 工具,README 明确要求先确保其已安装:

工具用途
avro-tools将 Avro IDL(.avdl)编译为协议文件(.avpr)并进一步生成 Java 代码
thrift将 Thrift schema(.thrift)编译为 Java 代码

在 macOS 上可通过 Homebrew 一键安装:

$ brew install thrift avro-tools

在 Linux 发行版上,可分别通过包管理器安装(如apt-get install thrift-compiler)或从 Avro / Thrift 官方发行版获取对应版本,安装后需确保avro-toolsthrift命令位于PATH中。工具版本应与项目测试环境保持一致,因为不同版本生成的 Java 代码存在 API 差异。

深入 Avro 测试协议:parquet-compat.avdl

parquet-compat.avdl 是 Avro 侧测试协议的完整定义,其用途在文件头注释中写得很清楚:"This is a test protocol for testing parquet-avro compatibility."协议声明了 Java namespaceorg.apache.spark.sql.execution.datasources.parquet.test.avro,与 Thrift 侧的 namespace 相互独立,避免生成类冲突。

整个协议按类型覆盖维度被精心拆分为多个 record,每个 record 对应一类兼容性场景:

枚举与嵌套类型

enum Suit { SPADES, HEARTS, DIAMONDS, CLUBS } record ParquetEnum { Suit suit; } record Nested { array<int> nested_ints_column; string nested_string_column; }

Suit枚举用于验证 Avro 枚举类型到 Parquet 底层存储(以整数表示枚举值)的映射;Nestedrecord 则验证嵌套结构化类型在 Parquet 中的展开方式。

基本类型与可空类型

record AvroPrimitives { boolean bool_column; int int_column; long long_column; float float_column; double double_column; bytes binary_column; string string_column; } record AvroOptionalPrimitives { union { null, boolean } maybe_bool_column; union { null, int } maybe_int_column; // ... maybe_long / maybe_float / maybe_double / maybe_binary / maybe_string }

AvroPrimitives覆盖 Avro 的全部标量类型到 Parquet 的类型映射;AvroOptionalPrimitives则通过union { null, T }形态构造"可空列",用于验证 optional 字段在 Parquet Schema 中的 nullable 语义。

数组、嵌套数组与 Map

record AvroNonNullableArrays { array<string> strings_column; union { null, array<int> } maybe_ints_column; } record AvroArrayOfArray { array<array<int>> int_arrays_column; } record AvroMapOfArray { map<array<int>> string_to_ints_column; }

这三个 record 分别验证:非空数组、可空数组、数组的数组(二维数组)、Map 的 value 为数组等复杂组合类型,覆盖 Parquet 中 repeated group 与 logical type 的交互。

综合兼容协议

record ParquetAvroCompat { array<string> strings_column; map<int> string_to_int_column; map<array<Nested>> complex_column; }

ParquetAvroCompat是集大成者:字符串数组、整型 Map、以及 value 为Nested结构数组的复杂 Map,用于端到端验证 Spark SQL 对 parquet-avro 写入的复杂嵌套 Schema 的读取能力。

深入 Thrift 测试 Schema:parquet-compat.thrift

parquet-compat.thrift 定义了 Thrift 侧的测试结构,namespace 为org.apache.spark.sql.execution.datasources.parquet.test.thrift,与 Avro 侧完全隔离。核心 structParquetThriftCompat的设计比 Avro 侧更细:

标量类型与 required 语义

struct ParquetThriftCompat { 1: required bool boolColumn; 2: required byte byteColumn; 3: required i16 shortColumn; 4: required i32 intColumn; 5: required i64 longColumn; 6: required double doubleColumn; 7: required binary binaryColumn; 8: required string stringColumn; 9: required Suit enumColumn

这里覆盖了 Thrift 独有的整型精度类型bytei16,它们对应 Parquet 中INT32上的不同逻辑类型标注,是验证 Thrift→Parquet 类型映射精确度的关键字段。

optional 字段族

10: optional bool maybeBoolColumn; 11: optional byte maybeByteColumn; // ... maybeShort / maybeInt / maybeLong / maybeDouble 16: optional binary maybeBinaryColumn; 17: optional string maybeStringColumn; 18: optional Suit maybeEnumColumn;

optional字段族与 Avro 侧的union { null, T }对应,二者共同覆盖"可空列"在两种序列化框架下的 Parquet 表现,用于对比验证 Spark 读取 optional 数据的一致性。

集合类型组合

19: required list<string> stringsColumn; 20: required set<i32> intSetColumn; 21: required map<i32, string> intToStringColumn; 22: required map<i32, list<Nested>> complexColumn; }

第 19~22 字段覆盖 Thrift 的listsetmapmap<k, list<Nested>>四层嵌套组合。其中set<i32>在 Parquet 中通常以 repeated group 表示,complexColumn则验证"Map 的 value 是结构体数组"这种复杂 repeated 结构的兼容性。Thrift 侧的Suit枚举与Nested结构分别对应 Avro 侧的同名定义,保证两套 Schema 的对比维度对齐。

测试侧如何使用这些生成产物

生成的 Java 类被 Parquet 兼容性测试用例直接引用。测试基类位于 ParquetCompatibilityTest.scala,它扩展了QueryTestParquetTest,提供两类核心辅助能力:

  1. Schema 读取与校验readParquetSchema(path)通过ParquetFileReader.readAllFootersInParallel读取 Parquet 文件 footer 中的 Schema,默认过滤_前缀的隐藏文件;logParquetSchema则将 parquet-avro 写入文件的 Schema 打印到日志,便于人工核对类型映射。
  2. 底层 Parquet 写入辅助:提供writeBinaryDatamakePointWkb/makeLineStringWkb/makePolygonWkb等工具方法,直接通过ValuesWriter构造 WKB 二进制数据,用于在测试中生成包含 Geometry 类型的数据。从该基类源码结构可以推断,其整体思路是"用生态工具(Avro/Thrift)写出文件 → 用 Spark SQL 读取并断言 Schema 与数据 → 反向验证 Spark 写入的数据可被生态工具读取",从而形成双向兼容性闭环。

实操:新增一个兼容性测试字段的完整流程

综合以上内容,当你在 Spark 中新增一个类型映射的兼容性验证时,标准操作路径如下:

  1. 修改 Schema 源:在 parquet-compat.avdl 或 parquet-compat.thrift 中新增 record / struct 或字段;
  2. 重新生成代码:在sql/core/src/test/下执行./scripts/gen-avro.sh./scripts/gen-thrift.sh,脚本会重建gen-java/并刷新avro/*.avpr协议文件;
  3. 核对产物:确认gen-java/下新增的 Java 类与 namespace 正确,git status中应同时包含 Schema 源与生成产物的变更;
  4. 补充测试用例:在ParquetCompatibilityTest的子类中新增用例,用生成的 POJO 写入 Parquet 文件,再通过spark.read.parquet(...)读取并断言数据类型与值;
  5. 运行测试:通过 Maven 或 sbt 运行sql/core模块下对应的 Parquet 测试套件,验证新字段在 Spark 读写两侧均兼容。

总结

sql/core/src/test/下的这套 Parquet 兼容性测试资产,通过"Schema 源(.avdl/.thrift)→ 生成脚本(gen-avro.sh/gen-thrift.sh)→ 入库的生成代码(.avpr/gen-java/)→ 测试基类(ParquetCompatibilityTest)"的完整链路,在避免构建期代码生成的同时,系统性地覆盖了 Avro 与 Thrift 到 Parquet 的基本类型、可空类型、枚举、数组、集合、Map 与多层嵌套组合的兼容性验证。理解这套机制,既是向 Spark 提交兼容性修复与新增测试的前提,也能帮助你更透彻地把握 Spark SQL 对 Parquet 复杂 Schema 的类型映射边界。

  • 大数据
  • 数据分析
  • 批处理
  • 流处理
  • 机器学习
  • 图计算

【免费下载链接】spark

Apache Spark - A unified analytics engine for large-scale data processing

项目地址:https://gitcode.com/gh_mirrors/sp/spark
点击查看免费下载
上一篇:Lottie-Android动画调试工具:OutlineMasksAndMattes使用指南
下一篇:Munder Difflin vs Vibe Kanban:管理AI Agent编队到底该用哪种看板?

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询