流数据还在写文件慢慢导?Kafka Connect 一条链路直接入湖【详解 OSS Tables 系列】

简介: 本文详解如何通过Kafka Connect Iceberg Sink Connector,将Kafka流数据实时写入阿里云OSS Tables数据湖。涵盖环境准备、Table Bucket/命名空间/表创建、Connector配置及权限设置,实现免运维、Exactly-Once的流式入湖。

流数据的归宿问题

在 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-aws-bundle-1.10.1.jar

匹配iceberg版本

提供 S3FileIO 实现及 REST Catalog SigV4 签名认证所需的 AWS SDK。

iceberg-aws-1.10.1.jar

匹配iceberg版本

提供 SigV4 签名和 S3FileIO 实现。版本需与 iceberg-aws-bundle 一致。

iceberg-parquet-1.10.1.jar

匹配iceberg版本

Parquet 文件格式写入支持。

hadoop-client-runtime-3.3.6.jar

3.3.6

Hadoop 运行时依赖(Iceberg 内部加载需要),版本可按需调整。

hadoop-client-api-3.3.6.jar

3.3.6

Hadoop API 依赖(Iceberg 内部加载需要),版本可按需调整。

failsafe-3.3.2.jar

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格式如下:

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

配置参数说明

参数

是否必填

说明

iceberg.catalog.type

固定为 rest,指定使用 REST Catalog。

iceberg.catalog.uri

REST Catalog 端点 URL。格式:

iceberg.catalog.warehouse

Table Bucket ARN。格式:acs:osstables:<Region>:<阿里云账号ID>:bucket/<Table Bucket名称>

iceberg.catalog.rest.sigv4-enabled

固定为 true,启用 SigV4 签名认证。

iceberg.catalog.rest.signing-name

固定为 osstables,OSS Tables 服务端点的 SigV4 签名服务名。

iceberg.catalog.io-impl

固定为 org.apache.iceberg.aws.s3.S3FileIO,使用 S3 协议访问 OSS 数据面。

iceberg.catalog.s3.endpoint

OSS 数据面端点。格式:

iceberg.catalog.s3.path-style-access

固定为 true,使用 Path-Style 访问模式。

权限配置

使用 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 级别

oss:CreateTableBucket

不允许

oss:GetTableBucket

允许

oss:ListTableBuckets

不允许

oss:CreateNamespace

允许

oss:GetNamespace

允许

oss:ListNamespaces

允许

oss:DeleteNamespace

允许

oss:DeleteTableBucket

允许

oss:PutTableBucketPolicy

不允许

oss:GetTableBucketPolicy

不允许

oss:DeleteTableBucketPolicy

不允许

oss:GetTableBucketMaintenanceConfiguration

允许

oss:PutTableBucketMaintenanceConfiguration

允许

oss:PutTableBucketEncryption

不允许

oss:GetTableBucketEncryption

不允许

oss:DeleteTableBucketEncryption

不允许

Table 级别

oss:GetTableMaintenanceConfiguration

允许

oss:PutTableMaintenanceConfiguration

允许

oss:PutTablePolicy

不允许

oss:GetTablePolicy

不允许

oss:DeleteTablePolicy

不允许

oss:CreateTable

允许

oss:GetTable

允许

oss:GetTableMetadataLocation

允许

oss:ListTables

允许

oss:RenameTable

允许

oss:UpdateTableMetadataLocation

允许

oss:GetTableData

允许

oss:PutTableData

允许

oss:GetTableEncryption

不允许

oss:PutTableEncryption

不允许

oss:DeleteTable

允许

Iceberg REST操作与权限映射

下表列出 Iceberg REST Catalog 各操作所需的 OSS Action:

Iceberg REST 操作

所需 OSS Action

getConfig

oss:GetTableBucket

listNamespaces

oss:ListNamespaces

createNamespace

oss:CreateNamespace

loadNamespaceMetadata

oss:GetNamespace

dropNamespace

oss:DeleteNamespace

listTables

oss:ListTables

createTable

oss:CreateTableoss:PutTableData

loadTable

oss:GetTableMetadataLocationoss:GetTableData

updateTable

oss:UpdateTableMetadataLocationoss:PutTableDataoss:GetTableData

dropTable

oss:DeleteTable

renameTable

oss:RenameTable

tableExists

oss:GetTable

namespaceExists

oss:GetNamespace

