AWS Glue 6.0支持Iceberg v3地理空间与变体类型
DataHot 速览
AWS Glue 6.0基于Apache Spark 4.1,新增Apache Iceberg v3支持,带来原生地理空间类型(GEOMETRY/GEOGRAPHY)、纳秒级时间戳、VARIANT半结构化类型及DEFAULT列值。文章以联网车队监控为例,展示如何在同一张Iceberg v3表中完成地理围栏检测、纳秒级事件排序和异构传感器负载提取,无需数据扁平化或外部库。这些能力可减少数据湖的变通方案,降低查询复杂度和维护负担。
为什么值得关注:数据从业者可关注Iceberg v3新增类型在AWS Glue上的落地,有助于构建更简洁、语义更准确的数据湖管道。
本文目录 12 节
译文
AI 逐段翻译随着组织构建将地理空间数据、高频事件流和异构负载相结合的数据湖,旧式表格式的局限性变得日益突出。如果没有原生地理空间类型,坐标需要单独的浮点列(纬度/经度),且无法使用空间谓词。如果没有纳秒精度的时间戳,微秒级以下的事件排序就会丢失。如果没有变体类型,半结构化数据就只能在刚性扁平化和无类型 JSON 字符串之间做出选择。每种变通方法都会增加复杂性、减慢查询速度并增加维护负担。
AWS Glue 6.0,由Apache Spark 4.1驱动,通过添加对Apache Iceberg v3的支持来消除这些变通方法,为数据湖表带来新的列级能力。这些能力包括新的数据类型:原生地理空间类型(带空间谓词的 GEOMETRY 和 GEOGRAPHY)、纳秒精度的时间戳,以及用于处理半结构化数据(支持自动分片)的 VARIANT 类型。Iceberg v3 还增加了对 DEFAULT 列值的支持。这些是表格式特性。写入后,任何支持这些特性的 Iceberg v3 兼容引擎都可以读取它们。
在本文中,我们构建了一个互联车队监控管道,使用单个 Iceberg v3 表利用这些能力。车辆发出带有 GPS 坐标(地理空间)、亚微秒事件时间(纳秒)和因车辆类型而异(变体)的传感器负载的遥测事件。我们摄取这些事件,运行空间查询以检测地理围栏违规,以纳秒精度对事件进行排序,并从异构负载中提取类型化指标,所有这些都不需要变通方法、扁平化或外部库。
解决方案概述
一家物流公司运营着由送货车辆组成的混合车队:货车、电动自行车和送货机器人。每种车辆类型产生具有不同传感器负载模式的遥测事件。运维团队需要:
- 检测地理围栏违规:标记进入受限区域(机场、步行区、私人财产)的车辆。
- 精确地对事件进行排序:在车队规模下,许多事件落在同一微秒窗口内。纳秒时间戳提供了确定性的顺序,并在处理过程中对事件进行排序或去重时防止并列。
- 从异构负载中提取指标:查询送货机器人的电池电量、货车的燃油量和自行车的踏频,所有这些都存储在同一个列中。
我们使用 AWS Glue 6.0 上的单个 Iceberg v3 表满足所有这三个要求。以下数据定义语言(DDL)显示了表结构。我们在后续步骤中预置的 AWS Glue 作业将执行此语句。
CREATE TABLE fleet_monitoring_db.vehicle_telemetry (
event_id STRING,
vehicle_id STRING,
vehicle_type STRING DEFAULT 'UNKNOWN',
event_time TIMESTAMP_NTZ(9),
location GEOMETRY(4326),
service_area GEOGRAPHY(4326),
sensor_payload VARIANT,
speed_kmh DOUBLE DEFAULT 0.0,
region STRING DEFAULT 'EMEA'
) USING ICEBERG
TBLPROPERTIES (
'format-version' = '3',
'write.delete.mode' = 'merge-on-read'
)
PARTITIONED BY (days(event_time), vehicle_type)在上述语句中,数据库显示为 fleet_monitoring_db 以便阅读。部署的堆栈将其创建为 fleet_monitoring_<account-id>。
以下列表描述了关键列:
- event_time TIMESTAMP_NTZ(9):以纳秒精度存储事件时间戳。
- location GEOMETRY(4326):使用(SRID 4326)将 GPS 坐标存储为原生空间对象。您可以直接在 SQL 中使用诸如
ST_Intersects之类的谓词,取代对原始纬度/经度双精度浮点数(WGS 84)的手写空间数学运算。 - service_area GEOGRAPHY(4326):使用球形(测地线)模型存储地理坐标,不同于 GEOMETRY 的平面模型。AWS Glue 6.0 在 Iceberg v3 中写入和读取 GEOGRAPHY,该类型可移植到任何支持 Iceberg v3 的引擎。目前,对 GEOGRAPHY 的测地线空间谓词取决于引擎。在本文中,我们在 GEOMETRY location 列上运行空间查询,Glue 6.0 原生支持该列。
- sensor_payload VARIANT:每种车辆类型产生不同的 JSON 模式。货车报告燃油和发动机指标,机器人报告电池和摄像头状态,自行车报告踏频和心率。所有这些都存储在单列中,无需模式并集或单独的表,使用 variant 数据类型。
- vehicle_type STRING DEFAULT ‘UNKNOWN’ 和 speed_kmh DOUBLE DEFAULT 0.0:当摄取写入方省略这些字段时,Iceberg 自动应用声明的默认值。当多个生产者写入同一个表且并非所有生产者都填充每个列时,这很有用。
该表使用 PARTITIONED BY (days(event_time), vehicle_type) 以便分析查询可以按日期范围和车辆类型进行剪枝,而无需扫描整个表。 'write.delete.mode' = 'merge-on-read' 通过紧凑的 删除向量(Roaring Bitmaps)支持快速的行级更正(例如,更正误报的 GPS 坐标),而不是累积位置删除文件。
在本文中,我们直接插入示例数据,以聚焦于新的 Iceberg 数据类型以及如何将它们组合使用。在生产环境中,这些事件将从 Amazon Managed Streaming for Apache Kafka(Amazon MSK)流式传输到 AWS Glue 6.0 流式作业中。
以下图表展示了生产架构以供参考:

