存储

AWS S3

师成师成· 更新于 2026-09-28· 阅读 89 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

Iceberg AWS 集成

Iceberg 通过 iceberg-aws 模块与多种 AWS 服务进行集成。本节介绍如何在 AWS 上使用 Iceberg。

启用 AWS 集成

从 0.11.0 起的所有版本中,iceberg-aws 模块已随 Spark 和 Flink 引擎运行时一同打包。但 AWS 客户端并未打包在内,以便你可以使用与应用程序相同的客户端版本。你需要提供 AWS v2 SDK,因为 Iceberg 依赖于它。你可以选择使用 AWS SDK bundle,也可以使用各个独立的 AWS 客户端包(Glue、S3、DynamoDB、KMS、STS),以尽可能减小依赖体积。

所有默认的 AWS 客户端均使用 Apache HTTP Client 进行 HTTP 连接管理。该依赖不属于 AWS SDK bundle 的一部分,需要单独添加。若要选用其他 HTTP 客户端库,例如 URL Connection HTTP Client,请参阅客户端定制一节了解详细信息。

AWS 模块的所有功能都可以通过自定义 catalog 属性加载,具体加载自定义 catalog 的方式请参阅各引擎的文档。以下是一些示例。

Spark

例如,要在 Spark 3.4(使用 Scala 2.12)中使用 AWS 功能(AWS 客户端打包在 iceberg-aws-bundle 中),可以这样启动 Spark SQL shell:

# start Spark SQL client shell
spark-sql --packages org.apache.iceberg:iceberg-spark-runtime-3.4_2.12:1.11.0,org.apache.iceberg:iceberg-aws-bundle:1.11.0 \
    --conf spark.sql.defaultCatalog=my_catalog \
    --conf spark.sql.catalog.my_catalog=org.apache.iceberg.spark.SparkCatalog \
    --conf spark.sql.catalog.my_catalog.warehouse=s3://my-bucket/my/key/prefix \
    --conf spark.sql.catalog.my_catalog.type=glue \
    --conf spark.sql.catalog.my_catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO

如你所见,在 shell 命令中,我们使用 --packages 指定包含所有相关 AWS 依赖的附加包 iceberg-aws-bundle。

Flink

要在 Flink 中使用 AWS 模块,你可以下载必要的依赖,并在启动 Flink SQL 客户端时指定它们:

# download Iceberg dependency
ICEBERG_VERSION=1.11.0
MAVEN_URL=https://repo1.maven.org/maven2
ICEBERG_MAVEN_URL=$MAVEN_URL/org/apache/iceberg

wget $ICEBERG_MAVEN_URL/iceberg-flink-runtime/$ICEBERG_VERSION/iceberg-flink-runtime-$ICEBERG_VERSION.jar

wget $ICEBERG_MAVEN_URL/iceberg-aws-bundle/$ICEBERG_VERSION/iceberg-aws-bundle-$ICEBERG_VERSION.jar

# start Flink SQL client shell
/path/to/bin/sql-client.sh embedded \
    -j iceberg-flink-runtime-$ICEBERG_VERSION.jar \
    -j iceberg-aws-bundle-$ICEBERG_VERSION.jar \
    shell

有了这些依赖,你就可以创建如下 Flink catalog:

CREATE CATALOG my_catalog WITH (
  'type'='iceberg',
  'warehouse'='s3://my-bucket/my/key/prefix',
  'catalog-type'='glue',
  'io-impl'='org.apache.iceberg.aws.s3.S3FileIO'
);

你也可以在 sql-client-defaults.yaml 中指定目录配置,以便预加载它:

catalogs:
  - name: my_catalog
    type: iceberg
    warehouse: s3://my-bucket/my/key/prefix
    catalog-type: glue
    io-impl: org.apache.iceberg.aws.s3.S3FileIO

Hive

要在 Hive 中使用 AWS 模块,您可以像 Flink 示例中那样下载必要的依赖项,然后将它们添加到 Hive 类路径中,或在 CLI 运行时添加 JAR 包:

add jar /my/path/to/iceberg-hive-runtime.jar;
add jar /my/path/to/aws/bundle.jar;

有了这些依赖,你就可以在命令行中于运行时注册 Glue Catalog 并在 Hive 中创建外部表:

SET iceberg.engine.hive.enabled=true;
SET hive.vectorized.execution.enabled=false;
SET iceberg.catalog.glue.type=glue;
SET iceberg.catalog.glue.warehouse=s3://my-bucket/my/key/prefix;

-- suppose you have an Iceberg table database_a.table_a created by GlueCatalog
CREATE EXTERNAL TABLE database_a.table_a
STORED BY 'org.apache.iceberg.mr.hive.HiveIcebergStorageHandler'
TBLPROPERTIES ('iceberg.catalog'='glue');

你也可以通过在 hive-site.xml 中设置上述配置来预加载 catalog。

目录(Catalog)

用户可以使用多种不同的方式基于 AWS 构建 Iceberg catalog。

Glue Catalog

Iceberg 支持使用 AWS Glue 作为 Catalog 的实现。使用时,Iceberg 的 namespace 会存储为一个 Glue Database,Iceberg 表会存储为一个 Glue Table,而每个 Iceberg 表版本都会存储为一个 Glue TableVersion。你可以通过将 catalog-impl 指定为 org.apache.iceberg.aws.glue.GlueCatalog,或者像上文启用 AWS 集成一节中所示那样将 catalog-type 设置为 glue,来开始使用 Glue catalog。关于加载 catalog 的更多细节,请参阅各个引擎的页面,例如 Spark 和 Flink。

Glue Catalog ID

