编程 Apache Fluss 深度拆解:当流存储把 Kafka、Redis、Flink State 与 Iceberg 揉进一个底座——湖仓原生实时架构全链路实战

2026-08-18 06:12:49 +0800 CST views 33

Apache Fluss 深度拆解:当流存储把 Kafka、Redis、Flink State 与 Iceberg 揉进一个底座——湖仓原生实时架构全链路实战

选题来源:GitHub Trending / Apache 顶级项目动态(2026-08-11 Apache Fluss 从孵化器毕业成为顶级项目)
适用读者:做过实时数仓、被"Kafka + Flink + Redis + Iceberg"四件套折磨过的后端 / 数据工程师


一、背景:实时数仓的"多系统税"

如果你在一家稍微有点规模的公司做过实时业务,大概率见过下面这套"标准答案":

  • 消息队列用 Kafka 做事件传输;
  • 流处理用 Flink(或 Spark Structured Streaming)做清洗、聚合、打宽;
  • 在线 KV用 Redis(或 DynamoDB、HBase)承接亚毫秒级点查,给 API 和特征服务喂数据;
  • 离线湖用 Iceberg / Hudi / Delta 落在 S3/OSS 上,喂训练、喂 BI、喂历史回放;
  • 同步层再叠一堆自研管道 + 新鲜度监控,把上面四个系统的数据"对齐"。

这套架构本身没毛病,业界跑了快十年。但它有一个被刻意忽视的隐性成本——每一条系统边界,都是数据分叉点

Kafka 里的订单流、Redis 里的订单快照、Iceberg 里的订单历史,三份数据的"此刻真值"永远差那么几秒到几分钟。你以为你在做实时,其实你在做"勉强不崩的准实时"。一旦某个同步管道悄悄挂了,Redis 里的库存和 Iceberg 里的库存就会出现肉眼不可见的偏差,直到一次对账失败或一次大促超卖才炸出来。

Apache 基金会在 2026-08-11 宣布:由阿里云捐赠的流式存储项目 Fluss 正式从孵化器毕业,成为 Apache 全球顶级项目(TLP),提案在委员会投票中全票通过。官方给 Fluss 的定位一句话概括——"Streaming Storage for Real-Time Analytics & AI",它要做的不是又一个消息队列,而是把上面那套"消息代理 + 在线 KV + 流处理状态后端 + 湖仓冷存"坍缩成同一个底座,让 Lakehouse 真正变成实时的。

这篇文章不吹概念,我们从架构、读写路径、存算分离、列式流分析、湖仓统一,一直拆到能跑的代码和性能调优清单。


二、核心概念:Fluss 到底是什么

2.1 一句话定义与六个能力支柱

Fluss 官方的自我定位是 lakehouse-native streaming storage(湖仓原生的流存储)。它和 Kafka 最大的区别在于词:Kafka 是 streaming transport(传输),Fluss 是 streaming storage(存储)。传输只管把字节从 A 搬到 B 且可重放;存储要回答"我现在这份数据的真值是什么、点查多快、分析怎么扫、冷了往哪沉"。

Fluss 把被拆散的能力重新收拢成六个支柱:

  1. Unified Architecture(统一架构):一套系统同时承担消息传输、点查、分析查询。架构基础是 PK Table 的"双表示"——同一份数据同时有 Log StoreKV Store 两种形态。
  2. Stream & Lakehouse Unification(流湖统一):实时层和批式层共享同一份 schema、同一份数据,靠 Tiering Service 下沉冷数据 + Union Read 跨冷热统一查询。
  3. Compute / Storage Separation(存算分离):计算无状态、秒级恢复,状态常驻在 Fluss leader 而不是 Flink slot 里,官方宣称相比 Kafka 拓扑最高可便宜 85%
  4. Columnar Streaming Analytics(列式流分析):Log 用 ARROW 列式格式,TabletServer 端做投影、谓词下推、分区裁剪,IO 与网络成本呈数量级下降。
  5. Feature & Context Stores(特征与上下文存储):行、列、向量三种数据形态同存一个 substrate,在线特征、RAG 上下文、分析查询都是同一张 PK Table 的不同视图。
  6. Ecosystem Openness(生态开放):Flink、Spark、Trino、StarRocks、Doris、DuckDB 都能读;冷层是 Iceberg / Paimon / Lance 等开放格式,无厂商锁定。

