五层递进,一台机器跑全栈

从 MVP 湖仓到 AI 接入层,5 个 Docker Compose 文件逐层叠加。15 分钟完成全栈部署。

AI/门户层Portal/MCP/Superset
全栈双平台层Flink/MySQL
实时采集层Kafka/CDC
离线中台层DS/JupyterLab
MVP湖仓层Iceberg/Spark/Trino/MinIO

五层架构,按需组装

每层独立可运行,按业务需求叠加。不是「全家桶强行安装」,而是「乐高式按需搭建」。

统一接入层:AI + 门户

.portal.yml

Portal (FastAPI + React) · MCP Server · Superset BI

API Key 认证 · 统一 Web UI

全栈双平台层:离线 + 实时

.flink.yml

Flink SQL Gateway · MySQL Business · 实时聚合

Flink CDC → Iceberg · 流批一体

实时采集层:消息队列

.streaming.yml

Kafka · Redpanda Console · 实时 CDC

MySQL → Kafka → MinIO 管道

离线中台层:调度 + 开发 + BI

.addons.yml

DolphinScheduler · JupyterLab · Supervisor

Streamlit 报表 · DuckDB 分析

MVP 湖仓层:存储 + 计算

docker-compose.yml

Apache Iceberg · Spark SQL · Trino · MinIO (S3)

PostgreSQL (REST Catalog) · 时间旅行 · ACID

一条数据,五层流转

从外部 MySQL 到 AI 问答,追踪数据的完整生命周期

MySQL
Kafka
Flink
Iceberg
MinIO
Spark ETL

ODS→DWD→DWS

Trino 查询

BI 看板

DuckDB 分析

Streamlit

🤖

AI Agent: "帮我分析销售趋势"

search_tables → NL2SQL → 结论

一份数据,多个引擎

Iceberg 表格式 + MinIO S3 对象存储,Spark/Trino/Flink/DuckDB 四引擎共享同一份 Parquet 文件。

Iceberg v2行级 DELETE/UPDATE时间旅行MERGE INTO分区裁剪元数据快照

MinIO S3

warehouse/*.parquet

Spark

写/改/删

Trino

只读查询

Flink

CDC写入

DuckDB

临时分析

写用 Spark,读用 Trino,各司其职

Spark SQL(写入)

-- ODS 全量快照,T-1 分区
INSERT OVERWRITE demo.gmall_dw.ods_order_info
PARTITION (pt = '20260629')
SELECT id, consignee, total_amount, order_status,
       user_id, payment_way, province_id, create_time
FROM mysql_gmall.order_info

Trino(查询)

-- 一条 SQL 跨 TPC-H + Iceberg
SELECT n.name AS nation,
       COUNT(DISTINCT o.orderkey) AS orders,
       ROUND(SUM(l.extendedprice), 2) AS revenue
FROM tpch.tiny.customer c
JOIN tpch.tiny.orders o   ON c.custkey = o.custkey
JOIN tpch.tiny.lineitem l ON o.orderkey = l.orderkey
JOIN tpch.tiny.nation n   ON c.nationkey = n.nationkey
GROUP BY n.name ORDER BY revenue DESC
维度Spark SQLTrino
职责数据入湖、ETL 写入联邦查询、BI 分析
写入✅ INSERT / MERGE / DELETE❌ 只读
速度批量处理,分钟级交互查询,秒级
起点SHOW DATABASESSHOW CATALOGS(7 个源)

不只是 Parquet,是完整的湖仓表格式

🕰

时间旅行

回到任意历史快照,数据永不丢失

FOR VERSION AS OF {snapshot_id}
✏️

行级更新

DELETE / UPDATE / MERGE INTO

Iceberg v2 + Copy-on-Write
📊

分区演进

不重写数据即可变更分区策略

Partition Evolution
🔍

元数据快照

秒级查表行数,无需全表扫描

snapshots 系统表

时间旅行动画演示

T1 快照

10 行数据

张三
李四
王五

T2 DELETE 孙七

9 行数据

张三
李四
王五

T3 时间旅行恢复

10 行数据

张三
李四
王五
孙七

一个 Portal,管理一切

数据源配置、文件浏览、SQL 查询、Schema 探索、血缘 DAG、Job 管理,6 个功能集成在同一个工作台。

6 个 Tab 无需切工具数据源配置自动同步 JupyterLabanalyze_table 一键输出401 自动跳转登录Token 过期优雅降级

Portal API (FastAPI)

/jupyter/* /trino/* /datasources/* /scheduler/* /jobs/* /auth/*

📁

文件管理

📡

数据源管理

🔍

SQL 查询

📋

Schema 浏览

🕸

血缘 DAG

🖥

Jobs 管理

JupyterLab · Trino · Spark · DolphinScheduler · Supervisor

PocketLakehouse 技术栈

14 个核心开源组件,协同工作

存储

MinIO

S3 兼容

Apache Iceberg

表格式

PostgreSQL

Catalog

MySQL

Business/DS

计算

Apache Spark

ETL

Trino

查询

Apache Flink

流处理

DuckDB

分析

工具链

JupyterLab

开发

DolphinScheduler

调度

Supervisor

进程

Streamlit

报表

接入层

FastAPI

Portal

React

前端

MCP Server

AI Agent

Kafka

消息