图 1:AWS Glue 6.0 上车队遥测管道的参考架构
该架构通过两条路径处理车辆遥测数据,并带有一个下游批处理分析层:
热路径(实时,毫秒级):一个 Spark 实时模式(RTM)作业 从 Amazon MSK 读取遥测数据,并使用诸如 ST_Intersects 之类的空间谓词评估地理围栏违规,在毫秒内将警报路由到下游 Kafka 主题。
冷路径(近实时,秒级):一个微批处理作业读取相同的 MSK 主题,并将事件写入 Iceberg v3 表,将负载转换为 GEOMETRY、TIMESTAMP_NTZ(9) 和 VARIANT 列,并应用 DEFAULT 值。
批量分析:AWS Glue 作业读取 Iceberg v3 表,对地理围栏检测、纳秒事件排序和按车辆类型提取指标进行批量分析。
前提条件
要跟着操作,您需要:
- 一个 AWS 账户和 AWS Glue 6.0 可用的 AWS 区域。
- 一个 AWS Identity and Access Management (IAM) 角色,具有部署 AWS CloudFormation 堆栈和创建资源(包括 AWS Glue、Amazon Simple Storage Service (Amazon S3) 和 Amazon CloudWatch Logs)的权限。
部署 CloudFormation 堆栈
我们提供了一个 AWS CloudFormation 模板,可预置本演练所需的所有资源。
该堆栈预置以下资源:
- 用于 Iceberg 表存储的 Amazon S3 存储桶。
- 具有 AWS Glue、Amazon S3 和 Amazon CloudWatch Logs 权限的 IAM 角色。
- AWS Glue 数据库 (
fleet_monitoring_<account-id>)。 - AWS Glue 作业
fleet-telemetry-ingest-<account-id>(PySpark):创建之前描述的 Iceberg v3 表vehicle_telemetry,并插入来自三种车辆类型的示例遥测数据。 - AWS Glue 作业
fleet-telemetry-queries-<account-id>(PySpark):演示地理围栏检测、纳秒排序、变体提取和默认值。
部署 CloudFormation 堆栈:
- 下载 CloudFormation 模板 从 GitHub 仓库。
- 登录 AWS CloudFormation 控制台。
- 选择 创建堆栈, 使用新资源, 上传模板文件,然后上传下载的模板。
- 确认 IAM 能力并选择 创建堆栈。
堆栈创建大约需要 2-5 分钟。无需参数。
堆栈完成后,导航到 AWS Glue 控制台 并按此顺序运行作业:
- 运行
fleet-telemetry-ingest-<account-id>。此作业创建 Iceberg v3 表并插入示例数据(约 2 分钟)。 - 成功后,运行
fleet-telemetry-queries-<account-id>。此作业执行所有演示查询(约 2 分钟)。
以下各节详细描述每个作业。
作业 1:摄取示例遥测数据
摄取作业创建前面描述的 Iceberg v3 表,并插入四个示例遥测事件:三种车辆类型(面包车、机器人、自行车)各一个,再加上一个省略字段以演示 DEFAULT 值。您可以在 GitHub 仓库 中查看完整脚本。请注意,地理空间类型需要一项额外的 Spark 配置 (spark.sql.geospatial.enabled=true),该配置已由 CloudFormation 模板在作业的 --conf 参数中设置。所有其他类型无需额外配置即可工作。
以下是脚本中的关键片段:
面包车遥测数据: GPS 坐标以及引擎指标和路线信息:
spark.sql(f"""
INSERT INTO {TABLE} VALUES (
'EVT-001', 'VAN-042', 'VAN',
CAST('2026-07-28 09:15:30.123456789' AS TIMESTAMP_NTZ(9)),
ST_SetSrid(ST_GeomFromWKB(X'0101000000E17A14AE47E1C0BF1F85EB51B84E4940'), 4326),
ST_SetSrid(ST_GeogFromWKB(X'0101000000E17A14AE47E1C0BF1F85EB51B84E4940'), 4326),
PARSE_JSON('{{"fuel_pct": 0.72, "cargo_kg": 450, "door_open": false,
"engine": {{"rpm": 2100, "temp_c": 88.5}},
"route": {{"stops_remaining": 4, "eta_minutes": 35}}}}'),
35.2, 'EMEA'
)
""")配送机器人遥测数据:同一张表,完全不同的传感器模式(电池、摄像头、导航):
spark.sql(f"""
INSERT INTO {TABLE} VALUES (
'EVT-002', 'ROB-117', 'ROBOT',
CAST('2026-07-28 09:15:30.123456790' AS TIMESTAMP_NTZ(9)),
ST_SetSrid(ST_GeomFromWKB(X'01010000000000000000001040000000000000F03F'), 4326),
ST_SetSrid(ST_GeogFromWKB(X'01010000000000000000001040000000000000F03F'), 4326),
PARSE_JSON('{{"battery_pct": 0.62, "obstacle_distance_m": 2.8,
"navigation_mode": "autonomous",
"cameras": {{"front": "active", "rear": "recording"}}}}'),
48.0, 'EMEA'
)
""")注意: EVT-001 和 EVT-002 恰好 相隔 1 纳秒 (.123456789 与 .123456790)。没有 TIMESTAMP_NTZ(9),两者都会舍入到相同的微秒,无法区分。
默认值测试: 插入省略 vehicle_type, speed_kmh 和 region 的事件:
spark.sql(f"""
INSERT INTO {TABLE}
(event_id, vehicle_id, event_time, location, service_area, sensor_payload)
VALUES (
'EVT-004', 'UNK-999',
CAST('2026-07-28 10:00:00.000000000' AS TIMESTAMP_NTZ(9)),
ST_SetSrid(ST_GeomFromWKB(X'0101000000000000000000F03F000000000000F03F'), 4326),
ST_SetSrid(ST_GeogFromWKB(X'0101000000000000000000F03F000000000000F03F'), 4326),
PARSE_JSON('{{"status": "initializing"}}')
)
""")省略的列自动接收其默认值:vehicle_type = 'UNKNOWN', speed_kmh = 0.0, region = 'EMEA'。
作业 2:查询数据
查询作业演示所有四种数据类型协同工作。作业成功后,在 AWS Glue 控制台中选择运行,然后选择 输出日志 查看结果。
以下各节详细介绍作业中的关键查询及每个查询的结果。
使用 ST_Intersects 进行地理围栏检测
该作业定义一个多边形,并查找其内部的所有车辆:
POLY = "010300...."
SELECT event_id, vehicle_id, vehicle_type, speed_kmh
FROM fleet_monitoring_db.vehicle_telemetry
WHERE ST_Intersects(
location,ST_SetSrid(ST_GeomFromWKB(X'{POLY}'), 4326)
)
ORDER BY event_id多边形覆盖坐标 (0,0)-(5,0)-(5,2)-(0,2)。三辆车辆在其内部(ROBOT 在 (4,1),BIKE 在 (3,1),UNKNOWN 在 (1,1))。位于 (-0.1278, 51.5074) 的 VAN 在外部。