每个 AWS 账户和每个 AWS 区域都有一个唯一的 Glue 元数据存储。默认情况下,GlueCatalog 会根据用户的默认 AWS 客户端凭证和区域设置来选择要使用的 Glue 元数据存储。你可以通过 glue.id catalog 属性指定 Glue catalog ID,以指向另一个 AWS 账户中的 Glue catalog。Glue catalog ID 即你的 AWS 账户数字 ID。如果 Glue catalog 位于不同的区域,你应当配置 AWS 客户端指向正确的区域,详见 AWS 客户端自定义。

跳过归档(Skip Archive)

AWS Glue 具有归档较旧表版本的能力,如有需要,用户可以将表回滚到任意历史版本。默认情况下,Iceberg Glue Catalog 会跳过对较旧表版本的归档。如果用户希望归档较旧的表版本,可以将 glue.skip-archive 设置为 false。请注意,对于向 Iceberg 表进行流式摄入的场景,将 glue.skip-archive 设置为 false 会快速生成大量的 Glue 表版本。更多细节请参阅 Glue 配额 和 UpdateTable API。

跳过名称校验(Skip Name Validation)

允许用户跳过表名和命名空间的名称校验。建议遵循 Glue 最佳实践,以确保操作与 Hive 兼容。此选项仅为那些已有使用非标准字符的既定约定的用户而添加。当数据库名称和表名校验被跳过时,无法保证所有下游系统都能支持这些名称。

乐观锁

默认情况下,Iceberg 使用 Glue 的乐观锁来处理对表的并发更新。在乐观锁机制下,每个表都有一个版本 ID。当用户获取表元数据时,Iceberg 会记录该表的版本 ID。只要服务端的版本 ID 保持不变,用户就可以更新该表。如果在你之前有其他人修改了该表,就会出现版本不匹配,导致更新失败。随后 Iceberg 会刷新元数据并检查是否存在冲突。如果没有提交冲突,操作将被重试。乐观锁保证了 Iceberg 表在 Glue 中的原子事务,同时也防止他人意外覆盖你的更改。

Info

请使用 AWS SDK 版本 >= 2.17.131 以利用 Glue 的乐观锁。如果 AWS SDK 版本低于 2.17.131,则仅使用内存锁。为确保原子事务,你需要配置 DynamoDb Lock Manager。

数据仓库位置

与其他所有 catalog 实现类似,warehouse 是一个必需的 catalog 属性,用于确定存储中数据仓库的根路径。默认情况下,由于使用了 S3FileIO,Glue 仅允许将数据仓库位置设置在 S3 上。若要将数据存储在其他本地或云存储中,可以通过设置 io-impl catalog 属性,将 Glue catalog 切换为使用 HadoopFileIO 或任何自定义 FileIO。有关此功能的详细信息,请参阅 自定义 FileIO 部分。

表位置

默认情况下,命名空间 my_ns 中的表 my_table 的根位置位于 my-warehouse-location/my-ns.db/my-table。此默认根位置可以在命名空间和表两个层级上进行更改。

若要为某个命名空间下的所有表使用不同的路径前缀,请使用 AWS 控制台或你喜欢的任意 AWS Glue 客户端 SDK 来更新对应 Glue 数据库的 locationUri 属性。例如,你可以将 my_ns 的 locationUri 更新为 s3://my-ns-bucket,那么之后任何新创建的表的默认根位置都会位于该新前缀之下。比如,新表 my_table_2 的根位置将位于 s3://my-ns-bucket/my_table_2。

若要为某张特定的表使用完全不同的根路径,请将 location 表属性设置为你期望的根路径值。例如,在 Spark SQL 中可以这样写:

CREATE TABLE my_catalog.my_ns.my_table (
    id bigint,
    data string,
    category string)
USING iceberg
OPTIONS ('location'='s3://my-special-table-bucket')
PARTITIONED BY (category);

对于支持 LOCATION 关键字的引擎(如 Spark),上面的 SQL 语句等价于:

CREATE TABLE my_catalog.my_ns.my_table (
    id bigint,
    data string,
    category string)
USING iceberg
LOCATION 's3://my-special-table-bucket'
PARTITIONED BY (category);

DynamoDB Catalog

Iceberg 支持使用 DynamoDB 表来记录和管理数据库与表的信息。

配置项

DynamoDB catalog 支持以下配置:

属性默认值说明
dynamodb.table-nameicebergDynamoDbCatalog 所使用的 DynamoDB 表名

内部表设计

DynamoDB 表按以下列进行设计:

列名键类型说明
identifier分区键string表标识符,例如 db1.table1;对于命名空间,则使用字符串 NAMESPACE
namespace排序键string命名空间名称。以 namespace 作为分区键、identifier 作为排序键创建了一个全局二级索引(GSI),不投影其他列
vstring行版本,用于乐观锁
updated_atnumber最后一次更新的时间戳(毫秒)
created_atnumber创建时间戳(毫秒)
p.<property_key>stringIceberg 定义的表属性,包括 table_type、metadata_location 和 previous_metadata_location,或命名空间属性

该设计具有以下优势:

  1. 这可以避免当同一 namespace 下的表存在大量写入流量时,因分区键位于表级别而可能产生的热点分区问题
  2. namespace 操作被聚集在单个分区中,以避免影响表的提交操作
  3. 列表操作使用排序键到分区键的反向 GSI,其余所有操作均为单行操作或单分区查询。目录中的任何操作都无需进行全表扫描
  4. 使用字符串类型的 UUID 版本字段 v 而非 updated_at,以避免两个进程在同一毫秒内提交
  5. catalog.renameTable 使用多行事务来确保幂等性
  6. 属性被展平为顶层列,这样用户可以在任意属性字段上创建自定义 GSI,从而自定义目录功能。例如,用户可以将 owner 信息存储为表属性 owner,并通过在 p.owner 列上添加 GSI 来按 owner 搜索表

RDS JDBC Catalog