2.2 PK Table 的"双表示"——这是整篇的题眼

传统消息队列里,一条消息进了 Kafka 就是一个 append-only 的 offset 记录,你没法"按主键秒级拿到用户 12345 的当前资料"。要做点查,你得另起一个 Redis 把同一份数据再写一遍。

Fluss 的 Primary-Key Table 把这两件事合并了:落盘时同时维护两份结构——

  • Log Store:顺序、可重放、带 offset 的流,给流式消费(等价于 Kafka 的 topic 能力);
  • KV Store:按主键组织的点查结构,给亚毫秒级 get(key)(等价于 Redis 的能力)。
        写入 (Append / Upsert)
                 │
                 ▼
        ┌─────────────────┐
        │  PK Table 写入   │
        └────────┬────────┘
                 │ 同一份数据,双写
        ┌────────┴────────┐
        ▼                 ▼
  ┌──────────────┐  ┌──────────────┐
  │  Log Store   │  │   KV Store   │
  │ (顺序/可重放) │  │ (主键/点查)  │
  │  offset 有序 │  │  亚毫秒 get  │
  └──────┬───────┘  └──────┬───────┘
         │                 │
   流式消费/回放        点查/维表 Join

这意味着:你不再需要在 Kafka 和 Redis 之间写同步管道。一份数据进来,既能当流回放,也能当 KV 点查,新鲜度天然一致,因为源是同一个。

2.3 它和 Kafka 的本质区别

一句话总结对比(官方也专门做了 Kafka 对比页):

维度Apache KafkaApache Fluss
定位流式传输(transport)流存储(storage)
点查不支持,需外挂 KVPK Table 内建亚毫秒点查
列式分析无,按字节流消费ARROW 列式 + 服务端裁剪
湖仓统一需外部 Connect / 自研Tiering + Union Read 内建
状态后端状态在 Flink 侧状态外置到 Fluss leader
向量/多模态Lance 集成,支持 RAG 上下文

选型判断:如果你只需要"大规模把事件从 A 搬到 B 再消费",Kafka 依旧是更轻、更成熟的选择;如果你要的是"实时分析 + 在线点查 + 湖仓一体 + AI 上下文",Fluss 是那块被缺了很久的共享底座。


三、架构分析:一个底座怎么把四件套坍缩掉

3.1 组件拓扑

一个最小 Fluss 集群由三类角色组成(外加你熟悉的 ZooKeeper 做协调):

  • CoordinatorServer:集群大脑,管元数据、路由、bucket 分配。客户端 bootstrap.servers 指向的就是它(默认端口 9123)。
  • TabletServer:真正扛读写的数据节点,做服务端投影/下推/裁剪的就是它;本地盘放热数据,冷数据落到远端对象存储。
  • ZooKeeper:集群协调与选主(和很多 Hadoop 系组件一样)。

数据来源(CDC、事件日志、IoT、点击流)先灌进 Fluss 热层(Coordinator + TabletServer),再通过 Tiering Service 把冷数据沉到 Iceberg / Paimon / Lance。

3.2 写入路径:一次写入,双结构落盘

当一条 UPSERT 进 PK Table:

  1. 先写 WAL / Log Store(顺序追加,保证持久化与可重放);
  2. 同步更新内存中的 KV Store(按主键定位,覆盖旧值);
  3. KV 的变更也通过 WAL 保护,故障后可由 Log 重建;
  4. 后台把"足够冷"的 KV 分区 / Log 段按 Tiering 规则上传到对象存储(Iceberg/Paimon 文件)。

对客户端来说,这就是一次写入;对系统来说,它同时喂饱了"流"和"KV"两种消费需求。

3.3 读取路径:三种读法,各走各的最快路径

  • 点查(PK Lookup)get(shop_id, user_id) 直接走 KV Store,亚毫秒返回,给 API / 特征服务 / 维表 Join 用。
  • 流读(Streaming Log):消费者从某个 offset 顺序读 Log Store,等价于 Kafka Consumer,可重放、可回溯。
  • 分析读(Columnar Analytics):下游 Flink/Spark/Trino/DuckDB 扫整表时,TabletServer 在 ARROW 列式格式上做服务端投影(只返回要的列)、谓词下推(只返回满足条件的行)、分区裁剪(跳过无关 bucket/分区)。三道裁剪叠加,IO 和网络是数量级下降,而不是线性下降。

