流数据的归宿问题
在 OSS Tables 的生态里,Kafka Connect 是入湖管道,把 Kafka 消息原样写入 OSS Tables,不加不改,只负责搬。Kafka 是数据管道的咽喉要道,但咽喉不是终点,事件流进去之后最终要落到数据湖。Kafka Connect 就是这条最短的路。
常见的做法是写入一个下游存储。写数据库、写搜索索引、写对象存储里的文件,各有各的理由。但当数据量大到一定程度,或者你需要把多路流汇聚起来做离线分析时,“数据湖”几乎是绕不开的终点。
问题出在“写入”这一步。传统方案要么自己写 Consumer 把消息拉出来转成 Parquet 文件塞进 OSS,要么搭一套专门的入湖管道,但不管哪种,都要处理 Schema 管理、文件滚动、Exactly-Once 语义、分区策略这些工程细节。
Kafka Connect 的 Iceberg Sink Connector 把这些问题封装了。而 OSS Tables 兼容 Iceberg REST Catalog 协议,意味着可以通过 Kafka Connect 的 Iceberg Sink Connector 将 Kafka 消息实时写入 OSS Tables 中的表,实现流式数据入湖。
步骤一:环境准备
下载依赖 JAR 包
将以下 JAR 包放入 Kafka Connect 的插件目录(plugin.path 指定的路径)。
JAR包 |
版本要求 |
说明 |
匹配iceberg版本 |
提供 S3FileIO 实现及 REST Catalog SigV4 签名认证所需的 AWS SDK。 |
|
匹配iceberg版本 |
提供 SigV4 签名和 S3FileIO 实现。版本需与 iceberg-aws-bundle 一致。 |
|
匹配iceberg版本 |
Parquet 文件格式写入支持。 |
|
3.3.6 |
Hadoop 运行时依赖(Iceberg 内部加载需要),版本可按需调整。 |
|
3.3.6 |
Hadoop API 依赖(Iceberg 内部加载需要),版本可按需调整。 |
|
3.3.2 |
Iceberg SnapshotProducer 运行时需要此依赖,缺少会导致 ClassNotFoundException。 |
步骤二:创建Table Bucket
在开始写入数据之前,需要创建 Table Bucket 和 Namespace。可以使用 ossutil 创建,AWS CLI 也支持。
方式一:使用ossutil
1. 安装或升级 ossutil
请安装ossutil 2.3.0以上版本,如已安装 ossutil,可执行以下命令升级到最新版本:
ossutil update -f
2. 配置凭证
执行 ossutil config 命令,按提示输入 AccessKey ID、AccessKey Secret 和 Region。
3. 创建 Table Bucket
ossutil tables-api create-table-bucket --name {table bucket名称} --endpoint http://{endpint} --region {region}
命令执行成功后,返回结果中包含 Table Bucket ARN,请记录该值。
4. 创建 Namespace
ossutil tables-api create-namespace --table-bucket-arn {Table Bucket ARN} --namespace {Namespace名称} --endpoint http://{endpint}
重要:Namespace 和 Table 名称不能包含连字符(-),可使用下划线(_),这是因为名称会用于 SQL 语句中的标识符。
5. 创建 Table
您可以选择以下任一方式创建 Iceberg 表:
- 通过其他计算引擎创建(如 Spark)。
- 通过 ossutil 创建:先将表 schema 保存为 JSON 文件,再调用
create-table。
以下示例的 schema 文件 schema.json 定义了 3 个字段:
{ "iceberg": { "schema": { "fields": [ {"name": "event_id", "type": "string", "required": true}, {"name": "event_time", "type": "string"}, {"name": "event_type", "type": "string"} ] } } }
基于 schema 文件创建 Table:
ossutil tables-api create-table --table-bucket-arn <Table Bucket ARN> --namespace <Namespace名称> --name <表名称> --format ICEBERG --metadata file://schema.json --endpoint http://{endpint}
方式二:使用 AWS CLI
OSS Tables 兼容 S3 Tables API,也可以使用 AWS CLI 管理 Table Bucket。
1. 安装 AWS CLI
curl "https://awscli.amazonaws.com/awscli-exe-linux-x86_64.zip" -o "awscliv2.zip" unzip awscliv2.zip sudo ./aws/install
2. 配置凭证
执行 aws configure 命令,按提示输入 AccessKey ID、AccessKey Secret 和 Region。
3. 创建 Table Bucket
aws s3tables --endpoint http://{endpint} create-table-bucket --region {region} --name {table bucket名称}
命令执行成功后,返回结果中包含 Table Bucket ARN。
4. 创建 Namespace
aws s3tables --endpoint http://{endpoint} create-namespace --table-bucket-arn {Table Bucket ARN} --namespace {namespace名称}
5. 创建 Table
- 通过其他计算引擎(如 Spark)创建表
- 使用 AWS CLI 创建。使用 AWS CLI 时,先将完整的入参保存为 JSON 文件
create-table.json,再调用create-table。
{ "tableBucketARN": "<Table Bucket ARN>", "namespace": "<Namespace名称>", "name": "<表名称>", "format": "ICEBERG", "metadata": { "iceberg": { "schema": { "fields": [ {"name": "event_id", "type": "string", "required": true}, {"name": "event_time", "type": "string"}, {"name": "event_type", "type": "string"} ] } } } }
aws s3tables --endpoint http://{endpoint} create-table --cli-input-json file://{文件路径}
6. 管理后台维护任务
OSS Tables 支持自动执行 Iceberg 表的后台维护(如文件清理、文件合并等),通过 AWS CLI 可以查询和配置维护任务。
查询 Table 维护任务状态:
aws s3tables get-table-maintenance-job-status \ --table-bucket-arn="{bucketArn}" \ --namespace="{namespace名称}" \ --name="{表名}"
配置 Bucket 级维护策略(文件清理):
aws s3tables put-table-bucket-maintenance-configuration \ --table-bucket-arn {tableArn} \ --type icebergUnreferencedFileRemoval \ --value '{"status":"enabled","settings":{"icebergUnreferencedFileRemoval":{"unreferencedDays":4,"nonCurrentDays":10}}}'
配置 Table 级维护策略(小文件合并):
aws s3tables put-table-maintenance-configuration \ --table-bucket-arn {bucketArn} \ --type icebergCompaction \ --namespace {namespace名称} \ --name {表名} \ --value='{"status":"enabled","settings":{"icebergCompaction":{"targetFileSizeMB":256}}}'
步骤三:配置 Kafka Connect
OSS Tables 提供 Iceberg REST Catalog 端点,Kafka Connect 通过 Iceberg Sink Connector 连接该端点写入数据。Endpoint格式如下:
- 内网:
https://{region}-internal.oss-tables.aliyuncs.com/iceberg - 外网:
https://{region}.oss-tables.aliyuncs.com/iceberg
OSS Tables 提供S3FileIO访问OSS数据面使用的访问端点,Spark 通过该端点访问表数据。Endpoint格式如下:
Connector配置
创建 Iceberg Sink Connector 时,指定 Connector 类为 org.apache.iceberg.connect.IcebergSinkConnector,并配置以下属性:
# --- Iceberg Catalog (REST) --- iceberg.catalog.type: rest iceberg.catalog.uri: https://{region}-internal.oss-tables.aliyuncs.com/iceberg iceberg.catalog.rest.sigv4-enabled: true iceberg.catalog.rest.signing-region: <Region> iceberg.catalog.warehouse: <Table Bucket ARN> iceberg.catalog.rest.signing-name: osstables iceberg.catalog.rest.access-key-id: <AccessKey ID> iceberg.catalog.rest.secret-access-key: <AccessKey Secret> # --- Force S3FileIO (catalog returns oss:// but storage is S3-compatible) --- iceberg.catalog.io-impl: org.apache.iceberg.aws.s3.S3FileIO # --- S3FileIO 存储配置 --- iceberg.catalog.s3.endpoint: https://oss-{region}-internal.aliyuncs.com iceberg.catalog.s3.access-key-id: <AccessKey ID> iceberg.catalog.s3.secret-access-key: <AccessKey Secret> iceberg.catalog.s3.path-style-access: true iceberg.catalog.client.region: <Region> # --- 数据格式转换 --- key.converter: org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable: false value.converter: org.apache.kafka.connect.json.JsonConverter value.converter.schemas.enable: false
重要:如果使用较高版本的 AWS SDK(2.20+),可能出现签名错误:aws-chunked encoding is not supported with the specified x-amz-content-sha256 value。此时需要在 Kafka Connect 的 Java 启动参数中添加以下 JVM 选项:
-Daws.requestChecksumCalculation=when_required -Daws.responseChecksumValidation=when_required
配置参数说明
参数 |
是否必填 |
说明 |
|
是 |
固定为 |
|
是 |
REST Catalog 端点 URL。格式: |
|
是 |
Table Bucket ARN。格式: |
|
是 |
固定为 |
|
是 |
固定为 |
|
是 |
固定为 |
|
是 |
OSS 数据面端点。格式: |
|
是 |
固定为 |
权限配置
使用 RAM 用户或 STS 临时凭证访问 OSS Tables 时,需确保对应身份具备所需的操作权限。
资源定义
- Table Bucket ARN:
acs:osstables:<Region>:<阿里云账号ID>:bucket/<bucket_name> - Table ARN:
acs:osstables:<Region>:<阿里云账号ID>:bucket/<bucket_name>/table/<table_id>
Action 定义
下表列出 OSS Tables 支持的 Action,及其是否支持跨账号授权:
分类 |
Action |
跨账号访问 |
Table Bucket 级别 |
|
不允许 |
|
允许 |
|
|
不允许 |
|
|
允许 |
|
|
允许 |
|
|
允许 |
|
|
允许 |
|
|
允许 |
|
|
不允许 |
|
|
不允许 |
|
|
不允许 |
|
|
允许 |
|
|
允许 |
|
|
不允许 |
|
|
不允许 |
|
|
不允许 |
|
Table 级别 |
|
允许 |
|
允许 |
|
|
不允许 |
|
|
不允许 |
|
|
不允许 |
|
|
允许 |
|
|
允许 |
|
|
允许 |
|
|
允许 |
|
|
允许 |
|
|
允许 |
|
|
允许 |
|
|
允许 |
|
|
不允许 |
|
|
不允许 |
|
|
允许 |
Iceberg REST操作与权限映射
下表列出 Iceberg REST Catalog 各操作所需的 OSS Action:
Iceberg REST 操作 |
所需 OSS Action |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|