PostgreSQL迁移Snowflake时pgSQL时长计算函数转Snowflake UDF咨询
PostgreSQL迁移Snowflake:停车进出事件时长计算存储过程实现
1. 涉及表结构
三张核心表的建表语句如下:
-- events:进出事件表 CREATE TABLE IF NOT EXISTS events ( id bigint NOT NULL autoincrement start 1 increment 1 PRIMARY KEY, odb_created_at timestamp without time zone NOT NULL, event_time timestamp without time zone NOT NULL, device_type integer NOT NULL, event_type integer NOT NULL, ticket_type integer NOT NULL, card_nr character varying(100), count integer DEFAULT 1 NOT NULL, manufacturer character varying(200), carpark_id bigint ); -- durations:进出时长结果表 CREATE TABLE IF NOT EXISTS durations ( id bigint NOT NULL autoincrement start 1 increment 1 PRIMARY KEY, odb_created_at timestamp without time zone NOT NULL, event_id_arrival bigint, event_id_departure bigint, event_time_arrival timestamp without time zone, event_time_departure timestamp without time zone, card_nr character varying(100), ticket_type integer, duration integer, manufacturer character varying(200), carpark_id bigint ); -- properties:配置表 create or replace TABLE PROPERTIES ( PROP_KEY VARCHAR(80) NOT NULL, PROP_VALUE VARCHAR(250), primary key (PROP_KEY) );
2. 业务规则说明
- events表归集所有车辆进出事件:
device_type=1为进场事件,device_type=2为出场事件 - 已完成计算的事件不可重复计算:进场事件id存入durations表
event_id_arrival字段,出场事件id存入event_id_departure字段 - properties表存储计算配置:
DURATION.LIMIT.DAYS为时间窗口天数,DURATION.LIMIT.DATE为计算起始阈值,仅计算大于该阈值的事件
3. 测试样例数据
events表测试插入语句:
INSERT INTO public.events (id, odb_created_at, event_time, device_type, event_type, ticket_type, card_nr, count, manufacturer, carpark_id) VALUES(188160996, '2021-10-02 04:28:26.338', '2021-10-01 09:14:41.32', 1, 2, 11, '03998988030897300007782', 1, 'XX', 1852); INSERT INTO public.events (id, odb_created_at, event_time, device_type, event_type, ticket_type, card_nr, count, manufacturer, carpark_id) VALUES(188160790, '2021-10-02 04:28:26.248', '2021-10-01 09:31:10.94', 2, 2, 11, '03998988030897300007782', 1, 'XX', 1852); INSERT INTO public.events (id, odb_created_at, event_time, device_type, event_type, ticket_type, card_nr, count, manufacturer, carpark_id) VALUES(188146489, '2021-10-02 04:26:55.069', '2021-10-01 10:03:01.57', 1, 2, 500, '01479804030429500089598', 1, 'XX', 1563); INSERT INTO public.events (id, odb_created_at, event_time, device_type, event_type, ticket_type, card_nr, count, manufacturer, carpark_id) VALUES(188146069, '2021-10-02 04:26:54.852', '2021-10-01 11:49:58.45', 2, 2, 500, '01479804030429500089598', 1, 'XX', 1563); INSERT INTO public.events (id, odb_created_at, event_time, device_type, event_type, ticket_type, card_nr, count, manufacturer, carpark_id) VALUES(188161161, '2021-10-02 04:28:26.372', '2021-10-01 18:44:33.62', 1, 2, 11, '03998988030897300007782', 1, 'XX', 1852); INSERT INTO public.events (id, odb_created_at, event_time, device_type, event_type, ticket_type, card_nr, count, manufacturer, carpark_id) VALUES(188160950, '2021-10-02 04:28:26.329', '2021-10-01 18:45:51.903', 2, 2, 11, '03998988030897300007782', 1, 'XX', 1852); INSERT INTO public.events (id, odb_created_at, event_time, device_type, event_type, ticket_type, card_nr, count, manufacturer, carpark_id) VALUES(188161227, '2021-10-02 04:28:26.374', '2021-10-01 23:21:18.58', 1, 2, 11, '04139733030897300003136', 1, 'XX', 1852); INSERT INTO public.events (id, odb_created_at, event_time, device_type, event_type, ticket_type, card_nr, count, manufacturer, carpark_id) VALUES(188160974, '2021-10-02 04:28:26.334', '2021-10-01 23:24:03.29', 2, 2, 11, '04139733030897300003136', 1, 'XX', 1852); INSERT INTO public.events (id, odb_created_at, event_time, device_type, event_type, ticket_type, card_nr, count, manufacturer, carpark_id) VALUES(188239864, '2021-10-03 04:24:43.345', '2021-10-02 06:49:55.97', 1, 2, 11, '01719400030897300061410', 1, 'XX', 1852); INSERT INTO public.events (id, odb_created_at, event_time, device_type, event_type, ticket_type, card_nr, count, manufacturer, carpark_id) VALUES(188239649, '2021-10-03 04:24:43.308', '2021-10-02 07:02:08.72', 2, 2, 11, '01719400030897300061410', 1, 'XX', 1852);
4. 完整Snowflake实现
注意:Snowflake中UDF仅支持只读操作,涉及表写入/更新需使用存储过程,以下为JavaScript存储过程实现,完全兼容原pg函数逻辑:
CREATE OR REPLACE PROCEDURE CALCULATEDURATION() RETURNS VARCHAR LANGUAGE JAVASCRIPT AS $$ // 1. 读取配置参数 var param_sql = ` SELECT MAX(CASE WHEN PROP_KEY = 'DURATION.LIMIT.DAYS' THEN PROP_VALUE::INTEGER END) AS limit_days, MAX(CASE WHEN PROP_KEY = 'DURATION.LIMIT.DATE' THEN PROP_VALUE::TIMESTAMP END) AS limit_date FROM PROPERTIES `; var param_res = snowflake.execute({sqlText: param_sql}); param_res.next(); var limit_days = param_res.getColumnValue('LIMIT_DAYS'); var limit_date = param_res.getColumnValue('LIMIT_DATE'); // 2. 匹配进出事件对,过滤已计算/不符合条件的事件 var match_sql = ` INSERT INTO durations ( odb_created_at, event_id_arrival, event_id_departure, event_time_arrival, event_time_departure, card_nr, ticket_type, duration, manufacturer, carpark_id ) WITH unprocessed_events AS ( SELECT * FROM events WHERE manufacturer = 'XX' AND event_type = 2 AND device_type IN (1,2) AND event_time >= :1 AND id NOT IN ( SELECT COALESCE(event_id_arrival, event_id_departure) FROM durations WHERE event_id_arrival IS NOT NULL OR event_id_departure IS NOT NULL ) ) SELECT CURRENT_TIMESTAMP() AS odb_created_at, event_id_arrival, event_id_departure, event_time_arrival, event_time_departure, card_nr, ticket_type, duration, manufacturer, carpark_id FROM unprocessed_events MATCH_RECOGNIZE ( PARTITION BY card_nr, carpark_id ORDER BY event_time MEASURES ARR.id AS event_id_arrival, DEP.id AS event_id_departure, ARR.event_time AS event_time_arrival, DEP.event_time AS event_time_departure, DATEDIFF('second', ARR.event_time, DEP.event_time) AS duration, ARR.ticket_type AS ticket_type, ARR.manufacturer AS manufacturer, ARR.carpark_id AS carpark_id ONE ROW PER MATCH PATTERN (ARR DEP) DEFINE ARR AS device_type = 1, DEP AS device_type = 2 ) `; snowflake.execute({sqlText: match_sql, binds: [limit_date]}); // 3. 更新配置表的DURATION.LIMIT.DATE var update_sql = ` UPDATE PROPERTIES SET PROP_VALUE = ( SELECT DATEADD('day', -:1, MAX(event_time))::VARCHAR FROM events WHERE event_time >= :2 ) WHERE PROP_KEY = 'DURATION.LIMIT.DATE' `; snowflake.execute({sqlText: update_sql, binds: [limit_days, limit_date]}); return '计算完成,配置已更新'; $$;
存储过程调用方式
CALL CALCULATEDURATION();
内容的提问来源于stack exchange,提问作者Djabone
相关产品推荐
相关产品推荐