Iceberg 还支持 JDBC Catalog,它使用关系数据库中的表来管理 Iceberg 表。你可以配置 JDBC Catalog 与 AWS RDS 等关系数据库服务配合使用。请参阅JDBC 集成页面了解使用 JDBC Catalog 的指南和示例。如需了解配置 JDBC Catalog 并使用 IAM 身份验证的更多细节,请参阅此 AWS 文档。

如何选择 Catalog?

在众多可用选项中,以下是在为应用选择合适的 Catalog 时的指导建议:

  1. 如果你的组织已有 Glue 元数据存储,或者计划使用包括 Glue、Athena、EMR、Redshift 和 LakeFormation 在内的 AWS 分析生态系统,Glue 目录能提供最简便的集成方式。
  2. 如果你的应用需要频繁更新表,或者有很高的读写吞吐量(例如流式写入),Glue 和 DynamoDB 目录可以通过乐观锁提供最佳性能。
  3. 如果你希望对目录中的表实施访问控制,Glue 表可以作为 IAM 资源进行管理,而 DynamoDB 目录中的表只能通过条目级权限来管理,后者要复杂得多。
  4. 如果你希望根据表属性信息查询表而无需扫描整个目录,DynamoDB 目录允许你为任意属性字段构建二级索引,并提供高效的查询性能。
  5. 如果你既想利用 DynamoDB 目录的优势,同时又需要连接 Glue,可以启用 DynamoDB 流与 Lambda 触发器,异步地将 DynamoDB 目录中的表信息更新到 Glue 元数据存储中。
  6. 如果你的组织已经在 RDS 中维护了关系型数据库,或者使用无服务器版 Aurora 来管理表,JDBC 目录能提供最简便的集成方式。

DynamoDb 锁管理器

HadoopCatalog 或 HadoopTables 可以使用 Amazon DynamoDB,使得每次提交时,目录先通过一张辅助 DynamoDB 表获取锁,然后再安全地修改 Iceberg 表。对于基于文件系统的目录而言,这是必需的,因为 S3 等存储不提供文件写入互斥,需要以此保证事务的原子性。

此功能需要设置以下与锁相关的目录属性:

  1. 将 lock-impl 设置为 org.apache.iceberg.aws.dynamodb.DynamoDbLockManager。
  2. 将 lock.table 设置为你希望使用的 DynamoDB 表名。如果 DynamoDB 中不存在该名称的锁表,则会创建一张新表,并将计费模式设为按请求付费。

还可以使用其他与锁相关的目录属性来调整锁的行为,例如心跳间隔。更多详情请参阅锁目录属性。

S3 FileIO

Iceberg 允许用户通过 S3FileIO 将数据写入 S3。GlueCatalog 默认使用该 FileIO,其他 catalog 则可以通过 io-impl catalog 属性加载此 FileIO。

渐进式分段上传

S3FileIO 实现了一种自定义的渐进式分段上传算法来上传数据。数据文件按分段并行上传,每一段一旦准备就绪就立即上传,且该文件分段的上传一旦完成就立即删除。这样可以在上传过程中获得最大的上传速度,同时使本地磁盘占用降至最低。以下是用户可以针对该功能调整的配置项:

属性默认值说明
s3.multipart.num-threads系统中可用的处理器数量用于将分段上传到 S3 的线程数(在所有输出流之间共享)
s3.multipart.part-size-bytes32MB分段上传请求中单个分段的大小
s3.multipart.threshold1.5以分段大小的倍数表示的阈值,超过该阈值时,将由单次 put object 请求上传切换为分段上传
s3.staging-dirjava.io.tmpdir 属性的值用于存放临时文件的目录

S3 服务端加密

S3FileIO 支持全部 3 种 S3 服务端加密模式:

  • SSE-S3:当您使用由 Amazon S3 托管密钥的服务端加密(SSE-S3)时,每个对象都会使用唯一密钥进行加密。作为额外的保护措施,它还会使用一个主密钥对密钥本身进行加密,且该主密钥会定期轮换。Amazon S3 服务端加密使用目前可用的最强分组密码之一——256 位高级加密标准(AES-256)——来加密您的数据。
  • SSE-KMS:使用存储在 AWS 密钥管理服务(AWS KMS)中的客户主密钥(CMK)进行服务端加密(SSE-KMS)与 SSE-S3 类似,但使用该服务需要额外的收益和费用。针对 CMK 的使用设有独立的权限控制,可为 Amazon S3 中的对象提供额外的未授权访问防护。SSE-KMS 还会为您提供审计追踪,显示您的 CMK 何时被使用以及由谁使用。此外,您可以创建和管理客户托管的 CMK,也可以使用专属于您、您的服务和您的区域的 AWS 托管 CMK。
  • DSSE-KMS:使用 AWS 密钥管理服务密钥的双层服务端加密(DSSE-KMS)与 SSE-KMS 类似,但它在对象上传到 Amazon S3 时会应用两层加密。DSSE-KMS 可用于满足那些要求对数据应用多层加密并能完全控制加密密钥的合规标准。
  • SSE-C:使用客户提供密钥的服务端加密(SSE-C)时,加密密钥由您管理,而 Amazon S3 负责在写入磁盘时进行加密,并在您访问对象时进行解密。

要启用服务端加密,请使用以下配置属性:

属性默认值说明
s3.sse.typenonenone、s3、kms、dsse-kms 或 custom
s3.sse.keykms 和 dsse-kms 类型为 aws/s3,其他类型为 null对于 kms 和 dsse-kms 类型,为 KMS 密钥 ID 或 ARN;对于 custom 类型,为自定义的 base-64 AES256 对称密钥。
s3.sse.md5null若 SSE 类型为 custom,则必须将此值设置为对称密钥的 base-64 MD5 摘要,以确保完整性。

S3 访问控制列表

S3FileIO 支持 S3 访问控制列表(ACL),以实现更细粒度的访问控制。用户可以通过设置 s3.acl 属性来选择 ACL 级别。有关更多详细信息,请参阅 S3 ACL 文档。

