第 1 篇:「认识 Fluss」—— 新一代流式存储系统概览
阅读本文你将了解:Fluss 是什么、为什么需要它、它与 Kafka 的本质区别、Fluss 的六大核心能力、以及在 Docker 中 5 分钟快速体验。
1.1 为什么需要 Fluss:实时数据栈的「多系统税」
构建一个典型的实时数据管道,你需要什么?
┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │ Kafka │───→│ Flink │───→│ Redis │ │ Iceberg │ │ HBase │ │ (传输层) │ │ (计算层) │ │ (在线存储) │ │ (离线存储) │ │ (查询服务) │ └──────────┘ └──────────┘ └──────────┘ └──────────┘ └──────────┘ │ │ │ │ │ └───────┬───────┴───────┬───────┴───────┬───────┘ │ │ │ │ │ ┌───────┴───────┐ ┌─────┴─────┐ ┌───────┴───────┐ ┌────────┴────────┐ │ Schema Registry│ │ 运维监控 │ │ 同步管道 │ │ CDC/Debezium │ └───────────────┘ └───────────┘ └───────────────┘ └─────────────────┘五套系统、四个集成边界、持续的工程税。每个边界都是数据可能发生漂移的地方。
Fluss 的答案是把这五层压缩成一个统一的流式存储基座:
┌────────────────────────────────────────────────────────────┐ │ Apache Fluss │ │ │ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌───────────┐ │ │ │ Streaming│ │ PK │ │ State │ │ Multi- │ │ │ │ Log │ │ Lookup │ │ Store │ │ Modal │ │ │ └──────────┘ └──────────┘ └──────────┘ └───────────┘ │ │ ↓ 一个基座、零同步边界、单一事实来源 ↓ │ ├────────────────────────────────────────────────────────────┤ │ Flink ←→ Spark ←→ Trino ←→ StarRocks ←→ DuckDB │ │ (所有引擎读写同一份数据,无需额外同步层) │ └────────────────────────────────────────────────────────────┘Fluss 替代了什么?
| 传统组件 | Fluss 如何替代 |
|---|---|
| Apache Kafka | Log Table:持久化、可回放的流式日志 |
| Redis / HBase | PK Table + KvStore:亚毫秒级主键查询 |
| RocksDB(Flink 状态) | Delta Join + Aggregation Merge Engine:状态外部化 |
| Iceberg / Paimon | Tiering Service:自动 Compaction 为 Parquet 写入湖仓 |
1.2 Fluss 的核心设计哲学
Fluss 的定位关键词是Streaming Storage——不是 Streaming Transport(传输),而是 Storage(存储)。
| Kafka(Streaming Transport) | Fluss(Streaming Storage) | |
|---|---|---|
| 数据形态 | 仅追加日志行(append-only rows) | 表(Tables):Log Table + PK Table |
| 读取方式 | 消费端按 offset 逐条拉取 | 服务端列裁剪、谓词下推、PK Lookup |
| 状态归属 | 消费端自行维护(RocksDB) | 服务端 KvStore,计算节点无状态 |
| 湖仓集成 | 外部 Connector Sink | 原生 Tiering + Union Read |
| CDC | Kafka Connect + Debezium | 原生$changelog/$binlog虚拟表 |
一句总结:Kafka 帮你把数据从 A 搬到 B,Fluss 帮你直接在原地查询和分析数据。
1.3 六大能力支柱全景解读
支柱 1:统一架构(Unified Architecture)
一套系统同时提供消息传输、KV 查询、OLAP 分析能力。
-- 注册 Fluss CatalogCREATECATALOG fluss_catalogWITH('type'='fluss','bootstrap.servers'='coordinator-server:9123');USECATALOG fluss_catalog;-- 创建 PK 表(同时是流式日志 + KV 索引 + 分析表)CREATETABLEuser_profile(user_idBIGINT,name STRING,city STRING,last_loginTIMESTAMP(3),PRIMARYKEY(user_id)NOTENFORCED)WITH('bucket.num'='8');-- 亚毫秒 PK 查询SELECT*FROMuser_profileWHEREuser_id=12345;-- 流式读取变更SELECT*FROMuser_profile$changelog;源码支撑:Fluss 的 PK 表在服务端同时维护 LogStore 和 KvStore,由org.apache.fluss.server.tablet.TabletService统一管理。
支柱 2:流式与湖仓统一(Stream & Lakehouse Unification)
冷热数据共享同一 Schema,统一查询入口。
┌─────────────────────────────────────────────┐ │ SQL: SELECT * FROM orders │ ├──────────────────┬──────────────────────────┤ │ Fluss Hot Tier │ Lakehouse Cold Tier │ │ (Arrow, 秒级) │ (Parquet, 分钟级) │ │ 保留 3 天 │ 保留 12 月 │ ├──────────────────┴──────────────────────────┤ │ Union Read 自动合并结果 │ └─────────────────────────────────────────────┘支柱 3:计算存储分离(Compute/Storage Separation)
状态从 Flink 的 RocksDB 移出,放到 Fluss 的 TabletServer 上。
// 传统方式:状态在 Flink 本地 RocksDB// 问题:Checkpoint 大、恢复慢(分钟级)、扩缩容受限// Fluss Delta Join:状态外部化CREATETABLEwide_table(order_idBIGINT,user_idBIGINT,amountDECIMAL(10,2),user_nameSTRING,user_citySTRING,PRIMARYKEY(order_id)NOTENFORCED)WITH('bucket.num' = '16');--Flinkjob 直接写外部化状态,无需维护本地RocksDBINSERTINTOwide_tableSELECTo.order_id,o.user_id,o.amount,u.name,u.cityFROMorders oLEFTJOINuser_profileFORSYSTEM_TIMEASOFo.ptimeASuONo.user_id=u.user_id;--↑ 状态在Fluss端,Flink任务是无状态的!支柱 4:列式流分析(Columnar Streaming Analytics)
基于 Apache Arrow 列式格式,服务端做列裁剪和谓词下推:
查询: SELECT user_id, amount FROM orders WHERE amount > 100 ┌─────────────────────────────────────────────────────┐ │ 200 列的表,只读取 2 列 → I/O 降低 ~99% │ │ WHERE amount > 100 → 谓词下推到 TabletServer 执行 │ └─────────────────────────────────────────────────────┘支柱 5:特征与上下文存储(Feature & Context Stores)
行、列、向量三种数据格式统一存储在 Fluss,通过不同视图访问。
源码支撑:org.apache.fluss.lake.LanceLakeFormat提供向量格式的湖仓存储支持。
支柱 6:生态开放性(Ecosystem Openness)
1.4 适用场景
| 场景 | Fluss 的价值 |
|---|---|
| 实时分析大屏 | 亚秒级数据新鲜度 + 列式查询,200 列的表只读需要的几列 |
| 实时特征存储 | PK Table 支持亚毫秒级特征查询,直接对接 ML 推理服务 |
| CDC 数据管道 | 原生$changelog,无需 Debezium + Schema Registry |
| 实时数仓 | Streaming Lakehouse 架构,流批统一查询 |
| 风控引擎 | 状态外部化,秒级故障恢复,计算弹性扩缩 |
| 客户 360 | 多源数据实时融合,一张 PK Table = 一个统一客户视图 |
1.5 快速体验:Docker Compose 一键启动
前置条件
- Docker & Docker Compose
- 至少 8 GB 可用内存
docker-compose.yml
version:'3.8'services:zookeeper:image:zookeeper:3.8ports:-"2181:2181"coordinator-server:image:apache/fluss:0.9.1command:coordinatorServerdepends_on:-zookeeperenvironment:-ZOOKEEPER_ADDRESS=zookeeper:2181ports:-"9123:9123"tablet-server:image:apache/fluss:0.9.1command:tabletServerdepends_on:-coordinator-serverenvironment:-ZOOKEEPER_ADDRESS=zookeeper:2181启动集群
docker-composeup-d# 验证集群状态docker-composelogs coordinator-server|grep"started"使用 Flink SQL Client 连接
# 启动带 Fluss Connector 的 Flink SQL Clientdockerrun-it--networkhostapache/fluss-quickstart-flink:0.9.1# 在 Flink SQL Client 中执行CREATE CATALOG fluss_catalog WITH('type'='fluss','bootstrap.servers'='localhost:9123');USE CATALOG fluss_catalog;CREATE DATABASE test_db;USE test_db;-- 创建你的第一张 PK 表 CREATE TABLE my_first_table(idBIGINT, name STRING, PRIMARY KEY(id)NOT ENFORCED)WITH('bucket.num'='4');-- 写入数据 INSERT INTO my_first_table VALUES(1,'Hello Fluss');-- 查询数据 SELECT * FROM my_first_table WHEREid=1;1.6 总结与下一篇预告
| 关键词 | 含义 |
|---|---|
| Streaming Storage | Fluss 的核心定位:不是传输工具,是存储基座 |
| PK Table | 同时是流式日志、KV 索引和分析表 |
| 外部化状态 | Flink 不再抱着 RocksDB,状态归 Fluss 管 |
| Union Read | 一个 SQL 查询同时覆盖实时数据和历史数据 |
下一篇我们将深入 Fluss 的分布式架构——CoordinatorServer 如何管理集群、TabletServer 如何存储数据、ZooKeeper 和 Remote Storage 各自扮演什么角色。理解这些底层机制,是后续进行表设计和性能优化的基础。
本文基于 Apache Fluss 0.9.1。项目 GitHub: https://github.com/apache/fluss