第 1 篇:「认识 Fluss」—— 新一代流式存储系统概览
第 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亚毫秒级主键查询RocksDBFlink 状态Delta Join Aggregation Merge Engine状态外部化Iceberg / PaimonTiering Service自动 Compaction 为 Parquet 写入湖仓1.2 Fluss 的核心设计哲学Fluss 的定位关键词是Streaming Storage——不是 Streaming Transport传输而是 Storage存储。KafkaStreaming TransportFlussStreaming Storage数据形态仅追加日志行append-only rows表TablesLog Table PK Table读取方式消费端按 offset 逐条拉取服务端列裁剪、谓词下推、PK Lookup状态归属消费端自行维护RocksDB服务端 KvStore计算节点无状态湖仓集成外部 Connector Sink原生 Tiering Union ReadCDCKafka Connect Debezium原生$changelog/$binlog虚拟表一句总结Kafka 帮你把数据从 A 搬到 BFluss 帮你直接在原地查询和分析数据。1.3 六大能力支柱全景解读支柱 1统一架构Unified Architecture一套系统同时提供消息传输、KV 查询、OLAP 分析能力。-- 注册 Fluss CatalogCREATECATALOG fluss_catalogWITH(typefluss,bootstrap.serverscoordinator-server:9123);USECATALOG fluss_catalog;-- 创建 PK 表同时是流式日志 KV 索引 分析表CREATETABLEuser_profile(user_idBIGINT,name STRING,city STRING,last_loginTIMESTAMP(3),PRIMARYKEY(user_id)NOTENFORCED)WITH(bucket.num8);-- 亚毫秒 PK 查询SELECT*FROMuser_profileWHEREuser_id12345;-- 流式读取变更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_idu.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 OpennessApache FlussFlink SQL/DataStreamSpark SQLTrinoStarRocksDuckDBIceberg V2/V3PaimonLance1.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.ymlversion:3.8services:zookeeper:image:zookeeper:3.8ports:-2181:2181coordinator-server:image:apache/fluss:0.9.1command:coordinatorServerdepends_on:-zookeeperenvironment:-ZOOKEEPER_ADDRESSzookeeper:2181ports:-9123:9123tablet-server:image:apache/fluss:0.9.1command:tabletServerdepends_on:-coordinator-serverenvironment:-ZOOKEEPER_ADDRESSzookeeper:2181启动集群docker-composeup-d# 验证集群状态docker-composelogs coordinator-server|grepstarted使用 Flink SQL Client 连接# 启动带 Fluss Connector 的 Flink SQL Clientdockerrun-it--networkhostapache/fluss-quickstart-flink:0.9.1# 在 Flink SQL Client 中执行CREATE CATALOG fluss_catalog WITH(typefluss,bootstrap.serverslocalhost: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.num4);-- 写入数据 INSERT INTO my_first_table VALUES(1,Hello Fluss);-- 查询数据 SELECT * FROM my_first_table WHEREid1;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