对象存储文件布局

S3 以及许多其他云存储服务会基于对象前缀对请求进行限流。使用传统 Hive 存储布局存储在 S3 中的数据可能会遇到 S3 请求限流,因为这些对象都存储在相同的文件路径前缀之下。

Iceberg 默认使用 Hive 存储布局,但可以切换为使用 ObjectStoreLocationProvider。使用 ObjectStoreLocationProvider 时,会为每个存储的文件生成一个确定性的哈希值,并将该哈希值直接追加在 write.data.path 之后。这确保了写入 S3 的文件能够均匀分布到 S3 存储桶中的多个前缀之下,从而最大限度地减少限流并最大化与 S3 相关的 IO 操作吞吐量。使用 ObjectStoreLocationProvider 时,让所有 Iceberg 表共享同一个 write.data.path 可以提升性能。

有关 S3 如何扩展 API QPS 的更多信息,请参阅 2018 年 re:Invent 大会的 Amazon S3 和 Amazon S3 Glacier 最佳实践主题演讲。在 53:39 处介绍了 S3 如何进行扩展/分区,在 54:50 处则讨论了在新分区创建之前需要等待的 30 到 60 分钟时间。

要使用 ObjectStorageLocationProvider,请在表属性中添加 'write.object-storage.enabled'=true。下面是使用 ObjectStorageLocationProvider 创建表的 Spark SQL 命令示例:

CREATE TABLE my_catalog.my_ns.my_table (
    id bigint,
    data string,
    category string)
USING iceberg
OPTIONS (
    'write.object-storage.enabled'=true,
    'write.data.path'='s3://my-table-data-bucket/my_table')
PARTITIONED BY (category);

随后,我们就可以向这张新表中插入一行数据。

INSERT INTO my_catalog.my_ns.my_table VALUES (1, "Pizza", "orders");

这样会将数据写入 S3,并在 write.object-storage.path 之后直接附加一个 20 位的 base2 哈希(01010110100110110010),确保对表的读取均匀分布到各个 S3 存储桶前缀上,从而提升性能。之前提供的 base64 哈希已更新为 base2,以便在 S3 通用目的存储桶上获得更好的自动扩缩容行为。

作为此次更新的一部分,我们还将熵值划分到多个目录中,以提高 Iceberg 孤儿文件清理过程的效率,因为目录被用作在各个 worker 之间分配工作的方式,从而加快遍历速度。从下面的示例中可以看到,我们将哈希拆分为深度为 3、每个 4 位的目录,并把哈希的最后一部分附加在末尾。

s3://my-table-data-bucket/my_ns.db/my_table/0101/0110/1001/10110010/category=orders/00000-0-5affc076-96a4-48f2-9cd2-d5efbc9f0c94-00001.parquet

请注意,ObjectStoreLocationProvider 的路径解析逻辑是先取 write.data.path,若未设置则回退到 <tableLocation>/data。

而对于 0.12.0 及之前的旧版本,逻辑如下:

  • 0.12.0 之前,必须设置 write.object-storage.path。
  • 0.12.0 版本中,依次取 write.object-storage.path、write.folder-storage.path,最后回退到 <tableLocation>/data。
  • 2.0.0 版本将移除 write.object-storage.path 和 write.folder-storage.path。

有关更多详细信息,请参阅 LocationProvider 配置章节。

我们还新增了一个表属性 write.object-storage.partitioned-paths,若将其设置为 false(默认为 true),则会从文件路径中省略分区值。Iceberg 在文件路径中并不需要这些值,将其设为 false 可以进一步减小 key 的大小。在这种情况下,我们还会将最后 8 位熵值直接附加到文件名中。设置该配置后,插入的 key 形如下面所示,注意 category=orders 已被移除:

s3://my-table-data-bucket/my_ns.db/my_table/1101/0100/1011/00111010-00000-0-5affc076-96a4-48f2-9cd2-d5efbc9f0c94-00001.parquet

S3 重试

遇到 S3 限流的工作负载应持续进行带指数退避的重试,以便在 S3 自动扩容的同时取得进展。为此我们提供以下配置项来调整 S3 重试行为。对于遇到限流并因重试次数耗尽而失败的工作负载,建议将重试次数设置为 32,以便让 S3 自动扩容。注意,对于吞吐量极高、而 S3 尚未对相应表完成扩容的工作负载,可能需要进一步提高重试次数。

属性默认值说明
s3.retry.num-retries5重试 S3 操作的次数。对于高吞吐量工作负载,建议设置为 32。
s3.retry.min-wait-ms2s重试 S3 操作的最小等待时间。
s3.retry.max-wait-ms20s重试 S3 读取操作的最大等待时间。

S3 强一致性

2020 年 11 月,S3 宣布所有读取操作均支持强一致性,Iceberg 也已更新以充分利用这一特性。IO 操作期间不再有冗余的一致性等待和检查,从而避免对性能造成负面影响。

Hadoop S3A FileSystem

重要

对于 S3 使用场景,推荐使用 S3FileIO,而非 S3A FileSystem(HadoopFileIO)。

在 S3FileIO 推出之前,许多 Iceberg 用户选择使用 HadoopFileIO,通过 S3A FileSystem 将数据写入 S3。如前几节所述,S3FileIO 采用了最新的 AWS 客户端和 S3 特性,在安全性和性能方面均得到优化。

S3FileIO 使用 s3:// URI scheme 写入数据,同时也兼容 S3A FileSystem 所使用的 scheme。这意味着,对于任何包含 s3a:// 或 s3n:// 文件路径的表 manifest,S3FileIO 依然能够读取它们。这一特性让人们可以轻松地从 S3A 切换到 S3FileIO。

