如何在Snowflake事务中获取当前事务的最后一个查询ID?
获取Snowflake事务内最后一个查询ID的解决方案
针对并行调用存储过程时,会话级last_query_id()无法准确获取当前事务内目标查询ID的问题,以下是几种可靠的解决思路:
方法1:在存储过程内捕获查询ID并返回
直接在存储过程执行完目标INSERT操作后,调用GET_LAST_QUERY_ID()捕获该查询的ID,通过输出参数传递给调用端。这种方式完全不受并行会话影响,因为查询ID是在事务内部实时捕获的。
存储过程定义
CREATE OR REPLACE PROCEDURE my_data_proc(input_param VARCHAR, OUT insert_query_id VARCHAR) RETURNS VARCHAR LANGUAGE SQL AS $$ BEGIN -- 显式启动事务(若存储过程未设置自动事务) BEGIN TRANSACTION; -- 执行删除步骤 DELETE FROM target_table WHERE filter_col = :input_param; -- 执行插入并捕获查询ID INSERT INTO result_table SELECT col1, col2 FROM source_table WHERE filter_col = :input_param; insert_query_id := GET_LAST_QUERY_ID(); -- 提交事务 COMMIT; RETURN 'PROCEDURE_EXECUTED_SUCCESS'; END; $$;
Python调用示例
import snowflake.connector # 初始化连接(每个并行调用建议使用独立连接) conn = snowflake.connector.connect( user='your_user', password='your_pwd', account='your_account', warehouse='your_wh', database='your_db', schema='your_schema' ) cursor = conn.cursor() try: # 调用存储过程,传入参数并接收输出的查询ID cursor.execute("CALL my_data_proc(%s, ?)", ('your_input_value',)) proc_result = cursor.fetchone() insert_qid = cursor.outputs[0] # 使用result_scan获取插入行数 cursor.execute(f"SELECT COUNT(*) FROM TABLE(RESULT_SCAN('{insert_qid}'))") inserted_rows = cursor.fetchone()[0] print(f"本次插入行数:{inserted_rows}") finally: cursor.close() conn.close()
方法2:直接返回插入行数(更简洁)
不需要依赖查询ID,在存储过程内用ROW_COUNT()函数直接获取最后一次DML操作的行数,通过输出参数返回。这是最直接的方案,完全绕开查询ID的问题。
存储过程定义
CREATE OR REPLACE PROCEDURE my_data_proc(input_param VARCHAR, OUT insert_count INTEGER) RETURNS VARCHAR LANGUAGE SQL AS $$ BEGIN BEGIN TRANSACTION; DELETE FROM target_table WHERE filter_col = :input_param; INSERT INTO result_table SELECT col1, col2 FROM source_table WHERE filter_col = :input_param; insert_count := ROW_COUNT(); -- 获取本次INSERT的行数 COMMIT; RETURN 'PROCEDURE_EXECUTED_SUCCESS'; END; $$;
Python调用示例
# 连接初始化部分同方法1 try: cursor.execute("CALL my_data_proc(%s, ?)", ('your_input_value',)) proc_result = cursor.fetchone() inserted_rows = cursor.outputs[0] print(f"本次插入行数:{inserted_rows}") finally: cursor.close() conn.close()
方法3:通过事务ID关联查询历史(备选)
如果必须通过事务维度查询历史,可在存储过程内获取当前事务ID,之后结合会话ID在QUERY_HISTORY_BY_SESSION()视图中过滤出目标查询。注意该方法存在查询历史延迟的可能,仅建议在前两种方法无法使用时采用。
存储过程定义
CREATE OR REPLACE PROCEDURE my_data_proc(input_param VARCHAR, OUT txn_id VARCHAR) RETURNS VARCHAR LANGUAGE SQL AS $$ BEGIN BEGIN TRANSACTION; -- 获取当前事务ID SELECT CURRENT_TRANSACTION() INTO txn_id; DELETE FROM target_table WHERE filter_col = :input_param; INSERT INTO result_table SELECT col1, col2 FROM source_table WHERE filter_col = :input_param; COMMIT; RETURN 'PROCEDURE_EXECUTED_SUCCESS'; END; $$;
Python调用示例
# 连接初始化部分同方法1 try: cursor.execute("CALL my_data_proc(%s, ?)", ('your_input_value',)) proc_result = cursor.fetchone() txn_id = cursor.outputs[0] # 查询当前会话内该事务的INSERT查询ID cursor.execute(""" SELECT QUERY_ID FROM TABLE(QUERY_HISTORY_BY_SESSION()) WHERE TRANSACTION_ID = %s AND QUERY_TYPE = 'INSERT' ORDER BY START_TIME DESC LIMIT 1 """, (txn_id,)) insert_qid = cursor.fetchone()[0] # 后续使用result_scan逻辑同方法1 finally: cursor.close() conn.close()
内容的提问来源于stack exchange,提问作者Kumar
相关产品推荐
相关产品推荐