这是 Fluss 在成本上"敢说便宜 85%"的底层原因。传统 Flink 作业的状态(Join 状态、聚合状态、维表缓存)全在 TaskManager 的内存/堆外里,作业一挂,状态恢复要靠 checkpoint 重放,又慢又占资源。

Fluss 把"流处理状态后端"也收进自己:状态外置到 Fluss leader,Flink 计算节点变成无状态、可秒级拉起。计算资源按峰谷弹性伸缩,存储独立扩缩,二者解耦。对云上按量付费的场景,这意味着你不用为了"保住状态"而常驻一大堆昂贵的 Flink slot。

3.5 湖仓统一:Tiering + Union Read

热数据在 TabletServer 本地盘,冷数据在对象存储上的 Iceberg/Paimon/Lance 文件。两层共享同一 schema,对查询引擎暴露成一张表

  • Tiering Service 负责把冷数据自动下沉、对历史分区做 compaction;
  • Union Read 让一次查询同时命中热层(实时)和冷层(历史),对外是单一数据源。

结果:你既不用为"实时性"单独维护 Kafka→Iceberg 同步管道,也不用为"历史回放"再起一套批管道。实时和批式读的是同一份真值。

3.6 多模态:Lance 向量与 RAG 上下文

第五根支柱很多人会忽略,但恰恰契合 2026 年的"实时 AI"叙事。Fluss 通过 Lance 集成把向量也变成同一 substrate 上的一种数据形态,行、列、向量同存。在线特征、RAG 检索上下文、流式分析都是同一张 PK Table 的不同视图——这正是它敢喊"Agentic Lake 最强底座"的技术底气:企业级 Agent 做实时决策时,所需的"秒级更新的上下文"可以直接从 Fluss 一个点拿到,而不是在五个系统间拼凑。


四、代码实战

下面所有可直接跑的示例都来自 Fluss 官方 0.9 文档(quickstart-flink),我做了连贯化整理。最后我会用一段从零实现的 Python 迷你版把"双表示 + 列式裁剪"讲透——注意那段是我自己的教学代码,不是 Fluss 客户端 API,目的是让你看清内核原理。

4.1 一键拉起:Docker Compose

最小可玩环境包含:RustFS(S3 兼容对象存储,做分层存储)、Fluss(CoordinatorServer + TabletServer + ZooKeeper)、Flink(JobManager + TaskManager + SQL Client)。

