Flink 算完的数据存哪?写进 OSS Tables,实时离线读同一张表【详解 OSS Tables 系列】
实时计算之后数据该去哪在 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.jar3.3.6Hadoop API 依赖Iceberg 内部加载需要版本可按需调整。hadoop-client-runtime-3.3.6.jar3.3.6Hadoop 运行时依赖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_CALCULATIONWHEN_REQUIRED重要如果使用较高版本的 AWS SDK2.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.requestChecksumCalculationwhen_required -Daws.responseChecksumValidationwhen_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 创建。方式一使用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 {bucketArn} --namespace {namespace名称} --name {表名称} --format ICEBERG --metadata file://{文件路径} --endpoint --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: {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 CatalogOSS Tables 提供 Iceberg REST Catalog 端点Flink 通过 Iceberg Connector 的 RESTCatalog 实现连接。Endpoint格式如下内网https://{region}-internal.oss-tables.aliyuncs.com/iceberg外网https://{region}.oss-tables.aliyuncs.com/icebergOSS Tables 提供S3FileIO访问OSS数据面使用的访问端点Flink 通过该端点访问表数据。Endpoint格式如下内网https://oss-{region}-internal.aliyuncs.com外网https://oss-{region}.aliyuncs.com创建Catalog在 Flink SQL 中执行以下语句创建 CatalogCREATE 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。格式内网https://{region}-internal.oss-tables.aliyuncs.com/iceberg外网https://{region}.oss-tables.aliyuncs.com/icebergwarehouse是Table Bucket ARN。格式acs:osstables:Region:阿里云账号ID:bucket/Table Bucket名称。rest.sigv4-enabled是固定为true启用 SigV4 签名认证。rest.signing-name是固定为osstablesOSS Tables 服务的 SigV4 签名服务名。rest.signing-region是SigV4 签名地域。例如cn-hangzhou。s3.endpoint是OSS 数据面端点。格式内网https://oss-{region}-internal.aliyuncs.com外网https://oss-{region}.aliyuncs.coms3.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 MB33554432避免产生过大的文件影响后续维护任务的效率。步骤四写入并查询数据完成建表后在 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-001和evt-002两条记录。如果建表成功但写入或查询失败请优先检查 S3FileIO Endpoint、凭证和数据面权限。权限配置使用 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

相关新闻

做线上模型的TAIR缓存技巧

做线上模型的TAIR缓存技巧

部署了一个线上模型, 可以用线上日志记录这个模型的预测结果, 然后把日志记录的预测结果直接灌入TAIR做成缓存啊, 不用离线用模型推理刷之后再灌入TAIR啊。 当然离线可以用大模型推理,而日志记录的是线上小模型的结果&#xf…

2026/8/26 18:53:33 阅读更多 →
算法(40):separate chaining-12.2

算法(40):separate chaining-12.2

Page 17(碰撞的现实)物理事实:即使哈希函数设计得再好,碰撞(两个不同键算出同一个索引)也必然会发生。这是由生日悖论(Birthday Problem)决定的——当插入的元素数量超过 ~√(πM/2)…

2026/8/26 15:37:53 阅读更多 →
离散制造工艺路线动态建模实践:基于JVS-APS三层柔性机制的技术实现

离散制造工艺路线动态建模实践:基于JVS-APS三层柔性机制的技术实现

本文以技术视角解析离散制造中PMC与车间协同失效的根因,聚焦工艺路线建模的结构性缺陷。通过JVS-APS平台的三层柔性机制(模板主干扩展属性人员事件),给出可落地的技术实现路径:包括模板配置策略、字段级扩展绑定方法、…

2026/8/26 15:33:33 阅读更多 →

最新新闻

Final2x:开源免费的多模型图片无损放大工具

Final2x:开源免费的多模型图片无损放大工具

Final2x:开源免费的多模型图片无损放大工具内置 Real-ESRGAN 等几十种 AI 模型,把模糊老照片、低分辨率图片放大变清晰,支持 Windows、macOS、Linux。📖 背景说明 Final2x 是一款开源的图片放大工具,内置了 Real-ESRGA…

2026/8/26 18:53:29 阅读更多 →
企微实战排查之“获取用户信息失败” - 排除后端因素

企微实战排查之“获取用户信息失败” - 排除后端因素

前言 1. 企业微信调试工具快捷键:Ctrl Alt Shift D / Command Shift Control D 2. 不使用开发者工具排查问题。通过 进程->接口->网络->secret->日志 五步定位 一、定位问题 1.1 确认进程和端口 预期:端口处于LISTEN # 进程 ps -…

2026/8/26 18:53:29 阅读更多 →
Linux open 函数 Flag 参数详解

Linux open 函数 Flag 参数详解

一、open 函数原型与核心认知int open(const char *pathname, int flags, mode_t mode);很多新手最大误区:误以为 mode 是文件固定权限、flags 可以随意填写。实际内核执行逻辑非常严格,两个参数分工完全不同。核心规则:flags:控制…

2026/8/26 18:53:29 阅读更多 →
Nacos‑Client 与 Nacos‑Server

Nacos‑Client 与 Nacos‑Server