如果出于某些原因你必须使用 S3A,请按照以下步骤操作:

  1. 要使用 S3A 存储数据,请将 warehouse 目录属性指定为 S3A 路径,例如 s3a://my-bucket/my-warehouse
  2. 对于 HiveCatalog,若也想使用 S3A 存储元数据,请将 Hadoop 配置属性 hive.metastore.warehouse.dir 指定为 S3A 路径。
  3. 将 hadoop-aws 添加为计算引擎的运行时依赖。
  4. 根据 hadoop-aws 文档 配置 AWS 相关设置(请务必确认版本,S3A 的配置会因所用版本不同而有很大差异)。

S3 写入校验和验证

为确保上传对象的完整性,可以通过将目录属性 s3.checksum-enabled 设置为 true 来开启 S3 写入的校验和验证。该属性默认关闭。

S3 标签

在写入和删除 S3 对象时,可以为其添加自定义标签。例如,要使用 Spark 3.5 写入 S3 标签,可以按如下方式启动 Spark SQL shell:

spark-sql --conf spark.sql.catalog.my_catalog=org.apache.iceberg.spark.SparkCatalog \
    --conf spark.sql.catalog.my_catalog.warehouse=s3://my-bucket/my/key/prefix \
    --conf spark.sql.catalog.my_catalog.type=glue \
    --conf spark.sql.catalog.my_catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO \
    --conf spark.sql.catalog.my_catalog.s3.write.tags.my_key1=my_val1 \
    --conf spark.sql.catalog.my_catalog.s3.write.tags.my_key2=my_val2

对于上述示例,S3 中的对象将携带标签 my_key1=my_val1 和 my_key2=my_val2 被保存。请注意,指定的写入标签仅在创建对象时保存。

当目录属性 s3.delete-enabled 设置为 false 时,对象不会被从 S3 中硬删除。预期的用法是配合 S3 删除标签使用,即为对象打上标签后,再通过 S3 生命周期策略 将其移除。该属性默认设置为 true。

通过 s3.delete.tags 配置,对象会在删除前被加上所配置的键值对标签。用户可以在存储桶级别配置基于标签的对象生命周期策略,以便将对象转换到不同的存储层级。例如,要在 Spark 3.5 中添加 S3 删除标签,可以使用以下命令启动 Spark SQL shell:

sh spark-sql --conf spark.sql.catalog.my_catalog=org.apache.iceberg.spark.SparkCatalog \
    --conf spark.sql.catalog.my_catalog.warehouse=s3://iceberg-warehouse/s3-tagging \
    --conf spark.sql.catalog.my_catalog.type=glue \
    --conf spark.sql.catalog.my_catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO \
    --conf spark.sql.catalog.my_catalog.s3.delete.tags.my_key3=my_val3 \
    --conf spark.sql.catalog.my_catalog.s3.delete-enabled=false

对于上述示例,S3 中的对象在删除前会先被保存标签:my_key3=my_val3。用户还可以使用目录属性 s3.delete.num-threads 来指定为 S3 对象添加删除标签时所使用的线程数。

当目录属性 s3.write.table-tag-enabled 和 s3.write.namespace-tag-enabled 设置为 true 时,S3 中的对象会保存标签:iceberg.table=<table-name> 和 iceberg.namespace=<namespace-name>。用户可以基于这些标签按命名空间或表来定义访问和数据保留策略。例如,若要使用 Spark 3.5 将表名和命名空间名写入 S3 标签,可以通过以下方式启动 Spark SQL shell:

sh spark-sql --conf spark.sql.catalog.my_catalog=org.apache.iceberg.spark.SparkCatalog \
    --conf spark.sql.catalog.my_catalog.warehouse=s3://iceberg-warehouse/s3-tagging \
    --conf spark.sql.catalog.my_catalog.type=glue \
    --conf spark.sql.catalog.my_catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO \
    --conf spark.sql.catalog.my_catalog.s3.write.table-tag-enabled=true \
    --conf spark.sql.catalog.my_catalog.s3.write.namespace-tag-enabled=true

有关标签限制的更多详细信息,请参阅用户自定义标签限制。

S3 访问点

访问点可以通过指定存储桶到访问点的映射来执行 S3 操作。这适用于多区域访问、跨区域访问、灾难恢复等场景。

要使用跨区域访问点,需要额外将目录属性 use-arn-region-enabled 设置为 true,以使 S3FileIO 能够发起跨区域调用;对于同区域或多区域访问点,则无需此设置。

例如,要在 Spark 3.5 中使用 S3 访问点,可以通过以下方式启动 Spark SQL shell:

spark-sql --conf spark.sql.catalog.my_catalog=org.apache.iceberg.spark.SparkCatalog \
    --conf spark.sql.catalog.my_catalog.warehouse=s3://my-bucket2/my/key/prefix \
    --conf spark.sql.catalog.my_catalog.type=glue \
    --conf spark.sql.catalog.my_catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO \
    --conf spark.sql.catalog.my_catalog.s3.use-arn-region-enabled=false \
    --conf spark.sql.catalog.my_catalog.s3.access-points.my-bucket1=arn:aws:s3::<ACCOUNT_ID>:accesspoint/<MRAP_ALIAS> \
    --conf spark.sql.catalog.my_catalog.s3.access-points.my-bucket2=arn:aws:s3::<ACCOUNT_ID>:accesspoint/<MRAP_ALIAS>

对于上面的示例,my-bucket1 和 my-bucket2 存储桶中的对象在执行所有 S3 操作时将使用 arn:aws:s3::<ACCOUNT_ID>:accesspoint/<MRAP_ALIAS> 访问点。

有关使用访问点的更多详情,请参阅在兼容的 Amazon S3 操作中使用访问点和示例笔记本。

S3 Access Grants

S3 Access Grants 可以使用 IAM 主体(IAM Principal)来授予对 S3 数据的访问权限。为了在 Iceberg 中启用 S3 Access Grants,你需要先将 S3 Access Grants 插件 jar 添加到类路径中,然后将 s3.access-grants.enabled 目录属性设置为 true。该插件在 Maven 上的列表页可以在此查看。