services:
  # RustFS: S3 兼容存储,承接 Fluss 的分层(tiered)冷数据
  rustfs:
    image: rustfs/rustfs:1.0.0-alpha.83
    ports:
      - "9000:9000"
      - "9001:9001"
    environment:
      - RUSTFS_ACCESS_KEY=rustfsadmin
      - RUSTFS_SECRET_KEY=rustfsadmin
      - RUSTFS_CONSOLE_ENABLE=true
    volumes:
      - rustfs-data:/data
    command: /data
  rustfs-init:
    image: minio/mc
    depends_on:
      - rustfs
    entrypoint: >
      /bin/sh -c "
      until mc alias set rustfs http://rustfs:9000 rustfsadmin rustfsadmin; do
        echo 'Waiting for RustFS...'; sleep 1;
      done;
      mc mb --ignore-existing rustfs/fluss;
      "

  # Fluss 集群:CoordinatorServer + TabletServer + ZooKeeper
  coordinator-server:
    image: apache/fluss:0.9.1-incubating
    command: coordinatorServer
    depends_on:
      - zookeeper
      - rustfs-init
    environment:
      - |
        FLUSS_PROPERTIES=
        zookeeper.address: zookeeper:2181
        bind.listeners: FLUSS://coordinator-server:9123
        remote.data.dir: s3://fluss/remote-data
        s3.endpoint: http://rustfs:9000
        s3.access-key: rustfsadmin
        s3.secret-key: rustfsadmin
        s3.region: us-east-1
        s3.path-style-access: true
        s3.assumed.role.arn: arn:aws:iam::000000000000:role/rustfsadmin
        s3.assumed.role.sts.endpoint: http://rustfs:9000
  tablet-server:
    image: apache/fluss:0.9.1-incubating
    command: tabletServer
    depends_on:
      - coordinator-server
    environment:
      - |
        FLUSS_PROPERTIES=
        zookeeper.address: zookeeper:2181
        bind.listeners: FLUSS://tablet-server:9123
        data.dir: /tmp/fluss/data
        remote.data.dir: s3://fluss/remote-data
        s3.endpoint: http://rustfs:9000
        s3.access-key: rustfsadmin
        s3.secret-key: rustfsadmin
        s3.region: us-east-1
        s3.path-style-access: true
        s3.assumed.role.arn: arn:aws:iam::000000000000:role/rustfsadmin
        s3.assumed.role.sts.endpoint: http://rustfs:9000
  zookeeper:
    restart: always
    image: zookeeper:3.9.2

  # Flink 集群(镜像里已打包 Fluss connector + flink-faker + S3 支持)
  jobmanager:
    image: apache/fluss-quickstart-flink:1.20-0.9.1-incubating
    ports:
      - "8083:8081"
    command: jobmanager
    environment:
      - |
        FLINK_PROPERTIES=
        jobmanager.rpc.address: jobmanager
  taskmanager:
    image: apache/fluss-quickstart-flink:1.20-0.9.1-incubating
    depends_on:
      - jobmanager
    command: taskmanager
    environment:
      - |
        FLINK_PROPERTIES=
        jobmanager.rpc.address: jobmanager
        taskmanager.numberOfTaskSlots: 10
        taskmanager.memory.process.size: 2048m
        taskmanager.memory.framework.off-heap.size: 256m
  sql-client:
    image: apache/fluss-quickstart-flink:1.20-0.9.1-incubating
    depends_on:
      - jobmanager
    command: /opt/sql-client/sql-client
    environment:
      - |
        FLINK_PROPERTIES=
        jobmanager.rpc.address: jobmanager
        rest.address: jobmanager

volumes:
  rustfs-data:

启动:

docker compose up -d
docker compose ps
# 打开 http://localhost:8083 看 Flink Web UI
# 打开 http://localhost:9001 (rustfsadmin/rustfsadmin) 看分层存储桶

进入 SQL Client:docker compose run sql-client

-- 1) 注册 Fluss Catalog
CREATE CATALOG fluss_catalog WITH (
  'type' = 'fluss',
  'bootstrap.servers' = 'coordinator-server:9123'
);
USE CATALOG fluss_catalog;

-- 2) 建主键表(注意 bucket.num 控制分桶并行度)
CREATE TABLE fluss_order (
  `order_key`   BIGINT,
  `cust_key`    INT NOT NULL,
  `total_price` DECIMAL(15, 2),
  `order_date`  DATE,
  `order_priority` STRING,
  `clerk`       STRING,
  `ptime` AS PROCTIME(),
  PRIMARY KEY (`order_key`) NOT ENFORCED
);

CREATE TABLE fluss_customer (
  `cust_key`    INT NOT NULL,
  `name`        STRING,
  `phone`       STRING,
  `nation_key`  INT NOT NULL,
  `acctbal`     DECIMAL(15, 2),
  `mktsegment`  STRING,
  PRIMARY KEY (`cust_key`) NOT ENFORCED
);

-- 3) 把 CDC / 事件源持续灌进 Fluss(EXECUTE STATEMENT SET 批量提交)
EXECUTE STATEMENT SET
BEGIN
  INSERT INTO fluss_order   SELECT * FROM `default_catalog`.`default_database`.source_order;
  INSERT INTO fluss_customer SELECT * FROM `default_catalog`.`default_database`.source_customer;
END;

这里 fluss_order / fluss_customer 就是 PK Table:进来的数据同时进入 Log Store(可被流式消费)和 KV Store(可被点查)。bucket.num 决定分桶数,直接影响点查并行度和热点分布,是后面性能优化的核心旋钮。

4.3 点查 + Lookup Join:一套数据两种吃法

正是因为 PK Table 内建 KV,Flink 的 Lookup Join(维表 Join)可以直接打在 Fluss 上,亚毫秒拿到客户 / 国家维度,把订单流实时打宽:

