Flink 算完的数据存哪?写进 OSS Tables,实时离线读同一张表【详解 OSS Tables 系列】

简介: OSS Tables 是阿里云推出的 Iceberg 数据湖服务,原生兼容 REST Catalog 协议,让 Flink 无需额外部署元数据服务即可直连建表、写入与查询,实现 Exactly-Once 流式入湖,简化实时计算后的数据归宿问题。

实时计算之后,数据该去哪

在 OSS Tables 的生态里,Flink 是实时加工厂,对流数据做清洗、聚合之后再落表,算完即存。Flink 擅长实时 ETL、流式聚合这些活,这没什么争议。但作业跑完之后,处理过的数据落在哪里?OSS Tables 就是那个终点,后面的引擎读的是同一张表。


当前很多团队的实践是,Flink 处理完直接写回 Kafka 或者写入 MySQL、ES 之类的在线系统。这对实时服务来说没问题,但如果你还需要对处理后的数据做离线分析、回溯历史状态、或者给 BI 和机器学习用,就需要一个更“重”的存储来兜底。


数据湖是最常见的选择。问题是,Flink 写数据湖通常需要自己管理状态一致性、文件提交、Checkpoint 对齐等一系列细节。Iceberg 社区提供了 Flink Sink,但你还得有一个 Catalog 服务来管理表的元数据,传统方案是 Hive Metastore,或者另起一个 REST Catalog 服务。


OSS Tables 把这层省掉了。它原生兼容 Iceberg REST Catalog 协议,Flink 通过标准的 Iceberg Flink Runtime 就能直连。建表用 Flink SQL,写入用 Flink SQL,查询也用 Flink SQL。不需要额外部署元数据服务,Flink 作业的 Checkpoint 机制和 Iceberg 的事务语义天然对齐,保证 Exactly-Once。


本文介绍通过 Iceberg REST Catalog(兼容标准 Iceberg REST 协议、无需额外部署)从 Flink 访问 OSS Tables,使用 Flink SQL 进行建表、写入和查询操作,实现流式或批量数据入湖。本文以 Apache Flink 1.20 为例。

步骤一:环境准备

下载依赖JAR包

将以下 JAR 包放入 Flink 的 $FLINK_HOME/lib 目录,或在提交作业时通过 -C 参数指定。

JAR包

版本要求

说明

iceberg-flink-runtime-1.20-1.10.1.jar

匹配iceberg版本

Iceberg 的 Flink 运行时集成包。请根据 Flink 版本选择对应的包(如 Flink 1.20 对应 iceberg-flink-runtime-1.20)。

iceberg-aws-bundle-1.10.1.jar

匹配iceberg版本

提供 S3FileIO 实现及 REST Catalog SigV4 签名认证所需的 AWS SDK。版本需与 Runtime 包一致。

hadoop-client-api-3.3.6.jar

3.3.6

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

hadoop-client-runtime-3.3.6.jar

3.3.6

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

配置访问凭证

配置环境变量

Iceberg REST Catalog 使用 SigV4 签名认证,S3FileIO 访问数据面也需要凭证。推荐通过环境变量统一传递,在启动 Flink 之前设置以下环境变量:

说明环境变量名使用 AWS_ 前缀,是因为 Iceberg 的 SigV4 签名模块和 S3FileIO 复用 AWS SDK 的标准凭证链。实际填入的是您阿里云账号的 AccessKey ID 和 AccessKey Secret。

export AWS_ACCESS_KEY_ID=<阿里云AccessKey ID>
export AWS_SECRET_ACCESS_KEY=<阿里云AccessKey Secret>
export AWS_REGION=<地域,例如cn-hangzhou>
export AWS_DEFAULT_REGION=<地域,例如cn-hangzhou>
# 可选,使用STS临时凭证时配置
export AWS_SESSION_TOKEN=<阿里云STS TOKEN>
# 关闭STREAMING-UNSIGNED-PAYLOAD-TRAILER分块上传编码,OSS暂不支持
export AWS_REQUEST_CHECKSUM_CALCULATION=WHEN_REQUIRED

重要如果使用较高版本的 AWS SDK(2.20+),写入数据时可能出现签名错误:aws-chunked encoding is not supported with the specified x-amz-content-sha256 value。此时需要在 Flink 配置文件 conf/config.yaml 中添加以下 JVM 参数:

env.java.opts.all: "-Daws.requestChecksumCalculation=when_required -Daws.responseChecksumValidation=when_required"

通过配置项传递

除通过环境变量传递凭证外,也可以在创建 Catalog 时通过配置项显式传递凭证。该方式更适用于多 Catalog 场景,或不便为 Flink 进程统一注入环境变量的情况。在步骤三创建 Catalog 的 WITH 子句中增加以下配置项:

-- REST Catalog SigV4 签名凭证
  'rest.access-key-id'     = '<阿里云AccessKey ID>',
  'rest.secret-access-key' = '<阿里云AccessKey Secret>',
  -- 可选,使用STS临时凭证时配置
  'rest.session-token'     = '<阿里云STS TOKEN>',
  -- S3FileIO 数据面凭证
  's3.access-key-id'       = '<阿里云AccessKey ID>',
  's3.secret-access-key'   = '<阿里云AccessKey Secret>',
  'client.region'          = '<地域,例如cn-hangzhou>',
  -- 可选,使用STS临时凭证时配置
  's3.session-token'       = '<阿里云STS TOKEN>'

