← 全部文章

在本地搭建 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》。

image

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

image-20250330105204780

图片来源:《AI风暴来袭:2024年数据平台的演进、挑战与机遇》

https://mp.weixin.qq.com/s/Iz2VPtZWF-ByZTfbcduvlQ

02 部署过程

1. 蓝图

Data & AI 数据平台整体示意图

image

参考《新一代DA平台架构的设计原则与演进思路》

https://mp.weixin.qq.com/s/ukqztUYXwDE9jCONvzEn8A

https://mp.weixin.qq.com/s/RLzJRSTRq-BuUdoGqdwP8Q

本文涉及的离线湖仓的部署示意图

image

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
配置对象存储

界面上配置,连接到对象存储

image

(

https://chatgpt.com/c/67d8cf18-3f24-8012-adbe-5c94eb954428

Enable path style access(启用路径样式访问)

  • 关闭时(默认),使用虚拟托管方式访问,如https://bucket-name.s3.amazonaws.com

  • 开启后,使用路径样式访问,如https://s3.amazonaws.com/bucket-name

  • MinIO 这类 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

https://airflow.apache.org/docs/apache-airflow/stable/administration-and-deployment/production-deployment.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/

  1. 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
  1. Download the latest DEB package.
  2. 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️⃣配置代理

方便拉取镜像

image

2️⃣启用 kubernetes

https://docs.docker.com/desktop/features/kubernetes/

选择 kubeadm 的方式创建 kubernetes 集群

image

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")

界面示例

image

image

(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 ,配置 源 和 目标 。

image

4️⃣源PostgreSQL

image

5️⃣目标Iceberg

image

6️⃣结果示例

1192529 条记录,12 分钟32秒,平均 1585 条/秒。

期间 Airbyte 的 Pod 崩溃重启,但最终数据全部同步到 Iceberg 。

image

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

image

  • 数据查询

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

image

  • 数据看板

数据来源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年数据看上去缺少很多)并且状态为"已上映" 。

最终的数据看板

image

彩蛋

image

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

3. 微信实验平台Iceberg湖仓一体架构改造

https://mp.weixin.qq.com/s/5E2nq51qA8qo_WFelykU5g