- 数据湖
- 大数据
- 数据存储
【免费下载链接】iceberg
Apache Iceberg
Apache Iceberg™ 是面向大规模分析型数据湖的开源表格式,而 Apache Spark™ 是 Iceberg 最常用的计算引擎之一。本指南以仓库中的 spark-quickstart.md 为骨架,完整讲解如何用 Docker Compose 一键拉起带 Iceberg 的本地 Spark 集群,并逐步演示建表、写数、读数的完整流程,最后深入说明如何为已有 Spark 环境添加 Iceberg Catalog。读完本文,你将掌握一套可直接复现的 Iceberg + Spark 本地开发环境,以及 Spark SQL、Spark-Shell、PySpark 三种接口下的等价操作。
环境准备:你需要哪些工具
开始之前,请确认本机已安装:
- Docker CLI:容器运行时,用于拉取并启动镜像;
- Docker Compose CLI:用于按 YAML 定义编排多个容器。
两种方式都需要 Docker CLI,因为 Compose 依赖 Docker 引擎执行容器。
方式一:用 Docker Compose 快速启动全套环境
官方推荐的最快上手方式,是使用一份 docker-compose 文件拉起一个"开箱即用"的本地集群。该方案基于社区镜像tabulario/spark-iceberg,镜像内已经内置了一个配置好 Iceberg Catalog 的本地 Spark 集群,你无需手动下载任何 Iceberg jar 或修改 Spark 配置。
编写 docker-compose.yml
将下面的内容保存为docker-compose.yml:
services: spark-iceberg: image: tabulario/spark-iceberg container_name: spark-iceberg build: spark/ networks: iceberg_net: depends_on: - rest - minio volumes: - ./warehouse:/home/iceberg/warehouse - ./notebooks:/home/iceberg/notebooks/notebooks environment: - AWS_ACCESS_KEY_ID=admin - AWS_SECRET_ACCESS_KEY=password - AWS_REGION=us-east-1 ports: - 8888:8888 - 8080:8080 - 10000:10000 - 10001:10001 rest: image: apache/iceberg-rest-fixture container_name: iceberg-rest networks: iceberg_net: ports: - 8181:8181 environment: - AWS_ACCESS_KEY_ID=admin - AWS_SECRET_ACCESS_KEY=password - AWS_REGION=us-east-1 - CATALOG_WAREHOUSE=s3://warehouse/ - CATALOG_IO__IMPL=org.apache.iceberg.aws.s3.S3FileIO - CATALOG_S3_ENDPOINT=http://minio:9000 minio: image: quay.io/minio/minio container_name: minio environment: - MINIO_ROOT_USER=admin - MINIO_ROOT_PASSWORD=password - MINIO_DOMAIN=minio networks: iceberg_net: aliases: - warehouse.minio ports: - 9001:9001 - 9000:9000 command: ["server", "/data", "--console-address", ":9001"] mc: depends_on: - minio image: quay.io/minio/mc container_name: mc networks: iceberg_net: environment: - AWS_ACCESS_KEY_ID=admin - AWS_SECRET_ACCESS_KEY=password - AWS_REGION=us-east-1 entrypoint: | /bin/sh -c " until (/usr/bin/mc alias set minio http://minio:9000 admin password) do echo '...waiting...' && sleep 1; done; /usr/bin/mc rm -r --force minio/warehouse; /usr/bin/mc mb minio/warehouse; /usr/bin/mc policy set public minio/warehouse; tail -f /dev/null " networks: iceberg_net:这份编排文件包含 4 个角色,各司其职:
| 服务 | 镜像 | 作用 |
|---|---|---|
spark-iceberg | tabulario/spark-iceberg | 内置 Iceberg 的 Spark 集群,暴露 Notebook(8888)、Spark 端口等 |
rest | apache/iceberg-rest-fixture | 运行 Iceberg REST Catalog 服务,Spark 通过它管理元数据 |
minio | quay.io/minio/minio | S3 兼容对象存储,充当数据文件仓库 |
mc | quay.io/minio/mc | MinIO 客户端,负责在启动时初始化warehouse存储桶 |
容器角色与源码印证
rest容器对应的镜像构建文件位于 docker/iceberg-rest-fixture/Dockerfile,其中可以看到:
- 默认后端 Catalog 是
org.apache.iceberg.jdbc.JdbcCatalog,通过CATALOG_URI=jdbc:sqlite:/tmp/iceberg_catalog.db?journal_mode=WAL落在本地 SQLite; - 默认监听
REST_PORT=8181,并配置了健康检查(curl --fail http://localhost:$REST_PORT/v1/config)。
容器启动入口 docker/iceberg-rest-fixture/entrypoint.sh 最终执行org.apache.iceberg.rest.RESTCatalogServer这个主类。在源码 open-api/src/testFixtures/java/org/apache/iceberg/rest/RESTCatalogServer.java 中可以看到默认端口常量REST_PORT_DEFAULT = 8181,服务通过 Jetty 启动,并将请求路由到RESTServerCatalogAdapter,实际后端 Catalog 则由CatalogUtil.buildIcebergCatalog依据配置动态构建。
值得说明的是,rest服务读取的是一组以CATALOG_前缀命名的环境变量。在 open-api/src/testFixtures/java/org/apache/iceberg/rest/RCKUtils.java 中定义了这套转换规则:CATALOG_前缀被去掉,双下划线__替换为连字符-,单下划线_替换为点.,名称统一转为小写。例如:
CATALOG_WAREHOUSE=s3://warehouse/→warehouse=s3://warehouse/CATALOG_IO__IMPL=org.apache.iceberg.aws.s3.S3FileIO→io-impl=org.apache.iceberg.aws.s3.S3FileIOCATALOG_S3_ENDPOINT=http://minio:9000→s3.endpoint=http://minio:9000
这也是上例中CATALOG_IO__IMPL使用双下划线的根本原因——它对应的是 Iceberg 的io-impl属性,即数据文件读写由org.apache.iceberg.aws.s3.S3FileIO实现(该类位于 aws/src/main/java/org/apache/iceberg/aws/s3/S3FileIO.java),配合s3.endpoint指向容器内的 MinIO 服务(地址http://minio:9000),从而让 REST Catalog 的元数据与数据文件全部落在本地对象存储上。
mc容器的 entrypoint 负责在 MinIO 就绪后执行初始化:先轮询等待 MinIO 可用,然后清理并创建warehouse存储桶,并设置为公开策略。spark-iceberg容器通过depends_on保证在rest与minio之后启动,并将宿主机./warehouse目录挂载进容器,方便直接查看落盘的数据文件。
启动集群
在保存好docker-compose.yml的目录下执行:
docker-compose up首次启动会拉取相关镜像,之后可以打开四个入口:
| 入口 | 命令 / 地址 | 说明 |
|---|---|---|
| SparkSQL | docker exec -it spark-iceberg spark-sql | 以 SQL 交互方式操作 |
| Spark-Shell | docker exec -it spark-iceberg spark-shell | Scala 交互式环境 |
| PySpark | docker exec -it spark-iceberg pyspark | Python API |
| Jupyter Notebook | http://localhost:8888 | 图形化 Notebook 环境 |
!!! note 除了上述三种 CLI 入口,集群还自带了 Notebook 服务,浏览器访问http://localhost:8888即可使用,适合边写边跑实验。
创建你的第一张 Iceberg 表
以demo.nyc.taxis为例:demo是 Catalog 名,nyc是数据库名,taxis是表名。这是 Iceberg 标准的三段式命名(catalog.database.table)。
创建数据库
如果数据库尚不存在,先创建它:
=== "SparkSQL"
```sql CREATE DATABASE IF NOT EXISTS demo.nyc; ```=== "Spark-Shell"
```scala spark.sql("CREATE DATABASE IF NOT EXISTS demo.nyc") ```=== "PySpark"
```py spark.sql("CREATE DATABASE IF NOT EXISTS demo.nyc") ```创建分区表
使用CREATE TABLE ... PARTITIONED BY显式声明分区列(这里按vendor_id分区):
=== "SparkSQL"
```sql CREATE TABLE demo.nyc.taxis ( vendor_id bigint, trip_id bigint, trip_distance float, fare_amount double, store_and_fwd_flag string ) PARTITIONED BY (vendor_id); ```=== "Spark-Shell"
```scala import org.apache.spark.sql.types._ import org.apache.spark.sql.Row val schema = StructType( Array( StructField("vendor_id", LongType,true), StructField("trip_id", LongType,true), StructField("trip_distance", FloatType,true), StructField("fare_amount", DoubleType,true), StructField("store_and_fwd_flag", StringType,true) )) val df = spark.createDataFrame(spark.sparkContext.emptyRDD[Row],schema) df.writeTo("demo.nyc.taxis").create() ```=== "PySpark"
```py from pyspark.sql.types import DoubleType, FloatType, LongType, StructType,StructField, StringType schema = StructType([ StructField("vendor_id", LongType(), True), StructField("trip_id", LongType(), True), StructField("trip_distance", FloatType(), True), StructField("fare_amount", DoubleType(), True), StructField("store_and_fwd_flag", StringType(), True) ]) df = spark.createDataFrame([], schema) df.writeTo("demo.nyc.taxis").create() ```Spark-Shell 与 PySpark 使用了 DataFrame API 的writeTo("...").create()语法:先构造一个与目标表同构的空 DataFrame,再调用create()完成建表。这与 SQL 的CREATE TABLE等价。
Iceberg 的 Catalog 支持完整的 SQL DDL 能力,除建表外还包括:
CREATE TABLE ... AS SELECT:用查询结果直接建表;ALTER TABLE:变更表结构(如演进 schema、修改分区);DROP TABLE:删除表。
更完整的 DDL 语法参见 Spark DDL 文档。
向表中写入数据
建表完成后即可写入数据:
=== "SparkSQL"
```sql INSERT INTO demo.nyc.taxis VALUES (1, 1000371, 1.8, 15.32, 'N'), (2, 1000372, 2.5, 22.15, 'N'), (2, 1000373, 0.9, 9.01, 'N'), (1, 1000374, 8.4, 42.13, 'Y'); ```=== "Spark-Shell"
```scala import org.apache.spark.sql.Row val schema = spark.table("demo.nyc.taxis").schema val data = Seq( Row(1: Long, 1000371: Long, 1.8f: Float, 15.32: Double, "N": String), Row(2: Long, 1000372: Long, 2.5f: Float, 22.15: Double, "N": String), Row(2: Long, 1000373: Long, 0.9f: Float, 9.01: Double, "N": String), Row(1: Long, 1000374: Long, 8.4f: Float, 42.13: Double, "Y": String) ) val df = spark.createDataFrame(spark.sparkContext.parallelize(data), schema) df.writeTo("demo.nyc.taxis").append() ```=== "PySpark"
```py schema = spark.table("demo.nyc.taxis").schema data = [ (1, 1000371, 1.8, 15.32, "N"), (2, 1000372, 2.5, 22.15, "N"), (2, 1000373, 0.9, 9.01, "N"), (1, 1000374, 8.4, 42.13, "Y") ] df = spark.createDataFrame(data, schema) df.writeTo("demo.nyc.taxis").append() ```这里有一个实用技巧:DataFrame API 方式直接通过spark.table("demo.nyc.taxis").schema从已建好的表上复用 schema,避免手工重复定义字段类型。写入动作由append()完成,对应 SQL 中的INSERT INTO。
从表中读取数据
读取只需直接引用 Iceberg 表名即可:
=== "SparkSQL"
```sql SELECT * FROM demo.nyc.taxis; ```=== "Spark-Shell"
```scala val df = spark.table("demo.nyc.taxis").show() ```=== "PySpark"
```py df = spark.table("demo.nyc.taxis").show() ```Iceberg 表的读取对 Spark 完全透明,spark.table与 SQLSELECT均直接可用。查询下推、分区裁剪、Snapshot 隔离等能力由 Iceberg 的 Spark 集成层自动处理,更多查询特性可参考 Spark 查询文档。
为已有 Spark 添加 Catalog
Catalog 配置原理
Iceberg 支持多种 Catalog 后端来跟踪表,例如 JDBC、Hive Metastore、AWS Glue 等。Catalog 全部通过spark.sql.catalog.(catalog_name)前缀下的属性进行配置。在 Spark 侧,Iceberg 提供了两个 Catalog 实现类(参见 spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/SparkCatalog.java 及其 类文档):
org.apache.iceberg.spark.SparkCatalog:支持hive、hadoop、rest、glue、jdbc、nessie六种类型;org.apache.iceberg.spark.SparkSessionCatalog:给 Spark 内置 Catalog 增加 Iceberg 表支持,非 Iceberg 表则委托给内置 Catalog 处理。
SparkCatalog的初始化逻辑最终调用CatalogUtil.buildIcebergCatalog(name, options, conf)来构建底层 Iceberg Catalog 实例,因此 Spark 配置的属性会被透传为 Iceberg Catalog 属性。
CLI 方式配置(以 Hadoop 路径型 Catalog 为例)
下面的配置创建了一个名为local的基于路径的 Catalog,表数据存放在$PWD/warehouse下,同时为 Spark 内置 Catalog(spark_catalog)接入 Iceberg 支持:
spark-sql --packages org.apache.iceberg:iceberg-spark-runtime-{{ sparkVersionMajor }}:{{ icebergVersion }}\ --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \ --conf spark.sql.catalog.spark_catalog=org.apache.iceberg.spark.SparkSessionCatalog \ --conf spark.sql.catalog.spark_catalog.type=hive \ --conf spark.sql.catalog.local=org.apache.iceberg.spark.SparkCatalog \ --conf spark.sql.catalog.local.type=hadoop \ --conf spark.sql.catalog.local.warehouse=$PWD/warehouse \ --conf spark.sql.defaultCatalog=local其中{{ sparkVersionMajor }}与{{ icebergVersion }}是版本占位符,需要替换为实际版本。就本仓库支持的 Spark 版本而言,参见 gradle/libs.versions.toml 中的spark35 = "3.5.9"、spark40 = "4.0.4"、spark41 = "4.1.3"、spark42 = "4.2.0",对应的 runtime 坐标形如iceberg-spark-runtime-3.5、iceberg-spark-runtime-4.1。
上述参数的含义:
| 参数 | 说明 |
|---|---|
spark.sql.extensions | 注册 Iceberg 的 Spark 扩展,启用 Iceberg 专属 SQL 语法与优化 |
spark.sql.catalog.spark_catalog | 将 Spark 内置 Catalog 替换为SparkSessionCatalog,使其能同时处理 Iceberg 与非 Iceberg 表 |
spark.sql.catalog.spark_catalog.type | 底层类型,hive表示通过 Hive Metastore 跟踪元数据 |
spark.sql.catalog.local | 注册名为local的SparkCatalog |
spark.sql.catalog.local.type | hadoop表示基于文件系统的路径型 Catalog |
spark.sql.catalog.local.warehouse | Catalog 的仓库根目录 |
spark.sql.defaultCatalog | 将默认 Catalog 设为local |
spark-defaults.conf 方式配置
同样的配置也可以写入 Spark 的spark-defaults.conf,效果完全一致,适合持久化配置:
spark.jars.packages org.apache.iceberg:iceberg-spark-runtime-{{ sparkVersionMajor }}:{{ icebergVersion }} spark.sql.extensions org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions spark.sql.catalog.spark_catalog org.apache.iceberg.spark.SparkSessionCatalog spark.sql.catalog.spark_catalog.type hive spark.sql.catalog.local org.apache.iceberg.spark.SparkCatalog spark.sql.catalog.local.type hadoop spark.sql.catalog.local.warehouse $PWD/warehouse spark.sql.defaultCatalog local切换默认 Catalog
!!! note 如果 Iceberg Catalog 未被设为默认 Catalog(未配置spark.sql.defaultCatalog),需要手动切换到它:执行USE local;。
其他 Catalog 类型
除hadoop外,type还可取hive、rest、glue、jdbc、nessie。以 REST Catalog 为例,只需把type改为rest并指定uri:
spark.sql.catalog.rest_prod = org.apache.iceberg.spark.SparkCatalog spark.sql.catalog.rest_prod.type = rest spark.sql.catalog.rest_prod.uri = http://localhost:8080本文开头的 Docker 方案中,spark-iceberg容器实际上就是通过类似方式连接rest容器的 REST Catalog 与 MinIO 对象存储。更多 Catalog 配置参数(如cache-enabled、cache.expiration-interval-ms、table-default.*、table-override.*等)详见 Spark 配置文档。
进阶:把 Iceberg 装进已有的 Spark 环境
如果你已经有一套 Spark 环境,无需使用 Docker,直接用--packages参数即可在会话级引入 Iceberg:
=== "SparkSQL"
```sh spark-sql --packages org.apache.iceberg:iceberg-spark-runtime-{{ sparkVersionMajor }}:{{ icebergVersion }} ```=== "Spark-Shell"
```sh spark-shell --packages org.apache.iceberg:iceberg-spark-runtime-{{ sparkVersionMajor }}:{{ icebergVersion }} ```=== "PySpark"
```sh pyspark --packages org.apache.iceberg:iceberg-spark-runtime-{{ sparkVersionMajor }}:{{ icebergVersion }} ```!!! note 如果想在 Spark 安装目录中全局引入 Iceberg(所有会话默认可用),可以把 Iceberg Spark runtime 的 jar 放进 Spark 的jars目录,runtime 可从 Releases 页面 下载。使用--packages时记得同时带上spark.sql.extensions与spark.sql.catalog.*配置(见上文 CLI 示例),否则 Iceberg 语法与 Catalog 不会生效。
下一步:继续深入
完成上述步骤后,你已经掌握了 Iceberg + Spark 的核心操作闭环:起环境 → 建库建表 → 写数据 → 读数据 → 配 Catalog。接下来可以继续探索:
- Spark DDL 文档:完整的建表、ALTER、DROP 语法;
- Spark 配置文档:Catalog 参数、读写选项、运行时配置优先级;
- Spark 查询文档:时间旅行、Snapshot 读取等高级查询能力;
- Spark 写入文档:upsert、merge、流式写入等数据写入模式;
- Spark 维护文档:过期 Snapshot 清理、数据文件合并等表维护操作。
仓库中的 docker/iceberg-rest-fixture 与 docker/iceberg-flink-quickstart 还提供了其他引擎侧的快速开始环境,可以作为横向参考。
- 数据湖
- 大数据
- 数据存储
【免费下载链接】iceberg
Apache Iceberg
相关推荐
ESP-IDF USB HAL(esp_hal_usb)深度解析:USB 控制器与 PHY 的硬件抽象层
ESP IDF USB HAL(esp_hal_usb)深度解析:USB 控制器与 PHY 的硬件抽象层 ESP IDF 中的 esp_hal_usb 组件为所
数据湖大数据数据存储Apache Iceberg + Flink 快速上手实战:基于 Docker Compose 的 REST Catalog 环境搭建与 SQL 读写全流程
Apache Iceberg + Flink 快速上手实战:基于 Docker Compose 的 REST Catalog 环境搭建与 SQL 读写全流程 A
数据湖大数据数据存储Envoy Dynamic Modules 健康检查器:用共享库自定义上游健康检查的完整指南
Envoy Dynamic Modules 健康检查器:用共享库自定义上游健康检查的完整指南 导读 Envoy 的 Dynamic Modules 健康检查器(
数据湖大数据数据存储
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考