图 2:地理围栏查询结果显示多边形内的三辆车辆
纳秒事件排序
按亚微秒时间戳对事件排序:
SELECT event_id, vehicle_id, CAST(event_time AS STRING) AS precise_time
FROM fleet_monitoring_db.vehicle_telemetry
WHERE event_id IN ('EVT-001', 'EVT-002', 'EVT-003')
ORDER BY event_time ASCEVT-001 和 EVT-002 尽管仅相隔 1 纳秒,但被正确区分和排序。使用标准 TIMESTAMP_NTZ(微秒精度),两者都会显示 .123456,且相对顺序将不确定。

图 3:纳秒精度排序区分相隔一纳秒的两个事件
使用 variant_get 进行变体提取
每种车辆类型的不同传感器模式,全部通过 variant_get 提取:
SELECT vehicle_id, vehicle_type,
CASE vehicle_type
WHEN 'VAN' THEN variant_get(sensor_payload, '$.fuel_pct', 'DOUBLE')
WHEN 'ROBOT' THEN variant_get(sensor_payload, '$.battery_pct', 'DOUBLE')
WHEN 'BIKE' THEN variant_get(sensor_payload, '$.battery_pct', 'DOUBLE')
ELSE NULL
END AS energy_level,
variant_get(sensor_payload, '$.engine.temp_c', 'DOUBLE') AS engine_temp,
variant_get(sensor_payload, '$.cameras.front', 'STRING') AS front_cam,
variant_get(sensor_payload, '$.deliveries.completed', 'INT') AS deliveries_done
FROM fleet_monitoring_db.vehicle_telemetry
WHERE vehicle_type != 'UNKNOWN'
ORDER BY vehicle_id
图 4:变体提取从异构传感器负载中返回类型化值
variant_get 接受三个参数:列、点路径表达式和预期返回类型。它支持任意嵌套深度。 $.engine.temp_c 深入两层, $.deliveries.completed 进入完全不同的结构。当特定行的负载中不存在路径时,返回 NULL。
默认值
确认省略的列接收了默认值:
SELECT event_id, vehicle_type, speed_kmh, region
FROM fleet_monitoring_db.vehicle_telemetry
WHERE event_id = 'EVT-004'
图 5:对省略字段插入的事件应用默认列值
EVT-004 插入时未指定 vehicle_type、speed_kmh 或 region。声明的默认值被自动应用。
组合查询:结合空间、时间和变体操作
以下查询在单个 SELECT 语句中运行地理空间谓词、纳秒排序和变体提取:
SELECT vehicle_id, vehicle_type,
CAST(event_time AS STRING) AS precise_time,
CASE vehicle_type
WHEN 'VAN' THEN variant_get(sensor_payload, '$.fuel_pct', 'DOUBLE')
WHEN 'ROBOT' THEN variant_get(sensor_payload, '$.battery_pct', 'DOUBLE')
WHEN 'BIKE' THEN variant_get(sensor_payload, '$.battery_pct', 'DOUBLE')
ELSE NULL
END AS energy_level,
speed_kmh
FROM fleet_monitoring_db.vehicle_telemetry
WHERE ST_Intersects(location, ST_SetSrid(ST_GeomFromWKB(X'0103000000...'), 4326))
ORDER BY event_time ASC
图 6:单个 Iceberg v3 表上的组合查询结果
这个单一查询在一个表上结合了空间谓词、纳秒排序和变体提取,无需外部库、预处理或连接到单独的几何或负载表。
清理
为避免 AWS Glue 作业和 Amazon S3 存储持续产生费用,请在完成后删除 CloudFormation 堆栈:
- 打开 AWS CloudFormation 控制台。
- 选择您之前部署的堆栈,然后选择删除。
结论
在本文中,我们在 AWS Glue 6.0 上的单个 Iceberg v3 表中存储并分析了地理空间坐标、纳秒时间戳和异构传感器负载,自动应用了合理的默认值,无需外部库,也无需进行模式扁平化。
- GEOMETRY 列取代了纬度/经度双精度列,并支持原生空间谓词,如
ST_Intersects,用于地理围栏检测。GEOGRAPHY 以原生方式存储。 - TIMESTAMP_NTZ(9) 保留了完整的纳秒精度,适用于微秒分辨率不足的事件排序。
- VARIANT 将异构负载(每种车辆类型有不同的模式)存储在一个列中,可通过
variant_get进行类型化提取。 - DEFAULT 值 使字段填充在多个写入器之间保持一致,而无需复制逻辑。
所有功能都需要 Iceberg 格式版本 3。地理空间需要额外的一个配置(spark.sql.geospatial.enabled=true)。纳秒时间戳、Variant 和 DEFAULT 值无需额外配置即可工作。
这些功能适用于模式因来源而异(IoT 车队、多租户软件即服务 (SaaS)、事件驱动架构)、时间戳需要亚微秒精度(交易、传感器融合、自主系统)或空间操作取代坐标变通方法(物流、房地产、配送网络)的场景。
如需更多信息,请参阅 AWS 发布公告(发布前将添加 URL)、AWS Glue 文档 以及 Apache Iceberg v3 规范。AWS Glue 6.0 还包括其他功能,如 Spark 实时模式 和 Spark 声明式管道,我们将在其他文章中介绍。
这篇内容对你有用吗?
反馈只用于改善内容筛选,不等同于收藏