在本地搭建 Data Lakehouse 实验环境
使用开源技术栈在本地搭建离线 Data Lakehouse 实验环境,记录从基础组件部署到湖仓连通性验证的过程。
使用开源技术栈在本地搭建 Data Lakehouse 实验环境,可以对一些数据领域的概念有个具体的理解,也可以理解技术的演变。
Data Lakehouse(数据湖仓)从数据处理时效方面来看,有离线、流式湖仓,这篇文章记录的是离线湖仓的搭建过程。技术栈的选择,主要基于个人的经验和兴趣,另外从可观测性考虑要有 Web UI。
注意:
这不是一篇推荐使用数据湖仓架构的文章,企业使用什么技术是要因地制宜的。
01 背景资料
在1980年代后期,Bill Inmon 提出了“数据仓库(Data Warehouse)”的概念。
在2011年,Pentaho 首席技术官的 James Dixon 正式创造了“数据湖(Data Lake)”一词。
在2021年,Databricks 发布了一篇介绍“数据湖仓(Data Lakehouse)”概念的论文《Lakehouse: A New Generation of Open Platforms that Unify Data Warehousing and Advanced Analytics》。

图片来源:Databricks 官网 https://www.databricks.com/discover/data-warehouse https://www.databricks.com/discover/data-lakes https://www.databricks.com/glossary/data-lakehouse

图片来源:《AI风暴来袭:2024年数据平台的演进、挑战与机遇》
02 部署过程
1. 蓝图
Data & AI 数据平台整体示意图

参考《新一代DA平台架构的设计原则与演进思路》
本文涉及的离线湖仓的部署示意图

