You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Clickhouse中如何基于历史行计算库存差值并实时/定期存储?

问题描述

我用ClickHouse存储5000万+产品的stocks数据,表结构如下:

create table stocks
(
    check_id UUID,
    product_id Int32,
    size_id UInt32,
    size_name Nullable(String),
    size_original_name Nullable(String),
    price_retail UInt32,
    price_discount UInt32,
    price_total UInt32,
    quantity UInt16,
    warehouse_id UInt32,
    created_at DateTime default now()
)
engine = MergeTree PRIMARY KEY (check_id, product_id, size_id, warehouse_id)
ORDER BY (check_id, product_id, size_id, warehouse_id, created_at)
SETTINGS index_granularity = 8192;

示例数据:

INSERT INTO default.stocks (check_id, product_id, size_id, size_name, size_original_name, price_retail, price_discount, price_total, quantity, warehouse_id, created_at) VALUES ('0d3b10c0-c831-4d4f-8000-0011647feeee', 56519877, 103180735, '', '0', 590, 0, 590, 8, 146666, '2023-05-02 01:00:00');
INSERT INTO default.stocks (check_id, product_id, size_id, size_name, size_original_name, price_retail, price_discount, price_total, quantity, warehouse_id, created_at) VALUES ('21303340-d2ad-4ed5-8000-00117981cccc', 119518232, 212006721, '', '0', 12558, 40, 7534, 2, 119261, '2023-05-02 01:00:00');
INSERT INTO default.stocks (check_id, product_id, size_id, size_name, size_original_name, price_retail, price_discount, price_total, quantity, warehouse_id, created_at) VALUES ('0d3b10c0-c831-4d4f-8000-0011647fe692', 56519877, 103180735, '', '0', 590, 0, 590, 10, 146666, '2023-05-01 19:22:31');
INSERT INTO default.stocks (check_id, product_id, size_id, size_name, size_original_name, price_retail, price_discount, price_total, quantity, warehouse_id, created_at) VALUES ('21303340-d2ad-4ed5-8000-00117981c344', 119518232, 212006721, '', '0', 12558, 40, 7534, 4, 119261, '2023-05-01 19:22:31');

我的需求是:按product_id, size_id, warehouse_id分组,按created_at DESC排序,计算当前quantity与上一条记录的差值,将该值存储为新列或写入orders表。

我尝试过的方法及遇到的问题:

  • 带POPULATE的物化视图:
    CREATE MATERIALIZED VIEW orders_mv POPULATE AS
    SELECT check_id, product_id, size_id, warehouse_id, orders, created_at
    FROM (SELECT check_id,
                 product_id,
                 size_id,
                 warehouse_id,
                 quantity,
                 leadInFrame(quantity) OVER (PARTITION BY product_id, size_id, warehouse_id ORDER BY product_id ASC, size_id ASC, warehouse_id ASC, created_at DESC ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING) - quantity AS orders,
                 row_number()          OVER (PARTITION BY product_id, size_id, warehouse_id ORDER BY product_id ASC, size_id ASC, warehouse_id ASC, created_at ASC) AS rn,
                 created_at
          FROM stocks)
    WHERE rn > 1
      AND orders > 0
    
    这个视图对已有数据处理有效(1亿行耗时30秒),但插入新数据时无法计算orders值——因为物化视图只能处理插入的批次数据,无法访问历史的库存记录。
  • 应用层预计算:每秒需要查询25000+产品的最新库存数据,速度太慢,无法满足需求。
  • INSERT嵌套SELECT:单条插入可行,但批量插入用UNION ALL会触发max_ast相关错误,无法批量操作。

现在需要可行的实时或定期计算并存储该差值的方案。

可行解决方案

方案1:库存快照表+定时增量计算(近实时)

