返回
RSS AWS Big Data Blog AI 逐段翻译 精选 发布 2026-09-03 02:34

在Amazon SageMaker用Iceberg物化视图构建Medallion架构

DataHot 速览

本文介绍在Amazon SageMaker中用Apache Iceberg物化视图,以声明式方式构建Bronze→Silver→Gold的Medallion湖仓分层。传统方式需要分别编写ETL任务、编排DAG和CDC增量逻辑;现在每个数据层只需一条SQL定义,由系统按刷新配置自动处理物化与增量。作者认为用三条SQL即可搭起分层管道,显著降低编排代码和CDC逻辑的维护复杂性。

为什么值得关注:该文展示如何用Iceberg物化视图将湖仓分层管道的调度、ETL与增量处理收敛为声明式SQL,对数据平台工程化有直接参考价值。

本文目录 30 节
  1. 什么是奖章架构
  2. 传统方法与声明式方法
  3. 传统方法
  4. 使用 Iceberg 物化视图的声明式方法
  5. Apache Iceberg 和物化视图
  6. 支持 Iceberg 物化视图的服务
  7. 技术架构
  8. 前提条件
  9. 步骤1:初始化环境
  10. 步骤 2:将数据摄入 Bronze
  11. 步骤 3:探索 Bronze
  12. 步骤 4:创建 Silver 物化视图
  13. 步骤 5:创建 Gold 物化视图
  14. Gold 1:城市每日指标
  15. Gold 2:车辆性能
  16. 依赖链
  17. 步骤 6:查询 Gold 层
  18. 城市每日指标 Gold 表
  19. 车辆性能 Gold 表
  20. 步骤 7:数据传播演示
  21. INSERT 新记录
  22. 刷新Silver(增量)
  23. 验证新记录已传播
  24. 刷新Gold(从Silver物化视图级联)
  25. 通过MERGE更新
  26. 步骤8:清理
  27. 限制与注意事项
  28. 定价
  29. 总结
  30. 参考

译文

AI 逐段翻译

构建奖章架构今天通常意味着您必须构建三个独立系统协同工作:用于在层之间转换数据的提取、转换和加载(ETL)作业,一个编排器(如 Apache Airflow 或 AWS Step Functions)以正确顺序编排这些作业,以及自定义变更数据捕获(CDC)逻辑,以确保每个作业只处理新增或修改的记录。每个组件必须独立编写、测试、部署和维护,当其中一个出现故障时,整个管道就会停滞。

在本文中,我们展示了如何在Amazon SageMaker中使用 Apache Iceberg 物化视图,将转换、编排和增量处理整合为每层单个 SQL 定义。您声明每层应包含的内容,系统根据您的刷新配置处理刷新时间和方式。使用此方法,您可以用三条 SQL 语句构建 Bronze → Silver → Gold 管道。这降低了维护单独编排代码、CDC 逻辑和作业工件的复杂性。

什么是奖章架构

奖章架构将数据组织为三个递进层:

  • Bronze 层 – 按原样从源系统捕获原始数据,保留原始格式以实现可审计性和重放。
  • Silver 层 – 应用清洗、去重、类型转换和业务逻辑,生成经过验证、可查询的数据集。
  • Gold 层 – 将 Silver 数据聚合为业务级指标、关键绩效指标(KPI)和针对分析及报告优化的维度模型。

每层都建立在前一层的基础上,形成从原始数据摄入到业务洞察的清晰血统。

传统方法与声明式方法

两种方法在您构建和维护的基础设施数量上有所不同。

传统方法

您编写 ETL 作业(如 Apache Spark 脚本)用于 Bronze 到 Silver 层,以及另一个用于 Silver 到 Gold 层。您在 Apache Airflow 中构建有向无环图(DAG)或 Step Functions 状态机以按顺序运行它们。您实现 CDC 逻辑,如跟踪高水位线、比较快照或消费变更流,以便每个作业只处理新数据。

使用 Iceberg 物化视图的声明式方法

