Spark 内置 Avro 数据源完全指南:读写、Schema 演进与 to_avro/from_avro 实战 Spark 内置 Avro 数据源完全指南读写、Schema 演进与 to_avro/from_avro 实战【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark导读Apache Spark 自 2.4 版本起通过spark-avro外部模块为 Spark SQL 提供 Apache Avro 格式的内置读写能力。本文以当前仓库中的 sql-data-sources-avro.md 为骨架结合 AvroOptions.scala 等源码实现系统讲解 Avro 数据源的部署方式、DataFrame API、to_avro/from_avro函数、全部数据源选项与 SQL 配置、Avro 与 Spark SQL 类型映射规则以及循环引用字段的处理方案。读完本文你将能够在 Spark 应用中直接读写 Avro 文件并借助 Avro 与 Kafka 的配合构建流式数据管道。一、spark-avro 模块外部但内置spark-avro是一个内置但外部built-in but external的模块它随 Spark 发布并受官方维护源码位于仓库的connector/avro目录但默认不会被打进spark-submit或spark-shell的运行时 classpath 中。因此使用前必须显式地把它添加为应用依赖。从源码结构看该模块的主逻辑集中在connector/avro/src/main/scala/org/apache/spark/sql/avro/表达式与 Schema 转换实现如 AvroDataToCatalyst.scala以及org/apache/spark/sql/v2/avro/DataSource V2 的读写实现如 AvroDataSourceV2.scala测试则覆盖了 AvroSuite.scala、AvroFunctionsSuite.scala 等大量用例。1.1 通过 --packages 部署与任何 Spark 应用一样使用spark-submit提交应用时把spark-avro_{Scala 二进制版本}及其依赖直接通过--packages加入即可./bin/spark-submit --packages org.apache.spark:spark-avro_{{site.SCALA_BINARY_VERSION}}:{{site.SPARK_VERSION_SHORT}} ...其中{{site.SCALA_BINARY_VERSION}}是 Scala 二进制版本号如2.12/2.13{{site.SPARK_VERSION_SHORT}}是 Spark 版本号如4.0.0实际使用时需替换为具体版本。例如./bin/spark-submit --packages org.apache.spark:spark-avro_2.13:4.0.0 \ your-application.jar1.2 在 spark-shell 中实验想在spark-shell里快速实验同样使用--packages./bin/spark-shell --packages org.apache.spark:spark-avro_{{site.SCALA_BINARY_VERSION}}:{{site.SPARK_VERSION_SHORT}} ...关于提交含外部依赖的应用的更多细节请参见 Application Submission Guide应用提交指南。如果使用自己的spark-avrojar 构建也可以用--jars方式部署详见下文与 Databricks spark-avro 的兼容性一节及 Advanced Dependency Management。二、Load 与 SaveDataFrame API 读写 Avro由于spark-avro是外部模块DataFrameReader和DataFrameWriter上没有.avro()快捷方法。读取/写入 Avro 数据时必须把数据源format指定为avro或全限定名org.apache.spark.sql.avro。2.1 Pythondf spark.read.format(avro).load(examples/src/main/resources/users.avro) df.select(name, favorite_color).write.format(avro).save(namesAndFavColors.avro)2.2 Scalaval usersDF spark.read.format(avro).load(examples/src/main/resources/users.avro) usersDF.select(name, favorite_color).write.format(avro).save(namesAndFavColors.avro)2.3 JavaDatasetRow usersDF spark.read().format(avro).load(examples/src/main/resources/users.avro); usersDF.select(name, favorite_color).write().format(avro).save(namesAndFavColors.avro);2.4 Rdf - read.df(examples/src/main/resources/users.avro, avro) write.df(select(df, name, favorite_color), namesAndFavColors.avro, avro)实战提示上述示例中的examples/src/main/resources/users.avro正是当前仓库中真实存在的示例数据文件可直接在仓库根目录下验证仓库内还提供了对应的 Avro Schema 文件examples/src/main/resources/user.avsc记录名User含name与可空的favorite_color两个字段。三、to_avro() 与 from_avro()列级 Avro 编解码Avro 包提供了两个函数用于在列级别完成 Avro 二进制与 Spark SQL 数据之间的转换to_avro把某一列编码为 Avro 格式的二进制from_avro()把 Avro 二进制数据解码为列。两个函数都只做一列到另一列的变换输入/输出的 SQL 数据类型既可以是复杂类型也可以是原始类型。使用 Avro 记录作为列在读写 Kafka 这类流式数据源时尤为有用Kafka 的每条 key-value 记录都会被附加元数据如进入 Kafka 的时间戳、offset 等。如果承载数据的value字段是 Avro 格式可以用from_avro()把数据提取出来进行补全、清洗后再写回 Kafka 或输出到文件to_avro()可以把 struct 转成 Avro 记录在写 Kafka 前把多个列重新编码到单个列非常实用。3.1 Python 示例Kafka 流管道from pyspark.sql.avro.functions import from_avro, to_avro # from_avro requires Avro schema in JSON string format. jsonFormatSchema open(examples/src/main/resources/user.avsc, r).read() df spark\ .readStream\ .format(kafka)\ .option(kafka.bootstrap.servers, host1:port1,host2:port2)\ .option(subscribe, topic1)\ .load() # 1. Decode the Avro data into a struct; # 2. Filter by column favorite_color; # 3. Encode the column name in Avro format. output df\ .select(from_avro(value, jsonFormatSchema).alias(user))\ .where(user.favorite_color red)\ .select(to_avro(user.name).alias(value)) query output\ .writeStream\ .format(kafka)\ .option(kafka.bootstrap.servers, host1:port1,host2:port2)\ .option(topic, topic2)\ .start()3.2 Scala 示例import org.apache.spark.sql.avro.functions._ // from_avro requires Avro schema in JSON string format. val jsonFormatSchema new String(Files.readAllBytes(Paths.get(./examples/src/main/resources/user.avsc))) val df spark .readStream .format(kafka) .option(kafka.bootstrap.servers, host1:port1,host2:port2) .option(subscribe, topic1) .load() // 1. Decode the Avro data into a struct; // 2. Filter by column favorite_color; // 3. Encode the column name in Avro format. val output df .select(from_avro($value, jsonFormatSchema) as $user) .where(user.favorite_color \red\) .select(to_avro($user.name) as $value) val query output .writeStream .format(kafka) .option(kafka.bootstrap.servers, host1:port1,host2:port2) .option(topic, topic2) .start()3.3 Java 示例import static org.apache.spark.sql.functions.col; import static org.apache.spark.sql.avro.functions.*; // from_avro requires Avro schema in JSON string format. String jsonFormatSchema new String(Files.readAllBytes(Paths.get(./examples/src/main/resources/user.avsc))); DatasetRow df spark .readStream() .format(kafka) .option(kafka.bootstrap.servers, host1:port1,host2:port2) .option(subscribe, topic1) .load(); // 1. Decode the Avro data into a struct; // 2. Filter by column favorite_color; // 3. Encode the column name in Avro format. DatasetRow output df .select(from_avro(col(value), jsonFormatSchema).as(user)) .where(user.favorite_color \red\) .select(to_avro(col(user.name)).as(value)); StreamingQuery query output .writeStream() .format(kafka) .option(kafka.bootstrap.servers, host1:port1,host2:port2) .option(topic, topic2) .start();3.4 R 示例# from_avro requires Avro schema in JSON string format. jsonFormatSchema - paste0(readLines(examples/src/main/resources/user.avsc), collapse ) df - read.stream( kafka, kafka.bootstrap.servers host1:port1,host2:port2, subscribe topic1 ) # 1. Decode the Avro data into a struct; # 2. Filter by column favorite_color; # 3. Encode the column name in Avro format. output - select( filter( select(df, alias(from_avro(value, jsonFormatSchema), user)), column(user.favorite_color) red ), alias(to_avro(user.name), value) ) write.stream( output, kafka, kafka.bootstrap.servers host1:port1,host2:port2, topic topic2 )3.5 SQL 示例to_avro/from_avro也以 SQL 函数的形式暴露配合NAMED_STRUCT、变量和MAP()使用CREATE TABLE t AS SELECT NAMED_STRUCT(u, NAMED_STRUCT(member0, member0, member1, member1)) AS s FROM VALUES (1, NULL), (NULL, a) tab(member0, member1); DECLARE avro_schema STRING; SET VARIABLE avro_schema { type: record, name: struct, fields: [{ name: u, type: [int,string] }] }; SELECT TO_AVRO(s, avro_schema) AS RESULT FROM t; SELECT FROM_AVRO(result, avro_schema, MAP()).u FROM ( SELECT TO_AVRO(s, avro_schema) AS RESULT FROM t); DROP TEMPORARY VARIABLE avro_schema; DROP TABLE t;实现注记从源码看from_avro由 AvroDataToCatalyst.scala 中的表达式实现其prettyName即为from_avro并支持doGenCode代码生成路径to_avro对应 CatalystDataToAvro.scala。解析 Avro Schema 时实际使用的选项如稳定 Union 标识符、递归深度等直接取自AvroOptions。四、数据源选项Data Source OptionAvro 数据源选项可通过两种途径设置DataFrameReader/DataFrameWriter的.option方法from_avro函数的options参数。从 AvroOptions.scala 的源码可以看到所有选项名avroSchema、avroSchemaUrl、recordName、recordNamespace、ignoreExtension、compression、mode、positionalFieldMatching、datetimeRebaseMode、enableStableIdentifiersForUnionType、stableIdentifierPrefixForUnionType、recursiveFieldMaxDepth等均以大小写不敏感的CaseInsensitiveMap解析且额外支持从 URL 读取 schema受spark.sql.avro.schema.url.allowedSchemes限制。各选项说明如下属性名默认值含义作用域引入版本avroSchemaNone用户以 JSON 格式提供的可选 Schema。① 读取 Avro 文件或调用from_avro时可设置为与真实 Avro Schema 兼容但不同的演进后 Schema反序列化结果将与演进后的 Schema 一致。例如设置一个带默认值新增列的演进 SchemaSpark 的读取结果也会包含该新列。注意配合from_avro使用时仍需把真实 Avro Schema 作为函数参数传入。② 写 Avro 时若期望的输出 Schema 与 Spark 转换出的默认 Schema 不一致例如某列期望是enum类型而默认转换结果是string可通过此选项指定。读、写与from_avro函数2.4.0recordNametopLevelRecord写入结果中的顶层记录名Avro 规范所必需。写2.4.0recordNamespace写入结果中的记录命名空间。写2.4.0ignoreExtensiontrue控制读取时是否忽略不带.avro扩展名的文件。开启后无论是否带.avro扩展名所有文件都会被加载。已废弃请改用通用数据源选项pathGlobFilter过滤文件名详见 sql-data-sources-generic-options.md。源码中该选项的默认值被标为 deprecated自 3.0 起。读2.4.0compressionsnappy写操作使用的压缩编解码器。目前支持uncompressed、snappy、deflate、bzip2、xz、zstandard。若未设置则取配置spark.sql.avro.compression.codec的值。写2.4.0modeFAILFASTfrom_avro函数的解析模式①FAILFAST遇到损坏记录时抛出异常。②PERMISSIVE损坏记录被当作 null 结果处理因此数据 Schema 会被强制为完全可空可能与用户提供的 Schema 不同。from_avro函数2.4.0datetimeRebaseModespark.sql.avro.datetimeRebaseModeInRead配置的值指定date、timestamp-micros、timestamp-millis逻辑类型的值从儒略历到格里高利历的 rebase 模式①EXCEPTION遇到两种历法下有歧义的古老日期/时间戳时读取失败。②CORRECTED不 rebase直接加载。③LEGACY把古老日期/时间戳从儒略历 rebase 到格里高利历。读与from_avro函数3.2.0positionalFieldMatchingfalse与avroSchema选项配合使用调整所提供的 Avro Schema 与 SQL Schema 的字段匹配方式。默认按字段名匹配、忽略位置设为true后按字段位置匹配。读与写3.2.0enableStableIdentifiersForUnionTypefalse设为true时Avro Schema 被反序列化为 Spark SQL SchemaUnion 类型被转换为字段名与各类型保持一致的 struct字段名转为小写如member_int、member_string。若两个用户自定义类型名或用户自定义类型名与内建类型名在忽略大小写后相同会抛出异常其他情况下字段名可唯一标识。读3.5.0stableIdentifierPrefixForUnionTypemember_启用enableStableIdentifiersForUnionType时用于配置 Avro Union 类型字段的前缀。读4.0.0recursiveFieldMaxDepth-1指定解析 Schema 时允许的递归层数上限。设为负数或 0 表示不允许递归字段设为 1 丢弃所有递归字段设为 2 允许递归一层设为 3 允许递归两层依此类推上限为 15源码中常量RECURSIVE_FIELD_MAX_DEPTH_LIMIT 15。超过限制会截断返回的 struct。示例见下文处理 Avro 字段的循环引用。读4.0.0五、SQL 配置ConfigurationAvro 相关配置可通过spark.conf.set设置或在 SQL 中用SET keyvalue命令设置属性名默认值含义引入版本spark.sql.legacy.replaceDatabricksSparkAvro.enabledtrue设为 true 时数据源 providercom.databricks.spark.avro被映射到内置但外部的 Avro 模块以保持向后兼容。注意该 SQL 配置已在 Spark 3.2 中废弃未来可能移除。2.4.0spark.sql.avro.compression.codecsnappy写 Avro 文件使用的压缩编解码器。支持uncompressed、deflate、snappy、bzip2、xz、zstandard。2.4.0spark.sql.avro.deflate.level-1deflate 编解码器写 Avro 文件时的压缩级别。合法值为 1~9含或 -1。默认 -1 对应当前实现中的 6 级。2.4.0spark.sql.avro.xz.level6xz 编解码器写 Avro 文件时的压缩级别。合法值为 1~9含。4.0.0spark.sql.avro.zstandard.level3zstandard 编解码器写 Avro 文件时的压缩级别。4.0.0spark.sql.avro.zstandard.bufferPool.enabledfalse设为 true 时写 Avro 文件启用 ZSTD JNI 库的缓冲池。4.0.0spark.sql.avro.datetimeRebaseModeInReadEXCEPTIONdate、timestamp-micros、timestamp-millis逻辑类型从儒略历到格里高利历的 rebase 模式①EXCEPTION读到两历法下有歧义的古老日期/时间戳时读取失败。②CORRECTED不 rebase按原样读取。③LEGACY读取 Avro 文件时把日期/时间戳从遗留混合历儒略历格里高利历rebase 到格里高利历。仅当 Avro 文件的写入方信息如 Spark、Hive未知时此配置才生效。3.0.0spark.sql.avro.datetimeRebaseModeInWriteEXCEPTION从格里高利历到儒略历的 rebase 模式①EXCEPTION写数据时遇到两历法下有歧义的古老日期/时间戳则失败。②CORRECTED不 rebase按原样写入。③LEGACY写 Avro 文件时把日期/时间戳从格里高利历 rebase 到遗留混合历儒略历格里高利历。3.0.0spark.sql.avro.filterPushdown.enabledtrue设为 true 时对 Avro 数据源启用谓词下推。3.1.0六、与 Databricks spark-avro 的兼容性这个内置的 Avro 数据源模块源自 Databricks 的开源仓库spark-avro并与之兼容。默认情况下SQL 配置spark.sql.legacy.replaceDatabricksSparkAvro.enabled处于启用状态数据源 providercom.databricks.spark.avro会被映射到这个内置 Avro 模块。对于那些在 Catalog 元数据中Provider属性为com.databricks.spark.avro的 Spark 表这种映射对于使用内置模块加载这些表是必不可少的。需要注意的是Databricks 的spark-avro中定义了隐式类AvroDataFrameWriter和AvroDataFrameReader来提供.avro()快捷函数而这个内置但外部的模块删除了这两个隐式类。请改用DataFrameWriter/DataFrameReader上的.format(avro)这更简洁也足够好用。如果你更倾向于使用自己构建的spark-avrojar 文件可以禁用spark.sql.legacy.replaceDatabricksSparkAvro.enabled配置并在部署应用时改用--jars选项。详见 Application Submission Guide 中的 Advanced Dependency Management。七、Avro → Spark SQL 类型映射7.1 基础类型与复杂类型当前 Spark 支持读取 Avro 记录下的所有原始类型与复杂类型Avro 类型Spark SQL 类型booleanBooleanTypeintIntegerTypelongLongTypefloatFloatTypedoubleDoubleTypestringStringTypeenumStringTypefixedBinaryTypebytesBinaryTyperecordStructTypearrayArrayTypemapMapTypeunion见下方说明7.2 Union 类型除上述类型外还支持读取union类型。以下三种视为基本union类型union(int, long)映射为 LongTypeunion(float, double)映射为 DoubleTypeunion(something, null)something 为任意受支持的 Avro 类型映射为与 something 相同的 Spark SQL 类型且 nullable 设为 true。所有其他 union 类型视为复杂类型映射为 StructType字段名为member0、member1…… 与 union 成员一一对应。这与 Avro 与 Parquet 之间转换时的行为一致。源码佐证这些映射规则在 AvroSuite.scala 中有专门测试例如union(int, long) is read as long、union(float, double) is read as double、union(float, double, null) is read as nullable double以及针对复杂 union 的字段名断言member0、member1等。7.3 Logical 类型还支持读取以下 Avro 逻辑类型Avro 逻辑类型Avro 类型Spark SQL 类型dateintDateTypetimestamp-millislongTimestampTypetimestamp-microslongTimestampTypetime-microslongTimeTypetimestamp-nanoslongTimestampType(p)p 为 7~9需spark.sql.timestampNanosTypes.enabledtruelocal-timestamp-nanoslongTimestampNTZType(p)p 为 7~9需spark.sql.timestampNanosTypes.enabledtruedecimalfixedDecimalTypedecimalbytesDecimalType目前读取时会忽略 Avro 文件中的 docs、aliases 及其他属性。八、Spark SQL → Avro 类型映射Spark 支持把所有 Spark SQL 类型写入 Avro。大多数类型的映射是直接的如 IntegerType 转换为 int但有几个特例如下表Spark SQL 类型Avro 类型Avro 逻辑类型ByteTypeintShortTypeintBinaryTypebytesDateTypeintdateTimestampTypelongtimestamp-microsTimeTypelongtime-microsTimestampType(p)p 为 7~9longtimestamp-nanosTimestampNTZType(p)p 为 7~9longlocal-timestamp-nanosDecimalTypefixeddecimal你还可以用avroSchema选项指定整个输出 Avro Schema从而把 Spark SQL 类型转换为其他 Avro 类型。以下转换默认不执行需要用户指定 Avro SchemaSpark SQL 类型Avro 类型Avro 逻辑类型BinaryTypefixedStringTypeenumTimestampTypelongtimestamp-millisDecimalTypebytesdecimal九、处理 Avro 字段的循环引用在 Avro 中循环引用发生在字段类型定义在其某个父记录中时即类型自引用。解析数据时这可能引发无限循环或其他异常行为。要读取包含循环引用 Schema 的 Avro 数据可以使用recursiveFieldMaxDepth选项指定解析 Schema 时允许的最大递归层数。默认情况下Spark Avro 数据源把recursiveFieldMaxDepth设为 -1不允许递归字段。需要时可设为 1~15。设为 1丢弃所有递归字段设为 2允许递归一次设为 3允许递归两次大于 15 不允许因为这可能导致性能问题甚至栈溢出。源码中该限制在 AvroOptions.scala 中由常量RECURSIVE_FIELD_MAX_DEPTH_LIMIT: Int 15强制约束。以下面的 Avro 消息为例其 SQL Schema 会根据recursiveFieldMaxDepth的取值而变化{ type: record, name: Node, fields: [ {name: Id, type: int}, {name: Next, type: [null, Node]} ] }上面定义的 Avro Schema会基于recursiveFieldMaxDepth的值转换为如下结构的 Spark SQL 列1: structId: int 2: structId: int, Next: structId: int 3: structId: int, Next: structId: int, Next: structId: int可以看到深度为 1 时递归字段Next被完全丢弃深度为 2 时Next展开一层深度为 3 时展开两层依此类推。十、小结部署spark-avro是内置但外部模块需通过--packages或--jars加入应用依赖API用.format(avro)读写或用to_avro/from_avro做列级编解码两者都适用于 Kafka 流管道场景选项与配置avroSchema支持 Schema 演进与自定义输出 SchemarecursiveFieldMaxDepth处理循环引用压缩、历法 rebase、谓词下推等均有对应选项或 SQL 配置详细定义见 AvroOptions.scala类型映射Avro 与 Spark SQL 的原始/复杂/逻辑类型映射规则如上表所示Union 有专门的简化与展开规则测试用例见 AvroSuite.scala兼容性默认映射 Databrickscom.databricks.spark.avroprovider但不提供.avro()快捷方法。相关阅读Spark SQL 编程指南、通用文件数据源选项、应用提交指南。【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考