一句话总结:Nacos‑Server 是独立运行的服务端;Nacos‑Client 是嵌入业务项目里的客户端 SDK,两者配合完成注册中心 配置中心能力。表格对比项Nacos‑ServerNacos‑Client角色服务端(注册 & 配置的中央服务器)客户…

2026/8/26 18:53:29 阅读更多 →
海思平台不依赖频道文件实现 DVB 直接播放

海思平台不依赖频道文件实现 DVB 直接播放

1. 频道文件与直接调谐的区别传统方案的流程是:搜索频道↓ 导出频道文件↓ 系统预置频道文件↓ 开机导入↓ 根据文件中的 PID 播放频道文件通常包含:频率;符号率;调制方式;节目号;视频 PID;音频…

2026/8/26 18:52:29 阅读更多 →
海思 DVB 播放超时与 C 语言同名函数冲突排查

海思 DVB 播放超时与 C 语言同名函数冲突排查

概要本文记录一次海思 DVB 直接播放失败的定位过程。程序最初出现 HI_UNF_TUNER_SetAttr failed,修正前端参数后又出现 HIADP_Search_GetAllPmt failed 和 DMX_DataRead timeout。进一步分析发现,unf_abs.c 与 hi_adp_demux.c 中存在同名函数&#xff0c…

2026/8/26 18:52:29 阅读更多 →

日新闻

Python random 模块常用函数详解:从入门到实战

Python random 模块常用函数详解:从入门到实战

目录 1. 引言2. 准备工作3. 基础随机函数4. 序列相关函数5. 随机种子与复现6. 实战案例7. 注意事项8. 常见问题与排查9. 总结 1. 引言 摘要: 本文系统介绍 Python 标准库 random 模块中最常用的随机数生成函数。内容涵盖基础随机函数(random()、unifor…

2026/8/26 0:00:40 阅读更多 →
《Microsoft Sql server 2008 Internals》读书笔记--第三章Databases and Database Files(2)

《Microsoft Sql server 2008 Internals》读书笔记--第三章Databases and Database Files(2)

《Microsoft Sql server 2008 Internals》索引目录: 《Microsoft Sql server 2008 Internals》读书笔记--目录索引 在上篇文章中,主要介绍了创建数据库的基本语法和FileGroup的初步知识。需要注意的是: 关于FileGroup 如果你的系统是用Raid设备直接存…

2026/8/26 1:18:18 阅读更多 →
政务AI智能体怎么建?三种模式、三步路径与四个误区

政务AI智能体怎么建?三种模式、三步路径与四个误区

政务AI智能体已经从概念试点阶段,转入了政务服务的常态化落地应用;在实际使用过程中,它能自主理解办事需求、辅助完成填报申报、开展材料预审,并联动多个系统协同作业,真正嵌入到政务办理的全流程当中。但在落地推进过…

2026/8/26 1:18:18 阅读更多 →

周新闻

[光学原理与应用-521]:对光的错误理解与纠偏

[光学原理与应用-521]:对光的错误理解与纠偏

首先光是一种能量的载体和形态,宏观上观察到的光是由无数个微观的光量子组成的,每个光子在产生的瞬间,其在真空的空间中以确定不变的速度沿着一个初始的方向一直向前,在微观层面,每个光量子的运动轨迹是以波函数所展现…

2026/8/26 14:45:33 阅读更多 →
SIP通话转接原理与REFER方法实战解析

SIP通话转接原理与REFER方法实战解析

1. 通话转接不是“挂断再拨号”,而是SIP会话的动态重定向你有没有遇到过这样的场景:客服坐席A正在和客户通电话,突然需要把这通对话无缝转给专家坐席B,客户完全感知不到中间的断连——既没听到忙音,也没被要求重新拨号…

2026/8/26 17:46:43 阅读更多 →
Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

1. 为什么选择Kolla-ansible来部署单节点OpenStack?如果你正在寻找一种能把OpenStack从“概念”快速变成“可用的实验环境”的方法,那么Kolla-ansible几乎是当前最主流、最省心的选择。我见过太多人卡在手动编译依赖、配置服务、处理版本冲突的泥潭里&am…

2026/8/26 14:46:37 阅读更多 →

月新闻

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南 【免费下载链接】BaiduNetdiskPlugin-macOS For macOS.百度网盘 破解SVIP、下载速度限制~ 项目地址: https://gitcode.com/gh_mirrors/ba/BaiduNetdiskPlugin-macOS 还在为百度网盘macOS版的龟速下…

2026/8/26 3:50:20 阅读更多 →
终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换 【免费下载链接】ncmdump 项目地址: https://gitcode.com/gh_mirrors/ncmd/ncmdump 还在为网易云音乐下载的NCM格式文件无法在其他播放器播放而烦恼吗?ncmdump解密工具帮你轻松解决这个困…

2026/8/26 17:46:39 阅读更多 →
HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

AgentCard 智能体卡片:为英语学习 App 打造桌面级学习助手适用平台:HarmonyOS 7.0 (API 26 Beta)一、引言 HarmonyOS 7.0(API 26 Beta)新增了 AgentCard 智能体卡片能力,这是继 HMAF(鸿蒙智能体框架&#x…

2026/8/26 1:24:05 阅读更多 →