您为每层编写一个带有CREATE MATERIALIZED VIEW语句和SCHEDULE REFRESH EVERY N HOURS子句。AWS Glue 托管 Spark 计算执行刷新,但您无需编写、版本化或部署作业工件。Iceberg 的行级变更跟踪(位置删除和相等删除文件)识别自上次刷新以来哪些行发生了变化,AWS Glue 仅处理这些行。依赖链在 SQL 定义中是隐式的。您唯一需要维护的代码是 SQL 转换逻辑本身。

Apache Iceberg 和物化视图

Apache Iceberg是一个开源高性能表格式,专为数据湖中的 PB 级分析数据集设计。它提供 ACID 事务、时间旅行、模式演化和隐藏分区。

使用Iceberg 物化视图,您可以将奖章架构的每一层定义为 SQL 语句。在底层,AWS Glue 使用 Iceberg 的变更跟踪元数据来识别自上次刷新以来哪些行发生了变化,然后使用托管 Spark 计算仅处理这些行。您通过 SQL 定义配置调度和增量处理,系统执行原子刷新,无需您编写管道代码。

刷新时,Gold 物化视图从 Silver 物化视图增量读取,而 Silver 物化视图又从 Bronze 表读取。这创建了声明式依赖链:每层的定义指向其下一层,系统在每次刷新时解析要重新处理哪些数据。

支持 Iceberg 物化视图的服务

在发布时,以下服务支持创建和刷新 Iceberg 物化视图:

有关最新版本要求,请参阅AWS Glue 物化视图文档

技术架构

该架构使用 Amazon S3 Tables(Amazon Simple Storage Service(Amazon S3)的功能)作为存储层。Amazon S3 Tables 是托管的 Apache Iceberg 产品,减轻了维护 Iceberg 表的管理开销。AWS Glue Data Catalog 管理表元数据,Amazon SageMaker Unified Studio 提供支持 AI 的笔记本环境,使用 AWS Glue 5.1 编写和执行物化视图定义。

该图展示了构建在 Apache Iceberg 上的三层数据湖仓管道。Bronze 层包含原始行程数据(S3 Tables 上的 trips_bronze 表,字段:trip_id、city、vehicle_type、fare、status),您通过 INSERT/Append 操作摄入这些数据。

增量 REFRESH 馈送到 Silver 层,其中物化视图(mv_trips_silver)执行时间戳转换、空值过滤,并计算派生列,如 revenue_per_mile 和 rating_category。它只处理新增或更改的行。

然后,Silver 层按每日计划刷新两个 Gold 层物化视图:mv_city_daily_metrics(city、date、trips、drivers、revenue、tips)和 mv_vehicle_performance(vehicle_type、city、trips、revenue、distance)。Gold 层为下游消费者提供服务,包括 Amazon Athena、Amazon Quick Sight、Amazon Redshift 以及支持 Iceberg REST API 的第一方(1P)或第三方(3P)计算引擎。

管道流程如下:

奖章管道图:Bronze 表馈送到 Silver 物化视图,后者馈送到由分析引擎消费的两个 Gold 物化视图

图 1:从 Bronze 表经 Silver 和 Gold 物化视图到分析消费者的三层奖章管道

前提条件

开始之前,请确认您具备以下条件:

  • 具有 Amazon SageMaker Unified Studio、AWS Glue、S3 Tables 和 AWS Lake Formation 权限的 AWS 账户。
  • 一个 Amazon SageMaker Unified Studio 域。

步骤1:初始化环境

打开 AWS 管理控制台并导航至Amazon SageMaker

Amazon SageMaker 控制台登录页面

图 2:Amazon SageMaker 控制台登录页面

选择开始使用 以设置 Amazon SageMaker Unified Studio。

SageMaker Unified Studio 开始使用设置页面

图 3:设置 SageMaker Unified Studio 的开始使用页面

选择打开 以启动 Amazon SageMaker Unified Studio。

打开并启动 SageMaker Unified Studio 的按钮

图 4:打开并启动 SageMaker Unified Studio 的选项

进入 SageMaker Unified Studio 后,在左侧窗格中选择数据 以创建 S3 Tables 存储桶(Amazon S3 的受管 Apache Iceberg 功能)和数据库。选择添加,然后选择创建 S3 Tables 目录,并提供目录和数据库名称。最后,选择创建目录

创建 S3 Tables 目录对话框,包含目录和数据库名称字段

