跳到主要内容
最后 更新

快速上手

本文用一张订单表演示 Row Binlog 与 Table Stream 的基本用法:把订单的新增、更新、删除记录下来,再不重不漏地消费到一张下游表。全程只需要一个 MySQL 客户端,大约 10 分钟。

前置条件

  • Doris 5.0.0 及以上版本。

  • 在所有 FE 的 fe.conf 中加入以下配置并重启 FE(两项均为非动态配置):

    enable_feature_binlog = true
    enable_table_stream = true
  • 一个能连接 Doris 的 MySQL 客户端。

流程总览

  1. 创建开启 Row Binlog 的 orders 表,写入第一批数据。
  2. orders 上创建 min_delta 类型的 Table Stream。
  3. 写入第二批数据,包含更新、删除、新增。
  4. 通过 Stream 查看变更,理解每一行的变更类型。
  5. INSERT INTO ... SELECT 把变更消费到下游表。
  6. 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 BinlogRow Binlog 只能在建表时开启,需要新建开启 Row Binlog 的表并导入数据,见 Row Binlog
LAST_CONSUMPTION_TIME 显示 -1该分区尚未消费过

下一步

  • 三种消费类型的区别、show_initial_rows 的含义、查询与消费的事务语义:Table Stream 基础
  • 按分区分批消费、快照读取、与维表关联、基表 DDL 对 Stream 的影响:Table Stream 进阶
  • Row Binlog 的属性、支持范围和限制:Row Binlog
  • 不建 Stream,直接按时间窗口读变更:增量查询