
数据湖大数据数据存储【免费下载链接】icebergApache Iceberg项目地址https://gitcode.com/gh_mirrors/icebe/iceberg点击查看免费下载Apache Iceberg 作为开放的表格式需要各类计算引擎围绕它构建完整的流批一体读写链路。RisingWave 是一个 Postgres 兼容的流式 SQL 数据库它通过内置的 Iceberg source 与 sink 连接器支持对 Iceberg 表进行批量读取与流式写入从而把实时流处理与湖上历史数据无缝衔接。本文以 site/docs/integrations/risingwave.md 为核心结合仓库内 Iceberg 规范与配置文档带你完整掌握用 RisingWave 创建 Iceberg sink写入与 source读取的 SQL 实战方法以及其中涉及的 Catalog、仓库位置与参数约束。RisingWave 与 Iceberg 的集成能力概览RisingWave 是一款面向实时事件流数据处理、分析与管理的 Postgres 兼容 SQL 数据库能够持续摄入并分析实时数据流将流数据与历史表进行连续 Join并对外提供新鲜的、一致性的查询结果。在与 Apache Iceberg 的集成上RisingWave 通过其内置的 source 与 sink 连接器实现了两类核心能力批量读取Batch Read通过 source 连接器将 Iceberg 表作为可查询的数据源接入 RisingWave支持实时分析场景下对湖表数据的即时查询流式写入Streaming Write通过 sink 连接器把 RisingWave 中表或物化视图的结果持续写入 Iceberg 表形成实时入库、湖上沉淀的典型链路。仓库的 集成目录 与 供应商列表 均将 RisingWave 列为 Iceberg 生态的重要流式计算接入方。官方对 RisingWave 的定位是云原生流式数据库可用于来自消息队列、数据库通过 Change Data Capture、数据湖和文件等来源的数据并对 Iceberg 表实施高效的文件合并等维护操作。表格式与仓库位置的硬性约束在开始使用前需要明确 RisingWave 接入 Iceberg 时当前版本的两项关键限制表格式Table Format仅支持Iceberg V2 表格式。V2 是 Iceberg 规范中引入行级更新与删除能力的版本。仓库中的 格式规范 明确指出Version 2 为分析型表的不可变文件增加了行级更新和删除能力具体通过两类删除文件实现位置删除Position deletes按数据文件路径加行号标记被删除的行编码在位置删除文件中等式删除Equality deletes按一个或多个列值如id 5标记被删除的行编码在等式删除文件中。这也意味着如果目标 Iceberg 表是 V1 格式不支持行级删除则无法满足 RisingWave 基于主键的 upsert 写入需求。关于格式版本仓库的 配置文档 显示format-version的默认值自 Iceberg 1.4.0 起即为2因此新建表默认即可满足该约束。仓库位置Warehouse Location仅支持S3 兼容对象存储作为 Iceberg 数据仓库的存储位置例如 MinIO、Ceph RGW、SeaweedFS 等 S3 兼容服务以及各云厂商的 S3 协议对象存储。在连接参数中仓库路径使用s3a://前缀的 Hadoop S3A 协议表达例如s3a://hummock001/demo。支持的 Catalog 类型RisingWave 的 Iceberg 连接器支持以下五种 Catalog覆盖了从独立部署到云托管的主流场景Catalog 类型说明rest通过 Iceberg REST Catalog 协议访问元数据适合与各类 REST Catalog 服务对接jdbc/sql通过 JDBC 驱动访问基于关系型数据库如 PostgreSQL存储元数据的 CatalogglueAWS Glue Catalog适用于 AWS 生态storage基于对象存储目录结构Hadoop Catalog 风格直接管理元数据hiveHive Metastore Catalog适用于已有 Hive 元数据体系的场景在下面的两个实战示例中均使用catalog.type storage即直接以对象存储上的目录布局来组织和发现表元数据无需额外部署元数据服务是最便于本地验证与入门的方式。实战一通过 Sink 将数据写入 Iceberg 表RisingWave 中向 Iceberg 表写入数据的方式是创建一个sink。下面的示例将已有表或物化视图t1中的数据写入 Iceberg 表t1示例中源表写为t1目标表名为t1原文档以此示意从现有表或物化视图rw_data写入 Iceberg 表t1的形态CREATE SINK sink_to_iceberg FROM t1 WITH ( connector iceberg, type upsert, primary_key id, database.name demo_db, table.name t1, catalog.name demo, catalog.type storage, warehouse.path s3a://hummock001/demo, s3.endpoint http://127.0.0.1:9301, s3.region us-east-1, s3.access.key hummockadmin, s3.secret.key hummockadmin );Sink 参数逐项说明参数取值示例作用connectoriceberg声明使用 Iceberg 连接器typeupsert写入类型结合primary_key实现按主键的行级 upsert依赖 Iceberg V2 的行级删除能力primary_keyid用于 upsert 判断的主键列database.namedemo_db目标 Iceberg 表的数据库命名空间名table.namet1目标 Iceberg 表名catalog.namedemoCatalog 名称自定义逻辑名catalog.typestorageCatalog 类型可取rest、jdbc/sql、glue、storage、hivewarehouse.paths3a://hummock001/demoIceberg 仓库根路径使用 S3A 协议表达s3.endpointhttp://127.0.0.1:9301S3 兼容存储的服务端点s3.regionus-east-1S3 区域s3.access.keyhummockadminS3 访问密钥 Access Keys3.secret.keyhummockadminS3 访问密钥 Secret Key自动建表从 RisingWave 2.1 起可以在 sink 配置中加入create_table_if_not_exists参数当目标 Iceberg 表不存在时自动创建省去先建表的步骤CREATE SINK sink_to_iceberg FROM t1 WITH ( connector iceberg, type upsert, primary_key id, create_table_if_not_exists true, database.name demo_db, table.name t1, catalog.name demo, catalog.type storage, warehouse.path s3a://hummock001/demo, s3.endpoint http://127.0.0.1:9301, s3.region us-east-1, s3.access.key hummockadmin, s3.secret.key hummockadmin );底层原理佐证upsert 写入之所以必须落在 V2 表格式上根源在于 Iceberg 的不可变文件模型——数据文件一经写入不再修改行级更新只能通过追加删除文件position deletes / equality deletes来表达这正是 格式规范 中 V2 引入的两类删除编码。仓库中 Flink 侧的写入实践也遵循同样的约定例如 Flink DDL 示例 在创建支持 upsert 的表时显式声明format-version2与 RisingWave 的 V2 约束互为印证。实战二通过 Source 从 Iceberg 表读取数据从 Iceberg 表读取数据的方式是创建一个source。下面的示例从 Iceberg 表t1读取数据CREATE SOURCE iceberg_t1_source WITH ( connector iceberg, s3.endpoint http://127.0.0.1:9301, s3.region us-east-1, s3.access.key hummockadmin, s3.secret.key hummockadmin, s3.path.style.access true, catalog.type storage, warehouse.path s3a://hummock001/demo, database.name demo_db, table.name t1 );Source 参数要点相比 sink 配置source 的关键差异在于新增s3.path.style.access true启用 S3 的路径风格path-style访问方式。对于 MinIO、SeaweedFS 等自建 S3 兼容服务路径风格访问是必须开启的选项否则默认的虚拟主机风格virtual-hosted-style解析会导致访问失败无需type/primary_key读取侧不需要 upsert 语义因此不必声明写入类型与主键无需catalog.namesource 侧通过catalog.type、warehouse.path、database.name、table.name即可定位目标表。source 创建完成后即可用标准 SQL 查询 Iceberg 表中的数据SELECT * FROM iceberg_t1_source;参数取值与组合建议结合仓库内 Iceberg 的配置约定补充几点在真实环境中的参数组合建议s3.endpoint与s3.path.style.access必须搭配使用本地/私有 S3 兼容存储如示例中的127.0.0.1:9301几乎都需要开启 path-style接入 AWS S3 官方服务时endpoint 使用区域默认端点通常无需开启warehouse.path的s3a://前缀对应 Hadoop S3A 文件系统协议仓库路径指向 Iceberg 元数据与数据文件的根目录需保证 RisingWave 运行环境对该桶具备读写权限Catalog 选择与元数据服务强相关storage类型无独立元数据服务、自包含于对象存储目录适合快速验证生产环境若有现成的 Hive Metastore、REST Catalog 或 AWS Glue应分别选择hive、rest、glue并补充对应的连接参数如 REST 端点、JDBC URL 等具体以 RisingWave 官方 Iceberg catalog 文档为准表格式一致性写入侧要求目标表为 V2 格式若目标表由其他引擎创建请确认其format-version为2。仓库的 配置文档 显示 Iceberg 1.4.0 之后新表默认即 V2一般无需额外处理。小结RisingWave 通过内置的 Iceberg source 与 sink 连接器为流批一体的数据链路提供了便捷的 SQL 入口一方面用CREATE SINK把实时计算结果以 upsert 语义持续写入 Iceberg V2 表另一方面用CREATE SOURCE把湖上历史表直接接入流式查询。使用时需牢牢把握两个前提——Iceberg V2 表格式与 S3 兼容对象存储仓库位置并根据目标元数据服务选择rest、jdbc/sql、glue、storage、hive五种 Catalog 之一。结合本文的完整参数示例与仓库规范佐证你可以在本地 S3 兼容存储上快速搭建起 RisingWave ↔ Iceberg 的读写验证环境。赞分享数据湖大数据数据存储【免费下载链接】icebergApache Iceberg项目地址https://gitcode.com/gh_mirrors/icebe/iceberg点击查看免费下载相关推荐Apache SeaTunnel Iceberg Sink 连接器实战指南CDC 写入、多表写入与 Schema 演进Apache SeaTunnel Iceberg Sink 连接器实战指南CDC 写入、多表写入与 Schema 演进 Apache Iceberg Sink数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Redis 连接器实战指南Source 读取与 Sink 写入全解析SeaTunnel Redis 连接器实战指南Source 读取与 Sink 写入全解析 SeaTunnel 内置的 Redis 连接器Connector数据集成ETL大数据批处理流处理变更数据捕获ScyllaDB 与 Apache Kafka 集成实战指南Sink 连接器与 CDC Source 连接器全解析ScyllaDB 与 Apache Kafka 集成实战指南Sink 连接器与 CDC Source 连接器全解析 Apache Kafka 以其高吞吐、可扩数据库分布式数据库后端大数据上一篇prek用 Rust 重写的 Git 钩子管理器把提交流程缩短 10 倍的完整指南下一篇Dorado未来路线图了解Oxford Nanopore碱基识别技术的最新发展创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考