欢迎光临
我们一直在努力

第 3 篇:「Fluss 表设计指南」—— 主键表与日志表的最佳实践

第 3 篇:「Fluss 表设计指南」—— 主键表与日志表的最佳实践

阅读本文你将了解: Log Table 与 Primary Key Table 的完整设计方法、分区与分桶策略的选择、Schema Evolution 最佳实践、以及三个生产级实战案例。


3.1 Log Table:仅追加场景的表设计

Log Table 适用于所有只需要追加写入的场景:日志采集、事件流、用户点击流、IoT 传感器数据等。

3.1.1 基本 DDL

— 切换到 Fluss Catalog
USE CATALOG fluss_catalog;
CREATE DATABASE ecommerce;
USE ecommerce;

— 创建一张日志表
CREATE TABLE click_events (
event_id BIGINT,
user_id BIGINT,
event_type STRING,
page_url STRING,
duration_ms INT,
event_time TIMESTAMP(3)
) WITH (
'bucket.num' = '8'
);

3.1.2 Log Table 的底层存储

Log Table 只使用 LogStore,不创建 KvStore:

click_events (Bucket 0)
├── LogTablet (tablet_id=0)
│ ├── 0000.log ← 实际数据
│ └── 0000.index ← 稀疏索引
├── LogTablet (tablet_id=1)
│ └── …
└── (没有 KvTablet)

源码支撑: org.apache.fluss.metadata.TableDescriptor 中,TableType.LOG 类型不会创建 KvTablet。

3.1.3 分区日志表

— 按日期分区的日志表(推荐做法)
CREATE TABLE click_events_partitioned (
event_id BIGINT,
user_id BIGINT,
event_type STRING,
page_url STRING,
duration_ms INT,
event_time TIMESTAMP(3),
event_date STRING — 分区列
) PARTITIONED BY (event_date) — Log Table 的分区列不必是主键子集
WITH (
'bucket.num' = '8'
);

分区的好处:

  • 数据清理:DROP PARTITION event_date='2026-01-01' 一键清除旧数据
  • 查询优化:WHERE event_date='2026-08-08' 只扫描特定分区
  • Lakehouse Compaction:按分区并行压缩,效率更高

3.2 Primary Key Table:可更新场景的表设计

PK Table 是 Fluss 的核心能力所在。一张 PK Table 同时是:

  • 流式日志:所有变更以 changelog 形式记录
  • KV 索引:通过 RocksDB 实现亚毫秒级点查询
  • 分析表:列式 Arrow 格式支持分析查询

3.2.1 基本 DDL

— 用户画像表
CREATE TABLE user_profile (
user_id BIGINT,
name STRING,
age INT,
city STRING,
balance DECIMAL(10, 2),
last_login TIMESTAMP(3),
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'bucket.num' = '16',
'table.merge-engine' = 'deduplicate' — 默认去重合并
);

— 写入/更新数据
INSERT INTO user_profile VALUES
(1, 'Alice', 28, 'Beijing', 1500.00, TIMESTAMP '2026-08-08 10:00:00');

— 更新
UPDATE user_profile SET balance = 2000.00 WHERE user_id = 1;

— 删除
DELETE FROM user_profile WHERE user_id = 1;

3.2.2 写入路径源码分析

用户执行: INSERT INTO user_profile VALUES (1, 'Alice', 28, …)


┌───────────────────────┐
│ Bucketing: hash(1) % 16 = 3 │
│ 路由到 Bucket 3 │
└───────────┬───────────┘


┌───────────────────────┐
│ TabletServer (Bucket 3 的 Leader) │
│ │
│ ┌──────────────────────┐ │
│ │ 1. LogTablet.append()│ ← WAL │
│ │ 写入日志段文件 │ │
│ └────────┬─────────────┘ │
│ │ │
│ ┌────────▼─────────────┐ │
│ │ 2. KvTablet.put() │ ← LSM │
│ │ 写入 RocksDB │ │
│ └──────────────────────┘ │
│ │
│ ┌──────────────────────┐ │
│ │ 3. ISR 副本同步 │ ← HA │
│ └──────────────────────┘ │
└───────────────────────┘

