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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 04:01:13