☰
Apache Iceberg + Spark 快速上手指南:从 Docker Compose 到表读写与 Catalog 配置
2026/9/25 2:55:12 网站建设 项目流程
  • 数据湖
  • 大数据
  • 数据存储

【免费下载链接】iceberg

Apache Iceberg

项目地址:https://gitcode.com/gh_mirrors/icebe/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-icebergtabulario/spark-iceberg内置 Iceberg 的 Spark 集群,暴露 Notebook(8888)、Spark 端口等
restapache/iceberg-rest-fixture运行 Iceberg REST Catalog 服务,Spark 通过它管理元数据
minioquay.io/minio/minioS3 兼容对象存储,充当数据文件仓库
mcquay.io/minio/mcMinIO 客户端,负责在启动时初始化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.S3FileIO
  • CATALOG_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

首次启动会拉取相关镜像,之后可以打开四个入口:

入口命令 / 地址说明
SparkSQLdocker exec -it spark-iceberg spark-sql以 SQL 交互方式操作
Spark-Shelldocker exec -it spark-iceberg spark-shellScala 交互式环境
PySparkdocker exec -it spark-iceberg pysparkPython API
Jupyter Notebookhttp://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.typehadoop表示基于文件系统的路径型 Catalog
spark.sql.catalog.local.warehouseCatalog 的仓库根目录
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

项目地址:https://gitcode.com/gh_mirrors/icebe/iceberg
点击查看免费下载
上一篇:Delta查询与变更追踪:高效同步数据的REST API模式终极指南
下一篇:Windows 95模拟器终极输入处理指南:鼠标捕获与键盘事件的高级技巧

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询