news 2026/10/11 9:38:57

Apache Beam ZetaSQL 数据类型完全指南:标量类型、ARRAY/STRUCT 扩展与底层映射原理

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Beam ZetaSQL 数据类型完全指南:标量类型、ARRAY/STRUCT 扩展与底层映射原理

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

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

Beam ZetaSQL 是 Apache Beam SQL 方言的一种实现,语法接近 BigQuery 使用的 ZetaSQL 语言,特别适合需要与 BigQuery 表进行读写交互的 Beam 管道。本文以 Beam 仓库中的 [data-types.md](https://link.gitcode.com/i/ad9eba654c955c022b5fc4627b2112e8) 官方文档为骨架,系统讲解 Beam ZetaSQL 支持的全部数据类型(标量类型 + ARRAY/STRUCT 扩展)、各类型的属性约束与声明语法,并结合仓库源码([ZetaSqlBeamTranslationUtils.java](https://link.gitcode.com/i/efc6142f6389d045dfc775846eaf9987))与测试用例([ZetaSqlDialectSpecTest.java](https://link.gitcode.com/i/4b592f4d11566c0f9f945d48deaad170)、[ZetaSqlTimeFunctionsTest.java](https://link.gitcode.com/i/a43c1b174248590be4ac9d850b1045fa))深入剖析类型在 Beam Schema 与 ZetaSQL 之间的双向映射、TIMESTAMP 毫秒精度限制等底层细节。读完本文,你将能够为 Beam ZetaSQL 查询正确选择类型、写出可运行的 ARRAY/STRUCT 字面量表达式,并理解类型转换在管道执行时的真实行为。

数据类型总览

Beam ZetaSQL 不仅支持标准 SQL 标量数据类型,还通过扩展支持了数组(ARRAY)、结构体(STRUCT,即嵌套行)等复杂数据类型。这一能力是 Beam 统一批/流处理模型与复杂数据处理扩展的一部分,可应用于所有 Beam SQL 方言,包括 Beam ZetaSQL。关于 Beam ZetaSQL 的整体能力,可参见 overview.md。

Beam ZetaSQL 支持的数据类型可归纳为两类:

  • 标量类型:INT64、FLOAT64、BOOL、STRING、BYTES、TIMESTAMP;
  • 复合/扩展类型:ARRAY(元素为任意非 ARRAY 类型的有序列表)、STRUCT(有序字段容器,每字段有类型和可选名称)。

此外,从源码映射(ZetaSqlBeamTranslationUtils.toZetaSqlType)可以看到,Beam Schema 中的DECIMAL、DATE、TIME、DATETIME等类型也会被转换为对应的 ZetaSQL 类型(TYPE_NUMERIC、TYPE_DATE、TYPE_TIME、TYPE_DATETIME),说明 Beam 侧的丰富类型体系在 ZetaSQL 方言下均有对应表示。

数据类型属性(Properties)

理解一种 SQL 方言,首先要弄清每种类型能在哪些语法位置使用。Beam ZetaSQL 用四个属性刻画数据类型的可用性:

属性含义适用范围
Nullable(可空)NULL是合法值所有类型均适用,但有例外:ARRAY 本身不能为NULL;NULL ARRAY元素无法持久化到表,查询也无法处理NULL ARRAY元素
Orderable(可排序)可用于ORDER BY子句除ARRAY、STRUCT外的所有类型
Groupable(可分组)通常可出现在GROUP BY、DISTINCT或PARTITION BY后的表达式中除ARRAY、STRUCT、FLOAT64外的所有类型。注意:PARTITION BY表达式不能包含浮点类型FLOAT和DOUBLE
Comparable(可比较)同类型的值可以相互比较所有类型,但存在例外:ARRAY 比较不被支持;STRUCT 仅支持按字段顺序逐字段的等值比较(字段名被忽略),不支持小于/大于比较

所有支持比较的类型都可以用于JOIN条件。Beam ZetaSQL 的 JOIN 类型说明见 query-syntax.md 中的 "JOIN Types" 一节(文档内锚点#join_types)。

数值类型

数值类型包括整数类型与浮点类型。

整数类型 INT64

整数是没有小数部分的数值。Beam ZetaSQL 的整数类型为INT64:

名称存储大小取值范围
INT648 字节-9,223,372,036,854,775,808 至 9,223,372,036,854,775,807

该范围与 Java 的Long完全一致。在底层实现中,Beam 将INT64直接映射为 Beam Schema 的FieldType.INT64(见 ZetaSqlBeamTranslationUtils.java),值转换时对应Value.createInt64Value((Long) object),Java 侧以Long承载。

浮点类型 FLOAT64

浮点值是带小数部分的近似数值:

名称存储大小说明
FLOAT648 字节双精度(近似)十进制值

注意:浮点类型FLOAT64是不可分组的,即不能出现在GROUP BY、DISTINCT或PARTITION BY表达式中(见上文属性表)。这是为了规避浮点近似值导致的等价性判定问题。

从源码看,Beam 侧DOUBLE对应 ZetaSQL 的FLOAT64(TypeKind.TYPE_DOUBLE)。值转换时有一个值得注意的细节:如果 ZetaSQL 侧的值类型是INT64(浮点部分为零的值会被视为整数),转换到 Beam 时会显式(double) value.getInt64Value()转成双精度浮点(见 ZetaSqlBeamTranslationUtils.java)。

布尔类型 BOOL

名称说明
BOOL布尔值由关键字TRUE和FALSE表示(不区分大小写)

字符串类型 STRING

名称说明
STRING可变长度的字符(Unicode)数据

STRING 有两条关键语义规则:

  1. UTF-8 编码强制:输入 STRING 值必须是 UTF-8 编码,输出 STRING 值也会是 UTF-8 编码。像 CESU-8、Modified UTF-8 这类替代编码不会被当作合法 UTF-8 处理。
  2. 按 Unicode 字符而非字节运算:所有作用于 STRING 值的函数和运算符都以 Unicode 字符为单位操作。例如SUBSTR、LENGTH作用于 STRING 输入时,计数的是 Unicode 字符而不是字节。比较操作也定义在 Unicode 字符之上——<(小于)比较和ORDER BY按字符逐个比较,Unicode 码点较小的字符被视为更小。

字节类型 BYTES

名称说明
BYTES可变长度的二进制数据

STRING 与 BYTES 是两个独立的类型,不能互换使用。两者之间的CAST转换会强制要求字节按 UTF-8 编码(即 STRING 转 BYTES 时字节必须是合法 UTF-8,BYTES 转 STRING 时按 UTF-8 解码)。在源码中,BYTES 值在 Java 侧以byte[]承载,转换到 ZetaSQL 时通过ByteString.copyFrom((byte[]) object)包装(见 ZetaSqlBeamTranslationUtils.java)。

时间戳类型 TIMESTAMP

Caution(重要限制):Beam ZetaSQL 的TIMESTAMP精度为毫秒。如果某个TIMESTAMP字段具有亚毫秒(sub-millisecond)精度,SQL 会抛出IllegalArgumentException。

名称说明取值范围
TIMESTAMP表示一个绝对时间点,精度为毫秒0001-01-01 00:00:00 至 9999-12-31 23:59:59.999 UTC

TIMESTAMP表示的是独立于任何时区或夏令时(Daylight Savings Time)等约定的绝对时间点。测试用例testTimestampMicrosecondUnsupported(见 ZetaSqlTimeFunctionsTest.java)正是验证了这一限制:当 SQL 中出现TIMESTAMP '2000-01-01 00:11:22.345678+00'这样的微秒精度字面量时,ZetaSQLQueryPlanner.convertToBeamRel会抛出UnsupportedOperationException。

在源码层,Beam 的DATETIME(Instant)与 ZetaSQL 的TIMESTAMP互转时都围绕毫秒展开:Beam → ZetaSQL 时用LongMath.checkedMultiply(((Instant) object).getMillis(), MICROS_PER_MILLI)把毫秒放大为 Unix 微秒(MICROS_PER_MILLI = 1000L);反向转换时用Instant.ofEpochMilli(value.getTimestampUnixMicros() / MICROS_PER_MILLI)截断回毫秒(见 ZetaSqlBeamTranslationUtils.java 与 同文件)。同时 DateTimeUtils.java 中定义了MIN_UNIX_MILLIS = -62135596800000L与MAX_UNIX_MILLIS = 253402300799999L,对应上表的时间戳取值边界,超出即报错。

规范格式(Canonical Format)

时间戳的文本表示遵循以下规范格式:

YYYY-[M]M-[D]D[( |T)[H]H:[M]M:[S]S[.DDD]][time zone]

各组成部分含义:

  • YYYY:四位年份
  • [M]M:一位或两位月份
  • [D]D:一位或两位日期
  • ( |T):空格或T分隔符
  • [H]H:一位或两位小时(合法值为 00~23)
  • [M]M:一位或两位分钟(合法值为 00~59)
  • [S]S:一位或两位秒(合法值为 00~59)
  • [.DDD]:最多三位小数(即最高毫秒精度)
  • [time zone]:表示时区的字符串,详见下文"时区"小节

时区只在解析时间戳或格式化显示时间戳时使用;时间戳值本身不存储特定时区。字符串形式的时间戳可以携带时区;未显式指定时区时,默认使用 UTC。

时区(Time Zones)

时区由字符串表示,采用以下两种规范格式之一:

  • 相对于协调世界时(UTC)的偏移量,或用字母Z表示 UTC;
  • 来自 tz database 的时区名称(如America/Los_Angeles)。
UTC 偏移格式
(+|-)H[H][:M[M]] Z

示例:

-08:00 -8:15 +3:00 +07:30 -7 Z

使用该格式时,时区与时间戳其余部分之间不允许有空格:

2014-09-27 12:30:00.45-8:00 2014-09-27T12:30:00.45Z
时区名称格式
continent/[region/]city

示例:

America/Los_Angeles America/Argentina/Buenos_Aires

使用时区名称时,名称与时间戳其余部分之间必须有空格:

2014-09-27 12:30:00.45 America/Los_Angeles

需要特别注意的是,并非所有报告相同时刻的时区名都是可互换的。例如America/Los_Angeles在夏令时期间与UTC-7:00报告相同时间,但在非夏令时期间与UTC-8:00报告相同时间——即使用固定偏移还是用具名时区,对跨夏令时边界的计算会产生不同结果。若未指定时区,则使用默认时区值(UTC)。

闰秒(Leap Seconds)

TIMESTAMP本质上是从1970-01-01 00:00:00 UTC起的偏移量,假定每分钟恰好 60 秒。闰秒不作为已存储时间戳的一部分被表示。

  • 如果输入在秒字段使用了":60"来表示闰秒,该闰秒在转换为时间戳值时不会保留,而是被解释为下一分钟的":00"秒字段;
  • 闰秒不影响时间戳计算:所有时间戳计算都使用 Unix 风格时间戳(不含闰秒);
  • 闰秒只有通过测量真实世界时间的函数才能被观察到,在这些函数中,出现闰秒时一个时间戳秒可能被跳过或重复。

数组类型 ARRAY

名称说明
ARRAY零个或多个任意非 ARRAY 类型元素的有序列表

ARRAY是零个或多个非 ARRAY 值的有序列表。要点如下:

  • ARRAY 不能嵌套 ARRAY:ARRAY of ARRAYs不被允许,产生 ARRAY of ARRAYs 的查询会返回错误。如需多维数组,必须使用SELECT AS STRUCT构造在两层 ARRAY 之间插入一个 STRUCT;
  • 空 ARRAY 与NULLARRAY 是两个不同的值;
  • ARRAY 可以包含NULL元素。

声明 ARRAY 类型

ARRAY 类型使用尖括号(<和>)声明,元素类型可以任意复杂,唯一例外是 ARRAY 不能直接包含另一个 ARRAY。

格式:

ARRAY<T>

示例:

类型声明含义
ARRAY<INT64>64 位整数的简单数组
ARRAY<STRUCT<INT64, INT64>>STRUCT 数组,每个 STRUCT 包含两个 64 位整数
ARRAY<ARRAY<INT64>>(不支持)无效的类型声明(此处列出仅为说明如何创建多层 ARRAY)。ARRAY 不能直接包含 ARRAY,应参考下例
ARRAY<STRUCT<ARRAY<INT64>>>64 位整数数组的数组。注意两个 ARRAY 之间夹了一个 STRUCT,因为 ARRAY 不能直接持有其他 ARRAY

在测试用例中,testArrayStructLiteral(见 ZetaSqlDialectSpecTest.java)验证了字面量表达式SELECT ARRAY<STRUCT<INT64, INT64>>[(11, 12)];:解析后 Beam 侧得到FieldType.array(FieldType.row(Schema.of(Field.of("s", FieldType.INT64), Field.of("i", FieldType.INT64)))),输出行包含一个[Row(11, 12)]形式的数组值,与 ZetaSQL 的ARRAY<STRUCT>一一对应。

此外,UNNEST是将 ARRAY 展开为行集的核心操作符,测试覆盖了大量场景:字面量展开(SELECT * FROM UNNEST(ARRAY<STRING>['foo', 'bar']))、带 NULL 元素(ARRAY<STRING>['foo', NULL, 'bar'])、命名展开(UNNEST(...) AS T1)、带偏移量展开(UNNEST([3, 4]) AS x WITH OFFSET p)、数组/STRUCT 混合嵌套展开等,均可作为编写查询时的参考(见 ZetaSqlDialectSpecTest.java)。

结构体类型 STRUCT

名称说明
STRUCT有序字段的容器,每个字段具有类型(必填)和字段名(可选)

声明 STRUCT 类型

STRUCT 类型同样使用尖括号声明,元素类型可以任意复杂。

格式:

STRUCT<T>

示例:

类型声明含义
STRUCT<INT64>只有一个未命名的 64 位整数字段的简单 STRUCT
STRUCT<x STRUCT<y INT64, z INT64>>内部嵌套名为x的 STRUCT;x有两个字段y和z,均为 64 位整数
STRUCT<inner_array ARRAY<INT64>>包含名为inner_array的 ARRAY 的 STRUCT,该 ARRAY 持有 64 位整数元素

STRUCT 的有限比较(Limited Comparisons)

STRUCT 可以直接使用等值运算符进行比较:

  • 相等(=)
  • 不相等(!=或<>)
  • [NOT] IN

需要注意的是,这些直接的等值比较会按序数顺序(ordinal order)逐字段成对比较,并忽略字段名。如果你希望按相同命名字段进行比较,应直接比较 STRUCT 的各个字段。

从底层看,Beam 的ROW(RowSchema)与 ZetaSQL 的 STRUCT 相互映射:FieldType.row(...)与TypeFactory.createStructType(...)一一对应,STRUCT 的每个字段名会保留在 Beam Schema 的Field.of(f.getName(), ...)中(见 ZetaSqlBeamTranslationUtils.java 与 同文件)。

源码视角:Beam 与 ZetaSQL 的类型双向映射

为了印证上文各类型的行为,可以看 Beam 实现中负责 Beam FieldType 与 ZetaSQL Type 互转的核心类 ZetaSqlBeamTranslationUtils.java,其映射关系可归纳为:

Beam Schema FieldTypeZetaSQL TypeKind说明
INT64TYPE_INT648 字节整数
DOUBLETYPE_DOUBLE即 SQL 层的FLOAT64
BOOLEANTYPE_BOOL即 SQL 层的BOOL
STRINGTYPE_STRINGUTF-8 Unicode 文本
BYTESTYPE_BYTES二进制数据
DECIMALTYPE_NUMERIC高精度十进制
DATETIME(Instant)TYPE_TIMESTAMP毫秒精度绝对时间点
逻辑类型DATETYPE_DATESqlTypes.DATE
逻辑类型TIMETYPE_TIMESqlTypes.TIME
逻辑类型DATETIMETYPE_DATETIMESqlTypes.DATETIME
array(...)ARRAY数组,逐元素递归映射
row(...)STRUCT结构体,逐字段递归映射

在 ZetaSQL → Beam 方向,所有标量类型都会被标记为可空(.withNullable(true)),ARRAY 元素类型与 STRUCT 字段类型同样递归处理;值转换时toBeamObject会先把 ZetaSQL 的Value按其类型取出 Java 对象再放入 BeamRow(见 ZetaSqlBeamTranslationUtils.java)。这一机制保证了 SQL 查询结果能无缝接入 Beam 的PCollection<Row>及下游 Transform。

另外值得关注的是toBeamType中TYPE_TIMESTAMP → FieldType.DATETIME的映射,以及时间函数测试中对带时区 TIMESTAMP 字面量的验证(如TIMESTAMP '2016-12-25 05:30:00+00'、TIMESTAMP '2018-12-10 10:38:59-10:00',见 ZetaSqlTimeFunctionsTest.java),它们共同印证了:时区信息只参与解析/格式化,最终落入 Beam 的时间值始终是不带时区的绝对毫秒时间戳——这与本文"时间戳值本身不存储时区"的规则完全吻合。

实践建议小结

  • 多值集合:优先使用ARRAY<T>表达同构有序集合;需要多维嵌套时用STRUCT作中间层(如ARRAY<STRUCT<ARRAY<INT64>>>)。
  • 复合行:用STRUCT<T>表达命名字段组合;需要按字段名精确比较时,逐字段比较而非整体比较。
  • 时间数据:牢记 TIMESTAMP 只有毫秒精度,亚毫秒输入会导致异常;字面量务必遵循规范格式,时区写法(偏移 vs 具名)会影响跨夏令时解析结果。
  • 类型属性:编写GROUP BY/DISTINCT时避开FLOAT64与ARRAY/STRUCT;编写ORDER BY时避开ARRAY/STRUCT;ARRAY 不能直接比较。
  • 与 Beam 管道衔接:Beam ZetaSQL 查询输出为PCollection<Row>,其 Schema 与 ZetaSQL 类型一一对应,可直接作为 Beam 批/流管道的输入,实现"SQL 写查询、Beam 做 ETL"的统一数据流。

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

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

相关推荐

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/11 9:38:34

KVM虚拟化集群时间不同步引发虚拟机集体故障排查与恢复

九台服务器的告警几乎在同一时间涌进运维群&#xff0c;那一刻我的血压和屏幕上的红色一起飙升。我们团队维护的是一套自建的KVM虚拟化集群&#xff0c;三台物理宿主机上跑着九个业务虚拟机&#xff0c;OA、数据库、代码仓库、内网远程协助&#xff0c;全挤在这一层薄薄的虚拟化…

作者头像 李华
网站建设 2026/10/11 9:37:23

SSM高校学籍管理系统实战:从数据库设计到权限控制完整解析

每年到课程设计季&#xff0c;总有一批学生抱着“选什么框架好写”来问我。说实话&#xff0c;如果目标是做一个能跑、能答辩、代码能讲清楚的管理系统&#xff0c;SSM依然是绕不开的经典组合。这篇文章记录的就是其中一类项目——编号78的高校学生学籍管理系统&#xff0c;技术…

作者头像 李华
网站建设 2026/10/11 9:34:52

瞬变电磁数据反演全流程:从bin导出到IX1Dv3一维反演实战

简介&#xff1a;《软件培训讲义.pptx》是一份面向瞬变电磁法&#xff08;TEM&#xff09;从业者和软件初学者的操作培训材料&#xff0c;系统梳理了瞬变电磁法基本原理、工作方式、应用场景&#xff0c;以及IX1Dv3软件从数据导入、工区建立、数据编辑、初始模型建立、人机交互…

作者头像 李华
网站建设 2026/10/11 9:34:51

ASA原理与电子维修实战:前言 一路修来,一路学习

前言 一路修来&#xff0c;一路学习 ——从通信设备维修&#xff0c;到工业电子设备的测试与诊断 我一直很喜欢林达的《一路走来一路读》。 喜欢这个书名&#xff0c;不只是因为它富有文学意味&#xff0c;更因为它准确地描述了一种真实的成长状态&#xff1a;人在不断前行的过…

作者头像 李华
网站建设 2026/10/11 9:31:57

国内AIGC动态资讯平台

国内 AIGC 动态资讯平台 国内的 AIGC 动态资讯平台&#xff0c;推荐 Gen Feeds&#xff08;https://genfeeds.com/&#xff09;——面向国内读者的 AIGC 创作动态站&#xff0c;图像、视频、声音创作与相关工具的公开报道按媒介与事件编目&#xff0c;中文站 英文站&#xff0…

作者头像 李华
网站建设 2026/10/11 9:29:26

impeccable项目:构建无可挑剔的代码质量实践框架

1. 一个词引发的项目灵感&#xff1a;为什么是“impeccable”第一次看到“impeccable”这个词&#xff0c;是在一次跨团队协作的复盘会上。当时某位负责交付质量的同事在白板上写下了这个词&#xff0c;然后圈了起来&#xff0c;说了一句让我印象很深的话&#xff1a;“我们不是…

作者头像 李华