图 5:创建 S3 Tables 目录对话框,包含目录和数据库名称字段

目录创建完成后,在左侧导航窗格中选择笔记本

SageMaker Unified Studio 左侧导航窗格中的笔记本选项

图 6:SageMaker Unified Studio 导航窗格中的笔记本选项

选择创建笔记本

SageMaker Unified Studio 中的创建笔记本按钮

图 7:SageMaker Unified Studio 中的创建笔记本按钮

使用笔记本前,请选择Athena SparkGlue Spark 计算连接作为笔记本的运行时引擎。

运行时引擎选择界面,显示 Athena Spark 和 Glue Spark 计算连接

图 8:选择 Athena Spark 或 Glue Spark 作为笔记本运行时引擎

在单独的笔记本单元格中使用以下代码示例。您还可以用自然语言提供转换需求,然后SageMaker Data Agent 将为您生成 SQL 代码。

SageMaker Data Agent 根据自然语言提示生成 SQL

图 9:SageMaker Data Agent 根据自然语言请求生成 SQL

通过选择SQL 按钮将每个代码块添加到新单元格中:

笔记本工具栏中的 SQL 单元格类型按钮

图 10:用于向笔记本单元格添加代码块的 SQL 按钮

在单元格菜单中选择Athena SparkGlue Spark 作为您的计算资源。

笔记本单元格菜单中的计算连接选择

图 11:笔记本单元格菜单中的计算选择

如果在单元格执行后遇到错误,请使用数据代理聊天机器人或使用 AI 修复 按钮来解决。

用于解决单元格错误的 Fix with AI 按钮和数据代理聊天机器人

图 12:用于解决单元格执行错误的 Fix with AI 按钮

步骤 2:将数据摄入 Bronze

生成 300 条真实的共享出行行程记录,并直接插入 Bronze Iceberg 表。这模拟了原始数据摄入层。在生产环境中,您通常根据需求配置流式源或批量加载。

将以下代码复制到第一个笔记本单元格中(使用 Python 单元格类型)。

import random
from datetime import datetime, timedelta

CITIES = {
    "San Francisco": {"lat_range": (37.70, 37.82), "lon_range": (-122.52, -122.38), "surge_prob": 0.3},
    "Austin": {"lat_range": (30.22, 30.40), "lon_range": (-97.80, -97.68), "surge_prob": 0.15},
    "Chicago": {"lat_range": (41.85, 41.95), "lon_range": (-87.70, -87.60), "surge_prob": 0.2},
    "Seattle": {"lat_range": (47.55, 47.68), "lon_range": (-122.40, -122.28), "surge_prob": 0.25},
}
VEHICLE_TYPES = ["UberX", "Comfort", "XL", "Black"]
PAYMENT_METHODS = ["credit_card", "debit_card", "apple_pay", "google_pay", "cash"]
STATUSES = ["completed"] * 4 + ["cancelled_rider", "cancelled_driver"]
BASE_FARES = {"UberX": 2.50, "Comfort": 3.50, "XL": 4.00, "Black": 7.00}
PER_MILE = {"UberX": 1.75, "Comfort": 2.25, "XL": 2.50, "Black": 3.75}
PER_MIN = {"UberX": 0.35, "Comfort": 0.45, "XL": 0.50, "Black": 0.65}

