AWS DynamoDB
Iceberg AWS 集成
Iceberg 通过 iceberg-aws 模块提供与各种 AWS 服务的集成。本节介绍如何在 AWS 上使用 Iceberg。
启用 AWS 集成
从 0.11.0 版本起,Spark 和 Flink 引擎运行时均已内置 iceberg-aws 模块。但 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 指定了额外的 iceberg-aws-bundle,其中包含了所有相关的 AWS 依赖。
Flink
要配合 Flink 使用 AWS 模块,你可以下载所需的依赖,并在启动 Flink SQL client 时指定它们:
# 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.S3FileIOHive
要在 Hive 中使用 AWS 模块,您可以按照与 Flink 示例类似的方式下载必要的依赖项,然后将它们添加到 Hive 的 classpath 中,或在 CLI 中运行时添加 jar 包:
add jar /my/path/to/iceberg-hive-runtime.jar;
add jar /my/path/to/aws/bundle.jar;有了这些依赖,你可以在 CLI 中注册 Glue 目录,并在运行时于 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)
用户可以选择多种不同的方式,使用 AWS 构建 Iceberg 目录。
Glue Catalog
Iceberg 支持将 AWS Glue 作为 Catalog 的实现。使用时,Iceberg 命名空间存储为 Glue Database,Iceberg 表存储为 Glue Table,每个 Iceberg 表版本存储为 Glue TableVersion。你可以将 catalog-impl 指定为 org.apache.iceberg.aws.glue.GlueCatalog,或将 catalog-type 设置为 glue,从而开始使用 Glue 目录,这与上文启用 AWS 集成一节中的示例相同。关于加载目录的更多详情,可参阅各个引擎的页面,例如 Spark 和 Flink。
Glue Catalog ID
每个 AWS 账户和每个 AWS 区域都有一个唯一的 Glue 元数据存储。默认情况下,GlueCatalog 会根据用户的默认 AWS 客户端凭据和区域设置来选择要使用的 Glue 元数据存储。你可以通过 glue.id 目录属性指定 Glue catalog ID,以指向位于不同 AWS 账户中的 Glue 目录。Glue catalog ID 就是你的数字形式 AWS 账户 ID。如果 Glue 目录位于不同的区域,你应当配置 AWS 客户端指向正确的区域,详见 AWS 客户端定制。
跳过归档
AWS Glue 能够归档较旧的表版本,如有需要,用户可以将表回滚到任意历史版本。默认情况下,Iceberg 的 Glue Catalog 会跳过对较旧表版本的归档。如果用户希望归档较旧的表版本,可以将 glue.skip-archive 设置为 false。请注意,对于向 Iceberg 表进行的流式摄取,将 glue.skip-archive 设置为 false 会迅速生成大量 Glue 表版本。更多详情请参阅 Glue 配额 和 UpdateTable API。
跳过名称校验
允许用户跳过表名和命名空间的名称校验。建议遵循 Glue 最佳实践,以确保操作与 Hive 兼容。此选项仅为那些已有使用非标准字符的既有约定的用户而添加。当跳过数据库名和表名校验时,无法保证所有下游系统都能支持这些名称。
乐观锁
默认情况下,Iceberg 使用 Glue 的乐观锁来处理对表的并发更新。通过乐观锁,每个表都有一个版本 ID。当用户获取表元数据时,Iceberg 会记录该表的版本 ID。只要服务端的版本 ID 保持不变,用户就可以更新该表。如果在你之前有其他人修改了该表,就会出现版本不匹配,导致更新失败。随后 Iceberg 会刷新元数据并检查是否存在冲突。如果没有提交冲突,操作将会重试。乐观锁保证了 Glue 中 Iceberg 表的原子事务,同时也能防止他人意外覆盖你的更改。
信息
请使用 AWS SDK 版本 >= 2.17.131 以使用 Glue 的乐观锁。如果 AWS SDK 版本低于 2.17.131,则仅使用内存锁。为确保原子事务,你需要配置 DynamoDb Lock Manager。
仓库位置
与其他所有目录实现类似,warehouse 是一个必需的目录属性,用于确定存储中数据仓库的根路径。默认情况下,由于使用 S3FileIO,Glue 仅允许将仓库位置设置在 S3 上。若要将数据存储在其他本地或云存储中,可以通过设置 io-impl 目录属性,将 Glue 目录切换为使用 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-name | iceberg | DynamoDbCatalog 使用的 DynamoDB 表名称 |
内部表设计
DynamoDB 表按如下列进行设计:
| 列 | 键 | 类型 | 描述 |
|---|---|---|---|
| identifier | 分区键 | string | 表标识符,例如 db1.table1;对于命名空间则使用字符串 NAMESPACE |
| namespace | 排序键 | string | 命名空间名称。以 namespace 为分区键、identifier 为排序键创建全局二级索引(GSI),不投影其他列 |
| v | string | 行版本,用于乐观锁 | |
| updated_at | number | 最后一次更新的时间戳(毫秒) | |
| created_at | number | 表创建的时间戳(毫秒) | |
| p.<property_key> | string | Iceberg 定义的表属性,包括 table_type、metadata_location 和 previous_metadata_location,或者命名空间属性 |
这种设计有以下优点:
- 由于分区键位于表级别,如果同一命名空间中的表存在大量写入流量,它可以避免潜在的热点分区问题
- 命名空间操作集中在单个分区中,以避免影响表的提交操作
- 列表操作使用从排序键到分区键的反向 GSI,其余所有操作均为单行操作或单分区查询。目录中的任何操作都无需进行全表扫描。
- 使用字符串 UUID 版本字段
v而非updated_at,以避免两个进程在同一毫秒内提交 catalog.renameTable使用多行事务以确保幂等性- 属性被展开为顶层列,以便用户可以在任意属性字段上创建自定义 GSI,从而定制目录。例如,用户可以将所有者信息存储为表属性
owner,并通过在p.owner列上添加 GSI 来按所有者搜索表。
RDS JDBC Catalog
Iceberg 还支持 JDBC catalog,它使用关系数据库中的表来管理 Iceberg 表。你可以配置使用 JDBC catalog,并配合 AWS RDS 等关系数据库服务。有关使用 JDBC catalog 的指南和示例,请参阅 JDBC 集成页面。有关使用 IAM 身份验证配置 JDBC catalog 的更多详细信息,请参阅此 AWS 文档。
如何选择目录?
在所有可用选项的基础上,我们为根据应用选择合适的目录提供以下指导:
- 如果您的组织已有 Glue 元数据存储,或计划使用包括 Glue、Athena、EMR、Redshift 和 LakeFormation 在内的 AWS 分析生态,Glue catalog 能提供最简单的集成方式。
- 如果您的应用需要频繁更新表,或者有较高的读写吞吐量(例如流式写入),Glue 和 DynamoDB catalog 通过乐观锁可提供最佳性能。
- 如果您希望对 catalog 中的表实施访问控制,Glue 表可以作为 IAM 资源进行管理,而 DynamoDB catalog 的表只能通过 条目级权限管理,后者要复杂得多。
- 如果您希望基于表属性信息进行查询而无需扫描整个 catalog,DynamoDB catalog 允许您为任意属性字段构建二级索引,并提供高效的查询性能。
- 如果您既想享受 DynamoDB catalog 的优势,又需要连接到 Glue,可以启用 DynamoDB 流与 Lambda 触发器,将 DynamoDB catalog 中的表信息异步更新到您的 Glue 元数据存储。
- 如果您的组织已经在 RDS 中维护了现有关系型数据库,或使用 无服务器 Aurora 来管理表,JDBC catalog 能提供最简单的集成方式。
DynamoDb 锁管理器
HadoopCatalog 或 HadoopTables 可以使用 Amazon DynamoDB,这样每次提交时,catalog 会先借助一张辅助 DynamoDB 表获取锁,然后再安全地修改 Iceberg 表。对于基于文件系统的 catalog 而言,这是必要的,因为 S3 等存储不提供文件写入互斥,需要以此保证事务的原子性。
此功能需要以下与锁相关的 catalog 属性:
- 将
lock-impl设置为org.apache.iceberg.aws.dynamodb.DynamoDbLockManager。 - 将
lock.table设置为您想要使用的 DynamoDB 表名。如果 DynamoDB 中不存在给定名称的锁表,则会创建一张新表,并将计费模式设置为按需计费。
还可以使用其他与锁相关的 catalog 属性来调整锁的行为,例如心跳间隔。更多详情请参阅锁相关的 catalog 属性。
S3 FileIO
Iceberg 允许用户通过 S3FileIO 将数据写入 S3。GlueCatalog 默认使用该 FileIO,其他目录可以通过 io-impl 目录属性加载该 FileIO。
渐进式分片上传
S3FileIO 实现了一种自定义的渐进式分片上传算法来上传数据。数据文件按分片并行上传,每个分片就绪后立即上传,每个文件分片上传完成后即被删除。这在上传期间提供了最大化的上传速度和最小化的本地磁盘占用。用户可以调整以下与该特性相关的配置:
| 属性 | 默认值 | 描述 |
|---|---|---|
| s3.multipart.num-threads | 系统可用的处理器数量 | 用于将分片上传到 S3 的线程数(在所有输出流之间共享) |
| s3.multipart.part-size-bytes | 32MB | 分片上传请求中单个分片的大小 |
| s3.multipart.threshold | 1.5 | 以分片大小的倍数表示的阈值,超过该阈值时从单个 put object 请求切换为分片上传 |
| s3.staging-dir | java.io.tmpdir 属性值 | 用于存放临时文件的目录 |
S3 服务端加密
S3FileIO 支持全部 3 种 S3 服务端加密模式:
- SSE-S3:使用由 Amazon S3 托管密钥的服务端加密(SSE-S3)时,每个对象都会使用唯一密钥进行加密。作为额外的安全措施,它会使用一个会定期轮换的主密钥对密钥本身进行加密。Amazon S3 服务端加密采用可用的最强大分组密码之一——256 位高级加密标准(AES-256)——来加密您的数据。
- SSE-KMS:使用存储在 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.type | none | none、s3、kms、dsse-kms 或 custom |
| s3.sse.key | kms 和 dsse-kms 类型为 aws/s3,其他类型为 null | 对于 kms 和 dsse-kms 类型,为 KMS 密钥 ID 或 ARN;对于 custom 类型,为自定义的 base-64 AES256 对称密钥。 |
| s3.sse.md5 | null | 如果 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 相关 I/O 操作的吞吐量。使用 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");这将在 write.object-storage.path 之后直接附加一个 20 位的 base2 哈希值(01010110100110110010)并将数据写入 S3,从而确保对表的读取均匀分布在各个 S3 存储桶前缀上,进而提升性能。之前提供的 base64 哈希已更新为 base2,以便在 S3 通用型存储桶上获得更好的自动扩缩容行为。
作为本次更新的一部分,我们还将熵值拆分到多个目录中,以提升 Iceberg 孤儿文件清理过程的效率,因为目录被用作将工作量划分给各个工作节点的手段,从而加快遍历速度。从下面的示例中可以看到,我们拆分哈希值,生成深度为 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将被移除。
更多详情请参阅位置提供者配置一节。
我们还新增了一个表属性 write.object-storage.partitioned-paths,将其设置为 false(默认为 true)时,文件路径中将不再包含分区值。Iceberg 在文件路径中并不需要这些值,将该配置设为 false 可以进一步减小键的大小。在这种情况下,我们还会将最终的 8 位熵直接附加到文件名上。设置该配置后,插入的键如下所示,注意 category=orders 已被移除:
s3://my-table-data-bucket/my_ns.db/my_table/1101/0100/1011/00111010-00000-0-5affc076-96a4-48f2-9cd2-d5efbc9f0c94-00001.parquetS3 重试
遇到 S3 限流(throttling)的工作负载应以持久重试的方式并配合指数退避(exponential backoff)来持续推进,同时由 S3 自动完成扩缩容。我们提供以下配置项,用于调整 S3 的重试行为。对于遇到限流且因重试次数耗尽而失败的工作负载,我们建议将重试次数设置为 32,以便 S3 有时间自动扩容。请注意,如果工作负载的吞吐量极高,而对应的表所使用的 S3 尚未完成扩容,可能需要进一步提高重试次数。
| 属性 | 默认值 | 说明 |
|---|---|---|
| s3.retry.num-retries | 5 | S3 操作的重试次数。对于高吞吐工作负载,建议设置为 32。 |
| s3.retry.min-wait-ms | 2s | 重试 S3 操作的最小等待时间。 |
| s3.retry.max-wait-ms | 20s | 重试 S3 读操作的最大等待时间。 |
S3 强一致性
2020 年 11 月,S3 宣布其所有读操作具备强一致性,Iceberg 也已更新以充分利用这一特性。IO 操作期间不再存在可能对性能产生负面影响的多余一致性等待与校验。
Hadoop S3A 文件系统
重要
对于 S3 使用场景,推荐使用 S3FileIO,而不是 S3A FileSystem(HadoopFileIO)。
在 S3FileIO 引入之前,许多 Iceberg 用户选择使用 HadoopFileIO,通过 S3A FileSystem 将数据写入 S3。如前面章节所述,S3FileIO 采用了最新的 AWS 客户端和 S3 特性,以实现更优的安全性和性能。
S3FileIO 使用 s3:// URI 协议写入数据,同时也兼容 S3A FileSystem 写入的协议。这意味着,对于包含 s3a:// 或 s3n:// 文件路径的任何表清单(manifest),S3FileIO 仍能正常读取。这一特性使得用户可以轻松地从 S3A 切换到 S3FileIO。
如果出于某些原因你必须使用 S3A,请按以下步骤操作:
- 要使用 S3A 存储数据,请将
warehousecatalog 属性指定为 S3A 路径,例如s3a://my-bucket/my-warehouse - 对于
HiveCatalog,若要同时使用 S3A 存储元数据,请将 Hadoop 配置属性hive.metastore.warehouse.dir指定为 S3A 路径。 - 将 hadoop-aws 添加为计算引擎的运行时依赖。
- 根据 hadoop-aws 文档 配置 AWS 相关设置(请确认版本,S3A 的配置方式会因所用版本不同而有很大差异)。
S3 写入校验和验证
为确保上传对象的完整性,可以通过将 catalog 属性 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 进行保存。用户还可以使用 catalog 属性 s3.delete.num-threads 来指定为 S3 对象添加删除标签时所使用的线程数。
当 catalog 属性 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 主体授予对 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 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.acceleration-enabled=true有关使用 S3 传输加速的更多详细信息,请参阅使用 Amazon S3 传输加速配置快速、安全的文件传输。
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.enabled | true | 控制 S3 异步客户端是否通过 CRT 创建 |
| s3.crt.max-concurrency | 500 | S3 CRT 客户端的最大并发数 |
库的其他特定配置按以下章节分别组织:
逻辑 IO 配置
| 属性 | 默认值 | 说明 | |
|---|---|---|---|
| s3.analytics-accelerator.logicalio.prefetch.footer.enabled | true | 控制是否启用页脚预取 | |
| s3.analytics-accelerator.logicalio.prefetch.page.index.enabled | true | 控制是否启用页索引预取 | |
| s3.analytics-accelerator.logicalio.prefetch.file.metadata.size | 32KB | 常规文件预取的元数据大小 | |
| s3.analytics-accelerator.logicalio.prefetch.large.file.metadata.size | 1MB | 大文件预取的元数据大小 | |
| s3.analytics-accelerator.logicalio.prefetch.file.page.index.size | 1MB | 常规文件预取的页索引大小 | |
| s3.analytics-accelerator.logicalio.prefetch.large.file.page.index.size | 8MB | 大文件预取的页索引大小 | |
| s3.analytics-accelerator.logicalio.large.file.size | 1GB | 将文件视为大文件的阈值 | |
| s3.analytics-accelerator.logicalio.small.objects.prefetching.enabled | true | 控制小对象的预取 | |
| s3.analytics-accelerator.logicalio.small.object.size.threshold | 3MB | 小对象预取的大小阈值 | |
| s3.analytics-accelerator.logicalio.parquet.metadata.store.size | 45 | Parquet 元数据存储的大小 | |
| s3.analytics-accelerator.logicalio.max.column.access.store.size | 15 | 列访问存储的最大大小 | |
| s3.analytics-accelerator.logicalio.parquet.format.selector.regex | `^.*.(parquet\ | par)$` | 用于识别 Parquet 文件的正则表达式 |
| s3.analytics-accelerator.logicalio.prefetching.mode | ROW_GROUP | 预取模式(可选值:OFF、ALL、ROW_GROUP、COLUMN_BOUND) |
物理 IO 配置
| 属性 | 默认值 | 说明 |
|---|---|---|
| s3.analytics-accelerator.physicalio.metadatastore.capacity | 50 | 元数据存储的容量 |
| s3.analytics-accelerator.physicalio.blocksizebytes | 8MB | 数据传输的块大小 |
| s3.analytics-accelerator.physicalio.readaheadbytes | 64KB | 预读取的字节数 |
| s3.analytics-accelerator.physicalio.maxrangesizebytes | 8MB | 范围请求的最大大小 |
| s3.analytics-accelerator.physicalio.partsizebytes | 8MB | 传输时各个分片的大小 |
| s3.analytics-accelerator.physicalio.sequentialprefetch.base | 2.0 | 顺序预取大小的基础系数 |
| s3.analytics-accelerator.physicalio.sequentialprefetch.speed | 1.0 | 顺序预取增长的速度系数 |
遥测配置
| 属性 | 默认值 | 说明 |
|---|---|---|
| s3.analytics-accelerator.telemetry.level | STANDARD | 遥测详细级别(有效值:CRITICAL、STANDARD、VERBOSE) |
| s3.analytics-accelerator.telemetry.std.out.enabled | false | 启用标准输出遥测输出 |
| s3.analytics-accelerator.telemetry.logging.enabled | true | 启用日志遥测输出 |
| s3.analytics-accelerator.telemetry.aggregations.enabled | false | 启用遥测聚合 |
| s3.analytics-accelerator.telemetry.aggregations.flush.interval.seconds | -1 | 刷新聚合遥测的间隔时间 |
| s3.analytics-accelerator.telemetry.logging.level | INFO | 遥测日志级别 |
| s3.analytics-accelerator.telemetry.logging.name | com.amazon.connector.s3.telemetry | 遥测日志记录器名称 |
| s3.analytics-accelerator.telemetry.format | default | 遥测输出格式(有效值:json、default) |
对象客户端配置
| 属性 | 默认值 | 说明 |
|---|---|---|
| s3.analytics-accelerator.useragentprefix | null | 添加到 S3 请求中 User-Agent 字符串的自定义前缀 |
S3 双栈
S3 双栈允许客户端通过双栈端点访问 S3 存储桶。当客户端请求双栈端点时,如果可能,存储桶 URL 会解析为 IPv6 地址,否则回退到 IPv4。
要使用 S3 双栈,我们需要将 s3.dualstack-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.dualstack-enabled=true有关使用 S3 双栈的更多详情,请参阅在 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.arn | null,需要用户输入 | 要切换到的角色的 ARN,例如 arniam::123456789:role/myRoleToAssume |
| client.assume-role.region | null,需要用户输入 | 除 STS 客户端外,所有 AWS 客户端都将使用给定的区域,而不是默认的区域链 |
| client.assume-role.external-id | null | 可选的外部 ID |
| client.assume-role.timeout-sec | 1 小时 | 每次切换角色会话的超时时间。超时结束后,将通过 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-1HTTP 客户端配置
AWS 客户端支持两种 HTTP 客户端类型:URL Connection HTTP Client 和 Apache HTTP Client。默认情况下,AWS 客户端使用 Apache HTTP 客户端与服务进行通信。该 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 Clientapache: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-ms | null | 可选的 套接字超时,单位为毫秒 |
| http-client.urlconnection.connection-timeout-ms | null | 可选的 连接超时,单位为毫秒 |
用户可以使用 catalog 属性覆盖默认值。例如,要在启动 spark shell 时为 URL Connection HTTP Client 配置套接字超时,可以添加:
--conf spark.sql.catalog.my_catalog.http-client.urlconnection.socket-timeout-ms=80Apache HTTP 客户端配置
Apache HTTP 客户端提供以下可配置属性:
| 属性 | 默认值 | 描述 |
|---|---|---|
| http-client.apache.socket-timeout-ms | null | 可选的 套接字超时(毫秒) |
| http-client.apache.connection-timeout-ms | null | 可选的 连接超时(毫秒) |
| http-client.apache.connection-acquisition-timeout-ms | null | 可选的 获取连接超时(毫秒) |
| http-client.apache.connection-max-idle-time-ms | null | 可选的 连接最大空闲时间(毫秒) |
| http-client.apache.connection-time-to-live-ms | null | 可选的 连接存活时间(毫秒) |
| http-client.apache.expect-continue-enabled | null,默认禁用 | 可选的 true/false 设置,用于控制是否启用 expect continue |
| http-client.apache.max-connections | null | 可选的 最大连接数(整数) |
| http-client.apache.tcp-keep-alive-enabled | null,默认禁用 | 可选的 true/false 设置,用于控制是否启用 tcp keep alive |
| http-client.apache.use-idle-connection-reaper-enabled | null,默认启用 | 可选的 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。
评论
登录后参与评论
KnowForge