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的物化视图:
这个视图对已有数据处理有效(1亿行耗时30秒),但插入新数据时无法计算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 > 0orders值——因为物化视图只能处理插入的批次数据,无法访问历史的库存记录。 - 应用层预计算:每秒需要查询25000+产品的最新库存数据,速度太慢,无法满足需求。
- INSERT嵌套SELECT:单条插入可行,但批量插入用
UNION ALL会触发max_ast相关错误,无法批量操作。
现在需要可行的实时或定期计算并存储该差值的方案。
方案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);同步最新库存:用物化视图将
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;定时计算差值写入订单表:用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作为数据管道:
创建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';创建流处理物化视图:实时计算每个库存变更的差值并写入
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:插入阶段预计算差值(批量友好)
如果允许调整原表结构,可在插入时直接计算差值,避免后续复杂处理:
给
stocks表新增orders字段:ALTER TABLE stocks ADD COLUMN orders UInt16 DEFAULT 0;批量插入时结合快照表计算差值:用
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

