快速上手
本文用一张订单表演示 Row Binlog 与 Table Stream 的基本用法:把订单的新增、更新、删除记录下来,再不重不漏地消费到一张下游表。全程只需要一个 MySQL 客户端,大约 10 分钟。
前置条件
-
Doris 5.0.0 及以上版本。
-
在所有 FE 的
fe.conf中加入以下配置并重启 FE(两项均为非动态配置):enable_feature_binlog = true
enable_table_stream = true -
一个能连接 Doris 的 MySQL 客户端。
流程总览
- 创建开启 Row Binlog 的
orders表,写入第一批数据。 - 在
orders上创建min_delta类型的 Table Stream。 - 写入第二批数据,包含更新、删除、新增。
- 通过 Stream 查看变更,理解每一行的变更类型。
- 用
INSERT INTO ... SELECT把变更消费到下游表。 - 在
information_schema.table_stream_consumption中查看消费进度。
第 1 步:创建开启 Row Binlog 的表
Row Binlog 只能在建表时开启。这里创建一张 Unique Key Merge-on-Write 表,并打开 binlog.need_historical_value,这样更新和删除时会记录变更前的值:
CREATE DATABASE IF NOT EXISTS demo;
USE demo;
CREATE TABLE orders (
order_id BIGINT,
status VARCHAR(16),
amount DECIMAL(10, 2)
)
UNIQUE KEY(order_id)
DISTRIBUTED BY HASH(order_id) BUCKETS 1
PROPERTIES (
"replication_num" = "1",
"enable_unique_key_merge_on_write" = "true",
"binlog.enable" = "true",
"binlog.format" = "ROW",
"binlog.need_historical_value" = "true"
);
写入第一批数据:
INSERT INTO orders VALUES
(1, 'created', 100.00),
(2, 'created', 200.00),
(3, 'created', 300.00);
第 2 步:创建 Table Stream
在 orders 表上创建一个 Stream,之后通过它读取和消费 orders 的变更:
CREATE STREAM orders_stream ON TABLE orders
PROPERTIES (
"type" = "min_delta",
"show_initial_rows" = "false"
);
type = min_delta:输出两次消费之间每个 key 的净变化。show_initial_rows = false:创建 Stream 之前已经存在的 3 行不作为变更输出,只关心之后的变化。
此时查询 Stream 没有任何数据:
SELECT * FROM orders_stream;
Empty set
第 3 步:写入第二批变更
这一批操作覆盖更新、删除、新增,以及"新增后又删除"四种情况:
-- 订单 1 状态更新(Unique Key 表写入同 key 即更新)
INSERT INTO orders VALUES (1, 'paid', 100.00);
-- 订单 2 删除
DELETE FROM orders WHERE order_id = 2;
-- 新订单 4
INSERT INTO orders VALUES (4, 'created', 400.00);
-- 新订单 5,随后又被删除
INSERT INTO orders VALUES (5, 'created', 500.00);
DELETE FROM orders WHERE order_id = 5;
第 4 步:查看变更
通过 Stream 查询,并带上两个虚拟列:__DORIS_STREAM_CHANGE_TYPE_COL__ 表示变更类型,__DORIS_STREAM_SEQUENCE_COL__ 表示这次变更的提交时间戳(TSO)。虚拟列不包含在 SELECT * 中,需要显式写出。
SELECT order_id, status, amount,
__DORIS_STREAM_CHANGE_TYPE_COL__ AS change_type,
__DORIS_STREAM_SEQUENCE_COL__ AS change_tso
FROM orders_stream
ORDER BY order_id, change_type DESC;
+----------+---------+--------+---------------+--------------------+
| order_id | status | amount | change_type | change_tso |
+----------+---------+--------+---------------+--------------------+
| 1 | created | 100.00 | UPDATE_BEFORE | 469067680972800000 |
| 1 | paid | 100.00 | UPDATE_AFTER | 469067680972800000 |
| 2 | created | 200.00 | DELETE | 469067681287372800 |
| 4 | created | 400.00 | APPEND | 469067681628160003 |
+----------+---------+--------+---------------+--------------------+
对照第 3 步的操作:
| 第 3 步的操作 | Stream 中的输出 |
|---|---|
| 更新订单 1 | 一对 UPDATE_BEFORE(更新前的值)和 UPDATE_AFTER(更新后的值) |
| 删除订单 2 | 一条 DELETE,携带删除前的值 |
| 新增订单 4 | 一条 APPEND(新 key) |
| 新增订单 5 后又删除 | 两次消费之间净变化为空,min_delta 不输出 |
再执行一次同样的查询,结果完全相同:普通 SELECT 只读取变更,不会推进消费位点。
第 5 步:消费变更
用 INSERT INTO ... SELECT ... FROM <stream> 把变更写入下游表,这条语句在写入成功的同时推进消费位点,两者在同一个事务内完成:
CREATE TABLE orders_changes (
order_id BIGINT,
status VARCHAR(16),
amount DECIMAL(10, 2),
change_type VARCHAR(16),
change_tso BIGINT
)
DUPLICATE KEY(order_id)
DISTRIBUTED BY HASH(order_id) BUCKETS 1
PROPERTIES ("replication_num" = "1");
INSERT INTO orders_changes
SELECT order_id, status, amount,
__DORIS_STREAM_CHANGE_TYPE_COL__,
__DORIS_STREAM_SEQUENCE_COL__
FROM orders_stream;
消费之后再查询 Stream,已经没有未消费的变更:
SELECT COUNT(*) FROM orders_stream;
+----------+
| count(*) |
+----------+
| 0 |
+----------+
继续写入新的变更,Stream 只会返回这次消费之后发生的变化:
INSERT INTO orders VALUES (4, 'paid', 400.00);
SELECT order_id, status, __DORIS_STREAM_CHANGE_TYPE_COL__ AS change_type
FROM orders_stream ORDER BY change_type DESC;
+----------+---------+---------------+
| order_id | status | change_type |
+----------+---------+---------------+
| 4 | created | UPDATE_BEFORE |
| 4 | paid | UPDATE_AFTER |
+----------+---------+---------------+
第 6 步:查看消费进度
information_schema.table_stream_consumption 按分区展示每个 Stream 的消费位点和积压情况:
SELECT STREAM_NAME, UNIT, CONSUMPTION_STATUS, LAG, LAST_CONSUMPTION_TIME
FROM information_schema.table_stream_consumption
WHERE DB_NAME = 'demo' AND STREAM_NAME = 'orders_stream';
+---------------+--------+--------------------+-----------+-----------------------+
| STREAM_NAME | UNIT | CONSUMPTION_STATUS | LAG | LAST_CONSUMPTION_TIME |
+---------------+--------+--------------------+-----------+-----------------------+
| orders_stream | orders | 469067681628160003 | 262144000 | 1789351206000 |
+---------------+--------+--------------------+-----------+-----------------------+
| 列 | 含义 |
|---|---|
UNIT | 消费单元,即基表分区。orders 没有显式分区,只有一个与表同名的分区 |
CONSUMPTION_STATUS | 该分区已消费到的 TSO |
LAG | 基表分区最新提交的 TSO 与已消费 TSO 的差值,0 表示没有积压 |
LAST_CONSUMPTION_TIME | 最近一次消费的时间(毫秒时间戳),-1 表示尚未消费过 |
再执行一次第 5 步的 INSERT INTO orders_changes SELECT ...,LAG 会回到 0。
清理
DROP STREAM orders_stream;
DROP TABLE orders_changes;
DROP TABLE orders;
常见问题
| 问题 | 原因与处理 |
|---|---|
创建 Stream 报 Table Stream is experimental. Please set enable_table_stream=true to enable it. | FE 未开启 enable_table_stream,或修改 fe.conf 后没有重启 FE。按 前置条件 处理 |
| 刚创建 Stream 就查询,结果为空 | show_initial_rows = false 时,创建 Stream 之前已有的数据不作为变更输出。写入新的变更后再查询 |
SELECT * FROM orders_stream 看不到变更类型 | 虚拟列不包含在 SELECT * 中,需要显式写出 __DORIS_STREAM_CHANGE_TYPE_COL__、__DORIS_STREAM_SEQUENCE_COL__ |
| 多次查询 Stream,同样的变更一直在 | 普通 SELECT 只读取变更,不推进消费位点。用 INSERT INTO ... SELECT ... FROM orders_stream 消费后才不再返回 |
| 想给已有的表开启 Row Binlog | Row Binlog 只能在建表时开启,需要新建开启 Row Binlog 的表并导入数据,见 Row Binlog |
LAST_CONSUMPTION_TIME 显示 -1 | 该分区尚未消费过 |
下一步
- 三种消费类型的区别、
show_initial_rows的含义、查询与消费的事务语义:Table Stream 基础 - 按分区分批消费、快照读取、与维表关联、基表 DDL 对 Stream 的影响:Table Stream 进阶
- Row Binlog 的属性、支持范围和限制:Row Binlog
- 不建 Stream,直接按时间窗口读变更:增量查询