相关文章
|
15天前
|
存储 人工智能 弹性计算
AgenticFS 重磅发布:专为 AI Agent 打造的全新文件存储
阿里云发布AgenticFS——全球首款Agent Native文件存储,专为AI Agent场景设计。单系统支持百万级AgenticSpace,提供目录级权限隔离、独立配额与性能QoS,峰值挂卸载达10万QPS,兼顾安全、弹性与超大规模扩展能力。(239字)
127 0
|
12天前
|
存储 人工智能 运维
大模型竞争进入下半场,存储正在成为 AI Infra 的决胜场
AI时代,存储正从“容量底座”跃升为“智能引擎”。阿里云发布CPFS、KVCacheStore、AgenticFS及OSS/CDE Agent等新品,覆盖训练、推理、Agent全生命周期——百PB级吞吐、亚毫秒KV缓存、百万级Agent空间隔离、自然语言智能运维,全面支撑AI原生与Agent原生演进。
124 2
|
2月前
|
存储 人工智能 自然语言处理
阿里云盘企业版 CDE Agent 正式发布!企业网盘会思考、能办事,存储即智能
阿里云盘企业版推出CDE Agent,依托Qoder Cloud Agent实现网盘内文件智能处理与内容生成。文件存入即成“可对话、可执行”的智能资产,支持多模态检索、跨格式分析、自动总结与安全分享,全程数据不出域、权限原生继承、操作全程审计,让网盘升级为安全合规的AI智能工作空间。
428 0
|
3月前
|
存储 人工智能 运维
3 人团队零推广获 1.2 万用户:Matrees 如何用 OSS 向量 Bucket 低成本构建 AI 创作平台
Matrees 是 Z 世代创作者的 AI 虚拟世界平台,3 人团队几乎零推广获取 1.2 万用户。依托阿里云 OSS Vector Bucket 实现全托管向量检索,成本降 90%,让创作者专注构建虚拟世界。
388 1
|
3月前
|
存储 人工智能 运维
拍封面,识唱片:UNHEARD 携手阿里云向量 Bucket,用 AI 重新定义实体唱片发现体验
UNHEARD 是一款专为黑胶爱好者打造的数字工具,依托阿里云 OSS 向量 Bucket 与百炼多模态模型,首创“以图搜碟”功能——拍封面即识专辑。集海量唱片数据库、唱片店地图、垂直社区于一体,助力用户轻松发现、识别、管理实体唱片。
811 0
拍封面,识唱片:UNHEARD 携手阿里云向量 Bucket,用 AI 重新定义实体唱片发现体验
|
3月前
|
存储 运维 数据管理
告别“大海捞针”:OSS Vector Bucket 如何赋能媒资管理平台
在 AI 时代,媒资平台面临多模态数据爆炸式增长的管理挑战。阿里云 OSS Vector Bucket 提供统一向量存储与语义检索能力,支持 30 亿级素材秒级精准查找,打破数据孤岛,降低成本,助力内容创作提效降本。
377 11
|
3月前
|
存储 自然语言处理 运维
阿里云发布 OSS Agent:对象存储的下一个交互方式,是自然语言
阿里云正式发布 OSS Agent——基于通义大模型的云存储智能体,兼具 7×24 小时智能运维专家与非结构化数据管理平台双重能力,支持自然语言创建 Bucket、异常诊断、健康巡检、费用分析及“ Talk to Bucket ”语义检索,让数据管理更简单高效。
661 0
|
19天前
|
存储 人工智能 安全
神眸 x 阿里云-低功耗芯片结合云端,构建百万规模智能摄像头平台
神眸是研极微电子推出的AI智能摄像机品牌,搭载自研超低功耗AI芯片,支持太阳能/电池供电,专为无电无网场景设计。依托阿里云OSS+Tablestore双引擎,提供高可靠云存储、端云协同AI分析(事件检测、智能摘要、自然语言检索等),兼顾安全、智能与易用,已跻身国内监控线上市场TOP5。(239字)
105 1
|
23天前
|
SQL 分布式计算 对象存储
还在自建 Hive Metastore?Spark 直连 OSS Tables 就能跑 Iceberg 【详解 OSS Tables 系列】
本文详解Spark对接OSS Tables数据湖的极简方案:无需自建Hive Metastore,原生兼容Iceberg REST Catalog协议,仅需配置JAR包、凭证及Catalog参数,即可通过标准SQL完成建表、读写、分区、Time Travel等全功能操作。
|
21天前
|
弹性计算 Cloud Native Linux
阿里云轻量应用服务器怎么部署网站?流程复杂吗?
阿里云轻量应用服务器支持一键部署网站:选择WordPress、宝塔等应用镜像,开箱即用,无需手动配置环境,新手分钟级建站。支持ICP备案,含登录地址、账号密码等完整指引。阿里云轻量应用服务器官网:https://t.aliyun.com/U/dwftch
134 4

热门文章

最新文章