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类型
咨询问题
- 并行任务时事务中止的原因是什么?
- 如何解决该问题,确保并行执行时数据正确插入?
- 是否需要先将数据保存为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
相关产品推荐
相关产品推荐

