Apache Iceberg 与 Apache Hudi 虽然都叫数据湖表格式,但它们的容器化思路并不能直接复用同一套模板。Iceberg 把表状态收敛到 Catalog 和元数据文件中,计算侧几乎不需要长驻进程;Hudi 的更新、合并和清理逻辑则高度依赖时间线,一旦提交失败就可能出现重复数据或文件泄漏。把这两者放到 Kubernetes 里,首先要明确哪些组件属于无状态 API,哪些属于有状态批处理,否则 Pod 一重启,任务不会自动恢复,元数据也可能落到本地临时目录而丢失。

一、容器化前的架构差异:Iceberg 与 Hudi 的状态边界
Iceberg 的写入路径把元数据和数据文件严格分开。写入作业先把 Parquet 文件写到数据目录,接着生成清单文件和清单列表,最后原子性地提交一个 metadata.json 指针。Catalog 中只保存当前元数据文件的位置,并不会记录每一批数据的提交过程。这意味着即使写入作业所在 Pod 被强制终止,只要最后一个 metadata.json 没有替换成功,表状态就还停留在上一次有效快照,不会出现半提交。因此 Iceberg 非常适合做成短生命周期批处理,在 Kubernetes 中也无需为它维护专门的长驻协调进程。
Hudi 的情况则完全不同。它的所有状态变化都记录在表目录下的 .hoodie 时间线里,写入前先写入 requested 事件,执行中写入 inflight 事件,完成后才写入 commit 事件。清理、归档、压缩、聚簇这些后台动作同样要写时间线。如果 Pod 在 inflight 阶段突然重启,时间线上会残留未完成的事件,需要人为清理解锁。因此容器化 Hudi 时不能只考虑如何提交作业,还要把表服务拆成独立任务,并设计好失败重试路径。
下面的目录结构能直观展示两者的差异。Iceberg 的元数据集中在 metadata 目录,Hudi 的状态则分散在 .hoodie 下的多个事件文件中:
warehouse/demo/orders/metadata/v1.metadata.json warehouse/demo/orders/metadata/v2.metadata.json warehouse/demo/orders/metadata/snap-2345678901234567890.avro warehouse/demo/orders/data/ts_day=2024-01-01/part-00000.parquet warehouse/demo/orders/data/ts_day=2024-01-02/part-00000.parquet
s3a://lakehouse/hudi/orders/.hoodie/20240101120000.commit s3a://lakehouse/hudi/orders/.hoodie/20240101120000.inflight s3a://lakehouse/hudi/orders/.hoodie/20240101120000.requested s3a://lakehouse/hudi/orders/.hoodie/20240101121000.clean s3a://lakehouse/hudi/orders/.hoodie/archived/20240101121000.commit
二、构建适合容器的 Iceberg 镜像与 Catalog 配置
Iceberg 本身不是一个可执行服务,它是以 Spark、Flink 等计算引擎的库形式存在。所以容器化的第一步不是启动一个独立的 Iceberg 服务,而是把 Iceberg 运行时 JAR 和对象存储驱动打进计算引擎镜像。使用多层 Dockerfile 可以让依赖更新不影响基础镜像缓存。注意不要只把 JAR 放在 /tmp,因为很多基础镜像的 /tmp 权限或清理策略会导致运行期找不到依赖。
FROM apache/spark:3.5.1-scala2.12-java11-python3 USER root RUN mkdir -p /opt/spark/jars/iceberg COPY iceberg-spark-runtime-3.5_2.12-1.5.2.jar /opt/spark/jars/iceberg/ COPY aws-java-sdk-bundle-1.12.262.jar /opt/spark/jars/ COPY hadoop-aws-3.3.4.jar /opt/spark/jars/ USER spark
Catalog 选型在容器环境中尤其重要。传统 Hive Metastore 需要额外部署数据库和协调服务,不利于快速扩缩容。Iceberg 的 REST Catalog 把元数据操作封装成 HTTP 接口,可以用两个副本的 Deployment 直接提供服务。仓库地址、S3 端点和访问凭证全部通过环境变量或 Secret 注入,Pod 漂移后仍能访问同一份对象存储。下面是一个精简的 Deployment 配置:
apiVersion: v1
kind: Secret
metadata:
name: s3-credentials
type: Opaque
stringData:
AWS_ACCESS_KEY_ID: minioadmin
AWS_SECRET_ACCESS_KEY: minioadmin
---
apiVersion: apps/v1
kind: Deployment
metadata:
name: iceberg-rest-catalog
spec:
replicas: 2
selector:
matchLabels:
app: iceberg-rest-catalog
template:
metadata:
labels:
app: iceberg-rest-catalog
spec:
containers:
- name: rest-catalog
image: tabulario/iceberg-rest-catalog:1.5.2
env:
- name: CATALOG_WAREHOUSE
value: s3a://lakehouse/warehouse
- name: CATALOG_IO__IMPL
value: org.apache.iceberg.aws.s3.S3FileIO
- name: CATALOG_S3_ENDPOINT
value: http://minio.default.svc:9000
- name: CATALOG_S3_ACCESS_KEY_ID
valueFrom:
secretKeyRef:
name: s3-credentials
key: AWS_ACCESS_KEY_ID
- name: CATALOG_S3_SECRET_ACCESS_KEY
valueFrom:
secretKeyRef:
name: s3-credentials
key: AWS_SECRET_ACCESS_KEY
ports:
- containerPort: 8181
计算侧通过 Spark 配置指向这个 Catalog。下面这段 Scala 代码配置了 Iceberg 扩展、REST Catalog 地址和 S3 仓库。键名中的 spark.sql.catalog.iceberg 表示注册一个名为 iceberg 的 Catalog,后续 SQL 中写 iceberg.db.table 即可。因为在 Kubernetes 中 executor 会动态分配到不同节点,访问密钥不能写死在代码里,而应通过 Pod 模板或 IRSA 注入环境变量。
import org.apache.spark.sql.SparkSession
val spark = SparkSession.builder()
.appName("Iceberg On Kubernetes")
.config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
.config("spark.sql.catalog.iceberg", "org.apache.iceberg.spark.SparkCatalog")
.config("spark.sql.catalog.iceberg.type", "rest")
.config("spark.sql.catalog.iceberg.uri", "http://iceberg-rest-catalog.default.svc:8181")
.config("spark.sql.catalog.iceberg.warehouse", "s3a://lakehouse/warehouse")
.config("spark.sql.defaultCatalog", "iceberg")
.getOrCreate()
三、容器化 Hudi 的表服务与写入方式
Hudi 的依赖包比 Iceberg 更重,通常需要包含 hudi-spark-bundle、hudi-utilities-bundle 和 Hadoop 兼容库。不建议把 Iceberg 与 Hudi 的完整 bundle 打进同一个计算引擎镜像,因为两者可能依赖不同版本的 Avro、Parquet 或 AWS SDK,容易在类加载阶段出现冲突。应该为 Hudi 单独构建一个镜像,并在提交作业时通过 spark.kubernetes.container.image 指定。
FROM apache/spark:3.5.1-scala2.12-java11-python3 USER root RUN mkdir -p /opt/hudi && chown spark:spark /opt/hudi COPY hudi-spark-bundle_2.12-0.14.1.jar /opt/spark/jars/ COPY hudi-utilities-bundle_2.12-0.14.1.jar /opt/spark/jars/ COPY hudi-hadoop-mr-bundle-0.14.1.jar /opt/spark/jars/ USER spark
流式入湖可以直接使用 Hudi DeltaStreamer。它以常驻 Deployment 或批次 Job 的方式消费 Kafka、Pulsar 等数据源,并把数据写入 Hudi 表。下面的提交命令把 DeltaStreamer 运行在 Kubernetes 集群模式,目标表使用 COPY_ON_WRITE 类型,适合读多写少的场景。如果写入频率高且需要降低文件放大,可以改用 MERGE_ON_READ 表类型,并让压缩任务独立调度。
spark-submit \ --class org.apache.hudi.utilities.deltastreamer.HoodieDeltaStreamer \ --master k8s://https://kubernetes.default.svc \ --deploy-mode cluster \ --conf spark.kubernetes.container.image=myrepo/hudi-spark:0.14.1 \ --conf spark.hadoop.fs.s3a.access.key=minioadmin \ --conf spark.hadoop.fs.s3a.secret.key=minioadmin \ --conf spark.hadoop.fs.s3a.endpoint=http://minio.default.svc:9000 \ --conf spark.serializer=org.apache.spark.serializer.KryoSerializer \ /opt/spark/jars/hudi-utilities-bundle_2.12-0.14.1.jar \ --table-type COPY_ON_WRITE \ --source-class org.apache.hudi.utilities.sources.JsonKafkaSource \ --source-ordering-field ts \ --target-base-path s3a://lakehouse/hudi/orders \ --target-table orders \ --props /opt/hudi/config/kafka-source.properties \ --schemaprovider-class org.apache.hudi.utilities.schema.FilebasedSchemaProvider \ --op UPSERT
清理和归档不能省略。Hudi 的每次提交都会保留旧版本文件和日志,如果长期不做清理,对象存储成本会迅速上升。容器环境里建议把 HoodieCleaner 或压缩任务配置成 CronJob,例如每小时或每天运行一次,而不是让主写入进程背负全部维护工作。这样既能控制资源使用,也能在清理失败时单独重跑,不影响流式写入主链路。
spark-submit \ --class org.apache.hudi.utilities.HoodieCleaner \ --master k8s://https://kubernetes.default.svc \ --deploy-mode cluster \ --conf spark.kubernetes.container.image=myrepo/hudi-spark:0.14.1 \ --conf spark.hadoop.fs.s3a.access.key=minioadmin \ --conf spark.hadoop.fs.s3a.secret.key=minioadmin \ --conf spark.hadoop.fs.s3a.endpoint=http://minio.default.svc:9000 \ /opt/spark/jars/hudi-utilities-bundle_2.12-0.14.1.jar \ --target-base-path s3a://lakehouse/hudi/orders \ --cleaner-policy KEEP_LATEST_COMMITS \ --retain-commits 20
并发写入是另一个容器化后容易忽略的问题。Deployment 多副本虽然可以提升吞吐,但如果多个 DeltaStreamer 实例写同一张表,必须保证每个实例只消费互斥的 Kafka 分区,否则会在时间线上产生冲突。Hudi 默认的乐观并发控制在低冲突时可用,但不能把它当作强一致锁来设计。对于有严格去重要求的场景,更好的做法是单写者多读者,或者使用 Kafka 分区键把数据天然隔离到不同 Hudi 分区。
四、同集群调度与对象存储调优
如果把 Iceberg 和 Hudi 同时放进一个 Kubernetes 集群,数据底座最好统一使用 S3 兼容对象存储。容器 DNS 名称必须稳定,例如用 minio.default.svc:9000,不要写临时端口或宿主机 IP。两类表的 warehouse 要分开,Iceberg 仓库示例为 s3a://lakehouse/iceberg,Hudi 仓库示例为 s3a://lakehouse/hudi。如果混用同一个目录,元数据文件会和 .hoodie 时间线互相干扰,查询时会出现不一致。
任务调度上,Iceberg 的批处理作业和 Hudi 的流式入湖可以共享同一个 Spark operator 或 Volcano 调度器。需要注意的是 Spark executor 的 shuffle 临时目录。容器化后如果只使用默认的 rootfs,容量很快会耗尽,导致大任务频繁失败。应该通过 spark.local.dir 指向挂载的本地 SSD 或 emptyDir,并设置合理的 spark.kubernetes.executor.volumes.emptyDir 参数。
对象存储的参数对两种表格式都会产生明显影响。Iceberg 写入时建议把 write.distribution-mode 设为 hash,并控制目标文件大小,减少小文件数量。Hudi 则要关注 hoodie.parquet.small.file.limit、hoodie.filesystem.view.type 和归档间隔。对于 S3 这种最终一致性的存储,Hudi 的视图类型应选择基于内存的文件系统视图,而不是依赖目录遍历,避免频繁 LIST 请求拖慢写入。
监控方面,REST Catalog 可以直接通过就绪探针检查 HTTP 健康接口,DeltaStreamer 则需要同时监控 Kafka lag、提交成功率和 .hoodie 时间线堆积情况。Iceberg 的快照可以通过 SELECT * FROM iceberg.db.orders.snapshots 查看,Hudi 的提交可以通过 hudi-cli 或 Spark SQL 查询。只有把调度、存储、监控和清理链路都容器化,整个数据湖才算真正具备在 Kubernetes 上长期稳定运行的能力。