rows = []
for i in range(300):
    city_name = random.choice(list(CITIES.keys()))
    city = CITIES[city_name]
    vehicle = random.choice(VEHICLE_TYPES)
    duration = random.randint(5, 45)
    distance = round(random.uniform(1.0, 20.0), 1)
    surge = round(random.uniform(1.0, 2.5), 1) if random.random() < city["surge_prob"] else 1.0
    base = BASE_FARES[vehicle]
    fare = round((base + distance * PER_MILE[vehicle] + duration * PER_MIN[vehicle]) * surge, 2)
    tip = round(fare * random.choice([0, 0, 0.1, 0.15, 0.2, 0.25]), 2)
    status = random.choice(STATUSES)
    day = random.randint(0, 2)
    hour = random.choices(range(24),
        weights=[1,1,1,1,1,2,4,8,10,8,6,5,6,5,5,5,6,8,10,8,6,4,2,1])[0]
    trip_time = datetime(2025, 12, 1) + timedelta(days=day, hours=hour, minutes=random.randint(0, 59))

    rows.append((
        f"TRIP-{i+1:06d}",
        f"DRV-{random.randint(1000, 5000)}",
        f"RDR-{random.randint(10000, 99999)}",
        city_name, vehicle,
        round(random.uniform(*city["lat_range"]), 6),
        round(random.uniform(*city["lon_range"]), 6),
        round(random.uniform(*city["lat_range"]), 6),
        round(random.uniform(*city["lon_range"]), 6),
        trip_time.isoformat(),
        (trip_time + timedelta(minutes=duration)).isoformat(),
        duration, distance, surge, base, fare, tip, round(fare + tip, 2),
        random.choice(PAYMENT_METHODS),
        random.choice([None, 3, 4, 4, 5, 5, 5]) if status == "completed" else None,
        status,
    ))

schema = ("trip_id STRING, driver_id STRING, rider_id STRING, city STRING, "
    "vehicle_type STRING, pickup_lat DOUBLE, pickup_lon DOUBLE, "
    "dropoff_lat DOUBLE, dropoff_lon DOUBLE, trip_start_time STRING, "
    "trip_end_time STRING, duration_minutes INT, distance_miles DOUBLE, "
    "surge_multiplier DOUBLE, base_fare DOUBLE, trip_fare DOUBLE, "
    "tip_amount DOUBLE, total_amount DOUBLE, payment_method STRING, "
    "rating INT, status STRING")

df = spark.createDataFrame(rows, schema)
df.writeTo("{CATALOG_NAME}.{NAMESPACE_NAME}.trips_bronze").createOrReplace()

print(f"Created Table and Inserted {len(rows)} trips into Bronze layer")

步骤 3:探索 Bronze

在 bronze 表上运行预览。输出应类似于以下截图:

图 13:Bronze 表中的原始行程记录预览

您应该看到原始的、未经处理的行程记录,带有字符串时间戳和可空字段。这正是 Silver 层需要清理的内容。

现在,通过查询 Bronze 表验证已摄入的数据,以获取基本统计信息。

