流数据还在写文件慢慢导?Kafka Connect 一条链路直接入湖【详解 OSS Tables 系列】 📅 发布时间:2026/8/20 12:18:23 👁 浏览次数: 流数据的归宿问题在 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.jar3.3.6Hadoop 运行时依赖Iceberg 内部加载需要版本可按需调整。hadoop-client-api-3.3.6.jar3.3.6Hadoop API 依赖Iceberg 内部加载需要版本可按需调整。failsafe-3.3.2.jar3.3.2Iceberg SnapshotProducer 运行时需要此依赖缺少会导致 ClassNotFoundException。步骤二创建Table Bucket在开始写入数据之前需要创建 Table Bucket 和 Namespace。可以使用 ossutil 创建AWS CLI 也支持。方式一使用ossutil1. 安装或升级 ossutil请安装ossutil 2.3.0以上版本如已安装 ossutil可执行以下命令升级到最新版本ossutil update -f2. 配置凭证执行ossutil config命令按提示输入 AccessKey ID、AccessKey Secret 和 Region。3. 创建 Table Bucketossutil tables-api create-table-bucket --name {table bucket名称} --endpoint http://{endpint} --region {region}命令执行成功后返回结果中包含 Table Bucket ARN请记录该值。4. 创建 Namespaceossutil 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 文件创建 Tableossutil tables-api create-table --table-bucket-arn Table Bucket ARN --namespace Namespace名称 --name 表名称 --format ICEBERG --metadata file://schema.json --endpoint http://{endpint}方式二使用 AWS CLIOSS Tables 兼容 S3 Tables API也可以使用 AWS CLI 管理 Table Bucket。1. 安装 AWS CLIcurl https://awscli.amazonaws.com/awscli-exe-linux-x86_64.zip -o awscliv2.zip unzip awscliv2.zip sudo ./aws/install2. 配置凭证执行aws configure命令按提示输入 AccessKey ID、AccessKey Secret 和 Region。3. 创建 Table Bucketaws s3tables --endpoint http://{endpint} create-table-bucket --region {region} --name {table bucket名称}命令执行成功后返回结果中包含 Table Bucket ARN。4. 创建 Namespaceaws 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 ConnectOSS Tables 提供 Iceberg REST Catalog 端点Kafka Connect 通过 Iceberg Sink Connector 连接该端点写入数据。Endpoint格式如下内网https://{region}-internal.oss-tables.aliyuncs.com/iceberg外网https://{region}.oss-tables.aliyuncs.com/icebergOSS Tables 提供S3FileIO访问OSS数据面使用的访问端点Spark 通过该端点访问表数据。Endpoint格式如下内网https://oss-{region}-internal.aliyuncs.com外网https://oss-{region}.aliyuncs.comConnector配置创建 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 SDK2.20可能出现签名错误aws-chunked encoding is not supported with the specified x-amz-content-sha256 value。此时需要在 Kafka Connect 的 Java 启动参数中添加以下 JVM 选项-Daws.requestChecksumCalculationwhen_required -Daws.responseChecksumValidationwhen_required配置参数说明参数是否必填说明iceberg.catalog.type是固定为rest指定使用 REST Catalog。iceberg.catalog.uri是REST Catalog 端点 URL。格式内网https://{region}-internal.oss-tables.aliyuncs.com/iceberg外网https://{region}.oss-tables.aliyuncs.com/icebergiceberg.catalog.warehouse是Table Bucket ARN。格式acs:osstables:Region:阿里云账号ID:bucket/Table Bucket名称。iceberg.catalog.rest.sigv4-enabled是固定为true启用 SigV4 签名认证。iceberg.catalog.rest.signing-name是固定为osstablesOSS Tables 服务端点的 SigV4 签名服务名。iceberg.catalog.io-impl是固定为org.apache.iceberg.aws.s3.S3FileIO使用 S3 协议访问 OSS 数据面。iceberg.catalog.s3.endpoint是OSS 数据面端点。格式内网https://oss-{region}-internal.aliyuncs.com外网https://oss-{region}.aliyuncs.comiceberg.catalog.s3.path-style-access是固定为true使用 Path-Style 访问模式。权限配置使用 RAM 用户或 STS 临时凭证访问 OSS Tables 时需确保对应身份具备所需的操作权限。资源定义Table Bucket ARNacs:osstables:Region:阿里云账号ID:bucket/bucket_nameTable ARNacs:osstables:Region:阿里云账号ID:bucket/bucket_name/table/table_idAction 定义下表列出 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 ActionIceberg REST 操作所需 OSS ActiongetConfigoss:GetTableBucketlistNamespacesoss:ListNamespacescreateNamespaceoss:CreateNamespaceloadNamespaceMetadataoss:GetNamespacedropNamespaceoss:DeleteNamespacelistTablesoss:ListTablescreateTableoss:CreateTable、oss:PutTableDataloadTableoss:GetTableMetadataLocation、oss:GetTableDataupdateTableoss:UpdateTableMetadataLocation、oss:PutTableData、oss:GetTableDatadropTableoss:DeleteTablerenameTableoss:RenameTabletableExistsoss:GetTablenamespaceExistsoss:GetNamespace