1.数据更新概述

数据更新是指对具有相同 key 的数据记录中的 value 列进行修改。对于不同的数据模型,数据更新的处理方式有所不同:

主键(Unique)模型:主键模型是专门为数据更新设计的一种数据模型。
主键模型支持使用 UPDATE 语句进行少量数据更新,也支持通过导入方式进行批量更新。导入方式包括 Stream Load、Broker Load、Routine Load 和 Insert Into 等,所有导入操作都遵循“UPSERT”语义,即如果记录不存在则插入,存在则更新。更新操作支持整行更新和部分列更新,默认为整行更新。

聚合(Aggregate)模型:在聚合模型中,数据更新是一种特殊用法。当聚合函数设置为 REPLACE 或 REPLACE_IF_NOT_NULL 时,可以实现数据更新。聚合模型仅支持基于导入方式的更新,不支持使用 UPDATE 语句。
通过设置聚合函数为 REPLACE_IF_NOT_NULL,可以实现部分列更新的能力。

Doris 支持两种存储方式:Merge-on-Read(MoR,读时合并数据)和 Merge-on-Write(MoW,写时合并数据)。
从 Doris 2.1 版本开始,默认存储方式为 MoW,而 MoW 则提供了更好的分析性能。


(1) 2.0 版本仅在 Unique Key 的 Merge-on-Write 实现中支持部分列更新能力。
(2) 从 2.0.2 版本开始,支持使用 INSERT INTO 进行部分列更新。
(3) 不支持在有同步物化视图的表上进行部分列更新。
(4) 不支持在进行 Schema Change 的表上进行部分列更新。

2. 主键模型的 Update 更新

1.适用场景

  • 小范围数据更新:适用于更新少量数据的场景,例如修复某些记录中的错误字段,或更新某些字段的状态(如订单状态更新等)。
  • ETL 批量加工部分字段:适用于大批量更新某个字段,常见于 ETL 加工场景。注意:大范围数据更新仅适合低频调用

2.基本原理

通过 where 过滤逻辑,筛选当天分区。使用 Update 语句先读取整行数据,然后更新部分字段值,再写回,从而实现行级别更新。

3.示例

假设在金融风控场景中,存在如下结构的交易明细表:

CREATE TABLE transaction_details (
    transaction_id BIGINT NOT NULL,        -- 唯一交易编号
    user_id BIGINT NOT NULL,               -- 用户编号
    transaction_date DATE NOT NULL,        -- 交易日期
    transaction_time DATETIME NOT NULL,    -- 交易时间
    transaction_amount DECIMAL(18, 2),     -- 交易金额
    transaction_device STRING,             -- 交易设备
    transaction_region STRING,             -- 交易地区
    average_daily_amount DECIMAL(18, 2),   -- 最近 3 个月日均交易金额
    recent_transaction_count INT,          -- 最近 7 天交易次数
    has_dispute_history BOOLEAN,           -- 是否有拒付记录
    risk_level STRING                      -- 风险等级
)
UNIQUE KEY(transaction_id)
DISTRIBUTED BY HASH(transaction_id) BUCKETS 16
PROPERTIES (
    "replication_num" = "3",               -- 副本数量,默认 3
    "enable_unique_key_merge_on_write" = "true"  -- 启用 MOW 模式,支持合并更新
);

存在如下交易数据:

+----------------+---------+------------------+---------------------+--------------------+--------------------+--------------------+----------------------+--------------------------+---------------------+------------+
| transaction_id | user_id | transaction_date | transaction_time    | transaction_amount | transaction_device | transaction_region | average_daily_amount | recent_transaction_count | has_dispute_history | risk_level |
+----------------+---------+------------------+---------------------+--------------------+--------------------+--------------------+----------------------+--------------------------+---------------------+------------+
|           1001 |    5001 | 2024-11-24       | 2024-11-24 14:30:00 |             100.00 | iPhone 12          | New York           |               100.00 |                       10 |                   0 | NULL       |
|           1002 |    5002 | 2024-11-24       | 2024-11-24 03:30:00 |             120.00 | iPhone 12          | New York           |               100.00 |                       15 |                   0 | NULL       |
|           1003 |    5003 | 2024-11-24       | 2024-11-24 10:00:00 |             150.00 | Samsung S21        | Los Angeles        |               100.00 |                       30 |                   0 | NULL       |
|           1004 |    5004 | 2024-11-24       | 2024-11-24 16:00:00 |             300.00 | MacBook Pro        | high_risk_region1  |               200.00 |                        5 |                   0 | NULL       |
|           1005 |    5005 | 2024-11-24       | 2024-11-24 11:00:00 |            1100.00 | iPad Pro           | Chicago            |               200.00 |                       10 |                   0 | NULL       |
+----------------+---------+------------------+---------------------+--------------------+--------------------+--------------------+----------------------+--------------------------+---------------------+------------+