此外,我们还支持回退到 IAM 的配置,当 S3 Access Grants 无法授权你的 S3 调用时,可以回退到使用你的 IAM 角色(及其权限集)来访问 S3 数据。这可以通过 s3.access-grants.fallback-to-iam 布尔类型目录属性来实现。默认情况下,该属性设置为 false。

例如,要将 S3 Access Grants 集成到 Spark 3.5 中,你可以使用以下命令启动 Spark SQL shell:

spark-sql --conf spark.sql.catalog.my_catalog=org.apache.iceberg.spark.SparkCatalog \
    --conf spark.sql.catalog.my_catalog.warehouse=s3://my-bucket2/my/key/prefix \
    --conf spark.sql.catalog.my_catalog.catalog-impl=org.apache.iceberg.aws.glue.GlueCatalog \
    --conf spark.sql.catalog.my_catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO \
    --conf spark.sql.catalog.my_catalog.s3.access-grants.enabled=true \
    --conf spark.sql.catalog.my_catalog.s3.access-grants.fallback-to-iam=true

有关使用 S3 Access Grants 的更多详情,请参阅使用 S3 Access Grants 管理访问权限。

S3 跨区域访问

通过将 catalog 属性 s3.cross-region-access-enabled 设置为 true,可以开启 S3 跨区域存储桶访问。该属性默认关闭,以避免首次 S3 API 调用产生额外延迟。

例如,要在 Spark 3.5 中启用 S3 跨区域存储桶访问,可以通过以下命令启动 Spark SQL shell:

spark-sql --conf spark.sql.catalog.my_catalog=org.apache.iceberg.spark.SparkCatalog \
    --conf spark.sql.catalog.my_catalog.warehouse=s3://my-bucket2/my/key/prefix \
    --conf spark.sql.catalog.my_catalog.type=glue \
    --conf spark.sql.catalog.my_catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO \
    --conf spark.sql.catalog.my_catalog.s3.cross-region-access-enabled=true

更多详情,请参阅 Amazon S3 跨区域访问。

S3 加速

S3 加速可用于加速与 Amazon S3 之间的数据传输,在长距离传输较大对象时,速度可提升 50% 到 500%。

要使用 S3 加速,我们需要将 s3.acceleration-enabled 目录属性设置为 true,以启用 S3FileIO 发起加速的 S3 调用。

例如,要在 Spark 3.5 中使用 S3 加速,可以通过以下方式启动 Spark SQL shell:

spark-sql --conf spark.sql.catalog.my_catalog=org.apache.iceberg.spark.SparkCatalog \
    --conf spark.sql.catalog.my_catalog.warehouse=s3://my-bucket2/my/key/prefix \
    --conf spark.sql.catalog.my_catalog.type=glue \
    --conf spark.sql.catalog.my_catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO \
    --conf spark.sql.catalog.my_catalog.s3.acceleration-enabled=true

有关使用 S3 加速的更多详细信息,请参阅使用 Amazon S3 Transfer Acceleration 配置快速、安全的文件传输。

S3 分析加速器

Amazon S3 分析加速器库可帮助你加速应用程序对 Amazon S3 数据的访问。这一开源解决方案能够减少数据分析工作负载的处理时间和计算成本。

要使 S3 分析加速器库在 Iceberg 中生效,你可以将 s3.analytics-accelerator.enabled 目录属性设置为 true。默认情况下,该属性设置为 false。

例如,要在 Spark 中使用 S3 分析加速器,你可以通过以下方式启动 Spark SQL shell:

spark-sql --conf spark.sql.catalog.my_catalog=org.apache.iceberg.spark.SparkCatalog \
    --conf spark.sql.catalog.my_catalog.warehouse=s3://my-bucket2/my/key/prefix \
    --conf spark.sql.catalog.my_catalog.type=glue \
    --conf spark.sql.catalog.my_catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO \
    --conf spark.sql.catalog.my_catalog.s3.analytics-accelerator.enabled=true

分析加速器库(Analytics Accelerator Library)既可以配合 S3 CRT 客户端使用,也可以配合 S3AsyncClient 使用。该库推荐使用 S3 CRT 客户端,因为它具有更出色的连接池管理能力以及更高的下载吞吐量。

客户端配置

属性默认值说明
s3.crt.enabledtrue控制是否使用 CRT 创建 S3 异步客户端
s3.crt.max-concurrency500S3 CRT 客户端的最大并发数

库特有的其他配置按以下章节进行组织:

逻辑 IO 配置

属性默认值说明
s3.analytics-accelerator.logicalio.prefetch.footer.enabledtrue控制是否启用页脚预取
s3.analytics-accelerator.logicalio.prefetch.page.index.enabledtrue控制是否启用页索引预取
s3.analytics-accelerator.logicalio.prefetch.file.metadata.size32KB普通文件预取的元数据大小
s3.analytics-accelerator.logicalio.prefetch.large.file.metadata.size1MB大文件预取的元数据大小
s3.analytics-accelerator.logicalio.prefetch.file.page.index.size1MB普通文件预取的页索引大小
s3.analytics-accelerator.logicalio.prefetch.large.file.page.index.size8MB大文件预取的页索引大小
s3.analytics-accelerator.logicalio.large.file.size1GB将文件视为大文件的阈值
s3.analytics-accelerator.logicalio.small.objects.prefetching.enabledtrue控制小对象的预取
s3.analytics-accelerator.logicalio.small.object.size.threshold3MB小对象预取的大小阈值
s3.analytics-accelerator.logicalio.parquet.metadata.store.size45Parquet 元数据存储的容量大小
s3.analytics-accelerator.logicalio.max.column.access.store.size15列访问存储的最大容量
s3.analytics-accelerator.logicalio.parquet.format.selector.regex`^.*.(parquet\par)$`用于识别 Parquet 文件的正则表达式模式
s3.analytics-accelerator.logicalio.prefetching.modeROW_GROUP预取模式(有效值:OFF、ALL、ROW_GROUP、COLUMN_BOUND)