步骤:

  1. 创建最新库存快照表:用ReplacingMergeTree自动维护每个product_id, size_id, warehouse_id组合的最新库存,旧记录会被自动合并替换:

    CREATE TABLE stock_latest
    (
        product_id Int32,
        size_id UInt32,
        warehouse_id UInt32,
        quantity UInt16,
        created_at DateTime,
        check_id UUID
    )
    ENGINE = ReplacingMergeTree(created_at)
    PRIMARY KEY (product_id, size_id, warehouse_id)
    ORDER BY (product_id, size_id, warehouse_id, created_at);
    
  2. 同步最新库存:用物化视图将stocks的新数据自动同步到快照表:

    CREATE MATERIALIZED VIEW mv_stock_latest TO stock_latest AS
    SELECT product_id, size_id, warehouse_id, quantity, created_at, check_id
    FROM stocks;
    
  3. 定时计算差值写入订单表:用ClickHouse内置的cron任务(或外部调度工具如Airflow)定期执行增量计算,避免全量扫描:

    INSERT INTO orders
    SELECT 
        s.check_id,
        s.product_id,
        s.size_id,
        s.warehouse_id,
        sl.quantity - s.quantity AS orders,
        s.created_at
    FROM stocks s
    LEFT JOIN stock_latest sl 
        ON s.product_id = sl.product_id 
        AND s.size_id = sl.size_id 
        AND s.warehouse_id = sl.warehouse_id
        AND sl.created_at < s.created_at
    WHERE sl.quantity - s.quantity > 0
      AND s.created_at >= now() - INTERVAL 1 HOUR -- 按时间范围增量计算
      AND NOT EXISTS (SELECT 1 FROM orders o WHERE o.check_id = s.check_id); -- 避免重复写入
    

    可根据业务延迟需求调整时间间隔(比如5分钟、1小时),平衡实时性与资源消耗。

方案2:Kafka流引擎实现实时计算

如果需要严格实时处理,可结合Kafka作为数据管道:

  1. 创建Kafka引擎消费表:接收库存数据的实时流(应用层直接写入Kafka,或通过ClickHouse日志引擎同步stocks插入操作):

    CREATE TABLE stocks_kafka
    (
        check_id UUID,
        product_id Int32,
        size_id UInt32,
        quantity UInt16,
        warehouse_id UInt32,
        created_at DateTime
    )
    ENGINE = Kafka
    SETTINGS 
        kafka_broker_list = 'kafka:9092',
        kafka_topic_list = 'stocks_topic',
        kafka_group_name = 'clickhouse_orders_group',
        kafka_format = 'JSONEachRow';
    
  2. 创建流处理物化视图:实时计算每个库存变更的差值并写入orders表:

    CREATE MATERIALIZED VIEW mv_orders TO orders AS
    SELECT 
        check_id,
        product_id,
        size_id,
        warehouse_id,
        prev_quantity - quantity AS orders,
        created_at
    FROM (
        SELECT 
            check_id,
            product_id,
            size_id,
            warehouse_id,
            quantity,
            created_at,
            -- 按分组获取上一条记录的库存
            LAG(quantity) OVER (
                PARTITION BY product_id, size_id, warehouse_id 
                ORDER BY created_at DESC
            ) AS prev_quantity
        FROM stocks_kafka
    )
    WHERE prev_quantity IS NOT NULL 
      AND prev_quantity - quantity > 0;
    

    这种方式能实时处理每一条库存数据,利用Kafka的消息队列特性保证数据不丢失,同时维护分组的历史状态。

方案3:插入阶段预计算差值(批量友好)

如果允许调整原表结构,可在插入时直接计算差值,避免后续复杂处理:

  1. 给stocks表新增orders字段:

    ALTER TABLE stocks ADD COLUMN orders UInt16 DEFAULT 0;
    
  2. 批量插入时结合快照表计算差值:用INSERT ... SELECT替代UNION ALL,直接从快照表获取最新库存:

    -- 假设新数据存在临时表new_stocks_data中
    INSERT INTO stocks (check_id, product_id, size_id, size_name, size_original_name, price_retail, price_discount, price_total, quantity, warehouse_id, created_at, orders)
    SELECT 
        nd.check_id,
        nd.product_id,
        nd.size_id,
        nd.size_name,
        nd.size_original_name,
        nd.price_retail,
        nd.price_discount,
        nd.price_total,
        nd.quantity,
        nd.warehouse_id,
        nd.created_at,
        COALESCE(sl.quantity - nd.quantity, 0) AS orders
    FROM new_stocks_data nd
    LEFT JOIN stock_latest sl 
        ON nd.product_id = sl.product_id 
        AND nd.size_id = sl.size_id 
        AND nd.warehouse_id = sl.warehouse_id
    WHERE COALESCE(sl.quantity - nd.quantity, 0) > 0;
    

    这种方式把计算逻辑前置到插入阶段,利用快照表快速获取最新数据,避免了批量插入的AST错误,同时保证写入效率。

内容的提问来源于stack exchange,提问作者glmn

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.22 20:07:04