第 1 篇:「认识 Fluss」—— 新一代流式存储系统概览
2026/8/10 8:56:14 网站建设 项目流程

第 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 KafkaLog Table:持久化、可回放的流式日志
Redis / HBasePK Table + KvStore:亚毫秒级主键查询
RocksDB(Flink 状态)Delta Join + Aggregation Merge Engine:状态外部化
Iceberg / PaimonTiering 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
CDCKafka 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)

Apache Fluss

Flink SQL/DataStream

Spark SQL

Trino

StarRocks

DuckDB

Iceberg V2/V3

Paimon

Lance


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 StorageFluss 的核心定位:不是传输工具,是存储基座
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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询