重要该方式仅传递凭证。若未设置 AWS_REQUEST_CHECKSUM_CALCULATION 环境变量,仍需按上文说明在 conf/config.yaml 中通过 env.java.opts.all 添加 JVM 参数关闭分块上传编码。

步骤二:创建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 {bucketArn} --namespace {namespace名称} --name  {表名称} --format ICEBERG --metadata file://{文件路径} --endpoint --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": "{BucketArn}",
    "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}}}'

步骤三:配置Flink

配置Flink进程参数

以下 flink-conf.yaml 配置为可选项,用于显式指定 S3FileIO 的 Region 和凭证提供方式。已按步骤一配置环境变量(含 AWS_REGION)时无需配置:

参数

是否必填

说明

s3.endpoint.region

S3FileIO 使用的地域。例如 cn-hangzhou

fs.s3a.aws.credentials.provider

凭证提供方式。固定为 software.amazon.awssdk.auth.credentials.EnvironmentVariableCredentialsProvider,从环境变量读取。

fs.s3a.endpoint.region

Hadoop S3A 文件系统使用的地域。例如 cn-hangzhou

通过 Iceberg REST Catalog

OSS Tables 提供 Iceberg REST Catalog 端点,Flink 通过 Iceberg Connector 的 RESTCatalog 实现连接。Endpoint格式如下:

OSS Tables 提供S3FileIO访问OSS数据面使用的访问端点,Flink 通过该端点访问表数据。Endpoint格式如下:

创建Catalog

在 Flink SQL 中执行以下语句创建 Catalog:

CREATE CATALOG mycatalog WITH (
  'type'                 = 'iceberg',
  'catalog-impl'         = 'org.apache.iceberg.rest.RESTCatalog',
  'io-impl'              = 'org.apache.iceberg.aws.s3.S3FileIO',
  'uri'                  = 'https://{region}-internal.oss-tables.aliyuncs.com/iceberg',
  'warehouse'            = '<Table Bucket ARN>',
  'rest.sigv4-enabled'   = 'true',
  'rest.signing-name'    = 'osstables',
  'rest.signing-region'  = '<Region>',
  's3.endpoint'          = 'https://oss-{region}-internal.aliyuncs.com',
  's3.path-style-access' = 'false'
);

配置参数说明

参数

是否必填

说明

type

固定为 iceberg,使用 Iceberg Connector。

catalog-impl

固定为 org.apache.iceberg.rest.RESTCatalog,指定使用 REST Catalog。

io-impl

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

uri

REST Catalog 端点 URL。格式:

warehouse

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

rest.sigv4-enabled

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

rest.signing-name

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

rest.signing-region

SigV4 签名地域。例如 cn-hangzhou

s3.endpoint

OSS 数据面端点。格式:

s3.path-style-access

是否使用 Path-Style 访问模式。默认为 false

建表示例

CREATE TABLE IF NOT EXISTS mycatalog.<Namespace>.<表名> (
  event_id    STRING,
  event_time  STRING,
  event_type  STRING
) WITH (
  'format-version'                     = '2',
  'write.format.default'               = 'parquet',
  'write.target-file-size-bytes'       = '33554432',        -- 32 MB(不要 128MB)
  'write.parquet.row-group-size-bytes' = '8388608'          -- 8 MB
);

说明建议将 write.target-file-size-bytes 设置为 32 MB(33554432),避免产生过大的文件影响后续维护任务的效率。

步骤四:写入并查询数据

完成建表后,在 Flink SQL Client 中依次执行以下语句。建表成功表示元数据链路已连通;测试数据能够写入并查询返回,表示数据链路已连通。通过 table.dml-sync 等待写入作业完成后再执行查询,适用于使用 SQL 文件进行非交互式验证的场景。

写入测试数据

SET 'table.dml-sync' = 'true';
SET 'sql-client.execution.result-mode' = 'TABLEAU';
INSERT INTO mycatalog.<Namespace>.<表名> VALUES
  ('evt-001', '2026-07-28 10:00:00', 'click'),
  ('evt-002', '2026-07-28 10:01:00', 'view');

查询并验证结果

-- 切换为批执行模式。流模式下对非时间属性字段执行 ORDER BY 会报错
SET 'execution.runtime-mode' = 'batch';
SELECT event_id, event_time, event_type
FROM mycatalog.<Namespace>.<表名>
ORDER BY event_id;

查询结果应包含 evt-001evt-002 两条记录。如果建表成功但写入或查询失败,请优先检查 S3FileIO Endpoint、凭证和数据面权限。

权限配置

使用 RAM 用户或 STS 临时凭证访问 OSS Tables 时,需确保对应身份具备所需的操作权限。

资源定义

  • Table Bucket ARN:acs:osstables::<阿里云账号ID>:bucket/
  • Table ARN:acs:osstables::<阿里云账号ID>:bucket//table/

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


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

热门文章

最新文章