还在自建 Hive Metastore?Spark 直连 OSS Tables 就能跑 Iceberg 【详解 OSS Tables 系列】

简介: 本文详解Spark对接OSS Tables数据湖的极简方案:无需自建Hive Metastore,原生兼容Iceberg REST Catalog协议,仅需配置JAR包、凭证及Catalog参数,即可通过标准SQL完成建表、读写、分区、Time Travel等全功能操作。

一个老问题:Spark 读数据湖,到底需要几个组件

在 OSS Tables 的生态里,Spark 是全能主力,从建表、批量 ETL 到复杂分析都能做。但真动手的时候,你会发现“把数据存进去”和“让 Spark 能查出来”之间,还隔着不少工程活。

最典型的就是元数据管理。Spark 要读写湖里的表,需要一个 Catalog 服务来记录“哪些表存在、Schema 是什么、数据文件在哪”。传统方案是搭一套 Hive Metastore,配一个关系型数据库做后端。这套组合能跑,但运维成本高,而且每多一个组件就多一层故障面。

OSS Tables 做的事情很直接。它原生兼容 Apache Iceberg REST Catalog 协议,Spark 通过标准 Iceberg 客户端就能直连,不需要额外部署任何 Catalog 服务。建表、查询、写入、Time Travel,都是标准 SQL,代码里改一行 Catalog 配置就行。

这意味着:你的数据已经在 OSS 上了,Spark 只需要“接上” OSS Tables,就能组成完整的数据仓库解决方案。

下面用一个完整的端到端示例,演示如何用 Spark 接入 OSS Tables,完成建表、数据写入和时间旅行查询。

前提条件

  • 已安装Spark 3.5及以上版本,且Java版本为11或以上。
  • 已创建OSS Tables的Table Bucket。如未创建,请参见OSS Tables

步骤一:环境准备

下载依赖JAR包

将以下 JAR 包放入$SPARK_HOME/jars目录,或在启动时通过--jars参数指定。

JAR 包

说明

iceberg-spark-runtime-3.5_2.12-1.10.1.jar

Iceberg Spark Runtime 包,提供 Iceberg 的Spark集成能力。请根据Spark版本选择对应的包(如 Spark 3.5 对应iceberg-spark-runtime-3.5_2.12)。

iceberg-aws-bundle-1.10.1.jar

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

配置环境变量

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

说明

环境变量名使用 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>
# 可选,使用STS临时凭证时配置
export AWS_SESSION_TOKEN=<阿里云STS TOKEN>

重要:使用环境变量方式时,必须确保 Driver 和 Executor 都能读取到这些环境变量。在 YARN 模式下需要在集群各节点上注入环境变量;在 Kubernetes 模式下需要在 Pod 中注入。如果无法满足此条件,建议使用下方的 Spark 配置项方式显式指定凭证。

配置Spark 配置项

除了通过环境变量传递凭证外,也可以在 PySpark 或 Spark SQL 中通过 Spark Catalog 相关配置项显式传递凭证。该方式更适用于多 Catalog 场景,或不便设置环境变量的情况。

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

步骤二:配置Spark连接

  • OSS Tables提供Iceberg REST Catalog端点,Spark通过该端点管理表元数据。Endpoint格式如下:
  • OSS Tables 提供S3FileIO访问OSS数据面使用的访问端点,Spark 通过该端点访问表数据。Endpoint格式如下:

说明io-impl 请勿配置为 org.apache.iceberg.hadoop.HadoopFileIO

Iceberg 的设计理念是去 List 化,OSS Table Bucket 作为面向 Iceberg 深度优化的存储类型,禁止 List 操作。HadoopFileIO 基于对象存储的文件语义实现,为兼容通用场景会触发 List 操作,因此不适用于 OSS Table Bucket 的数据访问。

通过PySpark启动

以下示例以 Table Bucket 名称 my-data-lake、地域 cn-hangzhou 为例,对应的 Table Bucket ARN 为 acs:osstables:cn-hangzhou:{accountId}:bucket/my-data-lake