CREATE TABLE enriched_orders (
  `order_key`       BIGINT,
  `cust_key`        INT NOT NULL,
  `total_price`     DECIMAL(15, 2),
  `order_date`      DATE,
  `order_priority`  STRING,
  `clerk`           STRING,
  `cust_name`       STRING,
  `cust_phone`      STRING,
  `cust_acctbal`    DECIMAL(15, 2),
  `cust_mktsegment` STRING,
  `nation_name`     STRING,
  PRIMARY KEY (`order_key`) NOT ENFORCED
);

INSERT INTO enriched_orders
SELECT
  o.order_key, o.cust_key, o.total_price, o.order_date,
  o.order_priority, o.clerk,
  c.name, c.phone, c.acctbal, c.mktsegment,
  n.name
FROM fluss_order o
LEFT JOIN fluss_customer FOR SYSTEM_TIME AS OF o.ptime AS c
  ON o.cust_key = c.cust_key
LEFT JOIN fluss_nation   FOR SYSTEM_TIME AS OF o.ptime AS n
  ON c.nation_key = n.nation_key;

注意:在传统架构里,fluss_customer 这份维表你本该在 Redis 里再存一份;在 Fluss 里,它就是 PK Table 的 KV 视图,零额外同步。

4.4 从零实现一个迷你"双表示"流存储(教学版)

下面这段 Python 没有用任何外部依赖,目的是把 Fluss 的"双表示 + 列式服务端裁剪 + 分层"内核讲透。它不是 Fluss 客户端 API,请勿直接当生产代码用。

"""
mini_fluss.py —— 教学用迷你流存储,讲透 Fluss 的核心思想:
1) PK Table 双表示:Log Store(顺序可重放) + KV Store(主键点查)
2) 列式存储 + 服务端投影/谓词下推/分区裁剪
3) 分层:热层本地内存,冷层对象存储(模拟)
作者:三哥 / 程序员茄子(仅用于原理演示)
"""

from typing import Any, Dict, List, Optional


class LogStore:
    """顺序、带 offset、可重放的流(对应 Kafka 的 topic 能力)。"""

    def __init__(self) -> None:
        self._segments: List[Dict[str, Any]] = []

    def append(self, record: Dict[str, Any]) -> int:
        offset = len(self._segments)
        self._segments.append({"offset": offset, **record})
        return offset

    def read_from(self, offset: int) -> List[Dict[str, Any]]:
        """从指定 offset 顺序消费(流读 / 回放)。"""
        return self._segments[offset:]


class KvStore:
    """按主键组织的点查结构(对应 Redis 的点查能力)。"""

    def __init__(self) -> None:
        self._data: Dict[Any, Dict[str, Any]] = {}

    def upsert(self, key: Any, value: Dict[str, Any]) -> None:
        self._data[key] = value

    def get(self, key: Any) -> Optional[Dict[str, Any]]:
        return self._data.get(key)


class ColumnarTable:
    """极简列式存储:每列一个 list,演示服务端裁剪。"""

    def __init__(self, columns: List[str]) -> None:
        self.columns = columns
        self._cols: Dict[str, List[Any]] = {c: [] for c in columns}

    def insert_row(self, row: Dict[str, Any]) -> None:
        for c in self.columns:
            self._cols[c].append(row.get(c))

    def scan(self, project: List[str], predicate=None) -> List[Dict[str, Any]]:
        """
        服务端裁剪演示:
        - project:只返回需要的列(投影下推,省网络/IO)
        - predicate:只在服务端过滤满足条件的行(谓词下推)
        """
        n = len(self._cols[self.columns[0]])
        out: List[Dict[str, Any]] = []
        for i in range(n):
            row = {c: self._cols[c][i] for c in self.columns}
            if predicate and not predicate(row):
                continue
            out.append({c: row[c] for c in project})
        return out


class MiniFlussPkTable:
    """PK Table:一次写入,双表示落盘。"""

    def __init__(self, primary_key: str) -> None:
        self.pk = primary_key
        self.log = LogStore()
        self.kv = KvStore()
        self.cold: List[Dict[str, Any]] = []  # 模拟对象存储冷层

    def upsert(self, record: Dict[str, Any]) -> int:
        offset = self.log.append(record)
        self.kv.upsert(record[self.pk], record)
        return offset

    def lookup(self, key: Any) -> Optional[Dict[str, Any]]:
        """点查 —— 走 KV Store,亚毫秒。"""
        return self.kv.get(key)

    def stream(self, offset: int = 0) -> List[Dict[str, Any]]:
        """流读 —— 走 Log Store。"""
        return self.log.read_from(offset)

    def tier_to_cold(self) -> None:
        """分层:把冷数据批量下沉到对象存储(模拟 Tiering Service)。"""
        for rec in self.log._segments:
            self.cold.append(rec)


