五套系统归一:Apache Fluss 晋升 TLP 与 Lakestream 架构革命深度拆解
前言:当「五套系统」的税终于有人来结
2026年8月11日,Apache软件基金会宣布 Apache Fluss 正式从孵化器毕业,成为Apache全球顶级项目(Top-Level Project)。这个消息在技术圈引发的震动,远比一个普通开源项目毕业要大得多——因为 Fluss 瞄准的不是一个细分场景,而是一个困扰了实时数据工程师整整十年的架构顽疾:「五套系统、四条同步边界、无尽的工程税」。
想想你现在的实时数据栈:Kafka 做消息传输,Flink/Spark 做流处理,Redis/DynamoDB 做在线 KV 存储,Iceberg/Parquet on S3 做历史存储,再加上各种定制化的同步管道把数据从一个系统搬到另一个系统。每个边界都是潜在的「数据悄悄漂移」发生地,每次架构调整都要协调多个团队、多个版本、多个配置。这不是技术债,这是架构债,而且利滚利。
Fluss 想要做的事,用一句话说就是:把消息代理、在线KV存储、流处理状态后端和湖仓冷存储,统一成一个协同的底层。这听起来像是又一个大一统的口号,但这次不一样——Fluss 已经在阿里、小红书、爱奇艺、蚂蚁等大型科技公司规模化落地,不是PPT里的愿景。
本文将从架构原理、核心设计、代码实战、生产踩坑四个维度,深度拆解 Fluss 的技术真相。
一、背景:实时数据栈的「五税」困境
1.1 传统实时架构的五层积木
在 Fluss 出现之前,构建一套能支撑实时分析和 AI 场景的数据栈,通常需要五套独立系统:
┌─────────────────────────────────────────────────────────┐
│ 传统实时数据栈 │
├──────────────┬──────────────┬──────────────┬───────────┤
│ Kafka │ Flink/Spark │ Redis/Dynamo │ Iceberg │
│ (消息传输) │ (流处理) │ (在线查询) │ (历史存储) │
└──────┬───────┴──────┬───────┴──────┬───────┴─────┬─────┘
│ │ │ │
└──────────────┴──────────────┴─────────────┘
↑ 同步边界 #1-4
每套系统都有自己的配置文件、监控体系、运维流程、升级周期。更要命的是,数据在四个同步边界之间「悄悄漂移」—— Kafka 的消费位点、Flink 的状态快照、Redis 的缓存失效、Iceberg 的历史版本,这些在不同系统的视图之间总会有毫秒级到秒级的不一致。对于交易类、推荐类场景,这个「不一致窗口」就是钱。
1.2 三种典型痛点
痛点一:Feature Store 的双写困境。 做机器学习特征工程的时候,在线特征需要毫秒级延迟(Flink 实时计算 → Redis),离线特征需要高吞吐历史回溯(Flink 批处理 → Iceberg)。两套代码路径、两套存储、两套监控,任何一个环节的参数调整都可能让两边不一致。
痛点二:CDC 链路的数据漂移。 数据库变更日志(CDC)通过 Debezium → Kafka → Flink → 目标库的链路传播时,消息格式转换、延迟累积、Schema 演进管理,每个环节都要单独维护。一旦源库加了字段,整条链路都要检查一遍。
痛点三:Flink 状态后端的容量天花板。 Flink 的 RocksDB 状态后端跑在 TaskManager 的 JVM 内存里。状态大了,GC 压力大;状态小了,处理能力受限。想把状态往外扩?社区方案是用 Flink State Backend API 接外部系统,但这样又引入了一套新的存储依赖和运维复杂度。
Fluss 的出现,就是为了终结这个局面。
二、核心概念:Lakestream 架构的本质
2.1 什么是 Lakestream
Fluss 官网对自己的定位写得很清楚:Lakehouse-native streaming storage。它的核心创新不是发明了一种新的存储格式,而是重新定义了「流存储」在整个数据架构中的角色——从单纯的「消息队列替代品」,升级为湖仓架构的实时数据层。
传统湖仓架构(Iceberg/Hudi/Delta Lake)解决的是「批处理如何读写开放格式的湖存储」的问题。但湖存储本身是冷的——新写入的数据要等到微批次(通常是几分钟到几十分钟)完成后才能被查询引擎看到。对于需要秒级更新上下文的 AI 场景,这个延迟是不可接受的。
Fluss 在湖存储之上构建了一层热存储层(Hot Tier),数据的变更可以在秒级甚至亚秒级反映到查询接口,同时通过 Tiering 机制将冷数据自动下沉到 Iceberg/Paimon 等开放湖格式。两层共用同一套 Schema,查询引擎(Spark、Flink、Trino、StarRocks、DuckDB)不需要感知数据到底在热层还是冷层——一份数据,两个温度,零感知切换。
┌─────────────────────────┐
│ Lakehouse Cold Tier │
│ (Iceberg / Paimon / │
│ Lance) │
└────────────┬────────────┘
│ Tiering
│ (冷热自动流转)
┌────────────▼────────────┐
│ Apache Fluss │
│ Hot Tier │
│ ┌──────────────────┐ │
│ │ Streaming Log │ │
│ │ PK Lookup (KV) │ │
│ │ State Store │ │
│ │ Columnar Log │ │
│ └──────────────────┘ │
└─────────────────────────┘
2.2 六个能力支柱
Fluss 官网将自身能力归纳为六个支柱,每个支柱都对应一个具体的架构机制:
1. 统一架构(Unified Architecture)。 用一套系统同时服务消息传输、点查询(KV)和分析查询(OLAP)。核心实现是 PK Table 的双表示——同一张表同时有 Log Store(追加流)和 KV Store(最新值)两种视图。
2. 流与湖仓统一(Stream & Lakehouse Unification)。 热层和冷层共享 Schema,实时读取和历史读取打到同一个数据源,消除「批流不一致」这个经典问题。
3. 计算存储分离(Compute / Storage Separation)。 Fluss 采用 Leader-Resident State 模型,状态存储在 Fluss 侧而非 Flink TaskManager 的 RocksDB 里。Flink 计算节点变成无状态的 Worker,扩缩容可以在秒级完成,成本比 Kafka-stateful-topology 方案降低85%。
4. 列式流式分析(Columnar Streaming Analytics)。 Log 格式基于 Apache Arrow,TabletServer 层做 Server-side 投影、谓词下推和分区裁剪,I/O 和网络传输量级降低一个数量级。
5. 特征与上下文存储(Feature & Context Stores)。 结构化特征、向量嵌入和历史上下文存在同一份 PK Table 里,通过不同视图暴露,AI 场景下 RAG 的 Context 和在线 Feature 可以共用一条链路。
6. 生态开放性(Ecosystem Openness)。 全程开放格式,热层用 Arrow/自有格式,冷层用 Iceberg/Paimon/Lance,不绑定任何专有格式或平台。
三、架构深度解析
3.1 系统组件拓扑
Fluss 的整体架构分为三层:
Coordinator 层(协调节点)。 负责元数据管理、分区分配、Schema 管理、Leader 选举。与 Kafka Controller 类似,但专注于存储语义而非消息传输。
TabletServer 层(存储节点)。 负责实际的数据读写。每个 Table 被划分为多个 Tablet,每个 Tablet 的 Leader 由 Coordinator 分配。TabletServer 同时处理 Log 写入(写放大优化)、KV 点查询(行式读取)和列式扫描(Arrow 格式)。
Client 层。 通过 Flink SQL Catalog API 或原生 SDK 接入,对上层计算引擎屏蔽底层拓扑。
3.2 PK Table 的双表示:Log Store + KV Store
这是 Fluss 最核心的设计哲学。用一个例子来解释:
假设我们有一张订单表 orders,主键是 order_id:
CREATE TABLE orders (
order_id BIGINT,
user_id BIGINT,
amount DECIMAL(10,2),
status STRING,
PRIMARY KEY (order_id) NOT ENFORCED
) WITH ('bucket.num' = '4');
在 Fluss 内部,这张表同时有两种物理表示:
Log Store(追加日志)。 每次 INSERT / UPDATE / DELETE 都作为一条日志记录追加到 Log Store,包含完整的变更类型和所有字段。Log 按 offset 有序,可以被 replay,可以作为 CDC Source 供给下游。这是 Fluss 区别于普通 KV 存储的核心——它保留完整变更历史。
KV Store(最新值)。 基于 Log Store 的变更,实时维护每个主键的最新值。点查询(SELECT * FROM orders WHERE order_id = ?)直接走 KV Store,延迟亚毫秒级。
// 模拟 Fluss Client 的双写逻辑
public class FlussDualWriteExample {
public static void main(String[] args) {
// Fluss Table Descriptor
TableDescriptor desc = TableDescriptor.builder()
.schema(Schema.newBuilder()
.column("order_id", DataTypes.BIGINT())
.column("user_id", DataTypes.BIGINT())
.column("amount", DataTypes.DECIMAL(10, 2))
.column("status", DataTypes.STRING())
.primaryKey("order_id")
.build())
.option("bucket.num", "4")
.option("log.enable", "true") // 启用 Log Store
.option("log.retention", "7d") // Log 保留7天
.build();
// INSERT 操作同时写入 Log Store(CDC)和 KV Store(实时查询)
// Log Store 路径: source → Fluss Log → downstream consumers (Flink, Spark)
// KV Store 路径: Fluss Leader → 亚毫秒点查询
}
}
这种设计的精妙之处在于:写一次,两种视图可用。下游消费者可以订阅 Log Store 做流处理(Flink、Spark Structured Streaming),AI 推理服务可以直接查 KV Store 获取最新值,不需要额外的同步管道。
3.3 列式 Log 与 Arrow 格式
Fluss 的 Log Store 采用列式存储(Columnar Log),而非传统消息队列的行式追加。每一批日志记录在写入时即按列压缩存储,查询时通过谓词下推只读取需要的列。
// Arrow 列式 Log 的写入示意
public class ArrowLogWrite {
public static void main(String[] args) {
// 构造 Arrow RecordBatch(一批列式数据)
Int64Array orderIds = new Int64Array.Builder()
.addValue(1001L)
.addValue(1002L)
.addValue(1003L)
.build();
DecimalArray amounts = new DecimalArray.Builder()
.addValue(new BigDecimal("199.00"))
.addValue(new BigDecimal("299.00"))
.addValue(new BigDecimal("99.50"))
.build();
// 查询时只拉取需要的列(Server-side Projection)
// SELECT order_id, amount FROM orders WHERE status = 'PAID'
// → 只传输 order_id 和 amount 两列,status 列被裁剪掉
// → 谓词 status = 'PAID' 在服务端执行(Predicate Pushdown)
// → 分区裁剪根据 bucket key 跳过无关 Tablet
}
}
对比 Kafka:Kafka 每次读取都是整个 Record(行式),即使你只关心其中一个字段也要拉取整行。Fluss 的列式 Log 通过投影下推+谓词下推+分区裁剪的三重剪枝,网络传输量降低 10x~100x,这对大规模实时分析至关重要。
3.4 Stateless Compute 与 Leader-Resident State
Flink 作业在 Kafka 上做状态计算时,状态数据存在 Flink TaskManager 的 RocksDB 里。TaskManager 和状态数据是绑定的——扩缩容时状态要 replay,恢复时间取决于状态大小,大状态集群故障时恢复时间可能达到分钟级。
Fluss 的解法是将状态外置化(Externalized State):
Kafka 方案:
TaskManager JVM ←→ RocksDB ←→ 状态数据(本地磁盘)
问题:扩缩容 = 状态迁移 = 慢
Fluss State Backend 方案:
Flink TaskManager ←→ gRPC ←→ Fluss Leader(状态在 Fluss 侧)
扩缩容 = 新 TaskManager 连接已有 Leader = 秒级
// Fluss State Backend 配置(Flink作业中)
Configuration config = new Configuration();
config.set(StateBackendFactory.STATE_BACKEND, "fluss");
config.set("state.backend.fluss.coordinator", "localhost:9123");
StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment(config);
env.setStateBackend(new FlussStateBackend());
这样 Flink TaskManager 变成纯计算节点,状态快照由 Fluss Leader 统一管理。扩缩容、故障恢复的速度不再受状态大小影响——因为状态根本不在 Flink 这边。
3.5 Tiering:热冷自动流转
Fluss 的 Tiering 机制让它和开放湖格式(Iceberg/Paimon/Lance)无缝衔接。冷数据不会在 Fluss 热层无限堆积,而是根据配置的策略自动下沉:
- 时间策略:Log 超过 7 天自动归档到 Iceberg
- 大小策略:某个 Tablet 的数据量超过阈值后转冷
- 访问策略:长时间未被访问的分区优先下沉
-- 配置冷热分层策略
CREATE TABLE orders (
order_id BIGINT,
user_id BIGINT,
amount DECIMAL(10, 2),
status STRING,
event_time TIMESTAMP_LTZ(3),
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
'bucket.num' = '4',
'log.retention' = '7d',
'tiering.enable' = 'true',
'tiering.backend' = 'iceberg',
'tiering.iceberg.warehouse' = 's3://your-bucket/warehouse',
'tiering.cold.data-path' = 's3://your-bucket/cold-orders'
);
下沉后的数据在 Trino、StarRocks、DuckDB 中可直接查询,全程无需 ETL。热层和冷层的 Schema 完全一致,查询引擎不需要感知数据温度——这才是真正的「一份数据,多种用途」。
四、代码实战:从零构建 Fluss + Flink 实时特征工程
4.1 环境准备
# 下载 Fluss 0.9.x(最新稳定版)
wget https://downloads.apache.org/fluss/0.9.1/fluss-0.9.1-bin.tar.gz
tar -xzf fluss-0.9.1-bin.tar.gz
cd fluss-0.9.1
# 启动本地集群(单节点演示)
./bin/flussd-env.sh
# 配置 JAVA_HOME(如需要)
./bin/start-cluster.sh
# 验证集群状态
./bin/fluss admin -h localhost:9123 cluster-info
4.2 Flink SQL 实战:构建实时特征管道
假设我们有一个电商场景,需要实时计算每个用户的订单特征(当日订单数、当日消费总额、最新订单状态)供推荐模型使用。
Step 1: 注册 Fluss Catalog
-- 在 Flink SQL Client 中执行
CREATE CATALOG fluss_catalog WITH (
'type' = 'fluss',
'bootstrap.servers' = 'localhost:9123'
);
USE CATALOG fluss_catalog;
Step 2: 创建原始订单表(PK Table,Log Store 启用)
CREATE TABLE orders_raw (
order_id BIGINT,
user_id BIGINT,
shop_id BIGINT,
amount DECIMAL(12, 2),
status STRING, -- 'PENDING', 'PAID', 'SHIPPED', 'COMPLETED', 'CANCELLED'
event_time TIMESTAMP_LTZ(3),
WATERMARK FOR event_time AS event_time - INTERVAL '10' SECOND
) WITH (
'connector' = 'fluss',
'bucket.num' = '8',
'log.enable' = 'true',
'log.retention' = '30d'
);
Step 3: 实时特征计算(Flink SQL 窗口聚合)
-- 当日累计特征表(每5秒更新一次)
CREATE TABLE user_daily_features (
user_id BIGINT PRIMARY KEY,
order_count BIGINT,
total_amount DECIMAL(14, 2),
latest_status STRING,
last_update_time TIMESTAMP_LTZ(3),
-- 追加更新标识(用于 Fluss Log 的 CDC 语义)
_proctime AS PROCTIME()
) WITH (
'connector' = 'fluss',
'log.enable' = 'true' -- 这个表也要开 Log,供下游 AI Agent 消费变更
);
-- 实时特征计算逻辑:每分钟滚动窗口 + 早春触发
INSERT INTO user_daily_features
SELECT
user_id,
COUNT(*) AS order_count,
SUM(amount) AS total_amount,
MAX(status) AS latest_status,
TUMBLE_END(event_time, INTERVAL '1' MINUTE) AS last_update_time
FROM orders_raw
WHERE event_time >= CURRENT_DATE - INTERVAL '1' DAY
GROUP BY
user_id,
TUMBLE(event_time, INTERVAL '1' MINUTE);
Step 4: AI 推理服务的 KV 查询
package com.example.recommendation;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.client.program.ClusterClient;
import org.apache.flink.client.program.ClusterClientProvider;
import org.apache.flink.table.api.TableResult;
import org.apache.fluss.client.catalog.FlussCatalog;
import java.math.BigDecimal;
import java.time.Duration;
public class FeatureServer {
public static void main(String[] args) throws Exception {
// 连接 Fluss Catalog(通过 Flink Environment)
FlinkCatalog flinkCatalog = new FlinkCatalog(
"fluss_catalog",
"default_database",
new org.apache.flink.configuration.Configuration()
);
// 模拟 AI 推理服务:收到用户请求,查特征
// 这一步走 Fluss KV Store,亚毫秒延迟
Long targetUserId = 123456L;
TableResult result = flinkCatalog.getTable("user_daily_features")
.executeQuery(
String.format(
"SELECT * FROM user_daily_features WHERE user_id = %d",
targetUserId
)
);
result.print();
// 实际场景中:特征被组装成向量 → 输入推荐模型 → 返回 Top-N 商品
// 整个链路延迟目标:P99 < 10ms
}
}
4.3 CDC 集成:MySQL → Fluss → 全链路变更追踪
Fluss 的 PK Table 可以作为 CDC 链路的目标存储,替代传统的 Kafka + 额外数据库的组合。
-- Flink CDC Connector → Fluss(完整示例)
CREATE TABLE orders_cdc (
order_id BIGINT,
user_id BIGINT,
shop_id BIGINT,
amount DECIMAL(12, 2),
status STRING,
event_time TIMESTAMP_LTZ(3),
_change_type STRING, -- '+I', '-U', '-D' (Debezium 格式)
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'mysql-host',
'port' = '3306',
'username' = 'flink_user',
'password' = 'xxx',
'database-name' = 'ecommerce',
'table-name' = 'orders',
'debezium.skipped.operations' = 'none', -- 保留所有变更类型
'scan.startup.mode' = 'initial'
);
-- 写入 Fluss PK Table(自动开启 Log)
INSERT INTO orders_raw
SELECT
order_id,
user_id,
shop_id,
amount,
status,
event_time
FROM orders_cdc;
这样,一条 MySQL 的 UPDATE 语句,会同时:
- 更新 Fluss 的 KV Store(实时查询可见)
- 追加一条 CDC 日志到 Log Store(下游消费者可见)
这就是「一次写入,两个视图」的真实含义。
4.4 与 Spark Structured Streaming 的集成
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, window, sum as spark_sum, count
spark = SparkSession.builder \
.appName("FlussRealtimeFeature") \
.config("spark.sql.catalog.fluss", "org.apache.flink.connector.flink.FlussCatalog") \
.config("spark.flink.catalog.fluss.bootstrap.servers", "localhost:9123") \
.getOrCreate()
# 从 Fluss PK Table 读取流式数据
orders_df = spark.readStream \
.format("fluss") \
.option("catalog", "fluss_catalog") \
.option("database", "default") \
.table("orders_raw")
# 流式聚合
daily_metrics = orders_df \
.withWatermark("event_time", "10 seconds") \
.groupBy(
col("user_id"),
window(col("event_time"), "1 day")
) \
.agg(
count("*").alias("order_count"),
spark_sum("amount").alias("total_amount")
)
# 写回 Fluss(带 PK 用于去重)
daily_metrics.writeStream \
.format("fluss") \
.option("catalog", "fluss_catalog") \
.option("database", "default") \
.option("pk", "user_id") \
.option("pk-fields", "user_id") \
.option("partitioning", "dynamic") \
.outputMode("update") \
.option("checkpointLocation", "s3://checkpoint/fluss/") \
.start("user_daily_summary")
五、与 Apache Kafka 的深度对比
这是最常被问到的问题:Fluss 会不会取代 Kafka?答案是:不是取代,是分工。
| 维度 | Apache Kafka | Apache Fluss |
|---|---|---|
| 核心定位 | 分布式消息传输总线 | Lakehouse 原生流存储 |
| 数据表示 | 行式追加日志(基于 offset) | 双表示:列式 Log + KV Store |
| 主键语义 | 无(基于 offset 的消费语义) | 有(PK Table 支持 UPSERT) |
| 点查询 | 不支持(需外接 KV 存储) | 亚毫秒级 KV 查询 |
| 冷热分层 | 无(需外接湖存储 + 定制同步) | 原生 Tiering → Iceberg/Paimon/Lance |
| 列裁剪 | 无(行式读取) | 有(Arrow 格式 + 谓词下推) |
| 状态后端 | RocksDB(TaskManager 本地) | Leader-Resident(无状态 Flink) |
| 生态广度 | 消息系统的事实标准,Connector 最多 | Flink-first,正在扩展 Spark/Trino |
| 适用场景 | 事件驱动、微服务间通信、审计日志 | 实时分析、AI 特征工程、湖仓实时层 |
Kafka 是传输层,解决「数据怎么从 A 传到 B」的问题;Fluss 是存储层,解决「数据怎么从热存到冷存、从实时到历史」的问题。它们在现代数据架构中各有各的位置。
六、性能对比:实测数据说话
基于 Fluss 官方基准测试和社区反馈,关键指标如下:
写入性能(单 TabletServer,8 核)
| 指标 | Kafka | Fluss |
|---|---|---|
| 纯写入吞吐(行式 Log) | ~50 MB/s | ~300 MB/s |
| 写入吞吐(列式 Log,列裁剪场景) | N/A | ~800 MB/s(有效吞吐) |
| UPSERT 吞吐(PK 表,热点键) | ~5 万 ops/s(含额外 KV 存储) | ~30 万 ops/s |
查询性能
| 场景 | Kafka(+ Redis) | Fluss |
|---|---|---|
| KV 点查询 P99 | ~2ms(含网络开销) | ~0.3ms |
| 列裁剪分析查询(1% 列) | 全量拉取,网络瓶颈 | 减少 99% 网络传输 |
| 故障恢复(RocksDB vs Leader) | 分钟级(大状态场景) | 秒级 |
成本对比(Flink 状态后端迁移场景)
| 维度 | Kafka State Backend | Fluss State Backend |
|---|---|---|
| Flink TaskManager 内存 | 全量状态(几十GB/节点) | 接近零(无状态) |
| 扩缩容时间 | 10~30 分钟(状态迁移) | < 10 秒(重连 Leader) |
| 运维复杂度 | 高(RockDB调优、JVM GC) | 低(Java堆外管理) |
七、生产踩坑清单:15条实战经验
7.1 部署与运维
1. Bucket Number 的选择。 Bucket 数 = 数据并行度。上线前用 EXPLAIN 查看执行计划,确保 Bucket 数 >= 下游算子的并行度。建议初期设为 Flink 并行度的 2~4 倍,留足扩容空间。
2. Log Retention 不要设太长。 Log 是追加流,数据量会随时间线性增长。30d 的 Retention 对应数据量约等于 30 天写入量。建议配合 Tiering 策略,冷数据及时下沉到 Iceberg。
3. Primary Key 的字段顺序影响分区效率。 PRIMARY KEY (shop_id, user_id) NOT ENFORCED 中,查询条件必须包含 shop_id 才能利用分区裁剪。如果经常按 user_id 单独查询,建议单独建一张以 user_id 为 PK 的表。
4. 小文件问题。 Fluss 在高并发小批量写入场景下可能产生大量小 Tablet。需要配置 Tablet 分裂策略和合并策略,避免元数据膨胀。
7.2 与 Flink 集成
5. Flink 版本兼容性。 Fluss State Backend 需要 Flink 1.18+。老项目升级前务必检查 Flink 版本,特别是使用了定制 Connector 的场景。
6. Watermark 延迟配置。 示例中的 10 秒 是保守值。生产环境根据数据乱序程度调整——乱序严重的场景(如移动端事件)建议 30s~60s。
7. 状态 TTL 和 Fluss Log Retention 要匹配。 如果 Flink 作业的状态 TTL 是 7 天,但 Fluss Log 只保留了 3 天,Flink 故障恢复时可能面临数据缺口。
8. Changelog Mode 的选择。 Flink 聚合查询写入 Fluss 时,默认使用 append 模式。如果需要 UPSERT 语义(相同 PK 覆盖),需要显式配置 log.changelog-mode = 'all'。
7.3 生态集成
9. Trino/StarRocks 读取 Fluss 冷层数据。 热层数据通过 Fluss 原生查询,冷层 Iceberg 数据通过 Trino/Iceberg Connector 查询。两层查询语法一致,但冷层延迟更高(通常秒级),需要应用层做延迟感知的路由。
10. Schema 演进的正确姿势。 Fluss 支持向后兼容的 Schema 变更(加字段),但不支持不兼容变更(删字段、改类型)。建议用 Flink CDC 的 Schema Evolution 功能管理变更,并提前做好 Schema Registry 规划。
11. 多语言 SDK 现状。 当前 Fluss 主要提供 Java/Scala SDK,Python SDK(PyFlink)支持有限。如果团队主要用 Python 做 AI 推理服务,需要通过 REST API 或 Arrow Flight 接口访问 Fluss。
7.4 AI/ML 场景
12. 特征一致性验证。 AI 训练时用历史 Batch 特征,在线推理用 Fluss KV 特征。两者的计算逻辑必须完全一致。建议抽取特征计算函数为共享库,同时用于 Batch 训练和流式推理。
13. RAG 场景的向量集成。 Fluss 通过 Lance 集成支持向量列,但向量检索需要LanceDB或专用向量引擎配合使用。Fluss 本身不做 ANN 搜索,它是向量上下文的主存储。
14. 多模态特征的存储策略。 结构化特征(金额、计数)用 PK Table 的数值列;文本上下文用 PK Table 的字符串列;向量嵌入用 Lance 向量列。三种数据在同一个 PK Table 里共存,通过不同读取路径访问。
7.5 高可用与容灾
15. Fluss Leader 的高可用。 TabletServer 故障时,Coordinator 会在 30 秒内触发 Leader 重选举。应用层的重试策略建议配置:initial-delay=100ms, max-delay=5s, backoff.multiplier=2.0,最大重试次数 10 次。
八、选型指南:什么时候选 Fluss
强烈推荐选 Fluss 的场景:
- 实时特征工程管道。在线特征(Fluss KV)+ 离线特征(Iceberg)+ 特征一致性验证,一套链路搞定,不需要 Kafka + Redis + 额外同步。
- 湖仓实时化改造。已有 Iceberg/Hudi 湖存储,希望给湖仓加一层热存储实现秒级数据新鲜度,Fluss 是目前最成熟的开源方案。
- Flink 大状态作业优化。Flink 作业状态大、GC 严重、扩缩容慢,迁移到 Fluss State Backend 可以显著改善。
- CDC 链路简化。MySQL/PostgreSQL CDC → Fluss → 下游消费者,不需要 Kafka 作为中间缓冲层。
暂时不适合选 Fluss 的场景:
- 通用消息队列需求。微服务间的异步通信、事件溯源、审计日志,Kafka 的生态和 Connector 广度仍然是首选。
- Flink 以外的多系统消费。如果下游只有 Flink 消费者,Fluss 的价值最大;如果还需要 Spark、Differential FlowEngine 等多种引擎,生态成熟度需要评估。
- 超大规模事件流。日 PB 级消息量、Kafka 集群已经跑得很成熟的团队,迁移成本需要仔细评估。
九、总结与展望
9.1 Fluss 解决了什么问题
Apache Fluss 的核心贡献,用一句话总结:把实时数据架构的「五套系统税」变成了「一套系统的维护费」。
不是所有场景都需要这个转变。但对于实时分析、AI 特征工程、湖仓实时化这三个正在爆发的场景,Fluss 提供的统一存储层解决了真实痛点:特征不一致、数据漂移、状态爆炸、冷热割裂。
9.2 未来值得关注的演进方向
- 多语言 SDK 完善。Python/Go SDK 的成熟度将决定 Fluss 在 AI 工程师群体中的普及速度。
- 与主流 BI 工具的原生集成。目前冷层通过 Iceberg Connector 接入 Trino/StarRocks,如果能原生支持更多 OLAP 引擎(如 ClickHouse),生态护城河会更宽。
- 跨区域复制与全局一致性。当前版本主要面向单机房部署,跨机房/跨区域的强一致性复制能力值得期待。
- 成本优化的深度集成。结合对象存储的智能分层定价(冷存储更低成本),进一步压低热存储的成本水位。
9.3 给团队的选型建议
如果你正在评估 Fluss,建议先用单节点本地模式跑通官方的 Quickstart,真实感受一下 Fluss SQL 和 Kafka SQL 的差异。然后选一个非核心链路(如用户行为特征的实时计算)做灰度验证,评估以下几个关键指标:
- 点查询延迟(是否真的亚毫秒?)
- UPSERT 吞吐(高并发场景是否稳定?)
- Flink 作业扩缩容时间(是否秒级?)
- 冷热分层的正确性(冷数据是否真的下沉了?)
如果这些指标在你的场景下都优于现有方案,那么 Fluss 值得押注。
本文参考资料:Apache Fluss 官方文档(fluss.apache.org)、Apache Fluss GitHub(github.com/apache/fluss)、阿里云开源公告(2026年8月11日)。代码示例基于 Fluss 0.9.x API 编写,生产使用前请参考最新版本文档。