按照如下风控规则来更新每日所有交易记录的风险等级:
(1)有拒付记录,风险为 high。
(2)在高风险地区,风险为 high。
(3)交易金额异常(超过日均 5 倍),风险为 high。
(4)最近 7 天交易频繁: a. 交易次数 > 50,风险为 high。 b. 交易次数在 20-50 之间,风险为 medium。
(5)非工作时间交易(凌晨 2 点到 4 点),风险为 medium。
(6)默认风险为 low。

UPDATE transaction_details
SET risk_level = CASE
    -- 有拒付记录或在高风险地区的交易
    WHEN has_dispute_history = TRUE THEN 'high'
    WHEN transaction_region IN ('high_risk_region1', 'high_risk_region2') THEN 'high'

    -- 突然异常交易金额
    WHEN transaction_amount > 5 * average_daily_amount THEN 'high'

    -- 最近 7 天交易频率很高
    WHEN recent_transaction_count > 50 THEN 'high'
    WHEN recent_transaction_count BETWEEN 20 AND 50 THEN 'medium'

    -- 非工作时间交易
    WHEN HOUR(transaction_time) BETWEEN 2 AND 4 THEN 'medium'

    -- 默认风险
    ELSE 'low'
END
WHERE transaction_date = '2024-11-24';

更新之后的数据为

+----------------+---------+------------------+---------------------+--------------------+--------------------+--------------------+----------------------+--------------------------+---------------------+------------+
| transaction_id | user_id | transaction_date | transaction_time    | transaction_amount | transaction_device | transaction_region | average_daily_amount | recent_transaction_count | has_dispute_history | risk_level |
+----------------+---------+------------------+---------------------+--------------------+--------------------+--------------------+----------------------+--------------------------+---------------------+------------+
|           1001 |    5001 | 2024-11-24       | 2024-11-24 14:30:00 |             100.00 | iPhone 12          | New York           |               100.00 |                       10 |                   0 | low        |
|           1002 |    5002 | 2024-11-24       | 2024-11-24 03:30:00 |             120.00 | iPhone 12          | New York           |               100.00 |                       15 |                   0 | medium     |
|           1003 |    5003 | 2024-11-24       | 2024-11-24 10:00:00 |             150.00 | Samsung S21        | Los Angeles        |               100.00 |                       30 |                   0 | medium     |
|           1004 |    5004 | 2024-11-24       | 2024-11-24 16:00:00 |             300.00 | MacBook Pro        | high_risk_region1  |               200.00 |                        5 |                   0 | high       |
|           1005 |    5005 | 2024-11-24       | 2024-11-24 11:00:00 |            1100.00 | iPad Pro           | Chicago            |               200.00 |                       10 |                   0 | high       |
+----------------+---------+------------------+---------------------+--------------------+--------------------+--------------------+----------------------+--------------------------+---------------------+------------+

3.聚合模型的导入更新

1.聚合模型的部分列更新

Aggregate 表主要在预聚合场景使用而非数据更新的场景使用,但也可以通过将聚合函数设置为 REPLACE_IF_NOT_NULL 来实现部分列更新效果。

CREATE TABLE order_tbl (
  order_id int(11) NULL,
  order_amount int(11) REPLACE_IF_NOT_NULL NULL,
  order_status varchar(100) REPLACE_IF_NOT_NULL NULL
) ENGINE=OLAP
AGGREGATE KEY(order_id)
COMMENT 'OLAP'
DISTRIBUTED BY HASH(order_id) BUCKETS 1
PROPERTIES (
"replication_allocation" = "tag.location.default: 1"
);

直接通过inser into写入即可

INSERT INTO order_tbl (order_id, order_status) values (1,'待发货');

2.部分列更新使用注意

Aggregate Key 模型在写入过程中不做任何额外处理,所以写入性能不受影响,与普通的数据导入相同。但是在查询时进行聚合的代价较大,典型的聚合查询性能相比 Unique Key 模型的 Merge-on-Write 实现会有 5-10 倍的下降。

由于 REPLACE_IF_NOT_NULL 聚合函数仅在非 NULL 值时才会生效,因此用户无法将某个字段值修改为 NULL 值。

Logo

DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。

更多推荐