实时分析、实时上下⽂与AI的数据基座 版权与法律声明 版权 Copyright © 2026 The Apache Software Foundation。本⽂档内容基于ApacheLicense, Version 2.0发布(http://www.apache.org/licenses/LICENSE-2.0)。 商标 Apache®、Apache Fluss、Fluss、the Apache feather logo与Apache Incubatorproject logo均为The Apache Software Foundation在美国及/或其他国家的注册商标或商标。 本⽂档中提及的Apache Flink、Apache Kafka、Apache Spark、Apache Iceberg、Apache Paimon、Apache Hudi、Apache Arrow同为The Apache SoftwareFoundation的商标。 其他产品名、Logo与品牌归各⾃所有者所有。 孵化器免责声明 中⽂:本⽩⽪书为社区实践总结与使⽤指南,⾮Apache软件基⾦会(ASF)的官⽅出版物,内容不代表ASF或孵化器PMC的⽴场或背书。Apache Fluss当前处于ASF孵化器孵化阶段。 English: Apache Fluss is an effort undergoing incubation at The Apache SoftwareFoundation (ASF), sponsored by the Apache Incubator. Incubation is required of allnewly accepted projects until a further review indicates that the infrastructure,communications, and decision making process have stabilized in a mannerconsistent with other successful ASF projects. While incubation status is not necessarily a reflection of the completeness or stability of the code, it does indicatethat the project has yet to be fully endorsed by the ASF. ThisdocumentisnotanofficialApacherelease. ⽬录 摘要致读者序章基于ApacheFluss的流式湖仓基于Apache Fluss的Streamhouse实践Apache Fluss的⻆⾊与定位基于Apache Fluss的Streamhouse核⼼价值第⼀章现有基础设施为何⼒不从⼼1.1⾛向实时AI基础设施的五段路1.2实时智能真正需要的基础设施1.3多系统税:碎⽚化架构的隐性代价第⼆章集群即⼀个逻辑系统2.1集群拓扑2.2存储引擎与⽇志格式2.3写路径、读路径与副本机制第三章两种表,同⼀存储底座3.1⽇志表:仅追加的事件流存储3.2主键表:流与表的统⼀存在3.3 Changelog:⼀等公⺠第四章三级存储,同⼀底座4.1同⼀底座中的三级存储4.2热层数据的两条归档路径4.3开放湖表格式:冷层⽬标的选择4.4联合读取:打破边界第五章复合裁剪:列裁剪、谓词下推与分区裁剪5.1服务端列裁剪5.2谓词下推⽀持5.3分区裁剪与复合效应5.4⽹络传输代价的回响第六章状态归位:Fluss作为状态底座6.1存储计算分离:状态外置6.2 Delta Join:⽆状态的双流Join 6.3聚合合并引擎6.4状态外移后的恢复语义6.5流式Lookup Join6.6部分更新:⽆需Join的数据打宽6.7消除训练-服务偏差6.8 Fluss作为AI系统的上下⽂存储6.9虚拟表:同⼀状态的多重视图6.10同⼀底座上的实时实体画像附录A三级存储运维参考A.1三级存储⼀览A.2写⼊流⽔线各阶段A.3配置参考附录B关键术语中英对照后记 摘要 本⽩⽪书系统阐述Apache Fluss的设计理念与架构实践:以⼀份数据、⼀个底座,承载实时分析、特征⼯程与AI上下⽂⼯程等多类⼯作负载。 过去⼗年,企业为构建实时智能系统,往往需要拼接消息队列、KV服务、OLAP引擎、特征存储与向量库等多类系统。每类系统各司其职,但系统之间的边界形成了⼤量数据搬运与⼀致性维护成本,即所谓“多系统税(Multiple Systems Tax)”。系统越多,数据在边界越容易出现偏差,⼯程团队的精⼒被迫从业务创新转向基础设施⼀致性维护。Apache Fluss提供了另⼀种解决思路:作为⼀款流式优先、湖仓原⽣的列式存储引擎,将事件传输、KV服务、列式分析与AI上下⽂等原本割裂的能⼒收敛到统⼀的存储底座。Fluss通过主键表实现流与表的统⼀,通过分层服务(TieringService)贯通热数据层与冷数据层,通过Union Read实现湖流透明读取,通过Delta Join与聚合合并引擎将计算回归⽆状态。 本书⾯向架构师、研发⼯程师与AI应⽤开发⼯程师,系统介绍Apache Fluss的整体架构、表模型、分层存储、裁剪机制与状态语义。读者阅读完毕后,应能独⽴判断Fluss是否契合其业务场景,并据此规划⾃身的实时数据底座。 作者 作者:Giannis Polyzos,Ververica⾸席流处理架构师,⻓期专注于流处理与湖仓存储的⼯程实践。深度参与Apache Fluss与Apache Flink社区建设,关注实时分析、特征⼯程与AI原⽣数据基础设施的演进路线。 翻译:徐榜江(雪尽),阿⾥云⾼级技术专家,Flink PMC Member & FlussCommitter。 致读者 数据基础设施的范式正在改变。⼆⼗年前,问题是“如何让数据在夜间被批量处理”;⼗年前,问题变成“如何让事件在分钟级被消费”;今天,问题已经升级为“如何让⼀份数据同时服务于流式分析、机器学习与⽣成式AI”。每⼀次需求的演进,都引⼊了新的系统: Kafka⽤于流式传输、Redis⽤于在线服务、Iceberg⽤于离线归档、向量数据库⽤于语义检索。能⼒不断叠加,却从未真正收敛。本⽩⽪书要讨论的,正是这种收敛。 本书将沿⼀条主线展开:“⼀个底座承载⼀切”为什么可⾏、为什么必要。我们从企业实时基础设施的演进路径出发,分析多系统税的形成原因;进⽽深⼊介绍ApacheFluss的整体架构、存储引擎、双表模型与三级存储;接着讨论数据裁剪如何在存储侧⽽⾮计算侧消除冗余开销;最后回到状态归属问题——状态应当归属于存储层,⽽⾮由流计算引擎承担。 读完本书,您未必认同每⼀个判断,但我们希望提供⼀组可以被讨论、被反驳的论据,为您下⼀次架构选型提供参考。 序章基于ApacheFluss的流式湖仓 1基于apacheFluss的Streamhouse实践 Ververica Streamhouse是⼀个基于Apache Fluss的统⼀流式湖仓平台(即社区惯称的"湖流⼀体"架构),其设计围绕⼀条核⼼原则:⼀份数据,服务于所有⼯作负载。流式写⼊、低延迟分析、批处理、机器学习特征与AI上下⽂都在同⼀份数据之上完成,避免了传统架构层层叠加的多系统税。 它将以下组件整合为⼀个统⼀的运营底座: 开源流式存储层——Apache Fluss开放湖仓层——Apache Paimon、Apache Iceberg、Lance统⼀的流批处理引擎——VERA(基于ApacheFlink的企业版)声明式数据管道——物化表(MaterializedTables)⾃动调优服务——Autopilot感知新鲜度的⼯作流调度器●●●●●● 上图为基于Apache Fluss的Streamhouse平台架构。Apache Fluss处于该平台的流式存储层,作为整个平台的低延迟列式存储底座。 1ApacheFluss的⻆⾊与定位 Apache Fluss是Ververica Streamhouse的低延迟列式流存储层。存储层之上有物化表、VERA引擎与⼯作流调度器,它们都直接读写Fluss;存储层之下是搭建在对象存储上的开放湖仓格式,由分层服务(Tiering Service)接收Fluss的⽇志段,并通过Union Read(联合读取)向上层计算引擎透明可读。 Fluss将原本需要消息队列(Kafka)、KV存储(Redis)与OLAP热存储(StarRocks、Doris、ClickHouse)三类系统协同承担的能⼒收敛到⼀个系统中,并以三种访问模式共享同⼀份底层数据: 1.KV点查:服务于在线特征与异步Lookup Join(异步维表点查Join);2.流式⽇志读取:服务于CDC消费与下游加⼯;3.列式批量扫描:服务于分析查询与训练数据构建。下⽂将依次介绍Apache Fluss如何实现上述能⼒。 2基于ApacheFluss的Streamhouse核⼼价值 基于Apache Fluss的Streamhouse为上层⼯程团队带来的核⼼价值,可以归纳为以下6个⽅⾯。 每⼀项核⼼价值都对应着第⼆章⾄第六章将要深⼊介绍的某项具体架构或技术特性,并对应⼀种可量化的收益: 统⼀架构:⼀个系统同时服务消息队列、应⽤开发、分析查询与AI场景,整合此前由消息队列、KV存储与OLAP引擎分别承担的⻆⾊。其架构基础是第⼆、三章介绍的主键表的双重表示——⽇志(历史变更)+ Leader上的RocksDB KV(快照)。● 湖流⼀体:实时层与批处理层共享同⼀份数据,元数据⾃动同步。其架构基础是第四章介绍的分层服务(Tiering Service)与Union Read:Fluss的热数据层与冷数据层共享同⼀逻辑Schema,查询可在统⼀存储底座上⽆缝完成。存算分离:存储与计算解耦,使计算层精简、弹性、⽆状态、恢复迅速;相⽐同等规模的Kafka集群,成本最⾼可降低85%。其架构基础是第六章介绍的⽆状态计算模型——状态外置到Fluss集群,⽽不再驻留于每个Flink作业的TaskManager中,恢复时间被显著压缩,RPO与RTO得以解耦。列式分析:在压缩列式数据上⾼效执⾏流式查询,基于Apache Arrow格式实现服务端列裁剪与谓词下推,显著降低⽹络传输代价。其架构基础是第五章介绍的复合裁剪机制——列裁剪、谓词下推与分区裁剪叠加作⽤,可实现数量级的数据缩减。特征与上下⽂存储:⾏式、列式、向量等多模态数据的访问,以及在线特征服务与RAG上下⽂,都在同⼀存储底座上完成。其架构基础是第六章介绍的统⼀底座——特征存储、上下⽂存储与实时实体画像统⼀收敛到⼀张主键表,由不同视图按需访问。⽣态开放:开放共享的存储底座,可被Apache Flink、Apache Spark、Trino、StarRocks、DuckDB等计算引擎直接读取;项⽬以Apache协议开源,⽆⼚商锁定。其架构基础是Streamhouse在数据冷层提供对开放湖仓格式(ApacheIceberg、Apache Paimon、Lance)的⽀持,以及Apache Fluss数据热层的原⽣连接器。●●●●● 第⼀章现有基础设施为何⼒不从⼼ 基础设施的演进,常常是把新能⼒叠加到既有系统之上,最终得到的不是⼀个统⼀的底座,⽽是⼀堆系统的堆叠。 1.1⾛向实时AI基础设施的五段路 1.1.1流式写⼊与变更数据捕获 演进路径⼏乎总是从同⼀个起点开始:业务数据库不断膨胀,批量ETL难以承载⽇益增⻓的数据量,团队由此开始引⼊消息中间件。 ⼯程师部署Kafka,将应⽤事件(点击流、交易、⽤户⾏为)以流的⽅式写⼊数据湖;与此同时,变更数据捕获(Change Data Capture,