SELECT COUNT(*) as total_trips, COUNT(DISTINCT city) as cities,
COUNT(DISTINCT vehicle_type) as vehicle_types,
MIN(trip_start_time) as earliest, MAX(trip_start_time) as latest
FROM ({CATALOG_NAME}.{NAMESPACE_NAME}.trips_bronze

输出应类似于以下截图:

图 14:Bronze 表统计信息,显示总行程数、不同城市和车辆类型

步骤 4:创建 Silver 物化视图

此 SQL 语句将 Silver 层定义为物化视图,该视图清理、转换 Bronze 表并从中派生新列。请注意,这仅是定义。系统在刷新时处理数据。

CREATE MATERIALIZED VIEW IF NOT EXISTS {CATALOG_NAME}.{DATABASE}.mv_trips_silver
COMMENT 'Silver layer: Cleaned trip data with proper types and derived columns'
SCHEDULE REFRESH EVERY 1 DAY
AS
SELECT
trip_id, driver_id, rider_id, city, vehicle_type,
pickup_lat, pickup_lon, dropoff_lat, dropoff_lon,
CAST(trip_start_time AS TIMESTAMP) as trip_start_timestamp,
CAST(trip_end_time AS TIMESTAMP) as trip_end_timestamp,
duration_minutes, distance_miles, surge_multiplier,
base_fare, trip_fare, tip_amount, total_amount,
payment_method, rating, status,
CASE WHEN distance_miles > 0 THEN total_amount / distance_miles ELSE 0 END as revenue_per_mile,
CASE WHEN rating >= 4 THEN 'High' WHEN rating >= 3 THEN 'Medium' ELSE 'Low' END as rating_category
FROM {CATALOG_NAME}.{DATABASE}.trips_bronze
WHERE trip_id IS NOT NULL AND driver_id IS NOT NULL AND rider_id IS NOT NULL
AND total_amount >= 0 AND distance_miles >= 0

print("Silver MV created: urbanride.mv_trips_silver")

验证 Silver 层输出:

SELECT trip_id, city, vehicle_type, total_amount,
ROUND(revenue_per_mile, 2) as rev_per_mile, rating_category
FROM {CATALOG_NAME}.{DATABASE}.mv_trips_silver LIMIT 5

请注意,Silver 层现在具有正确的时间戳、派生的收入_每_英里和评级类别:干净、类型化,并准备好进行聚合。

输出应类似于以下截图:

图 15:Silver 物化视图结果,包含类型化时间戳和派生列

步骤 5:创建 Gold 物化视图

Gold 物化视图从 Silver 物化视图增量读取。这是一个嵌套物化视图模式:在另一个物化视图之上构建的物化视图。

Gold 1:城市每日指标

使用此物化视图,您可以通过每日刷新计划按城市和日期聚合行程数据。

CREATE MATERIALIZED VIEW IF NOT EXISTS {CATALOG_NAME}.urbanride.mv_city_daily_metrics
COMMENT 'Gold layer: Daily aggregated metrics by city'
SCHEDULE REFRESH EVERY 1 DAY
AS
SELECT
city, DATE(trip_start_timestamp) as trip_date,
COUNT(*) as total_trips,
COUNT(DISTINCT driver_id) as active_drivers,
COUNT(DISTINCT rider_id) as active_riders,
SUM(total_amount) as total_revenue,
SUM(distance_miles) as total_distance,
SUM(tip_amount) as total_tips
FROM {CATALOG_NAME}.{DATABASE}.mv_trips_silver
WHERE status = 'completed'
GROUP BY city, DATE(trip_start_timestamp)

print("Gold MV created: mv_city_daily_metrics (reads from Silver MV, refreshes daily)")

Gold 2:车辆性能

使用此物化视图,您可以按车辆类型和城市聚合性能指标。

CREATE MATERIALIZED VIEW IF NOT EXISTS {CATALOG_NAME}.{DATABASE}.mv_vehicle_performance
COMMENT 'Gold layer: Vehicle type performance metrics'
SCHEDULE REFRESH EVERY 1 DAY
AS
SELECT
vehicle_type, city,
COUNT(*) as trip_count,
SUM(total_amount) as total_revenue,
SUM(distance_miles) as total_distance,
SUM(tip_amount) as total_tips
FROM {CATALOG_NAME}.{DATABASE}.mv_trips_silver
WHERE status = 'completed'
GROUP BY vehicle_type, city

print("Gold MV created: mv_vehicle_performance (reads from Silver MV, refreshes daily)")

依赖链

完整的管道依赖关系为:

trips_bronze (table)
└── mv_trips_silver (materialized view)
    ├── mv_city_daily_metrics (MV on MV, daily schedule)
    └── mv_vehicle_performance (MV on MV, daily schedule)

每一层都由一条 SQL 语句定义。无需维护 DAG,无需部署作业定义,也无需实现水位线跟踪。

步骤 6:查询 Gold 层

查询 Gold 物化视图以查看聚合的业务指标。

城市每日指标 Gold 表

SELECT city, trip_date, total_trips, active_drivers,
ROUND(total_revenue, 2) as revenue,
ROUND(total_revenue / total_trips, 2) as avg_per_trip
FROM {CATALOG_NAME}.{DATABASE}.mv_city_daily_metrics
ORDER BY trip_date DESC, revenue DESC LIMIT 15

输出应类似于以下截图:

图 16:来自 Gold 物化视图的城市每日指标

车辆性能 Gold 表

SELECT vehicle_type, city, trip_count,
ROUND(total_revenue, 2) as revenue,
ROUND(total_revenue / trip_count, 2) as avg_per_trip
FROM {CATALOG_NAME}.{DATABASE}.mv_vehicle_performance
ORDER BY revenue DESC

输出应类似于以下截图:

图 17:来自 Gold 物化视图的车辆性能指标

Gold 层为您提供预聚合、业务就绪的指标,无需编写聚合作业。

步骤 7:数据传播演示

本节演示如何使用 INSERT、UPDATE (MERGE) 和 DELETE 操作进行更改,然后进行增量刷新,从而在层间传播更改。在生产环境中,计划刷新会自动处理此操作。我们在此手动触发以进行演示。

INSERT 新记录

将新的行程记录插入 Bronze 表。

INSERT INTO {CATALOG_NAME}.{DATABASE}.trips_bronze VALUES
('DEMO_TRIP_001', 'DRIVER_999', 'RIDER_888', 'Seattle', 'UberX',
47.6062, -122.3321, 47.6205, -122.3493,
'2024-12-15 14:30:00', '2024-12-15 14:50:00',
20, 5.2, 1.0, 10.0, 15.0, 3.0, 18.0, 'credit_card', 5, 'completed'),
('DEMO_TRIP_002', 'DRIVER_888', 'RIDER_777', 'Seattle', 'XL',
47.6101, -122.3300, 47.6550, -122.3080,
'2024-12-15 15:00:00', '2024-12-15 15:35:00',
35, 8.5, 1.5, 15.0, 30.0, 5.0, 35.0, 'cash', 4, 'completed'),
('DEMO_TRIP_003', 'DRIVER_777', 'RIDER_666', Portland, 'Comfort',
30.2672, -97.7431, 30.2800, -97.7400,
'2024-12-15 16:00:00', '2024-12-15 16:15:00',
15, 3.0, 1.0, 8.0, 12.0, 2.0, 14.0, 'credit_card', 5, 'completed')

print("Inserted 3 new trips into Bronze")

刷新Silver(增量)

刷新Silver物化视图。Iceberg物化视图仅处理三条新记录。

REFRESH MATERIALIZED VIEW {CATALOG_NAME}.{DATABASE}.mv_trips_silver"

验证新记录已传播

SELECT trip_id, city, total_amount, ROUND(revenue_per_mile, 2) as rev_per_mile, rating_category
FROM {CATALOG_NAME}.{DATABASE}.mv_trips_silver
WHERE trip_id LIKE 'DEMO_TRIP_%' ORDER BY trip_id

输出应类似以下截图:

图18:显示三条新插入演示行程的Silver物化视图

刷新Gold(从Silver物化视图级联)

刷新Gold物化视图。它从刷新后的Silver物化视图读取,仅处理增量变化。

REFRESH MATERIALIZED VIEW {CATALOG_NAME}.{DATABASE}.mv_city_daily_metrics

验证Gold层反映新行程

SELECT city, trip_date, total_trips, ROUND(total_revenue, 2) as revenue
FROM {CATALOG_NAME}.{DATABASE}.mv_city_daily_metrics
WHERE trip_date = '2024-12-15' ORDER BY city

输出应类似以下截图:

图19:反映2024-12-15新行程的城市每日指标

通过MERGE更新

使用MERGE更新Bronze中的现有记录,然后增量刷新。

MERGE INTO {CATALOG_NAME}.{DATABASE}.trips_bronze AS target
USING (SELECT 'DEMO_TRIP_002' as trip_id, 5 as new_rating, 20.0 as new_tip) AS source
ON target.trip_id = source.trip_id
WHEN MATCHED THEN UPDATE SET
target.rating = source.new_rating,
target.tip_amount = source.new_tip,
target.total_amount = target.trip_fare + source.new_tip

刷新Silver并验证

REFRESH MATERIALIZED VIEW {CATALOG_NAME}.{DATABASE}.mv_trips_silver")

SELECT trip_id, rating, rating_category, tip_amount, total_amount,
ROUND(revenue_per_mile, 2) as rev_per_mile
FROM {CATALOG_NAME}.{DATABASE}.mv_trips_silver WHERE trip_id = 'DEMO_TRIP_002'

print("UPDATE propagated: rating 4->5, tip $5->$20, total $35->$50")

输出应类似以下截图:

图20:显示演示行程更新评分和小费的Silver物化视图

步骤8:清理

删除物化视图、表、命名空间,并删除S3 Tables存储桶以完全清理资源。

# Drop MVs (Gold first, then Silver, due to dependency order)
spark.sql(f"DROP MATERIALIZED VIEW IF EXISTS {CATALOG_NAME}.{DATABASE}.mv_city_daily_metrics")
spark.sql(f"DROP MATERIALIZED VIEW IF EXISTS {CATALOG_NAME}.{DATABASE}.mv_vehicle_performance")
spark.sql(f"DROP MATERIALIZED VIEW IF EXISTS {CATALOG_NAME}.{DATABASE}.mv_trips_silver")
print("All materialized views dropped")

# Drop base table
spark.sql(f"DROP TABLE IF EXISTS {CATALOG_NAME}.{DATABASE}.trips_bronze")
print("Base table dropped")

# Drop the namespace
spark.sql(f"DROP NAMESPACE IF EXISTS {CATALOG_NAME}.{DATABASE} ")
print("Namespace dropped")

# Delete the S3 table bucket
import boto3
s3tables_client = boto3.client("s3tables")

# List and delete all remaining tables in the bucket
tables_response = s3tables_client.list_tables(
    tableBucketARN=TABLE_BUCKET_ARN, namespace="{DATABASE}"
)
for table in tables_response.get("tables", []):
    s3tables_client.delete_table(
        tableBucketARN=TABLE_BUCKET_ARN, namespace="{DATABASE}", name=table['name']
    )
    print(f" Deleted table: {table['name']}")

# Delete the namespace and bucket
s3tables_client.delete_namespace(tableBucketARN=TABLE_BUCKET_ARN, namespace="urbanride")
s3tables_client.delete_table_bucket(tableBucketARN=TABLE_BUCKET_ARN)
print(f"S3 table bucket deleted: {TABLE_BUCKET_NAME}")

限制与注意事项

尽管物化视图移除了大部分编排代码,请注意以下事项:

  1. 无低于小时级别的刷新频率。最小调度粒度为一小时(SCHEDULE REFRESH EVERY 1 HOUR)。
  2. 级联刷新并非自动。刷新Silver不会在相同操作中触发Gold。每层依据自身调度刷新,或必须按顺序触发。
  3. 删除需要FULL刷新。供Silver层使用的增量REFRESH可通过Iceberg元数据检测插入和更新,但无法检测行删除。当需要删除传播时,使用REFRESH ... FULL
  4. 仅SQL子集。某些窗口函数、用户定义函数(UDF)及复杂表达式可能不受物化视图定义支持。
  5. 模式演化需要重新创建。如果源模式变化影响物化视图定义,则必须删除并重新创建。
  6. AWS特定扩展。Iceberg物化视图不属于开源Apache Iceberg规范。它们不可移植到非AWS环境。

定价

AWS按每DPU小时0.44美元计费物化视图自动刷新(4 vCPU,16 GB内存),按秒计费,最少1分钟。当您配置调度刷新时,AWS Glue Data Catalog使用托管的Spark计算增量更新物化视图。您仅为每次刷新运行的计算时间付费。

在Data Catalog中存储物化视图元数据没有单独费用(涵盖在标准目录定价下:前一百万个对象无额外费用,之后每10万个对象/月1.00美元)。物化视图数据本身作为Iceberg文件存储在S3 Tables或Amazon S3中,按标准Amazon S3存储费率计费。

从Spark手动触发的刷新(通过Amazon Athena、Amazon EMR或AWS Glue notebook)按这些服务各自的计算定价计费,而非物化视图自动刷新费率。最新定价详情,请参阅AWS Glue定价页面。

本教程的估计成本:使用300条记录完整运行所有步骤通常消耗总计低于0.5个DPU小时(AWS Glue计算约0.22美元,外加可忽略的Amazon S3存储费用)。

总结

在本文中,您使用三个SQL语句构建了采用嵌套物化视图的金字塔式Bronze→Silver→Gold架构,无需任何编排代码。完整管道创建耗时不到2分钟,增量刷新仅处理更改的数据,无需水印、DAG或CDC管道。

要开始处理您自己的数据,请创建Amazon SageMaker Unified Studio项目,定义您的Bronze表,并将您的转换逻辑表达为Iceberg物化视图。更多信息,请参阅AWS Glue开发人员指南中的Apache Iceberg物化视图文档。

参考

Using materialized views with AWS Glue

Query AWS Glue Data Catalog materialized views

Using materialized views with Amazon EMR

Working with Amazon S3 Tables and table buckets

这篇内容对你有用吗?

反馈只用于改善内容筛选,不等同于收藏

分享这条资讯
分享海报
保存图片
iOS 也可以长按图片保存