编程 五套系统归一:Apache Fluss 晋升 TLP 与 Lakestream 架构革命深度拆解

2026-08-12 15:46:05 +0800 CST views 8

五套系统归一: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 完全一致,查询引擎不需要感知数据温度——这才是真正的「一份数据,多种用途」。

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

假设我们有一个电商场景,需要实时计算每个用户的订单特征(当日订单数、当日消费总额、最新订单状态)供推荐模型使用。

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 语句,会同时:

  1. 更新 Fluss 的 KV Store(实时查询可见)
  2. 追加一条 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 KafkaApache 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 核)

指标KafkaFluss
纯写入吞吐(行式 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 BackendFluss 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 分裂策略和合并策略,避免元数据膨胀。

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 的场景:

  1. 实时特征工程管道。在线特征(Fluss KV)+ 离线特征(Iceberg)+ 特征一致性验证,一套链路搞定,不需要 Kafka + Redis + 额外同步。
  2. 湖仓实时化改造。已有 Iceberg/Hudi 湖存储,希望给湖仓加一层热存储实现秒级数据新鲜度,Fluss 是目前最成熟的开源方案。
  3. Flink 大状态作业优化。Flink 作业状态大、GC 严重、扩缩容慢,迁移到 Fluss State Backend 可以显著改善。
  4. CDC 链路简化。MySQL/PostgreSQL CDC → Fluss → 下游消费者,不需要 Kafka 作为中间缓冲层。

暂时不适合选 Fluss 的场景:

  1. 通用消息队列需求。微服务间的异步通信、事件溯源、审计日志,Kafka 的生态和 Connector 广度仍然是首选。
  2. Flink 以外的多系统消费。如果下游只有 Flink 消费者,Fluss 的价值最大;如果还需要 Spark、Differential FlowEngine 等多种引擎,生态成熟度需要评估。
  3. 超大规模事件流。日 PB 级消息量、Kafka 集群已经跑得很成熟的团队,迁移成本需要仔细评估。

九、总结与展望

9.1 Fluss 解决了什么问题

Apache Fluss 的核心贡献,用一句话总结:把实时数据架构的「五套系统税」变成了「一套系统的维护费」

不是所有场景都需要这个转变。但对于实时分析、AI 特征工程、湖仓实时化这三个正在爆发的场景,Fluss 提供的统一存储层解决了真实痛点:特征不一致、数据漂移、状态爆炸、冷热割裂。

9.2 未来值得关注的演进方向

  1. 多语言 SDK 完善。Python/Go SDK 的成熟度将决定 Fluss 在 AI 工程师群体中的普及速度。
  2. 与主流 BI 工具的原生集成。目前冷层通过 Iceberg Connector 接入 Trino/StarRocks,如果能原生支持更多 OLAP 引擎(如 ClickHouse),生态护城河会更宽。
  3. 跨区域复制与全局一致性。当前版本主要面向单机房部署,跨机房/跨区域的强一致性复制能力值得期待。
  4. 成本优化的深度集成。结合对象存储的智能分层定价(冷存储更低成本),进一步压低热存储的成本水位。

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 编写,生产使用前请参考最新版本文档。

推荐文章

JavaScript中的常用浏览器API
2024-11-18 23:23:16 +0800 CST
禁止调试前端页面代码
2024-11-19 02:17:33 +0800 CST
File 和 Blob 的区别
2024-11-18 23:11:46 +0800 CST
程序员茄子在线接单