# ---------- 跑一个端到端示例 ----------
if __name__ == "__main__":
    orders = MiniFlussPkTable(primary_key="order_key")

    orders.upsert({"order_key": 1, "cust_key": 100, "total_price": 99.5, "mkt": "AUTO"})
    orders.upsert({"order_key": 2, "cust_key": 200, "total_price": 12.0, "mkt": "BUILDING"})
    orders.upsert({"order_key": 3, "cust_key": 100, "total_price": 250.0, "mkt": "AUTO"})

    # 点查:等价于 Redis GET,无需外挂 KV
    print("点查 order_key=2 ->", orders.lookup(2))

    # 流读:等价于 Kafka Consumer,从 offset 0 重放
    print("流读全部 ->", orders.stream(0))

    # 列式分析:只投影 total_price,只保留 mkt=AUTO(服务端裁剪)
    ct = ColumnarTable(["order_key", "total_price", "mkt"])
    for r in orders.stream(0):
        ct.insert_row(r)
    cheap_auto = ct.scan(
        project=["order_key", "total_price"],
        predicate=lambda row: row["mkt"] == "AUTO",
    )
    print("列式裁剪结果(仅 AUTO 且只投影两列) ->", cheap_auto)

    # 分层:下沉冷数据
    orders.tier_to_cold()
    print("冷层记录数 ->", len(orders.cold))

运行输出会清晰展示:同一次 upsert 既进了 Log 也进了 KV;点查走 KV、流读取 Log、分析走列式裁剪——三个消费需求来自同一份真值。这就是 Fluss 把 Kafka + Redis 收进一个底座的直觉来源。

4.5 与 Kafka 的迁移对照

如果你现在用 Kafka + Redis + Flink,向 Fluss 收敛时的心智变化:

你现在做的Fluss 里怎么表达
Kafka topic 传事件Fluss Log Store(PK Table 的流视图)
Redis 存维表点查Fluss KV Store(同一 PK Table 的点查视图)
Flink 状态后端外置到 Fluss leader,Flink 变无状态
Kafka→Iceberg 同步管道Tiering Service 自动下沉 + Union Read
自研新鲜度监控单一数据源,不再需要跨系统对账

五、性能优化清单

Fluss 不是"装上就快",要把它喂好,下面几条是实战里最值钱的旋钮。

5.1 bucket.num 是点查与并行的第一旋钮

PK Table 按主键哈希分桶。桶数太少 → 单桶过热、点查热点、并行度上不去;桶数太多 → 元数据和小文件膨胀、KV 维护成本上升。经验值:单桶目标承载 百万级 ~ 千万级主键行,再按 QPS 反推。热点主键(比如大 V 用户)要靠主键设计打散(加 salt / 复合主键)而非单纯加桶。

5.2 服务端裁剪:把计算推到存储侧

这是 Fluss 相对"裸 Kafka + 你自己在 Flink 里全量拉"的最大收益点。确保查询走的是列裁剪 + 谓词下推路径:

  • SELECT 需要的列,别 SELECT *
  • 把过滤条件下推到存储(Flink/Trino/DuckDB 的 planner 会自动做,但你要避免在 UDF 里把整行拉出来再过滤);
  • 利用分区/桶裁剪,让无关 bucket 根本不被扫描。

5.3 点查优化:PK Table + 必要二级索引

亚毫秒点查依赖 KV Store 的主键结构。如果你的点查维度不是主键,需要在上游把"查询键→主键"的映射建模进表设计(例如用 (shop_id, user_id) 复合主键直接覆盖最常见的点查路径),而不是在查询时再做全表过滤。

