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

使用psycopg更新PostgreSQL表时遇OperationalError求助

解决PostgreSQL临时表+COPY更新时的"another command is already in progress"错误

问题场景

通过Python结合psycopg3,尝试用临时表+COPY命令高效批量更新PostgreSQL数据,代码运行时触发错误:

psycopg.OperationalError: sending query and params failed: another command is already in progress

环境:Windows 10、Python 3.12、psycopg 3.1.18,未使用多线程。

原代码:

conn = psycopg.connect(
    host="localhost", dbname=DB_NAME, user=DB_USER, password=DB_PASSWORD
)


with conn.cursor(name="wordfreq_cursor") as cur, conn.cursor() as ins_cur:
        ins_cur.execute("CREATE TEMP TABLE temp_frequency(id INTEGER NOT NULL, frequency FLOAT4) ON COMMIT DROP")
        cur.itersize = 20_000
        cur.execute("SELECT id, lang_code, word FROM etymology LIMIT 3000") # The limit is only for debugging
        i = 1 
        pbar = tqdm(total=20_000_000)
        with ins_cur.copy("COPY temp_frequency (id, frequency) FROM STDIN") as copy:
            for row in cur:
                id, lang_code, word = row
                if lang_code in langs: # langs is a set of string
                    frequency = zipf_frequency(word, lang_code)
                    copy.write_row((id, frequency))
                pbar.update(1)
            ins_cur.execute("UPDATE etymology e SET e.frequency = t.frequency FROM temp_frequency t WHERE e.id = t.id")
            i += 1
        conn.commit()

完整错误栈:

Traceback (most recent call last):
  File "c:\Users\hanne\Documents\Programme\ultimate-dictionary-api\ebook_dictionary_creator\ebook_dictionary_creator\add_wordfreq_to_db.py", line 75, in <module>
    temp_table_solution()
  File "c:\Users\hanne\Documents\Programme\ultimate-dictionary-api\ebook_dictionary_creator\ebook_dictionary_creator\add_wordfreq_to_db.py", line 70, in temp_table_solution
    ins_cur.execute("UPDATE etymology e SET e.frequency = t.frequency FROM temp_frequency t WHERE e.id = t.id")
  File "C:\Users\hanne\Documents\Programme\ultimate-dictionary-api\ebook_dictionary_creator\.venv\Lib\site-packages\psycopg\cursor.py", line 732, in execute
    raise ex.with_traceback(None)
psycopg.OperationalError: sending query failed: another command is already in progress
  0%|                                                                                  | 3000/20000000 [00:04<7:26:13, 746.90it/s]
PS C:\Users\hanne\Documents\Programme\ultimate-dictionary-api> & C:/Users/hanne/Documents/Programme/ultimate-dictionary-api/ebook_dictionary_creator/.venv/Scripts/python.exe c:/Users/hanne/Documents/Programme/ultimate-dictionary-api/ebook_dictionary_creator/ebook_dictionary_creator/add_wordfreq_to_db.py
  0%|                                                                                                | 0/20000000 [00:00<?, ?it/s]Traceback (most recent call last):
  File "c:\Users\hanne\Documents\Programme\ultimate-dictionary-api\ebook_dictionary_creator\ebook_dictionary_creator\add_wordfreq_to_db.py", line 75, in <module>
    temp_table_solution()
  File "c:\Users\hanne\Documents\Programme\ultimate-dictionary-api\ebook_dictionary_creator\ebook_dictionary_creator\add_wordfreq_to_db.py", line 64, in temp_table_solution
    for row in cur:
  File "C:\Users\hanne\Documents\Programme\ultimate-dictionary-api\ebook_dictionary_creator\.venv\Lib\site-packages\psycopg\server_cursor.py", line 332, in __iter__
    recs = self._conn.wait(self._fetch_gen(self.itersize))
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "C:\Users\hanne\Documents\Programme\ultimate-dictionary-api\ebook_dictionary_creator\.venv\Lib\site-packages\psycopg\connection.py", line 969, in wait
    return waiting.wait(gen, self.pgconn.socket, timeout=timeout)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "C:\Users\hanne\Documents\Programme\ultimate-dictionary-api\ebook_dictionary_creator\.venv\Lib\site-packages\psycopg\waiting.py", line 228, in wait_select
    s = next(gen)
        ^^^^^^^^^
  File "C:\Users\hanne\Documents\Programme\ultimate-dictionary-api\ebook_dictionary_creator\.venv\Lib\site-packages\psycopg\server_cursor.py", line 173, in _fetch_gen
    res = yield from self._conn._exec_command(query, result_format=self._format)
          ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "C:\Users\hanne\Documents\Programme\ultimate-dictionary-api\ebook_dictionary_creator\.venv\Lib\site-packages\psycopg\connection.py", line 467, in _exec_command
    self.pgconn.send_query_params(command, None, result_format=result_format)
  File "psycopg_binary\\pq/pgconn.pyx", line 276, in psycopg_binary.pq.PGconn.send_query_params
psycopg.OperationalError: sending query and params failed: another command is already in progress

错误原因

PostgreSQL的单个连接同一时间只能处理一个未完成的命令,原代码存在两个冲突点:

  1. 在COPY操作的上下文管理器内部执行UPDATE,此时COPY会话尚未结束,连接被占用,无法执行新的查询。
  2. 同时使用两个游标(服务器端游标cur和普通游标ins_cur),迭代cur获取数据时,连接处于查询状态,此时用同一连接的另一个游标执行操作会触发冲突。

修复后的代码

conn = psycopg.connect(
    host="localhost", dbname=DB_NAME, user=DB_USER, password=DB_PASSWORD
)

# 拆分游标使用时序,避免交叉操作
with conn.cursor(name="wordfreq_cursor") as cur:
    cur.itersize = 20_000
    cur.execute("SELECT id, lang_code, word FROM etymology LIMIT 3000")  # 调试用LIMIT
    pbar = tqdm(total=20_000_000)
    
    # 准备临时表并添加索引
    with conn.cursor() as ins_cur:
        ins_cur.execute("CREATE TEMP TABLE temp_frequency(id INTEGER NOT NULL, frequency FLOAT4) ON COMMIT DROP")
        ins_cur.execute("CREATE INDEX idx_temp_freq_id ON temp_frequency(id)")
    
    # 执行COPY写入临时表
    with conn.cursor() as copy_cur, copy_cur.copy("COPY temp_frequency (id, frequency) FROM STDIN") as copy:
        for row in cur:
            id, lang_code, word = row
            if lang_code in langs:  # langs是字符串集合
                frequency = zipf_frequency(word, lang_code)
                copy.write_row((id, frequency))
            pbar.update(1)
    
    # COPY完成后执行UPDATE
    with conn.cursor() as update_cur:
        update_cur.execute("UPDATE etymology e SET e.frequency = t.frequency FROM temp_frequency t WHERE e.id = t.id")
    
    conn.commit()

关键修改说明

  1. 时序拆分:先完成数据读取和COPY写入,退出COPY上下文后再执行UPDATE,确保同一时间连接只处理一个命令。
  2. 独立游标处理阶段:每个操作(创建临时表、COPY、UPDATE)使用独立游标,避免游标间状态冲突。
  3. 添加临时表索引:给temp_frequency的id字段加索引,大幅提升UPDATE时的关联匹配速度。

额外优化建议

  • 若处理超大规模数据,可将数据分成多个批次,每批次写入临时表后执行UPDATE并清空临时表,避免单批次占用过多内存。
  • 服务器端游标itersize可根据服务器配置调整,平衡内存占用和查询效率。

内容的提问来源于stack exchange,提问作者Pux

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 02:15:56