3.2.3 PK 表的分区约束

PK 表的分区列为必须为主键的子集——这是关键的约束条件:

— ✅ 正确:分区列 (country) 在联合主键中
CREATE TABLE orders (
order_id BIGINT,
country STRING,
amount DECIMAL(10, 2),
order_time TIMESTAMP(3),
PRIMARY KEY (order_id, country) NOT ENFORCED
) PARTITIONED BY (country)
WITH ('bucket.num' = '4');

— ❌ 错误:分区列 city 不在主键中
CREATE TABLE users_bad (
user_id BIGINT,
city STRING,
name STRING,
PRIMARY KEY (user_id) NOT ENFORCED
) PARTITIONED BY (city); — 这会报错!
— 错误信息:Partition column 'city' must be a subset of primary key

为什么有这个约束? 因为 PK 表需要保证主键全局唯一。如果分区列不是主键的一部分,同一条主键的数据可能落入不同分区,导致唯一性无法保证。


3.3 Bucket 数量设计

Bucket 数量是最影响性能的参数之一。

3.3.1 Bucket 数量选择公式

推荐 Bucket 数 = max(
TabletServer 节点数 × 2~4, — 保证负载均衡
目标 QPS / 单 Bucket 能力 — 满足并发需求
)

单 Bucket 能力参考:
– LogTablet 写入: ~50K records/sec
– KvTablet 点查: ~20K QPS
– KvTablet 写入: ~10K ops/sec

3.3.2 不同场景的推荐值

场景数据量推荐 Bucket 数理由
轻量维表 < 100万行 4-8 足够负载均衡,减少小文件
中等业务表 100万-1亿行 16-32 平衡并发度和资源消耗
大流量事件表 > 1亿行 64-128 高并发写入需要更多 Bucket
高 QPS 特征表 任意 按 QPS/20K 计算 保证查询并发度

3.3.3 调整 Bucket 数量

— 查看当前 Bucket 数
SHOW CREATE TABLE orders;

— 调整 Bucket 数(需要 Rescale 操作,可能有短暂的影响)
ALTER TABLE orders SET ('bucket.num' = '32');
— 注意:Rescale 会触发数据重分布,建议在低峰期执行


3.4 Merge Engine 配置

Fluss 提供三种合并引擎:

— 1. DEDUPLICATE(默认):按主键去重,保留最新值
CREATE TABLE t1 (
id INT, name STRING, PRIMARY KEY (id) NOT ENFORCED
) WITH ('table.merge-engine' = 'deduplicate');

— 2. PARTIAL_UPDATE:只更新指定列,其他列保留原值
CREATE TABLE user_tags (
user_id BIGINT,
tag1 STRING,
tag2 STRING,
tag3 STRING,
PRIMARY KEY (user_id) NOT ENFORCED
) WITH ('table.merge-engine' = 'partial-update');

INSERT INTO user_tags VALUES (1, 'vip', NULL, NULL);
INSERT INTO user_tags VALUES (1, NULL, 'tech', NULL);
INSERT INTO user_tags VALUES (1, NULL, NULL, 'gamer');
— 最终结果: (1, 'vip', 'tech', 'gamer')

— 3. AGGREGATION:聚合合并
CREATE TABLE user_stats (
user_id BIGINT,
total_spent DECIMAL(10, 2),
order_count INT,
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'table.merge-engine' = 'aggregation',
'fields.total_spent.aggregate-function' = 'sum',
'fields.order_count.aggregate-function' = 'sum'
);

Merge Engine 源码结构

org.apache.fluss.table.merge
├── MergeEngine.java # 合并引擎接口
├── DeduplicateMergeEngine.java # 去重合并
├── PartialUpdateMergeEngine.java # 部分更新
└── AggregationMergeEngine.java # 聚合合并


3.5 Schema Evolution

当前支持 ADD COLUMN 操作:

— 向已有表添加列
ALTER TABLE user_profile ADD COLUMN email STRING;
ALTER TABLE user_profile ADD COLUMN vip_level INT DEFAULT 0;

— 注意:
— 1. 新列对旧数据为 NULL
— 2. 可以使用 DEFAULT 设置默认值
— 3. 当前不支持 DROP COLUMN 和修改列类型

Schema Evolution 源码要点

// 简化自 org.apache.fluss.metadata.SchemaEvolution
public class SchemaEvolution {

public TableSchema apply(TableSchema current, AlterTableOp operation) {
if (operation instanceof AddColumnOp) {
AddColumnOp addCol = (AddColumnOp) operation;

// 校验:列名不能重复
if (current.hasColumn(addCol.getColumnName())) {
throw new SchemaException("Column already exists: " + addCol.getColumnName());
}

// 创建新 Schema(追加列)
return current.addColumn(addCol.getColumnName(), addCol.getDataType());
}
throw new UnsupportedOperationException("Only ADD COLUMN is supported");
}
}


3.6 实战案例

案例 1:电商订单表(PK Table)

CREATE TABLE orders (
order_id BIGINT,
user_id BIGINT,
product_id BIGINT,
amount DECIMAL(10, 2),
status STRING, — plced, paid, shipped, delivered
order_time TIMESTAMP(3),
update_time TIMESTAMP(3),
order_date STRING, — 分区列,格式 'yyyy-MM-dd'
PRIMARY KEY (order_id, order_date) NOT ENFORCED
) PARTITIONED BY (order_date)
WITH (
'bucket.num' = '32',
'table.merge-engine' = 'deduplicate'
);

— 订单状态更新
UPDATE orders SET status = 'paid', update_time = NOW()
WHERE order_id = 12345 AND order_date = '2026-08-08';

案例 2:用户画像表(PK Table + Partial Update)

CREATE TABLE user_profile (
user_id BIGINT,
name STRING,
age INT,
city STRING,
balance DECIMAL(12, 2),
total_orders INT,
last_login_time TIMESTAMP(3),
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'bucket.num' = '8',
'table.merge-engine' = 'partial-update'
);

— 多个上游系统独立更新各自字段,互不影响
— 用户系统更新基本信息
INSERT INTO user_profile(user_id, name, age, city)
VALUES (1, 'Alice', 28, 'Beijing');

— 账务系统只更新余额
INSERT INTO user_profile(user_id, balance)
VALUES (1, 2500.00);

— 推荐系统更新统计指标
INSERT INTO user_profile(user_id, total_orders)
VALUES (1, 42);

— 最终合并结果:(1, 'Alice', 28, 'Beijing', 2500.00, 42, …)

案例 3:实时指标表(Log Table + 分区)

CREATE TABLE realtime_metrics (
metric_name STRING,
metric_value DOUBLE,
tags STRING,
event_time TIMESTAMP(3),
dt STRING,
hour STRING
) PARTITIONED BY (dt, hour)
WITH (
'bucket.num' = '16'
);

— 每小时自动清理 7 天前的分区
— 通过 Lakehouse Compaction 将历史数据转为 Parquet


3.7 总结与下一篇预告

设计决策核心原则
选择表类型 仅追加 → Log Table;需要更新/查询 → PK Table
分区设计 PK 表分区列是主键子集;Log 表按日期分区
Bucket 数量 按节点数×4 和数据量综合计算,后续可调整
Merge Engine 去重用 deduplicate,多源写入用 partial-update,统计用 aggregation
Schema Evolution 当前仅支持 ADD COLUMN,加列使用 DEFAULT

下一篇我们将实战 Fluss + Flink 的完整集成——如何注册 Catalog、创建 Source/Sink、执行流式和批量查询。


本文基于 Apache Fluss 0.9.1 源码。项目 GitHub: https://github.com/apache/fluss

赞(0)
未经允许不得转载:171主机测评 » 第 3 篇:「Fluss 表设计指南」—— 主键表与日志表的最佳实践
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址