从Talend调用Snowflake存储过程提示作用域事务未完成回滚异常
问题原因分析
- 报错核心是存储过程内发起的作用域事务未完成,被系统主动回滚。出现运行差异的原因是不同环境的事务提交规则不同:
- Snowflake网页预览端默认开启自动提交(AUTOCOMMIT),存储过程执行完成后会自动提交所有未结束事务,因此运行正常。
- Talend的Snowflake JDBC连接默认关闭自动提交,外部始终处于开启事务的状态,存储过程内部执行的INSERT、UPDATE操作生成的作用域事务在存储过程执行结束后未被显式提交,因此触发报错。
- 存储过程中没有显式处理事务提交逻辑,也是触发问题的原因之一。
解决方案
方案1:开启Talend Snowflake连接的自动提交(最便捷)
在Talend的Snowflake连接配置中,找到自动提交(AUTOCOMMIT)选项,设置为true即可。配置生效后调用存储过程结束后会自动提交所有事务,和网页端运行逻辑完全一致。
方案2:修改存储过程,增加显式事务控制
在存储过程的JS逻辑中增加事务提交逻辑,同时补充异常捕获回滚逻辑提升鲁棒性,修改示例如下:
CREATE OR REPLACE PROCEDURE "CALCULATEDURATIONFORDESIGNAABACUS"() RETURNS VARCHAR(16777216) LANGUAGE JAVASCRIPT EXECUTE AS OWNER AS ' try { var query1 = ` INSERT INTO durations (id, odb_created_at, event_id_arrival, event_id_departure, event_time_arrival, event_time_departure, card_nr, ticket_type, duration, manufacturer, carpark_id) WITH cte AS ( SELECT e.id, e.card_nr, e.event_time, e.ticket_type, e.manufacturer, e.carpark_id, e.device_type, ROW_NUMBER() OVER (ORDER BY e.card_nr, e.carpark_id, e.event_time, e.device_type) AS rn FROM events e LEFT JOIN durations d ON d.event_id_arrival = e.id OR d.event_id_departure = e.id WHERE e.event_time >= (SELECT PROP_VALUE::timestamp FROM properties WHERE prop_key = ''DURATION.LIMIT.DATE'') AND e.device_type IN (1, 2) AND event_type = 2 AND e.manufacturer LIKE ''DESIGNA_ABACUS%'' AND d.id IS NULL ) SELECT durationseq.nextval, current_timestamp(), arrived_entry.id, departed_entry.id, arrived_entry.event_time, departed_entry.event_time, arrived_entry.card_nr, arrived_entry.ticket_type, timestampdiff(second, arrived_entry.event_time, departed_entry.event_time), arrived_entry.manufacturer, arrived_entry.carpark_id FROM (SELECT * FROM cte WHERE cte.device_type = 1) AS arrived_entry INNER JOIN (SELECT * FROM cte WHERE cte.device_type = 2) AS departed_entry ON arrived_entry.card_nr = departed_entry.card_nr AND arrived_entry.carpark_id = departed_entry.carpark_id AND arrived_entry.rn + 1 = departed_entry.rn `; snowflake.execute({ sqlText: query1 }); var query2 = "SELECT PROP_VALUE FROM properties WHERE prop_key = ''DURATION.LIMIT.DAYS''"; var stmt = snowflake.createStatement({ sqlText: query2 }); var resultSet = stmt.execute(); resultSet.next(); var prop_value = resultSet.getColumnValue(1); var query3 = ` UPDATE properties SET PROP_VALUE = ( SELECT dateadd(day, -1 * ${prop_value}, MAX(event_time)) FROM events WHERE event_time >= ( SELECT PROP_VALUE::timestamp FROM properties WHERE prop_key = ''DURATION.LIMIT.DATE'' ) ) WHERE PROP_KEY =''DURATION.LIMIT.DATE''; ` stmt = snowflake.createStatement({ sqlText: query3 }); stmt.execute(); // 新增显式提交 snowflake.execute({sqlText: "COMMIT"}); return ''true''; } catch (e) { // 异常时回滚事务 snowflake.execute({sqlText: "ROLLBACK"}); return ''failed: '' + e.message; } ';
方案3:在tSnowflakeRow中补充提交语句
如果不想修改存储过程和连接配置,可以在tSnowflakeRow的SQL输入框中,把调用语句和提交语句写在一起:
CALL "CALCULATEDURATIONFORDESIGNAABACUS"(); COMMIT;
注意:因为创建存储过程时使用了双引号包裹名称,调用时也需要携带双引号,避免大小写不匹配的问题。
内容的提问来源于stack exchange,提问作者Djabone
相关产品推荐
相关产品推荐