物理 IO 配置

属性默认值说明
s3.analytics-accelerator.physicalio.metadatastore.capacity50元数据存储的容量
s3.analytics-accelerator.physicalio.blocksizebytes8MB数据传输的块大小
s3.analytics-accelerator.physicalio.readaheadbytes64KB预读的字节数
s3.analytics-accelerator.physicalio.maxrangesizebytes8MB范围请求的最大大小
s3.analytics-accelerator.physicalio.partsizebytes8MB传输时各个分片的大小
s3.analytics-accelerator.physicalio.sequentialprefetch.base2.0顺序预取大小计算的基数因子
s3.analytics-accelerator.physicalio.sequentialprefetch.speed1.0顺序预取增长的速度因子

遥测配置

属性默认值说明
s3.analytics-accelerator.telemetry.levelSTANDARD遥测详细级别(有效值:CRITICAL、STANDARD、VERBOSE)
s3.analytics-accelerator.telemetry.std.out.enabledfalse启用标准输出遥测输出
s3.analytics-accelerator.telemetry.logging.enabledtrue启用日志遥测输出
s3.analytics-accelerator.telemetry.aggregations.enabledfalse启用遥测聚合
s3.analytics-accelerator.telemetry.aggregations.flush.interval.seconds-1刷新聚合遥测的间隔
s3.analytics-accelerator.telemetry.logging.levelINFO遥测日志级别
s3.analytics-accelerator.telemetry.logging.namecom.amazon.connector.s3.telemetry遥测日志记录器名称
s3.analytics-accelerator.telemetry.formatdefault遥测输出格式(有效值:json、default)

对象客户端配置

属性默认值说明
s3.analytics-accelerator.useragentprefixnull在 S3 请求的 User-Agent 字符串前添加的自定义前缀

S3 双栈

S3 双栈允许客户端通过双栈终端节点访问 S3 存储桶。当客户端请求双栈终端节点时,如果可能,存储桶 URL 将解析为 IPv6 地址,否则回退到 IPv4。

要使用 S3 双栈,需要将 s3.dualstack-enabled catalog 属性设置为 true,以使 S3FileIO 能够发起双栈 S3 调用。

例如,要在 Spark 3.5 中使用 S3 双栈,可以这样启动 Spark SQL shell:

spark-sql --conf spark.sql.catalog.my_catalog=org.apache.iceberg.spark.SparkCatalog \
    --conf spark.sql.catalog.my_catalog.warehouse=s3://my-bucket2/my/key/prefix \
    --conf spark.sql.catalog.my_catalog.type=glue \
    --conf spark.sql.catalog.my_catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO \
    --conf spark.sql.catalog.my_catalog.s3.dualstack-enabled=true

有关使用 S3 双栈(Dual-stack)的更多详情,请参阅从 AWS CLI 和 AWS SDK 使用双栈端点

AWS 客户端定制

许多组织都有自己的方式来配置 AWS 客户端,例如使用自定义的凭证提供器、访问代理、重试策略等。Iceberg 允许用户通过设置 client.factory 目录属性,插入自定义的 org.apache.iceberg.aws.AwsClientFactory 实现。

跨账户与跨区域访问

组织通常会为 Glue 元数据存储和 S3 存储桶设置一个集中的 AWS 账户,并让不同的团队使用各自的 AWS 账户和区域来访问这些资源。在这种情况下,需要使用跨账户 IAM 角色来访问这些集中资源。Iceberg 提供了一个 AWS 客户端工厂 AssumeRoleAwsClientFactory 来支持这一常见用例。它同时也可作为希望实现自定义 AWS 客户端工厂的用户的参考示例。

该客户端工厂提供以下可配置的目录属性:

属性默认值描述
client.assume-role.arnnull,需要用户输入要代入(assume)的角色的 ARN,例如 arniam::123456789:role/myRoleToAssume
client.assume-role.regionnull,需要用户输入除 STS 客户端外,所有 AWS 客户端都将使用给定的区域,而不是默认的区域链
client.assume-role.external-idnull可选的外部 ID
client.assume-role.timeout-sec1 小时每次代入角色会话的超时时间。超时结束后,将通过 STS 客户端获取一套新的角色会话凭证。

使用该客户端工厂时,会先使用默认凭证和区域初始化一个 STS 客户端以代入指定角色,然后再使用代入角色后的凭证和区域初始化 Glue、S3 和 DynamoDB 客户端来访问资源。以下是一个使用该客户端工厂启动 Spark shell 的示例:

spark-sql --packages org.apache.iceberg:iceberg-spark-runtime-3.4_2.12:1.11.0,org.apache.iceberg:iceberg-aws-bundle:1.11.0 \
    --conf spark.sql.catalog.my_catalog=org.apache.iceberg.spark.SparkCatalog \
    --conf spark.sql.catalog.my_catalog.warehouse=s3://my-bucket/my/key/prefix \
    --conf spark.sql.catalog.my_catalog.type=glue \
    --conf spark.sql.catalog.my_catalog.client.factory=org.apache.iceberg.aws.AssumeRoleAwsClientFactory \
    --conf spark.sql.catalog.my_catalog.client.assume-role.arn=arn:aws:iam::123456789:role/myRoleToAssume \
    --conf spark.sql.catalog.my_catalog.client.assume-role.region=ap-northeast-1

HTTP 客户端配置

AWS 客户端支持两种 HTTP 客户端类型:URL Connection HTTP Client 和 Apache HTTP Client。默认情况下,AWS 客户端使用 Apache HTTP Client 与服务进行通信。该 HTTP 客户端支持多种功能和自定义设置(例如 expect-continue 握手和 TCP KeepAlive),但代价是额外的依赖项和更多的启动延迟。相比之下,URL Connection HTTP Client 以最少的依赖项和最低的启动延迟为优化目标,但其支持的功能少于其他实现。