5.4 分层存储调优:热层本地盘,冷层对象存储

  • 热层(TabletServer 本地盘):用高 IOPS 本地 NVMe 放活跃 KV 与近期 Log,点查和最新流读都在这。
  • 冷层(S3/OSS/Iceberg/Paimon):调 Tiering 的触发阈值(数据年龄 / 大小),让"够冷"的分区自动下沉,并把 compaction 放到低峰期,避免和实时读写抢 IO。
  • Union Read 让你对冷热无感,但冷层访问延迟天然更高,实时看板要显式约束查询时间窗,别让热查询意外扫到冷分区。
  • Lookup Join 缓存:维表 Join 打开 Flink 的 lookup cache,减少重复点查;但缓存 TTL 要和数据新鲜度要求匹配,别为了性能牺牲一致性。
  • State Backend 外置:把 Join / 聚合状态放到 Fluss,而不是堆在 TaskManager,作业失败恢复从"重放 checkpoint"变成"无状态拉起 + 读 Fluss 状态",恢复时间从分钟级降到秒级。
  • 并行度对齐桶数:Flink 消费 Fluss 的并行度尽量与 bucket.num / TabletServer 数对齐,避免个别 subtask 成为瓶颈。

5.6 成本:存算分离的真红利

把"状态常驻在昂贵 Flink slot"改成"状态在 Fluss、Flink 无状态弹性伸缩",是 Fluss 宣称比 Kafka 拓扑便宜最高 85% 的来由。真实收益取决于你的峰谷比:峰谷越悬殊,把计算弹性出去省得越多;如果是常驻平稳流量,收益会打折。这点和所有存算分离架构一致——它不是免费的午餐,但是峰谷明显的实时业务的一剂良药


六、总结与展望

6.1 它解决的真问题

Fluss 不是在造一个新轮子,而是在填一个被行业习惯性忽略的坑:实时数据栈的"分叉税"。当一份数据的真值被拆在 Kafka、Redis、Iceberg 三处,你花在"对齐"上的工程力,往往超过花在"业务"上的。Fluss 用 PK Table 的双表示 + 湖仓统一 + 状态外置,把这几个系统重新收拢成一个 source of truth。

6.2 什么时候该用,什么时候别用

适合:实时分析 + 在线点查 + 湖仓一体叠加的场景;Agentic / 实时 AI 需要"秒级更新的统一上下文"的场景;峰谷明显的实时业务(吃存算分离红利)。

暂不建议:你只需要单纯的高吞吐事件传输、团队对 Kafka 生态极其熟练且不需要点查/湖仓统一——Kafka 更轻更稳;或者你需要的是严格 Exactly-Once 跨多系统的分布式事务(那是另一类问题,Fluss 不替你做全局 2PC)。

6.3 生态与里程碑

截至 2026-08-11,Fluss 已从 Apache 孵化器毕业成为顶级项目(TLP),已在阿里、小红书、爱奇艺、蚂蚁等大规模落地。开放生态是其关键卖点:热层由 Flink/Spark/Trino/StarRocks/Doris/DuckDB 直接读,冷层是 Iceberg/Paimon/Lance 开放格式,治理归 ASF,无厂商锁定。

6.4 结语

流式系统过去十年在"传输"上卷到了极致(Kafka 几乎成了事实标准),但"存储"这一层一直被将就——我们用一堆系统拼出了实时,却始终在为它们的边界买单。Fluss 的意义,是第一次把"流"和"存储"、"实时"和"湖仓"、"点查"和"分析"当成同一个问题的不同侧面来设计。Lakestream(湖流一体) 会不会成为下一个十年的默认范式,现在下结论还早;但至少,被 Kafka + Redis + Flink + Iceberg 四件套折磨过的工程师,多了一个值得认真评估的选择。

本文基于 Apache Fluss 0.9 官方文档与 2026-08-11 TLP 公告整理,代码示例除标注"教学版"的 Python 外均来自官方 quickstart。版本演进快,落地前请以你所用工版本的官方文档为准。

推荐文章

mysql int bigint 自增索引范围
2024-11-18 07:29:12 +0800 CST
JavaScript设计模式:发布订阅模式
2024-11-18 01:52:39 +0800 CST
404错误页面的HTML代码
2024-11-19 06:55:51 +0800 CST
一个有趣的进度条
2024-11-19 09:56:04 +0800 CST
Vue3中怎样处理组件引用?
2024-11-18 23:17:15 +0800 CST
php获取当前域名
2024-11-18 00:12:48 +0800 CST
程序员茄子在线接单