from pyspark.sql import SparkSession
spark = SparkSession.builder \
    .appName("OSS Tables Demo") \
    .config("spark.jars", "/path/to/iceberg-spark-runtime-3.5_2.12-1.10.1.jar,"
            "/path/to/iceberg-aws-bundle-1.10.1.jar") \
    .config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
    .config("spark.sql.catalog.oss_tables", "org.apache.iceberg.spark.SparkCatalog") \
    .config("spark.sql.catalog.oss_tables.catalog-impl", "org.apache.iceberg.rest.RESTCatalog") \
    .config("spark.sql.catalog.oss_tables.uri", "https://cn-hangzhou-internal.oss-tables.aliyuncs.com/iceberg") \
    .config("spark.sql.catalog.oss_tables.warehouse", "acs:osstables:cn-hangzhou:{accountId}:bucket/my-data-lake") \
    .config("spark.sql.catalog.oss_tables.rest.sigv4-enabled", "true") \
    .config("spark.sql.catalog.oss_tables.rest.signing-region", "cn-hangzhou") \
    .config("spark.sql.catalog.oss_tables.rest.signing-name", "osstables") \
    .config("spark.sql.catalog.oss_tables.io-impl", "org.apache.iceberg.aws.s3.S3FileIO") \
    .config("spark.sql.catalog.oss_tables.s3.endpoint", "https://oss-cn-hangzhou-internal.aliyuncs.com") \
    .getOrCreate()

通过spark-sql启动

spark-sql \
  --jars /path/to/iceberg-spark-runtime-3.5_2.12-1.10.1.jar,/path/to/iceberg-aws-bundle-1.10.1.jar \
  --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
  --conf spark.sql.catalog.oss_tables=org.apache.iceberg.spark.SparkCatalog \
  --conf spark.sql.catalog.oss_tables.catalog-impl=org.apache.iceberg.rest.RESTCatalog \
  --conf spark.sql.catalog.oss_tables.uri=https://cn-hangzhou-internal.oss-tables.aliyuncs.com/iceberg \
  --conf spark.sql.catalog.oss_tables.warehouse=acs:osstables:cn-hangzhou:{accountId}:bucket/my-data-lake \
  --conf spark.sql.catalog.oss_tables.rest.sigv4-enabled=true \
  --conf spark.sql.catalog.oss_tables.rest.signing-region=cn-hangzhou \
  --conf spark.sql.catalog.oss_tables.rest.signing-name=osstables \
  --conf spark.sql.catalog.oss_tables.io-impl=org.apache.iceberg.aws.s3.S3FileIO \
  --conf spark.sql.catalog.oss_tables.s3.endpoint=https://oss-cn-hangzhou-internal.aliyuncs.com

配置参数说明

参数

是否必填

说明

catalog-impl

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

uri

REST Catalog端点URL。格式:

warehouse

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

rest.sigv4-enabled

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

rest.signing-region

否(环境变量配置凭证时)

REST Catalog签名地域,需与Table Bucket所在地域一致,例如cn-hangzhou。使用AWS_REGION环境变量配置凭证时可省略;通过Spark配置项传递凭证时必填。

rest.signing-name

签名服务名,固定为osstables

io-impl

指定Iceberg读写底层数据文件的FileIO实现,固定为org.apache.iceberg.aws.s3.S3FileIO

s3.endpoint

S3FileIO访问OSS数据面使用的Endpoint,必须包含https://前缀。格式:

步骤三:使用SQL操作数据

连接成功后,使用标准SQL语句进行数据操作。

管理Namespace

Namespace(命名空间)用于对表进行逻辑分组,作用相当于数据库。

-- 查看现有Namespace
SHOW NAMESPACES IN oss_tables;
-- 创建Namespace
CREATE NAMESPACE oss_tables.my_namespace;
-- 删除Namespace(需先删除其中所有Table)
DROP NAMESPACE oss_tables.my_namespace;

建表与表管理

-- 创建非分区表
CREATE TABLE oss_tables.my_namespace.users (
    id BIGINT NOT NULL COMMENT '用户ID',
    name STRING COMMENT '用户名',
    email STRING COMMENT '邮箱',
    created_at TIMESTAMP COMMENT '创建时间'
) USING iceberg;
-- 创建分区表(按天分区)
CREATE TABLE oss_tables.my_namespace.events (
    id BIGINT NOT NULL,
    event_type STRING,
    data STRING,
    ts TIMESTAMP
) USING iceberg
PARTITIONED BY (days(ts));
-- 查看Namespace中的所有Table
SHOW TABLES IN oss_tables.my_namespace;
-- 查看表结构
DESCRIBE TABLE oss_tables.my_namespace.users;
-- 删除表(OSS Tables 要求必须带 PURGE 关键字,否则会报错:OSS Tables only supports dropping tables with purge enabled)
DROP TABLE oss_tables.my_namespace.users PURGE;

