【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
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:
| 名称 | 存储大小 | 取值范围 |
|---|---|---|
INT64 | 8 字节 | -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
浮点值是带小数部分的近似数值:
| 名称 | 存储大小 | 说明 |
|---|---|---|
FLOAT64 | 8 字节 | 双精度(近似)十进制值 |
注意:浮点类型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 有两条关键语义规则:
- UTF-8 编码强制:输入 STRING 值必须是 UTF-8 编码,输出 STRING 值也会是 UTF-8 编码。像 CESU-8、Modified UTF-8 这类替代编码不会被当作合法 UTF-8 处理。
- 按 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 FieldType | ZetaSQL TypeKind | 说明 |
|---|---|---|
INT64 | TYPE_INT64 | 8 字节整数 |
DOUBLE | TYPE_DOUBLE | 即 SQL 层的FLOAT64 |
BOOLEAN | TYPE_BOOL | 即 SQL 层的BOOL |
STRING | TYPE_STRING | UTF-8 Unicode 文本 |
BYTES | TYPE_BYTES | 二进制数据 |
DECIMAL | TYPE_NUMERIC | 高精度十进制 |
DATETIME(Instant) | TYPE_TIMESTAMP | 毫秒精度绝对时间点 |
逻辑类型DATE | TYPE_DATE | SqlTypes.DATE |
逻辑类型TIME | TYPE_TIME | SqlTypes.TIME |
逻辑类型DATETIME | TYPE_DATETIME | SqlTypes.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.
相关推荐
Apache Beam ZetaSQL 数据类型完整指南:标量类型、ARRAY、STRUCT 与比较规则详解
Apache Beam ZetaSQL 数据类型完整指南:标量类型、ARRAY、STRUCT 与比较规则详解 Apache Beam 的 ZetaSQL 方言提
大数据批处理流处理数据工程Apache Beam Calcite SQL 数据类型完全指南:标量、复合类型与 Java 运行时映射
Apache Beam Calcite SQL 数据类型完全指南:标量、复合类型与 Java 运行时映射 Beam Calcite SQL 是 Apache B
大数据批处理流处理数据工程Daft SQL 数据类型完全指南:SQL 类型到 DataType 的映射关系与底层解析原理
Daft SQL 数据类型完全指南:SQL 类型到 DataType 的映射关系与底层解析原理 Daft 作为面向 AI 与多模态负载的高性能数据引擎,其 SQ
大数据数据分析数据工程AI 应用
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考