2. Ubuntu 24.04
安装 Ubuntu 24.04 Desktop
3. Java 17
https://www.fosstechnix.com/how-to-install-openjdk-on-ubuntu-24-04-lts/
pan@pan-SER8:~$ sudo apt install openjdk-17-jdk
pan@pan-SER8:/opt/spark$ java --version
openjdk 17.0.14 2025-01-21
OpenJDK Runtime Environment (build 17.0.14+7-Ubuntu-124.04)
OpenJDK 64-Bit Server VM (build 17.0.14+7-Ubuntu-124.04, mixed mode, sharing)
# 查看java的安装路径
update-alternatives --config java
update-alternatives --config javac
选择 Java 17,主要考虑到 Spark 4.0 最低要求 Java 17 以及其它湖仓组件也支持Java 17 。
4. Python 3.12
Ubuntu 24.04 自带 Python 3.12
5. PostgreSQL
https://www.postgresql.org/download/linux/ubuntu/
pan@pan-SER8:$ sudo apt install postgresql
pan@pan-SER8:/opt/spark/spark-3.5.5$ psql --version
psql (PostgreSQL) 16.8 (Ubuntu 16.8-0ubuntu0.24.04.1)
6. Minio
https://min.io/open-source/download?platform=linux
# minio 目录
pan@pan-SER8:/opt$ sudo mkdir minio
[sudo] pan 的密码:
pan@pan-SER8:/opt$ sudo chown pan:pan minio
pan@pan-SER8:/opt$ cd minio/
# MinIO Server
pan@pan-SER8:/opt/minio$ wget https://dl.min.io/server/minio/release/linux-amd64/minio
pan@pan-SER8:/opt/minio$ chmod +x minio
# MinIO Client
pan@pan-SER8:/opt/minio$ wget https://dl.min.io/client/mc/release/linux-amd64/mc
pan@pan-SER8:/opt/minio$ chmod +x mc
# 数据目录
pan@pan-SER8:/$ sudo mkdir data
[sudo] pan 的密码:
pan@pan-SER8:/$ sudo chown pan:pan data
pan@pan-SER8:/data$ mkdir minio
# 运行
pan@pan-SER8:/opt/minio$ MINIO_ROOT_USER=molihua MINIO_ROOT_PASSWORD="Az&&&&09" ./minio server /data/minio --console-address ":9001"
# 配置MinIO别名,关联本地客户端与远程MinIO服务端
pan@pan-SER8:/opt/minio$ ./mc alias set myminio http://127.0.0.1:9000 molihua "Az&&&&09"
pan@pan-SER8:/opt/minio$ ./mc admin info myminio
# Web UI
http://192.168.31.72:9001/login
7. Spark & Iceberg
(1)单机 Standalone 部署模式 部署 Spark
https://spark.apache.org/downloads.html
# 安装文件校验
pan@pan-SER8:/opt/spark$ sha512sum spark-3.5.5-bin-hadoop3-scala2.13.tgz
a462656d2a87afcf81f24e4dea9a99e6666ed5c8ff7ba0cc9c936188498c549cf852d0ebb475325761d6cf645d1184d527be9a609503d587bb47d5532b533aa1 spark-3.5.5-bin-hadoop3-scala2.13.tgz
# 解压
pan@pan-Redmi-Book-14:/opt/spark$ tar -zxvf spark-3.5.2-bin-hadoop3.tgz
# 配置环境变量
sudo vim /etc/environment
JAVA_HOME="/usr/lib/jvm/java-17-openjdk-amd64"
SPARK_HOME="/opt/spark/spark-3.5.5"
PATH="/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin:/usr/games:/usr/local/games:/snap/bin:$SPARK_HOME/sbin:$SPARK_HOME/bin"
source /etc/environment
echo $JAVA_HOME
echo $SPARK_HOME
# 测试 Spark
# 启动 master
pan@pan-SER8:/opt/spark/spark-3.5.5$ ./sbin/start-master.sh
starting org.apache.spark.deploy.master.Master, logging to /opt/spark/spark-3.5.5/logs/spark-pan-org.apache.spark.deploy.master.Master-1-pan-SER8.out
pan@pan-SER8:/opt/spark/spark-3.5.5$ cat /opt/spark/spark-3.5.5/logs/spark-pan-org.apache.spark.deploy.master.Master-1-pan-SER8.out
Spark Command: /usr/lib/jvm/java-17-openjdk-amd64/bin/java -cp /opt/spark/spark-3.5.5/conf/:/opt/spark/spark-3.5.5/jars/* -Xmx1g org.apache.spark.deploy.master.Master --host pan-SER8 --port 7077 --webui-port 8080
========================================
25/03/17 20:16:32 INFO Utils: Successfully started service 'sparkMaster' on port 7077.
25/03/17 20:16:32 INFO Master: Starting Spark master at spark://pan-SER8:7077
25/03/17 20:16:32 INFO Master: Running Spark version 3.5.5
25/03/17 20:16:32 INFO JettyUtils: Start Jetty 0.0.0.0:8080 for MasterUI
25/03/17 20:16:32 INFO Utils: Successfully started service 'MasterUI' on port 8080.
25/03/17 20:16:32 INFO MasterWebUI: Bound MasterWebUI to 0.0.0.0, and started at http://192.168.31.72:8080
25/03/17 20:16:32 INFO Master: I have been elected leader! New state: ALIVE
# 启动 worker
#pan@pan-SER8:/opt/spark/spark-3.5.5$ ./sbin/start-worker.sh spark://192.168.31.72:7077
pan@pan-SER8:/opt/spark/spark-3.5.5$ ./sbin/start-worker.sh spark://pan-SER8:7077
Spark Command: /usr/lib/jvm/java-17-openjdk-amd64/bin/java -cp /opt/spark/spark-3.5.5/conf/:/opt/spark/spark-3.5.5/jars/* -Xmx1g org.apache.spark.deploy.worker.Worker --webui-port 8081 spark://pan-SER8:7077
========================================
25/03/17 20:28:24 INFO JettyUtils: Start Jetty 0.0.0.0:8081 for WorkerUI
25/03/17 20:28:24 INFO Utils: Successfully started service 'WorkerUI' on port 8081.
25/03/17 20:28:24 INFO WorkerWebUI: Bound WorkerWebUI to 0.0.0.0, and started at http://192.168.31.72:8081
25/03/17 20:28:24 INFO Worker: Connecting to master pan-SER8:7077...
25/03/17 20:28:24 INFO TransportClientFactory: Successfully created connection to pan-SER8/127.0.1.1:7077 after 14 ms (0 ms spent in bootstraps)
25/03/17 20:28:24 INFO Worker: Successfully registered with master spark://pan-SER8:7077
(2)Iceberg
https://iceberg.apache.org/docs/latest/spark-getting-started/
https://iceberg.apache.org/releases/
1.8.1 Spark 3.5_with Scala 2.13 runtime Jar
Spark 3.5.x 版本集成的 Iceberg 运行时依赖。用于在 Spark 作业中处理 Iceberg 表,支持各种数据操作,包括查询、插入、删除、更新等。
pan@pan-SER8:/opt/spark/spark-3.5.5$ cp ../iceberg-spark-runtime-3.5_2.13-1.8.1.jar ./jars/
(3)Iceberg AWS Integrations
https://iceberg.apache.org/docs/latest/aws/
https://iceberg.apache.org/releases/#downloads
1.8.1 aws-bundle Jar
Apache Iceberg的AWS(Amazon Web Services)集成包,它包含了与AWS服务(如Amazon S3)交互所需的所有类和资源。
pan@pan-SER8:/opt/spark/spark-3.5.5$ cp ../iceberg-aws-bundle-1.8.1.jar ./jars/
(4)Iceberg REST Catalog - Lakekeeper
考虑到 iceberg 官方的 REST catalog 实现 很久没有更新,故选择其它的 REST catalog 实现
https://github.com/tabular-io/iceberg-rest-image
https://github.com/databricks/iceberg-rest-image
https://github.com/lakekeeper/lakekeeper
https://docs.lakekeeper.io/getting-started/#option-4-binary
# 安装
pan@pan-SER8:/opt/spark$ tar -zxvf iceberg-catalog-x86_64-unknown-linux-gnu.tar.gz
pan@pan-SER8:/opt/spark$ ./iceberg-catalog --help
# 运行
pan@pan-SER8:/opt/spark$ export LAKEKEEPER__PG_DATABASE_URL_READ="postgres://molihua:Az&&&&09@127.0.0.1:5432/lakekeeper"
pan@pan-SER8:/opt/spark$ export LAKEKEEPER__PG_DATABASE_URL_WRITE="postgres://molihua:Az&&&&09@127.0.0.1:5432/lakekeeper"
pan@pan-SER8:/opt/spark$ export LAKEKEEPER__PG_ENCRYPTION_KEY="YOUR_ENCRYPTION_KEY"
# 9000端口,minio已用
pan@pan-SER8:/opt/spark$ export LAKEKEEPER__METRICS_PORT=8182
pan@pan-SER8:/opt/spark$ ./iceberg-catalog migrate
pan@pan-SER8:/opt/spark$ ./iceberg-catalog serve
# Web UI
http://192.168.31.72:8181/ui/
# 测试
# https://github.com/databricks/iceberg-rest-image
➜ ~ pyiceberg --uri http://localhost:8181 list
➜ ~ pyiceberg --uri http://localhost:8181 list nyc
➜ ~ pyiceberg --uri http://localhost:8181 describe --entity=table tpcds_iceberg.customer
# https://docs.lakekeeper.io/getting-started/#connect-compute
import pandas as pd
import pyspark
配置对象存储
界面上配置,连接到对象存储

(
https://chatgpt.com/c/67d8cf18-3f24-8012-adbe-5c94eb954428
Enable path style access(启用路径样式访问)
关闭时(默认),使用虚拟托管方式访问,如
https://bucket-name.s3.amazonaws.com开启后,使用路径样式访问,如
https://s3.amazonaws.com/bucket-nameMinIO 这类 S3 兼容存储通常需要开启该选项,以确保兼容性。
Enable alternative s3 protocols(启用其他 S3 协议,如 s3a、s3n)
启用后,可能允许使用 Hadoop 生态系统支持的
s3a://``s3n://协议,而不仅仅是s3://在 Spark、Hadoop 等大数据框架中,
s3a://比s3://提供更好的性能和兼容性。
Bucket Region(存储桶区域)us-east-1
在 AWS S3 上,存储桶区域指定存储的位置(如
us-east-1表示美国东部)。在 MinIO 或自建 S3 兼容存储时,通常可以随意设置,但某些应用程序可能需要该参数才能正确工作。
Enable Soft Deletion(启用软删除)
- 软删除表示删除的对象会进入回收站,可以在一定时间内恢复。
Enable STS(启用 STS 临时凭证)
- 启用 STS(Security Token Service)后,可使用临时访问凭证,而非固定的 Access Key 和 Secret Key。
Key Prefix(键前缀):
- 这个选项允许你指定一个键前缀,以便只访问存储桶中具有特定前缀的对象。
)
(5)Spark 配置
https://www.tabular.io/apache-iceberg-cookbook/getting-started-connect-rest-catalog/
https://docs.lakekeeper.io/getting-started/#connect-compute
# spark event 目录
pan@pan-SER8:/data$ mkdir -p /data/spark/spark-events
# 配置
# 配置 Iceberg REST Catalog 和 对象存储,备注有“新增”、“有变化”,是指 Iceberg 官方的 REST catalog 实现 和 Lakekeeper 的配置对比
pan@pan-SER8:/opt/spark/spark-3.5.5$ cp conf/spark-defaults.conf.template conf/spark-defaults.conf
pan@pan-SER8:/opt/spark/spark-3.5.5$ vim conf/spark-defaults.conf
spark.sql.extensions org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions
spark.sql.defaultCatalog demo
spark.sql.catalog.demo org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.demo.type rest
# 新增,与 spark.sql.catalog.demo.type 冲突
#spark.sql.catalog.demo.catalog-impl org.apache.iceberg.rest.RESTCatalog
# 有变化
#spark.sql.catalog.molihua.uri http://192.168.31.72:8181
spark.sql.catalog.demo.uri http://192.168.31.72:8181/catalog/
spark.sql.catalog.demo.io-impl org.apache.iceberg.aws.s3.S3FileIO
# 有变化
#spark.sql.catalog.demo.warehouse s3://warehouse
spark.sql.catalog.demo.warehouse warehouse
spark.sql.catalog.demo.s3.endpoint http://192.168.31.72:9000
#spark.sql.catalog.molihua.default-namespace=examples
spark.eventLog.enabled true
spark.eventLog.dir /data/spark/spark-events
spark.history.fs.logDirectory /data/spark/spark-events
spark.sql.catalogImplementation in-memory
# 字段名称
spark.sql.cli.print.header=true
(
对于"spark.sql.catalogImplementation in-memory",Spark会使用内存中的元数据存储来管理数据库、表和视图,而不会依赖外部的Hive Metastore或其他持久化存储。
如果你通过 `spark.sql.catalog.molihua.type=rest` 配置了REST Catalog并使用该Catalog来管理表,那么表的元数据将被存储在REST Catalog中,而不是在Spark的内存中。此时,`spark.sql.catalogImplementation=in-memory` 这个配置并不会影响 `molihua` Catalog,因为 `molihua` 是一个明确使用REST协议的**外部Catalog**。
)
8. Airflow
https://airflow.apache.org/docs/apache-airflow/stable/start.html
https://airflow.apache.org/docs/apache-airflow/stable/installation/index.html
Airflow 的扩展库还是要严格按约束文件中的版本安装,避免兼容性问题。
# python 虚拟环境
pan@pan-SER8:~$ cd /opt
pan@pan-SER8:/opt$ sudo mkdir airflow
pan@pan-SER8:/opt$ sudo chown -R pan:pan airflow
pan@pan-SER8:/opt$ python3 -m venv airflow
pan@pan-SER8:/opt$ cd airflow/
pan@pan-SER8:/opt/airflow$ source bin/activate
# 安装 airflow
(airflow) pan@pan-SER8:/opt/airflow$ export AIRFLOW_HOME=/opt/airflow
(airflow) pan@pan-SER8:/opt/airflow$ AIRFLOW_VERSION=2.10.5
(airflow) pan@pan-SER8:/opt/airflow$ PYTHON_VERSION="$(python -c 'import sys; print(f"{sys.version_info.major}.{sys.version_info.minor}")')"
(airflow) pan@pan-SER8:/opt/airflow$ CONSTRAINT_URL="https://raw.githubusercontent.com/apache/airflow/constraints-${AIRFLOW_VERSION}/constraints-${PYTHON_VERSION}.txt"
(airflow) pan@pan-SER8:/opt/airflow$ pip install "apache-airflow==${AIRFLOW_VERSION}" --constraint "${CONSTRAINT_URL}"
# 如果访问异常,约束文件可以手动下载
curl -x http://127.0.0.1:7890 -O https://raw.githubusercontent.com/apache/airflow/constraints-2.10.5/constraints-3.12.txt
(airflow) pan@pan-SER8:/opt/airflow$ pip install "apache-airflow==2.10.5" --constraint ./constraints-3.12.txt
# 运行之后在 AIRFLOW_HOME 生成 默认配置、默认数据库文件
(airflow) pan@pan-SER8:/opt/airflow$ airflow standalone
pan@pan-SER8:/opt/airflow$ ls -alh
-rw------- 1 pan pan 85K 3月 30 21:00 airflow.cfg
-rw-r--r-- 1 pan pan 1.2M 3月 30 21:03 airflow.db
-rw-rw-r-- 1 pan pan 4.7K 3月 30 21:00 webserver_config.py
(
If you want to run the individual parts of Airflow manually rather than using the all-in-one standalone command, you can instead run:
airflow db migrate
airflow users create \
--username admin \
--firstname Peter \
--lastname Parker \
--role Admin \
--email spiderman@superhero.org
airflow webserver --port 8080
airflow scheduler
)
# 配置后端数据库为 PostgreSQL
# https://airflow.apache.org/docs/apache-airflow/stable/howto/set-up-database.html#setting-up-a-postgresql-database
# 在 PostgreSQL 中创建用户、数据库
CREATE USER airflow WITH PASSWORD 'airflow';
create database airflow with owner airflow;
# 配置数据库连接
pan@pan-SER8:~/airflow$ vim airflow.cfg
sql_alchemy_conn = postgresql+psycopg2://airflow:airflow@127.0.0.1:5432/airflow
# 运行配置
# 如果缺少驱动,在约束文件查找并安装(约束文件的库版本是官方测试过的)
(airflow) pan@pan-SER8:/opt/airflow$ cat constraints-3.12.txt | grep psycopg
(airflow) pan@pan-SER8:/opt/airflow$ pip install psycopg2-binary==2.9.10
(airflow) pan@pan-SER8:/opt/airflow$ airflow db migrate
(airflow) pan@pan-SER8:/opt/airflow$ airflow users create \
--username admin \
--firstname huida \
--lastname pan \
--role Admin \
--email panhuida@qq.com
# 运行
(airflow) pan@pan-SER8:/opt/airflow$ sudo vim /etc/environment
AIRFLOW_HOME="/opt/airflow"
(airflow) pan@pan-SER8:/opt/airflow$ source /etc/environment
(airflow) pan@pan-SER8:/opt/airflow$ airflow webserver --port 8886
(airflow) pan@pan-SER8:/opt/airflow$ airflow scheduler
# Web UI
http://192.168.31.72:8886/
配置 PostgresSQL 连接示例
https://airflow.apache.org/docs/apache-airflow-providers-postgres/stable/index.html
https://www.astronomer.io/docs/learn/connections/postgres/
(airflow) pan@pan-SER8:/opt/airflow$ cat constraints-3.12.txt | grep apache-airflow-providers-postgres
#注意:要按constraints-3.12.tx的版本安装,安装 apache-airflow-providers-postgres 6.1.1 不会出现 postgres的连接类型,
(airflow) pan@pan-SER8:/opt/airflow$ pip install apache-airflow-providers-postgres==6.0.0
Successfully installed apache-airflow-providers-postgres-6.1.1 asyncpg-0.30.0
# 创建连接
# 通过界面配置或是命令行配置连接
(airflow) pan@pan-SER8:/opt/airflow$ airflow connections add 'postgres_default' \
--conn-type 'postgres' \
--conn-host '192.168.31.72' \
--conn-schema 'demo' \
--conn-login 'molihua' \
--conn-password 'Az&&&&09' \
--conn-port 5432
9. Superset
https://superset.apache.org/docs/installation/pypi
# 操作系统依赖
pan@pan-SER8:~$ sudo apt install build-essential libssl-dev libffi-dev python3-dev python3-pip libsasl2-dev libldap2-dev default-libmysqlclient-dev
# python 虚拟环境
pan@pan-SER8:~$ cd /opt
pan@pan-SER8:/opt$ sudo mkdir superset
pan@pan-SER8:/opt$ sudo chown -R pan:pan superset
pan@pan-SER8:/opt$ python3.12 -m venv superset
pan@pan-SER8:/opt$ cd superset/
pan@pan-SER8:/opt$ source bin/activate
# 安装 superset
# 稳定版(apache- superset==4.1.2)和Python 3.12有兼容性问题,尝试安装5.0.0rc1,安装成功
(superset) pan@pan-SER8:/opt/superset$ pip install apache-superset==5.0.0rc1
# 配置 SECRET_KEY
(superset) pan@pan-SER8:/opt/superset$ openssl rand -base64 42
YOUR_SUPERSET_SECRET_KEY
(superset) pan@pan-SER8:/opt/superset$ vim superset_config.py
SECRET_KEY = 'YOUR_SUPERSET_SECRET_KEY'
# 配置后端数据库为 PostgreSQL
https://superset.apache.org/docs/configuration/configuring-superset/#setting-up-a-production-metadata-database
# 在 PostgreSQL 中创建用户、数据库
CREATE USER superset WITH PASSWORD 'superset';
create database superset with owner superset;
# 配置数据库连接
(superset) pan@pan-SER8:/opt/superset$ vim superset_config.py
SQLALCHEMY_DATABASE_URI = postgresql+psycopg2://superset:superset@127.0.0.1:5432/superset
# 运行配置
(superset) pan@pan-SER8:/opt/superset$ export SUPERSET_CONFIG_PATH=/opt/superset/superset_config.py
(superset) pan@pan-SER8:/opt/superset$ export FLASK_APP=superset
# 安装后端数据库Pg驱动
(superset) pan@pan-SER8:/opt/superset$ pip install psycopg2-binary==2.9.10
(superset) pan@pan-SER8:/opt/superset$ superset db upgrade
# Create an admin user in your metadata database (use `admin` as username to be able to load the examples)
(superset) pan@pan-SER8:/opt/superset$ superset fab create-admin
# Load some data to play with
(superset) pan@pan-SER8:/opt/superset$ export http_proxy=http://127.0.0.1:7890
(superset) pan@pan-SER8:/opt/superset$ export https_proxy=http://127.0.0.1:7890
(superset) pan@pan-SER8:/opt/superset$ superset load_examples
# Create default roles and permissions
(superset) pan@pan-SER8:/opt/superset$ superset init
# To start a development web server on port 8088, use -p to bind to another port
# 运行
(superset) pan@pan-SER8:/opt/superset$ superset run -h 0.0.0.0 -p 8887 --with-threads --reload --debugger
10. Airbyte
Airbyte 在 v0.63.4 版本开始弃用对 Docker Compose 部署的支持,只支持 Kubernetes 部署。
https://docs.airbyte.com/release_notes/june_2024#announcements
在 Quick Start 文档中,Airbyte使用abctl进行本地部署,默认使用 Docker Desktop的 kind 方式创建 Kubernetes 集群。
使用 kind 的方式创建 kubernetes 集群,支持多节点,但其是通过单独的容器 kind-registry-mirror 拉取镜像,网络不通的话拉取镜像是个问题,故还是使用 Docker Desktop 的 kubeadm 方式创建 Kubernetes 集群(master和node组件都在一个节点),并通过 helm 部署 Airbyte
https://docs.airbyte.com/using-airbyte/getting-started/oss-quickstart
https://docs.airbyte.com/deploying-airbyte/
https://docs.docker.com/desktop/features/kubernetes/
(1)安装 Docker Desktop
https://docs.docker.com/desktop/setup/install/linux/ubuntu/
- Set up Docker’s package repository.
# 设置代理
export http_proxy=http://127.0.0.1:7890
export https_proxy=http://127.0.0.1:7890
env | grep proxy
## 这里的 -E 选项很重要,它会保留你设置的环境变量
sudo -E apt-get update
# Add Docker's official GPG key:
sudo apt-get update
sudo apt-get install ca-certificates curl
sudo install -m 0755 -d /etc/apt/keyrings
sudo curl -fsSL https://download.docker.com/linux/ubuntu/gpg -o /etc/apt/keyrings/docker.asc
sudo chmod a+r /etc/apt/keyrings/docker.asc
# Add the repository to Apt sources:
echo \
"deb [arch=$(dpkg --print-architecture) signed-by=/etc/apt/keyrings/docker.asc] https://download.docker.com/linux/ubuntu \
$(. /etc/os-release && echo "$VERSION_CODENAME") stable" | \
sudo tee /etc/apt/sources.list.d/docker.list > /dev/null
#sudo apt-get update
sudo -E apt-get update
- Download the latest DEB package.
- Install the package with apt as follows:
sudo apt-get update
#sudo apt-get install ./docker-desktop-<arch>.deb
pan@pan-Redmi-Book-14:~/Downloads$ sudo -E apt install ./docker-desktop-amd64.deb
By default, Docker Desktop is installed at/opt/docker-desktop.
(2)在 Docker Desktop 中启用 kubernetes
1️⃣配置代理
方便拉取镜像

2️⃣启用 kubernetes
https://docs.docker.com/desktop/features/kubernetes/
选择 kubeadm 的方式创建 kubernetes 集群

3️⃣安装 kubectl
https://kubernetes.io/zh-cn/docs/tasks/tools/install-kubectl-linux/
pan@pan-SER8:/opt/k8s$ curl -LO "https://dl.k8s.io/release/v1.32.2/bin/linux/amd64/kubectl"
pan@pan-SER8:/opt/k8s$ sudo install -o root -g root -m 0755 kubectl /usr/local/bin/kubectl
4️⃣安装 Helm
https://helm.sh/docs/intro/install/
Helm(https://github.com/helm/helm/releases)都提供了各种操作系统的二进制版本,这些版本可以手动下载和安装。
pan@pan-SER8:/opt/k8s$ tar -zxvf helm-v3.17.2-linux-amd64.tar.gz
pan@pan-SER8:/opt/k8s$ sudo mv linux-amd64/helm /usr/local/bin/helm
(3)安装 Airbyte
https://artifacthub.io/packages/helm/airbyte/airbyte
# 注意配置代理,注意代理不稳定
pan@pan-SER8:/opt/airbyte$ helm repo add airbyte https://airbytehq.github.io/helm-charts
# 如下命令没有执行成功,多是几次,或改为手工下载安装
#pan@pan-SER8:/opt/airbyte$ helm install demo-airbyte airbyte/airbyte --version 1.5.1 --create-namespace --namespace airbyte --debug
pan@pan-SER8:/opt/airbyte$ helm install demo-airbyte airbyte/airbyte --version 1.5.1 --create-namespace --namespace airbyte
NAME: demo-airbyte
LAST DEPLOYED: Thu Mar 27 09:56:54 2025
NAMESPACE: airbyte
STATUS: deployed
REVISION: 1
NOTES:
Get the application URL by running these commands:
echo "Visit http://127.0.0.1:8080 to use your application"
kubectl -n airbyte port-forward deployment/demo-airbyte-webapp 8080:8080
# 查看安装的 release
pan@pan-SER8:/opt/airbyte$ helm list -A
NAME NAMESPACE REVISION UPDATED STATUS CHART APP VERSION
demo-airbyte airbyte 1 2025-03-27 09:56:54.849031361 +0800 CST deployed airbyte-1.5.1 1.5.1
kubernetes-dashboard kubernetes-dashboard 1 2025-03-27 08:30:25.7844142 +0800 CST deployed kubernetes-dashboard-7.11.1
pan@pan-SER8:/opt/airbyte$ kubectl get pods -n airbyte
NAME READY STATUS RESTARTS AGE
airbyte-db-0 1/1 Running 0 4m49s
airbyte-minio-0 1/1 Running 0 4m49s
my-airbyte-airbyte-bootloader 0/1 Completed 0 4m48s
my-airbyte-connector-builder-server-5f9fc6bcbb-msfsm 1/1 Running 0 4m31s
my-airbyte-cron-7c8dfd4b4-jbgs7 1/1 Running 0 4m31s
my-airbyte-server-798ff4d969-nt9lq 1/1 Running 0 4m31s
my-airbyte-temporal-859ffcc49d-zpm54 1/1 Running 0 4m31s
my-airbyte-webapp-76bbc9b86f-wmr2d 1/1 Running 0 4m31s
my-airbyte-worker-b54df7887-rgrqf 1/1 Running 0 4m31s
my-airbyte-workload-api-server-f55df5d66-7cj8z 1/1 Running 0 4m31s
my-airbyte-workload-launcher-7b7fd47979-bbcs8 1/1 Running 0 4m31s
# Web UI
# 不使用代理
pan@pan-SER8:/opt/k8s$ unset HTTP_PROXY http_proxy HTTPS_PROXY https_proxy ALL_PROXY all_proxy
# 按 helm install 输出的操作
# pod
# pan@pan-SER8:/opt/k8s$ kubectl -n airbyte port-forward --address 192.168.31.72 my-airbyte-webapp-76bbc9b86f-wmr2d 8881:8080
# deployment
pan@pan-SER8:/opt/k8s$ kubectl -n airbyte port-forward --address 192.168.31.72 deployment/demo-airbyte-webapp 8881:8080
# airbyte不支持,kubernetes-dashboard 支持
# kubectl -n airbyte port-forward --address 192.168.31.72 svc/demo-airbyte-airbyte-webapp-svc 8881:8080
# Web UI
http://192.168.31.72:8881/
11. RisingWave
https://docs.risingwave.com/get-started/quickstart
原计划是安装DuckDB,但 DuckDB 目前不支持 Iceberg REST Catalog,先用 RisingWave 测试。
# curl -L https://risingwave.com/sh | sh
pan@pan-SER8:/opt/risingwave$ wget https://raw.githubusercontent.com/risingwavelabs/risingwave/main/scripts/install/install-risingwave.sh
# 默认是最新的版本,修改为v2.2.4
pan@pan-SER8:/opt/risingwave$ vim install-risingwave.sh
VERSION="v2.2.4"
pan@pan-SER8:/opt/risingwave$ sh install-risingwave.sh
pan@pan-SER8:/opt/risingwave$ ./risingwave
pan@pan-SER8:/opt/risingwave$ psql -h 0.0.0.0 -p 4566 -d dev -U root
pan@pan-SER8:/opt/risingwave$ ./risingwave --help
pan@pan-SER8:/opt/risingwave$ ./risingwave single --help
12. Jupyter
https://jupyterlab.readthedocs.io/en/latest/getting_started/installation.html#pip
# python 虚拟环境
pan@pan-SER8:~$ cd /opt
pan@pan-SER8:/opt$ sudo mkdir jupyter
pan@pan-SER8:/opt$ sudo chown -R pan:pan jupyter
pan@pan-SER8:/opt$ python3 -m venv jupyter
pan@pan-SER8:/opt$ cd jupyter/
pan@pan-SER8:/opt/jupyter$ source bin/activate
# 安装 jupyter
# 注意库之间的依赖
(jupyter) pan@pan-SER8:/opt/python/jupyter$ pip install jupyterlab
(jupyter) pan@pan-SER8:/opt/python/jupyter$ pip install pyspark
(jupyter) pan@pan-SER8:/opt/python/jupyter$ pip install pandas
# 运行
jupyter lab --ip=0.0.0.0 --port=8888
13. 测试 - 是否可以跑通
(1)使用 Spark SQL 进行湖仓连通测试
https://iceberg.apache.org/spark-quickstart/#creating-a-table
按 Iceberg 官方文档操作一遍 。
Iceberg Catalog 的对象层次结构为:catalog -> namespace(对应Spark 的 database)-> table 。
pan@pan-SER8:/opt/spark/spark-3.5.5$ ./bin/spark-sql
# 连接到 Spark 集群
# pan@pan-SER8:/opt/spark/spark-3.5.5$ ./bin/spark-sql --master spark://pan-SER8:7077
spark-sql ()> show catalogs;
spark-sql ()> use demo;
spark-sql ()> create database nyc;
spark-sql ()> use demo.nyc;
spark-sql (nyc)>
CREATE TABLE demo.nyc.taxis
(
vendor_id BIGINT,
trip_id BIGINT,
trip_distance FLOAT,
fare_amount DOUBLE,
store_and_fwd_flag STRING,
ts STRING
)
;
CREATE TABLE demo.ods.taxis
(
vendor_id bigint,
trip_id bigint,
trip_distance float,
fare_amount double,
store_and_fwd_flag string
)
PARTITIONED BY (vendor_id);
INSERT INTO demo.ods.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');
INSERT INTO demo.nyc.taxis
VALUES
(1, 1000371, 1.8, 15.32, 'N', '2024-01-01 9:15:23'),
(2, 1000372, 2.5, 22.15, 'N', '2024-01-02 12:10:11'),
(2, 1000373, 0.9, 9.01, 'N', '2024-01-01 3:25:15'),
(1, 1000374, 8.4, 42.13, 'Y', '2024-01-03 7:12:33');
spark-sql (nyc)> select * from taxis limit 5;
# 更新记录
spark-sql (nyc)> update taxis set ts='2025-03-18' where trip_id=1000373;
# 时间旅行用例
# 微信实验平台Iceberg湖仓一体架构改造
# https://mp.weixin.qq.com/s/5E2nq51qA8qo_WFelykU5g
# 测试时间旅行
spark-sql (nyc)> update taxis set ts='2025-03-29 10:14:30' where trip_id=1000371;
# 修改之前
spark-sql (nyc)> SELECT * FROM demo.nyc.taxis TIMESTAMP AS OF '2025-03-29 10:14:13';
vendor_id trip_id trip_distance fare_amount store_and_fwd_flag ts
1 1000371 1.8 15.32 N 2024-01-01 9:15:23
2 1000372 2.5 22.15 N 2024-01-02 12:10:11
1 1000374 8.4 42.13 Y 2024-01-03 7:12:33
2 1000373 0.9 9.01 N 2025-03-18
Time taken: 0.177 seconds, Fetched 4 row(s)
# 修改之后
spark-sql (nyc)> SELECT * FROM demo.nyc.taxis TIMESTAMP AS OF '2025-03-29 10:15:13';
vendor_id trip_id trip_distance fare_amount store_and_fwd_flag ts
2 1000373 0.9 9.01 N 2025-03-18
1 1000371 1.8 15.32 N 2025-03-29 10:14:30
2 1000372 2.5 22.15 N 2024-01-02 12:10:11
1 1000374 8.4 42.13 Y 2024-01-03 7:12:33
Time taken: 0.173 seconds, Fetched 4 row(s)
(2)使用 PySpark 进行湖仓连通测试
(jupyter) pan@pan-SER8:/opt/jupyter$ jupyter lab --ip=0.0.0.0 --port=8888
pyspark_demo.ipynb
# Set the PySpark environment variables
import os
# 使用单独安装的 Spark 环境,而不是 PySpark 自带的环境
os.environ['SPARK_HOME'] = "/opt/spark/spark-3.5.5"
# 使用 PySpark 的 Python 环境,而不是单独安装的 Spark 中的 Python 环境,避免 Python 版本不一致
os.environ['PYSPARK_PYTHON'] = 'python3'
# os.environ['PYSPARK_DRIVER_PYTHON'] = 'python'
os.environ['PYSPARK_DRIVER_PYTHON'] = 'jupyter'
os.environ['PYSPARK_DRIVER_PYTHON_OPTS'] = 'lab'
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("PySpark Demo").getOrCreate()
# 查询湖仓中的数据表
spark.sql(""" select * from demo.nyc.taxis """).show()
# 数据集来源 https://www.kaggle.com/datasets/asaniczka/tmdb-movies-dataset-2023-930k-movies
# 从CSV文件读取数据
# 指定第一行为标题,自动推断数据类型
df = spark.read.format("csv") \
.option("header", "true") \
.option("inferSchema", "true") \
.load("/data/dataset/TMDB_movie_dataset_v11.csv")
# 创建临时视图,可以用SQL查询
df.createOrReplaceTempView("tmdb_movies")
# 查看表结构
df.printSchema()
# 查看数据示例
df.show(3)
# 如果需要将数据保存为 Iceberg 表
df.write.mode("overwrite").saveAsTable("demo.nyc.tmdb_movies_table")
界面示例


(3)使用 Airbyte 同步数据 - 从 PostgreSQL 到 Iceberg
https://docs.airbyte.com/integrations/sources/postgres
https://docs.airbyte.com/integrations/destinations/iceberg
Airbyte 推荐 CDC 的方式同步数据(定时触发执行CDC同步)。
1️⃣PostgreSQL 配置逻辑复制
# 配置逻辑复制
sudo vim /etc/postgresql/16/main/postgresql.conf
wal_level = logical
sudo systemctl restart postgresql
# 为用户赋予权限
psql -h localhost -U postgres
ALTER USER molihua REPLICATION;
# 创建复制槽
psql -h localhost -U molihua -d demo
SELECT pg_create_logical_replication_slot('airbyte_slot', 'pgoutput');
# 设置表复制标识和创建发布
# 设置表复制标识
# https://www.postgresql.org/docs/current/logical-replication-publication.html
# https://zriyansh.medium.com/understanding-postgresql-replica-identity-the-complete-guide-92bc5c756056
ALTER TABLE public.tmdb_movies REPLICA IDENTITY DEFAULT;
# 创建发布
# 注意:这里不要使用for all tables创建数据发布者,因为这将造成发布者在以后不能增、删表。
# Debezium扩展默认发布所有的数据表(包含分区表)
# DROP PUBLICATION airbyte_publication;
CREATE publication airbyte_publication;
# 设置发布的事件类型
# 默认为 INSERT, UPDATE, DELETE, TRUNCATE
ALTER PUBLICATION airbyte_publication SET (publish = 'update, insert, delete');
# 增加发布数据表
# public.tmdb_movies 的数据来源 https://www.kaggle.com/datasets/asaniczka/tmdb-movies-dataset-2023-930k-movies
ALTER PUBLICATION airbyte_publication ADD TABLE public.tmdb_movies;
2️⃣Iceberg 创建 namespace
在目录 demo 下创建 命名空间 default(Airbyte 同步数据前会自动创建,但没有权限创建)。
pan@pan-SER8:/opt/spark/spark-3.5.5$ ./bin/spark-sql
spark-sql ()> use demo;
spark-sql ()> create database default;
3️⃣Airbyte 配置
创建 connection ,配置 源 和 目标 。

4️⃣源PostgreSQL

5️⃣目标Iceberg

6️⃣结果示例
1192529 条记录,12 分钟32秒,平均 1585 条/秒。
期间 Airbyte 的 Pod 崩溃重启,但最终数据全部同步到 Iceberg 。

Airbyte 同步数据到 Iceberg,数据表只有3个字段,其中 _airbyte_data 保存的是源数据表的数据。
pan@pan-SER8:/opt/spark/spark-3.5.5$ ./bin/spark-sql
spark-sql ()> use demo.default;
Response code
Time taken: 0.513 seconds
spark-sql (default)> show tables;
namespace tableName isTemporary
airbyte_raw_tmdb_movies
Time taken: 0.474 seconds, Fetched 1 row(s)
spark-sql (default)> show create table demo.default.airbyte_raw_tmdb_movies;
createtab_stmt
CREATE TABLE demo.default.airbyte_raw_tmdb_movies (
_airbyte_ab_id STRING,
_airbyte_emitted_at TIMESTAMP,
_airbyte_data STRING)
USING iceberg
LOCATION 's3://warehouse/0195f19f-73bd-7972-ad4a-9b5521820be1/0195f1b5-1647-7b62-b022-f3f6aca47cd6'
TBLPROPERTIES (
'current-snapshot-id' = '7396066456257038935',
'format' = 'iceberg/parquet',
'format-version' = '2',
'write.format.default' = 'parquet',
'write.parquet.compression-codec' = 'zstd')
Time taken: 1.09 seconds, Fetched 1 row(s)
(4)使用 RisingWave 访问 Iceberg 表,再通过 Superset 使用数据
原计划是使用 DuckDB访问 Iceberg,但DuckDB 当前版本(v1.2.1)前不支持 Iceberg REST Catalog(还在支持的路上),先用 RisingWave 测试。
1️⃣RisingWave 读 Iceberg 表
https://docs.risingwave.com/ingestion/getting-started/overview
https://docs.risingwave.com/ingestion/sources/iceberg
还是要把 overview 看一遍,熟悉核心概念,如 source、table 等。
RisingWave 目前只支持离线读取 iceberg 数据表。
在 RisingWave 创建读 Iceberg 表的 source,再基于source 创建 table,Superset 的 SQL IDE 支持
RisingWave 的 source 查询,但直接基于表的数据集定义只支持 RisingWave 的table(也可以基于
SQL IDE 查询的数据集制作图表) 。
# 连接 RisingWave
pan@pan-SER8:~$ psql -h 0.0.0.0 -p 4566 -d dev -U root
# 创建源
# tmdb_movies_table 在 “使用 PySpark 进行湖仓连通测试” 时导入
dev=>
CREATE SOURCE iceberg_nyc_tmdb_movies
WITH (
connector = 'iceberg',
database.name = 'nyc',
table.name = 'tmdb_movies_table',
catalog.name = 'demo',
catalog.type = 'rest',
catalog.uri = 'http://192.168.31.72:8181/catalog/',
warehouse.path = 'warehouse',
s3.endpoint = 'http://192.168.31.72:9000',
s3.access.key = 'YOUR_MINIO_ACCESS_KEY',
s3.secret.key = 'YOUR_MINIO_SECRET_KEY',
s3.region = 'us-east-1'
);
dev=> select * from iceberg_nyc_tmdb_movies limit 3;
# 创建单独的用户在 Superset 使用
dev=>
create user molihua with password 'Az&&&&09';
GRANT ALL PRIVILEGES ON DATABASE dev TO molihua;
-- PostgreSQL 15 requires additional privileges:
GRANT ALL ON SCHEMA public TO molihua;
GRANT ALL PRIVILEGES ON SOURCE iceberg_nyc_tmdb_movies TO molihua;
FLUSH;
2️⃣Superset 制作数据看板
- 配置 RisingWave 连接
https://superset.apache.org/docs/configuration/databases/#risingwave
# 先安装驱动
pip install sqlalchemy-risingwave
# 配置数据源
#risingwave://root:@127.0.0.1:4566/dev?sslmode=disable
risingwave://molihua:YOUR_PASSWORD@127.0.0.1:4566/dev?sslmode=disable

- 数据查询
可以使用 SQL IDE 进行数据探索

- 数据看板
数据来源Kaggle 上的数据集
Full TMDB Movies Dataset 2024 (1M Movies)
https://www.kaggle.com/datasets/asaniczka/tmdb-movies-dataset-2023-930k-movies
手动下载 Full TMDB Movies Dataset 2024 (1M Movies) 数据集 csv 文件,通过 PySpark 写入 Iceberg 表
demo.nyc.tmdb_movies_table,在 RisingWave 基于这个表创建 source、table ,Superset 基于
RisingWave 的source、table 制作数据看板。
数据处理
“电影时长分组” 在 Iceberg 表 demo.nyc.tmdb_movies_table 处理。
pan@pan-SER8:/opt/spark/spark-3.5.5$ ./bin/spark-sqlspark-sql ()> use demo.nyc;
spark-sql (nyc)># 添加 runtime_group 字段ALTER TABLE tmdb_movies_table ADD COLUMNS runtime_group STRING;
# 注意正则表达式的转义符UPDATE tmdb_movies_tableSET runtime_group = CASE WHEN runtime RLIKE '^\\d+$' THEN CASE WHEN CAST(runtime AS INT) > 0 AND CAST(runtime AS INT) < 60 THEN '(0, 1小时)' WHEN CAST(runtime AS INT) >= 60 AND CAST(runtime AS INT) < 120 THEN '[1小时, 2小时)' WHEN CAST(runtime AS INT) >= 120 AND CAST(runtime AS INT) < 240 THEN '[2小时, 4小时)' WHEN CAST(runtime AS INT) >= 240 AND CAST(runtime AS INT) < 480 THEN '[4小时, 8小时)' WHEN CAST(runtime AS INT) >= 480 AND CAST(runtime AS INT) < 720 THEN '[8小时, 12小时)' WHEN CAST(runtime AS INT) >= 720 AND CAST(runtime AS INT) < 960 THEN '[12小时, 16小时)' WHEN CAST(runtime AS INT) >= 960 AND CAST(runtime AS INT) < 1200 THEN '[16小时, 20小时)' WHEN CAST(runtime AS INT) >= 1200 AND CAST(runtime AS INT) < 1440 THEN '[20小时, 24小时)' WHEN CAST(runtime AS INT) >= 1440 THEN '[24小时, +)' ELSE 'invalid' END ELSE 'invalid'END;Time taken: 9.34 seconds
spark-sql (nyc)> SELECT title, runtime, runtime_group FROM tmdb_movies_table LIMIT 10;
原始语言、上映年份在 Superset 上处理(一般都是在数仓处理的,这里时测试Superset 功能)。
数据范围
数据看板的数据范围为 电影上映日期在1900~2023年(2024年数据看上去缺少很多)并且状态为"已上映" 。
最终的数据看板

彩蛋

https://movie.douban.com/subject/10744830/
喜欢这部电影的人也喜欢 · · · · · ·
03 附录
1.《Lakehouse: A New Generation of Open Platforms that Unify Data Warehousing and Advanced Analytics》
https://www.cidrdb.org/cidr2021/papers/cidr2021_paper17.pdf
2.《新一代DA平台架构的设计原则与演进思路》
https://mp.weixin.qq.com/s/ukqztUYXwDE9jCONvzEn8A
https://mp.weixin.qq.com/s/RLzJRSTRq-BuUdoGqdwP8Q