Apache Fluss:实时数仓的流式存储层,如何补上Kafka与数据湖之间的断层 做实时数仓的同行应该都有这种体会一提到实时数据第一个想到的是Kafka一提到分析查询马上又接一个数据湖或OLAP引擎。两套系统之间搬运数据、保证一致性几乎成了每天的常规操作。Apache Fluss就是冲着这个痛点来的它把自己定位成“流式存储”Streaming Storage——既提供Kafka那样的低延迟流式写入和订阅又提供数据库表那样的主键更新、点查和数据结构化语义。这篇文章我从“它解决了什么问题”开始讲然后会拆解它的核心表模型、一条数据从写入到可查的完整链路、和Flink的集成方式最后分享我从小集群跑通到压测调优踩过的一些坑。适合正在做实时数仓、实时湖仓选型或者想把Flink链路里的存储层再往前推一步的工程师参考。1. 实时湖仓的存储断层为什么Kafka和数据湖中间还缺一层1.1 经典实时链路的拧巴之处很多公司的实时链路长这样业务库和业务日志通过CDC或者SDK进入KafkaFlink消费Kafka做加工再写进Iceberg/Hudi/Paimon最后供OLAP查询。这套架构单独看每一环都很成熟但合在一起有个尴尬的问题——Kafka并不是一个“存储系统”。消息默认保留几天就被清理topic里的数据不能按主键更新你也没法对topic做点查。对Flink来说Kafka只是“暂时停靠的传送带”数据一旦从Kafka出去历史就被遗忘如果想再回查某个订单的最新状态就得到下游湖表里跑一遍扫描或者再维护一套在线存储。这也是很多团队最终走向“Lambda架构”的原因实时链路走Kafka Flink Redis/HBase离线链路走Kafka 定时任务 数据湖。两条链路处理同一份数据逻辑写两遍结果还经常对不上。我见过不少团队花大量时间在“实时结果和离线结果为什么差了3分钟”这类问题上根因往往不是计算逻辑错了而是存储层没有一个能同时满足实时消费和历史查询的统一出口。1.2 数据湖为什么也不能完全兜住数据湖表天然是为大吞吐分析设计的擅长批量小文件合并、事务提交、快照隔离但它的写入路径很重落数据通常是分钟级甚至小时级频繁提交小事务还会产生大量小文件给合并带来压力。所以数据湖适合当“结果仓库”不适合当“实时数据中转站”。于是实时场景里就出现了很典型的双写一份数据同时进Kafka供实时消费又进HBase/Redis提供点查再定时落湖供分析。链路过长成本高最要命的是三份数据很难保证强一致。Paimon这类流式湖格式已经解决了一部分问题它支持主键表、支持小文件自动合并让数据可以“近实时”地进湖。但Paimon依然把重点放在湖格式上写入路径仍然偏批式实时更新和点查不是它的强项。换句话说湖存储擅长“沉淀历史”不擅长“服务当下”。而实时业务里大量场景需要的是数据刚产生就能被查到而且这个“查”不是扫描全表而是像数据库一样按主键拿最新值。1.3 Fluss想填补的空档Fluss想做的事情用一句话概括把“流的写入速度”和“表的查询能力”放进同一个存储系统。它不再只是一个消息管道而是一个真正以表为单位的实时存储层。数据写进去之后既能像Kafka一样被订阅消费又能按主键更新、被直接查询后台还会自动把历史数据转成列存格式放上对象存储实现冷热分层。这样Flink链路里就不再需要为了“既要实时又要可查”去维护两套系统。这个定位听起来很理想但要落地得先理解它的两个核心表模型——这也是Fluss和Kafka最本质的区别。顺便说一句Fluss这个名字来自德语意思就是“河流”项目理念里强调数据像水一样不断流动同时又有河床可以沉淀和查询。这个比喻其实很准确河流是动态的河床是静态的Fluss想同时当好这两者。理解了这句话后面看它的架构就不会觉得绕。2. 两张表撑起整个设计Log Table和Primary Key Table2.1 表而不是topic是Fluss的第一印象用过Kafka的人第一次接触Fluss最直观的感受是它竟然是用建表的方式来管理实时数据。这是因为Fluss的数据模型是“表”不是“topic”。表意味着有字段、有类型、有主键、有约束Flink SQL可以像操作普通表一样去定义它、读写它。这个差异非常关键因为它把存储从“字节流”升级成了“结构化数据集”。这带来的直接好处是你在Flink SQL里写CREATE TABLE就能创建实时存储写INSERT INTO就能持续写入写SELECT就能查询写JOIN就能和别的表关联。整个过程几乎不需要理解底层文件格式和分区细节。对比Kafka需要维护topic、schema registry、序列化格式这一整套配套Fluss把概念大大简化了。2.2 Log Table纯追加事件流Log Table是Fluss里最像Kafka的部分只支持追加写入append-only不支持更新和删除适合存日志、埋点、订单流转事件这类不可变数据。它提供低延迟的实时订阅能力Flink消费起来和消费Kafka很接近。如果你只是需要一个比Kafka语义更丰富的“事件流”Log Table就够了。但它在Fluss里的定位更像是一个基础底座真正体现能力的是Primary Key Table。我自己的理解是Log Table的价值不在它本身多特殊而在于它让Fluss从一个“键值存储”扩展成了“也能处理事件流”的系统。一个存储只有log没有表就退化成了消息队列只有表没有log就退化成了普通数据库。Fluss两个都要所以同一个集群里既能跑流式ETL又能跑在线查询。2.3 Primary Key Table支持更新的流式表Primary Key Table的核心特征是带主键支持upsert支持点查。这意味着一笔订单从“已下单”变成“已支付”再变成“已发货”你不需要往表里追加三条流水而是可以用同一主键持续更新这条记录的最终状态。做实时数仓的同学看到这里应该会很兴奋——这正好是实时维表、实时宽表最需要的能力。以前我们用Kafka HBase拼出这套效果现在Fluss把“log订阅”和“KV点查”两种能力合到了一张表里。它甚至可以作为CDC链路的目标端把MySQL Binlog实时同步成一张可查询的实时表。实际操作中主键表的写入是幂等的同一条主键的数据反复写最后只保留最新版本。这个语义对“订单状态”这类不断变化的实体非常友好下游读到的永远是最新的业务事实不需要再去重或做窗口去重。同时它还保留了log的订阅能力这意味着下游Flink任务可以实时感知每一笔更新。Kafka做不到更新传统数据库做不到实时订阅Fluss主键表的组合能力正好卡在这个交叉点上。2.4 底层为什么能同时做到“流”和“表”从架构上看Fluss集群里主要有Coordinator Server和Tablet Server两类节点。Coordinator负责元数据和分桶调度Tablet Server负责实际数据读写。数据写入时会根据主键或指定分桶键做哈希路由到对应的bucketbucket先以日志log形式写内存和本地磁盘保证写入低延迟同时后台会把日志中的增量数据定期转换成列存文件上传到对象存储。这样一张主键表就同时拥有了日志的“可订阅性”和列存的“可分析性”。不理解这一层后面看数据写入链路就容易懵。这个设计让我想到HBase的LSM结构写入先走内存MemStore再异步刷盘成HFile。Fluss的思路类似但它走得更远——log部分承担了消息队列的功能列存文件则承担了分析存储的功能一份写入两用不需要像传统架构那样写两遍。这种“写入一次多种读法”的思想是理解Fluss性价比的关键。3. 一条订单数据从写入到可查的完整旅程3.1 先搭一个最小的集群环境要验证这套链路最快的方式是下一份Fluss发行版以Quickstart方式起一个单机伪集群。以当前版本的常见步骤是下载二进制包解压确认本机有JDK按官方文档启动依赖的协调服务再启动Fluss Server。启动完成后会看到服务监听在默认端口可以通过Web界面确认集群状态。整个流程十分钟内能跑起来但这里有个容易忽略的点Fluss会把bucket数据先落在本地磁盘再异步上传到远端对象存储所以本地数据盘的剩余空间和文件句柄数直接影响能否稳定运行。我用的是单机伪集群做功能验证生产环境至少需要3个Tablet Server节点来做副本容灾。这里提醒一下Quickstart里的默认配置主要用于“跑通”不是“压测”如果你一上来就拿着默认配置灌大数据大概率会遇到写入延迟飙升、刷盘任务积压的问题。第一次玩建议只导入几万条数据验证链路不要急着压测。3.2 用Flink SQL创建一张订单主键表环境就绪后我习惯用Flink SQL来做演示因为它最能体现表语义。先在Flink里注册Fluss的Catalog然后执行建表。比如CREATE CATALOG fluss_catalog WITH ( type fluss, bootstrap.servers localhost:9123 );然后USE CATALOG fluss_catalog; CREATE TABLE orders ( order_id BIGINT, user_id BIGINT, amount DOUBLE, order_status STRING, ts TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( bucket.num 4, bucket.key order_id );这里的核心参数是bucket.num和bucket.key。bucket.num决定数据被切成多少个分桶bucket.key决定按哪个字段哈希分桶。主键表按主键做更新所以通常把主键或高频查询字段设为分桶键让同一主键的数据落到同一个bucket避免跨桶查。实际生产里分桶数要和Flink写入并行度、查询并发一起评估。3.3 数据写入后发生了什么执行INSERT语句往orders表写数据时Flink的Fluss Connector会把每条记录按bucket.key哈希确定它属于哪个bucket然后发给对应的Tablet Server。Tablet Server把数据顺序追加到bucket的log里同时按配置做多副本复制保证可用性。这一步和写Kafka体验很接近毫秒级延迟。区别在于Fluss的bucket并不会一直只存log。后台会有类似compaction的任务把已经确认的消息从log转换成列存格式推到远端对象存储并更新元数据。这样一张表里最新数据在Tablet Server本地历史数据在对象存储查询引擎能通过元数据定位到正确的文件。这个“先log后列存”的过程可以类比成快递分拣中心快递到了先卸货到暂存区log保证你随时能取走分拣人员再按区域把快件装车运走列存上传到对象存储。暂存区空间有限所以要持续清走暂存区速度最快所以最新件你总能马上拿到。3.4 查询端能看到什么对客户端来说读Fluss表有两种姿势。第一种是流式订阅Flink消费Log Table或主键表的log部分拿到实时增量这对应实时计算场景。第二种是查询点查主键值直接命中bucket索引返回最新状态分析型查询则去扫描远端列存文件。一个容易被忽略的细节是主键表的点查返回的一定是基于log和列存合并后的最终状态而不是log里的中间状态。换句话说同一主键的多条更新最终只会看到一个最新版本。这个能力来源于写入时的版本管理和合并策略也是它能够替代“Kafka 数据库”组合的根本原因。我在测试时做了一个小实验连续写入三条同一order_id、不同status的记录然后立即点查返回的是最后一条的status。再用Flink消费这张表的log却能按顺序看到三条变更记录。这就非常直观——同一份数据流读拿到的是变更历史点查拿到的是最新快照两者互不干扰。这个特性对实时数仓的“多用途存储”诉求来说几乎是量身定做的。4. 和Flink的深度集成CDC入湖、实时查询、流批一体4.1 为什么生态上先对齐FlinkFluss的生态策略非常明确优先和Flink深度绑定。这和它出身于实时计算社区有关也和Flink正在从“计算引擎”走向“流式数据平台”的路线一致。对工程师来说好处是学习成本低——不用学一套全新的SQL方言Flink SQL怎么写Fluss就怎么写Flink的流批两用能力也能直接落到Fluss上。官方和维护团队一直在迭代Flink Connector建议使用前对照自己Flink版本选对应的连接器版本这是很多初次接入的人踩得最多的地方。我遇到过一个团队Flink用的是1.18Connector却下了最新版结果SQL提交报错排查了半天才定位到是jar包版本冲突。后来统一改成Flink 1.19 对应连接的稳定版本问题立刻消失。这类问题在Flink生态里太常见了不怪组件只怪版本管理太随意。4.2 场景一MySQL CDC实时同步到Fluss最常见的落地场景是把业务库实时同步成Fluss主键表。Flink SQL里可以这样组合先定义MySQL CDC源表再定义Fluss结果表然后一条INSERT语句启动持续同步。CREATE TABLE mysql_orders ( order_id BIGINT PRIMARY KEY, user_id BIGINT, amount DOUBLE, order_status STRING, ts TIMESTAMP(3) ) WITH ( connector mysql-cdc, hostname localhost, port 3306, username root, password ******, database-name app, table-name orders ); CREATE TABLE fluss_orders ( order_id BIGINT, user_id BIGINT, amount DOUBLE, order_status STRING, ts TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector fluss, bootstrap.servers localhost:9123 ); INSERT INTO fluss_orders SELECT * FROM mysql_orders;这条链路跑起来之后业务库每次变更Fluss里的订单状态都会实时更新。下游想查最新订单直接点查Fluss想订阅变更流Flink再消费Fluss整个过程不再需要额外维护一套HBase或者Redis。我拿一个模拟订单表做过验证MySQL里执行一条UPDATE大概一两秒内Fluss里的对应记录就已经是最新值延迟体感非常好。这背后是MySQL CDC的Binlog监听 Flink的实时写入 Fluss主键表的合并存储三者配合几乎没有多余环节。4.3 场景二Fluss实时缓冲 Paimon离线分析另一个我比较看好的组合是“Fluss Paimon”。Paimon是流式数据湖适合把大量历史数据以列存形式管好Fluss负责最热的实时部分。数据先以秒级延迟写入FlussFlink消费Fluss再定期或按分区提交到Paimon这样既保留了数据湖的离线分析能力又避免了高频小事务对湖表的冲击。相比“Kafka攒数据再落湖”的经典方案Fluss的优点是数据不必跨系统复制两份实时计算和湖表之间有了一个能查、能更新的中间层整个架构的语义更顺。这套组合有点类似“热存储 温存储”分层Fluss是热层对象存储是温层Paimon是最终的分析层。数据不必再睡在Kafka里等着被批任务扫走而是先进入一个真正可查询的实时存储让需要秒级结果的应用先跑起来湖表只负责历史沉淀。这个架构对“热数据服务线上查询、冷数据跑T1分析”的场景特别合适。4.4 场景三用主键表当Flink的实时维表Flink做流式关联时维表经常是MySQL或者HBase要么查询压力大要么数据更新不及时。如果维表数据也能实时从业务库同步到Fluss主键表Flink就可以直接对Fluss做关联查询拿到的还是最新的结果。这个用法看起来不起眼但很多实时宽表、实时风控场景恰恰卡在这一步。把维表的“新鲜度”和“可查性”同时解决是Fluss最实际的价值之一。我以前做实时宽表时最头疼的就是维表更新业务库改了商品分类实时任务不能立即感知导致宽表里分类字段是旧的。用Fluss做维表配合CDC同步这个问题几乎被消除。因为维表本身就是实时更新的Flink任务每次关联拿到的都是当前最新值。这个场景我觉得比“替换Kafka”更有说服力也是我向别人推荐Fluss时第一个强调的落地点。5. 选型对照Kafka、Paimon、Fluss各自解决什么问题5.1 一张表看清边界社区里经常有人问“Fluss是不是要取代Kafka”“Fluss和Paimon有什么区别”。我的回答是它们解决的问题有重叠但边界不同硬要分高下没意义。下面这个表是我自己整理的选型视角维度KafkaPaimon/IcebergApache Fluss核心定位消息管道流式数据湖流式存储数据模型Topic字节流湖表列存Log Table / Primary Key Table写入延迟毫秒级分钟级为主秒级至毫秒级更新能力不支持支持Merge/Upsert路径偏重主键表原生支持Upsert点查能力弱弱适合扫描主键表支持点查历史数据本地盘过期即删对象存储长期保留对象存储冷热分层与Flink集成成熟但只是管道深度集成偏批式写入原生深度集成流表一体从这个表能明显看出来Fluss和Kafka最核心的差异是从“管道”变成了“存储”和Paimon最核心的差异是从“结果仓库”变成了“实时服务层”。它不是要消灭谁而是填补时间维度上的空档数据产生的那一瞬间到它成为历史数据之间这段“最热”的时间以前没有好的存储方案。5.2 什么场景优先考虑Fluss总结下来这几类场景值得认真评估Fluss实时维表和实时宽表既需要实时更新又需要被高频点查。CDC实时同步的目标端想摆脱“MySQL - Kafka - 清洗 - 数据库”的多级同步。实时湖仓的缓冲层不想让高频实时写入直接冲击湖表又想保留查询能力。对Flink SQL有强依赖的团队可以用一套SQL体系覆盖实时加工和存储。我做选型时还会额外看一个指标这个存储能不能同时被“流任务”和“查询服务”两种角色直接使用。能就说明可以省掉一套中间存储不能说明它只是又增加了一环。Fluss在我目前测试的场景里是能胜任双角色的这也是我持续跟进它的主要原因。5.3 什么场景建议继续用老方案如果你的核心诉求只是“高吞吐消息广播”Kafka依然是生态最成熟的选择没必要为一两张表引入新组件。如果你只关心离线大规模分析对秒级实时性没有强诉求直接用Paimon、Iceberg更简单直接。Fluss不是万能存储它的价值集中在“实时 更新 可查”这个交叉点上。把这点想清楚比背下一堆参数更有用。另外还有一点要诚实Fluss的周边生态还在成长期监控告警、数据迁移工具、资深专家储备都不如Kafka那么丰富。如果你的团队没有足够的能力去读源码排查问题建议初期先在低风险场景试用不要一上来就放到核心交易链路上。这不是否定Fluss而是对线上稳定性负责。6. 我从零开始跑通Fluss的小集群踩坑与建议6.1 搭建时的三个前置准备我第一次跑Fluss的时候按官方Quickstart走还算顺但有几个前置条件容易坑到新手。第一JDK版本要对得上发行版说明里写的要求一定看清楚版本不对启动时会有很隐晦的报错。第二给Tablet Server的数据目录留足空间它是先落本地盘再异步上传远端存储的本地盘满了会直接卡写入。第三配置里远端对象存储的Endpoint、Bucket、AK/SK都要提前准备好不要在集群启动之后再改存储配置改起来很麻烦。我在本地测试时用的就是系统盘的默认目录结果跑了一晚上数据量起来后直接把磁盘占满整个集群写入卡死。后来单独挂了一块盘给数据目录再把对象存储的路径配好才稳定下来。这里也建议如果只是本地测试可以先把远端存储关掉或指向本地路径先把功能跑通再接入对象存储排查问题会更简单。6.2 我踩过的三个坑坑一Flink版本和Connector版本不匹配。这是最常见的问题症状是提交SQL时报ClassNotFound或者NoSuchMethodError其实不是代码问题就是jar版本不对。建议固定一个Flink版本不要去试最新版connector除非你的Flink也是最新。坑二分桶数设置和实际数据量不匹配。分桶太少写入并行度上不去出现热点分桶太多小文件、元数据开销大查询反而变慢。我个人的习惯是先按写入并行度的1至2倍设置观察一段时间再做调整。坑三把Fluss当Kafka用忽略了主键表的更新语义。Log Table确实可以当Kafka用但如果你创建的是主键表写入相同主键的多条记录会被合并而不是像Kafka一样全保留。想保留全量流水要用Log Table想维护最新状态才用Primary Key Table。这个语义区分清楚了能省很多排查时间。6.3 生产化之前的几个建议如果准备小规模生产试用我建议先把监控做起来重点看三类指标bucket的写入延迟、远端文件转换任务积压、对象存储的上传耗时。这三个指标能提前暴露存储瓶颈。其次Coordinator节点不需要堆很高规格但稳定性要求高不要和业务混部。最后容灾演练别省——把某个Tablet Server进程直接kill观察客户端是否自动切换、数据是否可恢复这类实验越早做越好。我个人实际跑下来最大的感受是别把Fluss当作Kafka的替代品来评估把它理解成“Flink生态里的可更新、可查询、可订阅状态存储”很多设计就都说得通了。它的定位和周边生态还在快速演进生产环境大规模替换成熟组件前建议先在一条业务链路上做小范围验证把监控、备份、恢复方案都补齐再逐步铺开。这套思路对任何一个还处于成长期的存储组件都适用。