有关配置的更多详情,请参见 URL Connection HTTP Client 配置 和 Apache HTTP Client 配置 章节。

HTTP 客户端的配置可以通过 catalog 属性进行设置。以下是可用配置的概览:

属性 默认值 描述

http-client.type apache HTTP 客户端的类型。
urlconnection:URL Connection HTTP Client
apache:Apache HTTP Client

http-client.proxy-endpoint null 用于 HTTP 客户端的可选代理端点。

http-client.proxy-use-system-property-values null,默认启用 可选的 true/false 设置,控制是否从 Java 系统属性(http.proxyHost、http.proxyPort、http.nonProxyHosts 等)中读取代理配置。

http-client.proxy-use-environment-variable-values null,默认启用 可选的 true/false 设置,控制是否从环境变量(HTTP_PROXY、HTTPS_PROXY、NO_PROXY 等)中读取代理配置。

URL Connection HTTP Client 配置

URL Connection HTTP Client 具有以下可配置属性:

属性默认值描述
http-client.urlconnection.socket-timeout-msnull可选的套接字超时,单位为毫秒
http-client.urlconnection.connection-timeout-msnull可选的连接超时,单位为毫秒

用户可以使用 catalog 属性覆盖默认值。例如,要在启动 spark shell 时配置 URL Connection HTTP Client 的套接字超时,可以添加:

--conf spark.sql.catalog.my_catalog.http-client.urlconnection.socket-timeout-ms=80

Apache HTTP 客户端配置

Apache HTTP 客户端提供以下可配置属性:

属性默认值说明
http-client.apache.socket-timeout-msnull可选的套接字超时,单位为毫秒
http-client.apache.connection-timeout-msnull可选的连接超时,单位为毫秒
http-client.apache.connection-acquisition-timeout-msnull可选的连接获取超时,单位为毫秒
http-client.apache.connection-max-idle-time-msnull可选的连接最大空闲时间,单位为毫秒
http-client.apache.connection-time-to-live-msnull可选的连接存活时间,单位为毫秒
http-client.apache.expect-continue-enablednull,默认禁用可选的 true/false 设置,用于控制是否启用 expect continue
http-client.apache.max-connectionsnull可选的最大连接数,整数值
http-client.apache.tcp-keep-alive-enablednull,默认禁用可选的 true/false 设置,用于控制是否启用 tcp keep alive
http-client.apache.use-idle-connection-reaper-enablednull,默认启用可选的 true/false 设置,用于控制是否使用 use idle connection reaper

用户可以使用 catalog 属性来覆盖默认值。例如,要在启动 spark shell 时配置 Apache HTTP Client 的最大连接数,可以添加:

--conf spark.sql.catalog.my_catalog.http-client.apache.max-connections=5

在 AWS 上运行 Iceberg

Amazon Athena

Amazon Athena 提供了一个无服务器查询引擎,可用于对 Iceberg 表执行读取、写入、更新和优化操作。更多详情请参见此处。

Amazon EMR

Amazon EMR 可以配置集成了 Spark(Spark 3 使用 EMR 6,Spark 2 使用 EMR 5)、Hive、Flink、Trino 的集群,这些组件均可运行 Iceberg。

从 EMR 6.5.0 版本开始,EMR 集群可以配置为自动安装必要的 Apache Iceberg 依赖,而无需使用引导操作(bootstrap action)。有关如何创建已安装 Iceberg 的集群,请参阅官方文档。

对于 6.5.0 之前的版本,你可以使用类似如下的引导操作来预安装所有必要的依赖:

#!/bin/bash

ICEBERG_VERSION=1.11.0
MAVEN_URL=https://repo1.maven.org/maven2
ICEBERG_MAVEN_URL=$MAVEN_URL/org/apache/iceberg
# NOTE: this is just an example shared class path between Spark and Flink,
#  please choose a proper class path for production.
LIB_PATH=/usr/share/aws/aws-java-sdk/


ICEBERG_PACKAGES=(
  "iceberg-spark-runtime-3.5_2.12"
  "iceberg-flink-runtime"
  "iceberg-aws-bundle"
)

install_dependencies () {
  install_path=$1
  download_url=$2
  version=$3
  shift
  pkgs=("$@")
  for pkg in "${pkgs[@]}"; do
    sudo wget -P $install_path $download_url/$pkg/$version/$pkg-$version.jar
  done
}

install_dependencies $LIB_PATH $ICEBERG_MAVEN_URL $ICEBERG_VERSION "${ICEBERG_PACKAGES[@]}"

AWS Glue

AWS Glue 提供了一项无服务器数据集成服务,可用于对 Iceberg 表执行读取、写入和更新操作。更多详情请参见此处。

AWS EKS

AWS Elastic Kubernetes Service (EKS) 可用于启动任意 Spark、Flink、Hive、Presto 或 Trino 集群,以配合 Iceberg 使用。

Amazon Kinesis

Amazon Kinesis Data Analytics 提供了一个运行全托管 Apache Flink 应用程序的平台。你可以将 Iceberg 包含在应用程序 Jar 中,并在该平台上运行。

AWS Redshift

AWS Redshift Spectrum 或 Redshift Serverless 支持查询在 AWS Glue Data Catalog 中登记的 Apache Iceberg 表。

Amazon Data Firehose

你可以使用 Firehose 将流式数据直接投递到 Amazon S3 中的 Apache Iceberg 表。借助此功能,你可以将单个流中的记录路由到不同的 Apache Iceberg 表,并自动对 Apache Iceberg 表中的记录执行插入、更新和删除操作。此功能需要使用 AWS Glue Data Catalog。

评论

登录后参与评论

正在加载评论…