数据写入与查询

-- 插入数据
INSERT INTO oss_tables.my_namespace.users VALUES
    (1, '张三', 'zhangsan@example.com', TIMESTAMP '2024-01-15 10:30:00'),
    (2, '李四', 'lisi@example.com', TIMESTAMP '2024-01-16 14:20:00'),
    (3, '王五', 'wangwu@example.com', TIMESTAMP '2024-01-17 09:15:00');
-- 全表查询
SELECT * FROM oss_tables.my_namespace.users;
-- 条件查询
SELECT * FROM oss_tables.my_namespace.users WHERE id = 2;
-- 聚合查询
SELECT COUNT(*) AS total FROM oss_tables.my_namespace.users;
-- 分组聚合
SELECT name, COUNT(*) AS cnt FROM oss_tables.my_namespace.users GROUP BY name;
-- 更新数据
UPDATE oss_tables.my_namespace.users SET name = '赵六' WHERE id = 3;
-- 删除数据
DELETE FROM oss_tables.my_namespace.users WHERE id = 1;
-- 查询验证
SELECT * FROM oss_tables.my_namespace.users ORDER BY id;

分区表操作

-- 插入分区数据
INSERT INTO oss_tables.my_namespace.events VALUES
    (1, 'click', '{"page": "home"}', TIMESTAMP '2024-01-15 10:30:00'),
    (2, 'view', '{"page": "product"}', TIMESTAMP '2024-01-15 11:00:00'),
    (3, 'click', '{"page": "detail"}', TIMESTAMP '2024-01-16 09:00:00');
-- 分区裁剪查询(仅扫描匹配分区)
SELECT * FROM oss_tables.my_namespace.events
WHERE ts >= TIMESTAMP '2024-01-15 00:00:00'
  AND ts < TIMESTAMP '2024-01-16 00:00:00';
-- 聚合统计
SELECT event_type, COUNT(*) AS cnt
FROM oss_tables.my_namespace.events
GROUP BY event_type;

时间旅行查询

Iceberg 支持时间旅行(Time Travel)查询,可读取历史某个时间点的数据快照。

-- 查看快照历史
SELECT snapshot_id, committed_at, operation
FROM oss_tables.my_namespace.users.snapshots;
-- 基于快照ID查询历史数据
SELECT * FROM oss_tables.my_namespace.users
VERSION AS OF <snapshot_id>;
-- 查询指定时间点的数据
SELECT * FROM oss_tables.my_namespace.users
TIMESTAMP AS OF TIMESTAMP '2024-01-16 00:00:00';
-- 查看数据文件分布
SELECT * FROM oss_tables.my_namespace.users.files;

注意事项

  • 版本要求:推荐使用Spark 3.5+和Iceberg 1.10.1。Spark版本与iceberg-spark-runtime JAR包的版本号需要匹配(如Spark 3.5对应iceberg-spark-runtime-3.5_2.12)。
  • 凭证配置:S3FileIO 与 REST Catalog sigv4 签名使用同一对 AccessKey ID 和 AccessKey Secret。推荐通过以下环境变量统一配置(同时覆盖 REST Catalog 签名和 S3FileIO 数据访问):
  • AWS_ACCESS_KEY_ID:阿里云AccessKey ID
  • AWS_SECRET_ACCESS_KEY:阿里云AccessKey Secret
  • AWS_REGION:Table Bucket所在地域,用于REST Catalog签名。设置此环境变量后,Catalog配置中可省略rest.signing-region
  • 表格式:OSS Tables目前仅支持Iceberg格式。建表时必须指定USING iceberg
  • 数据维护:OSS Tables内置文件合并、快照清理和未引用文件清理功能,无需在Spark中手动执行Iceberg Maintenance Procedures,详情请参见数据维护
  • OSS Tables Endpoint 支持:
相关文章
|
2月前
|
存储 人工智能 自然语言处理
阿里云盘企业版 CDE Agent 正式发布!企业网盘会思考、能办事,存储即智能
阿里云盘企业版推出CDE Agent,依托Qoder Cloud Agent实现网盘内文件智能处理与内容生成。文件存入即成“可对话、可执行”的智能资产,支持多模态检索、跨格式分析、自动总结与安全分享,全程数据不出域、权限原生继承、操作全程审计,让网盘升级为安全合规的AI智能工作空间。
421 0
|
27天前
|
人工智能 安全 前端开发
基于 AgentScope 构建金融级智能体底座实战
金融级 AI 原生智能体底座白皮书发布。
|
22天前
|
前端开发 Java 数据库连接
Spring Boot 详细简介!
Spring Boot 是什么?能干啥?
224 0
Spring Boot 详细简介!
|
20天前
|
人工智能 Cloud Native 算法
从 VDBBench 到 MMEB 榜首:阿里云 AI Search 从引擎到模型的全栈优化
当 Agent 开始进入生产环境,搜索正在从“找到相关文档”的工具,升级为持续供给高质量上下文的基础设施。面向这一变化,AI Search 也不再只是一个搜索引擎,而是贯穿模型理解、向量与全文召回、混合检索、融合重排和引擎执行的完整链路。近期,阿里云 AI Search 在 VDBBench 向量检索、Big5 与 Tantivy 传统检索,以及 MMEB 多模态文档检索等评测中取得一系列突破。这些结果背后,是阿里云从 Elasticsearch 自研引擎到多模态 Embedding 模型的全栈优化。 关键词: AI Search、Agent、FalconSeek、向量检索、混合检索、多模态
236 0
从 VDBBench 到 MMEB 榜首:阿里云 AI Search 从引擎到模型的全栈优化
|
1月前
|
人工智能 运维 Linux
凌晨告警不再慌!SysOM 巡检 Skill 一键锁定根因
凌晨两点被叫醒,还要花 40 分钟拼出根因?阿里云操作系统控制台发布的 SysOM 巡检 Skill,沉淀了内核专家的排查经验,37 秒即可生成报告,巡检发现问题后自动衔接诊断、精准定位根因。目前 SysOM 巡检 Skill 已开源,一行命令即可立即上手,欢迎体验。
|
2月前
|
存储 人工智能 Kubernetes
阿里云 AgentTeams 解读:当 Agent 开始真正在企业里干活
多 Agent 协作不只是任务并行,更是组织运转。从产品主创团队视角,聊聊 AgentTeams 在安全、协作、弹性、进化四个方向的设计思考。
|
3月前
|
存储 人工智能 运维
拍封面,识唱片:UNHEARD 携手阿里云向量 Bucket,用 AI 重新定义实体唱片发现体验
UNHEARD 是一款专为黑胶爱好者打造的数字工具,依托阿里云 OSS 向量 Bucket 与百炼多模态模型,首创“以图搜碟”功能——拍封面即识专辑。集海量唱片数据库、唱片店地图、垂直社区于一体,助力用户轻松发现、识别、管理实体唱片。
806 0
拍封面,识唱片:UNHEARD 携手阿里云向量 Bucket,用 AI 重新定义实体唱片发现体验
|
27天前
|
人工智能 前端开发 安全
词帆 CiFan:基于故事分级阅读与 AI 伴学的英语学习平台
词帆(CiFan)是面向乡村儿童与外语初学者的公益AI伴学平台,基于GGU分级阅读体系,提供极轻量Web端体验。支持点击查词、AI长难句大白话解析、趣味串词故事生成、艾宾浩斯复习及多模态学习统计,适配低配设备
97 3
词帆 CiFan:基于故事分级阅读与 AI 伴学的英语学习平台
|
18天前
|
存储 人工智能 安全
神眸 x 阿里云-低功耗芯片结合云端,构建百万规模智能摄像头平台
神眸是研极微电子推出的AI智能摄像机品牌,搭载自研超低功耗AI芯片,支持太阳能/电池供电,专为无电无网场景设计。依托阿里云OSS+Tablestore双引擎,提供高可靠云存储、端云协同AI分析(事件检测、智能摘要、自然语言检索等),兼顾安全、智能与易用,已跻身国内监控线上市场TOP5。(239字)

热门文章

最新文章