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

Apache Airflow并行任务PostgreSQL事务中止与死锁问题求助

问题描述

我用Apache Airflow自动化将数据加载至PostgreSQL数据库,工作流核心步骤为:通过API获取客户与订单数据、插入/更新客户记录、插入订单数据。串行执行时一切正常,但并行执行每周独立数据任务时,出现事务中止和死锁错误。

错误日志

事务中止错误

[2024-12-26T17:24:59.503+0000] {db_insert_customer.py:60} ERROR - Failed to insert/update customer 4507****0 - leticia.****@hotmail.com: current transaction is aborted, commands ignored until end of transaction block
Traceback (most recent call last):
  File "/opt/airflow/airflow_files/tasks/data_load/db_insert_customer.py", line 57, in db_insert_customer
    cursor.execute(query, (row['name'], row['email'], row['cpf'], row['phone']))
psycopg2.errors.InFailedSqlTransaction: current transaction is aborted, commands ignored until end of transaction block

死锁错误

[2024-12-26T18:59:11.974+0000] {db_insert_customer.py:66} ERROR - Failed to insert/update customer 030**** - anyn****4@gmail.com: deadlock detected
DETAIL:  Process 1456 waits for ShareLock on transaction 2230; blocked by process 1458.
Process 1458 waits for ShareLock on transaction 2227; blocked by process 1456.
HINT:  See server log for query details.
CONTEXT:  while inserting index tuple (103,67) in relation "customers"
[2024-12-26T18:59:11.982+0000] {db_insert_customer.py:66} ERROR - Failed to insert/update customer 030**** - anyn****4@gmail.com: current transaction is aborted, commands ignored until end of transaction block

工作流详情

  • 获取指定周的数据;
  • 使用UPSERT语句插入/更新客户数据以避免重复:
query = sql.SQL("""
      INSERT INTO app_schema.customers(name, email, cpf, phone)
      VALUES (%s, %s, %s, %s)
      ON CONFLICT (cpf) DO UPDATE
      SET email = EXCLUDED.email,
        phone = EXCLUDED.phone,
        name = EXCLUDED.name
    """)
  • 获取订单数据并插入orders表;
  • 每周数据独立,并行执行上述操作。

补充信息

  • 数据库:PostgreSQL
  • 驱动:psycopg2
  • customer_id为serial类型

咨询问题

  1. 并行任务时事务中止的原因是什么?
  2. 如何解决该问题,确保并行执行时数据正确插入?
  3. 是否需要先将数据保存为CSV再插入数据库以避免事务问题,该方案是否冗余?

解答

1. 并行任务事务中止的原因

事务中止是死锁触发的连锁反应:

  • 并行任务同时操作customers表时,两个任务各自持有部分锁,又互相等待对方释放锁,触发PostgreSQL的死锁检测,其中一个事务会被强制回滚。
  • 死锁回滚后,当前数据库连接的事务处于失败状态,如果不主动回滚或开启新事务,后续所有SQL命令都会被PostgreSQL拒绝,进而抛出current transaction is aborted错误。
  • 另外,单条UPSERT若遇到冲突(比如并行更新同一CPF),也可能导致单条语句失败,若未处理异常,整个事务会进入失败状态,后续操作也会报错。

2. 解决方法

(1)统一数据处理顺序,从根源避免死锁

死锁的核心是两个任务以相反顺序获取锁,因此要强制所有并行任务按相同顺序处理客户数据:

  • 将每个任务要处理的客户数据按cpf(或唯一键)排序后再执行UPSERT,这样所有任务都会按相同顺序加锁,避免循环等待。

(2)修复事务失败后的处理逻辑

在psycopg2中,一旦事务失败,必须手动回滚才能继续使用当前连接:

import psycopg2
import logging
import time

def process_customer(conn, row):
    query = """
        INSERT INTO app_schema.customers(name, email, cpf, phone)
        VALUES (%s, %s, %s, %s)
        ON CONFLICT (cpf) DO UPDATE
        SET email = EXCLUDED.email,
            phone = EXCLUDED.phone,
            name = EXCLUDED.name
    """
    cursor = conn.cursor()
    try:
        cursor.execute(query, (row['name'], row['email'], row['cpf'], row['phone']))
        conn.commit()
    except psycopg2.errors.DeadlockDetected:
        conn.rollback()
        # 延迟重试当前操作
        time.sleep(1)
        cursor.execute(query, (row['name'], row['email'], row['cpf'], row['phone']))
        conn.commit()
    except psycopg2.Error as e:
        conn.rollback()
        logging.error(f"Failed to process customer {row['cpf']}: {str(e)}")
    finally:
        cursor.close()
  • 注意:尽量缩小事务范围(比如单条或小批量提交),减少锁的持有时间,降低冲突概率。

(3)使用批量UPSERT优化性能

将多条UPSERT合并为批量语句,减少单条语句的执行次数,降低锁竞争:

INSERT INTO app_schema.customers(name, email, cpf, phone)
VALUES (%s, %s, %s, %s), (%s, %s, %s, %s), ...
ON CONFLICT (cpf) DO UPDATE
SET email = EXCLUDED.email,
    phone = EXCLUDED.phone,
    name = EXCLUDED.name;
  • 批量插入前同样要按cpf排序,避免死锁。

(4)调整PostgreSQL锁参数(可选)

适当提高deadlock_timeout(默认1秒),给事务更多时间完成,但这只是缓解手段,无法从根本解决死锁问题。

3. 关于CSV导入的方案

不需要先存CSV再导入,该方案属于冗余:

  • CSV导入(比如COPY命令)确实能提高批量插入效率,但并不能直接解决死锁问题——如果并行导入的CSV包含相同CPF的记录,依然会触发锁竞争和死锁。
  • 若将所有并行任务的数据合并成一个CSV再串行导入,会失去并行执行的意义,反而降低整体效率。
  • 直接优化当前的UPSERT逻辑和事务处理,比转CSV更高效、更直接。

内容的提问来源于stack exchange,提问作者giovanni simoes delsoto

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 08:52:23