SQL并行脚本死锁问题:设计排查与解决方案咨询
先聊聊你的场景:多段Python脚本并行跑,用pymssql操作三张关联表,结果碰到了死锁——哪怕开了autocommit=True和快照隔离都没用,甚至SQL Server没自动杀死锁进程,得手动终止Python脚本才行。咱们从设计问题到具体解决方案一步步捋:
一、现有设计的固有问题
你的死锁根源主要在几个设计细节上:
1. 子表缺少主键/有效索引
Nmap_Scans_Hosts和Nmap_Scans_Ports没有主键,连匹配更新条件的索引都没有。比如你更新Nmap_Scans_Hosts时用Host_Address和Tags过滤,没有索引的话SQL Server只能全表扫描,会加大范围的页锁甚至表锁——并行脚本同时操作时,很容易互相堵住。
2. 操作顺序不一致
不同脚本对表的访问顺序可能乱掉:比如脚本A先更新Hosts表再碰Ports表,脚本B反过来先操作Ports再操作Hosts,这就形成了经典的循环等待锁,直接触发死锁。
3. Autocommit的副作用
开autocommit=True会让每个SQL语句变成独立事务,锁的持有和释放频繁切换,反而增加了多个脚本抢锁的概率。而且这种碎片化的事务没法保证数据一致性,比如更新了旧数据但插入新数据失败,会导致数据状态混乱。
4. 快照隔离用错了场景
快照隔离解决的是读-写阻塞问题,让读操作不阻塞写、写不阻塞读,但你的死锁是写-写冲突——写操作还是会加排他锁,快照隔离管不了这种情况。
二、解决死锁的具体方案
1. 给子表补全索引/主键(最关键!)
先给两个子表加合适的索引,让SQL Server能快速定位要更新的行,只加行级锁:
针对Nmap_Scans_Hosts
你的更新条件是Host_Address + Tags,而且每条记录应该和Scan_ID唯一绑定,所以直接加复合主键或者唯一索引:
-- 方案1:设置复合主键(推荐,保证数据唯一性) ALTER TABLE nvlanteam.LAN.Nmap_Scans_Hosts ADD CONSTRAINT PK_Nmap_Scans_Hosts PRIMARY KEY (Scan_ID, Host_Address, Tags); -- 方案2:如果不想改主键,加覆盖更新条件的索引 CREATE NONCLUSTERED INDEX IX_Nmap_Scans_Hosts_HostTags ON nvlanteam.LAN.Nmap_Scans_Hosts (Host_Address, Tags) INCLUDE (Latest); -- 包含要更新的列,避免书签查找
针对Nmap_Scans_Ports
你的更新有两种场景:Host_Address + Port IS NULL和Host_Address + Port,所以加复合索引覆盖这两种情况:
-- 方案1:设置复合主键(唯一标识一条端口记录) ALTER TABLE nvlanteam.LAN.Nmap_Scans_Ports ADD CONSTRAINT PK_Nmap_Scans_Ports PRIMARY KEY (Scan_ID, Host_Address, Port); -- 方案2:加覆盖索引 CREATE NONCLUSTERED INDEX IX_Nmap_Scans_Ports_HostPort ON nvlanteam.LAN.Nmap_Scans_Ports (Host_Address, Port) INCLUDE (Latest);
2. 统一所有脚本的操作顺序
所有并行脚本必须严格按照相同的顺序访问表:比如先操作Nmap_Scans,再处理Nmap_Scans_Hosts,最后操作Nmap_Scans_Ports;针对同一IP的操作,必须先完成Hosts的所有更新插入,再碰Ports。这样就能避免循环等待锁的情况——所有脚本都按同一个顺序拿锁,不会出现A等B、B等A的死锁循环。
3. 关闭Autocommit,把相关操作放进一个事务
把针对同一次扫描的所有操作(更新旧的Latest=0、插入新扫描结果)放在一个事务里,这样锁的持有时间是从事务开始到提交,减少中间其他脚本插进来抢锁的机会,还能保证数据一致性:
import pymssql from time import sleep def run_scan_transaction(target, hosts, args): conn = pymssql.connect( server="你的服务器地址", user="用户名", password="密码", database="数据库名", autocommit=False # 关闭自动提交,手动管理事务 ) try: cursor = conn.cursor() # 第一步:处理主表Nmap_Scans cursor.execute( "UPDATE lan.nmap_scans SET latest = 0 WHERE latest = 1 AND target = %s", (target,) ) # 插入新扫描记录,注意用参数化查询避免SQL注入 cursor.execute( """INSERT INTO nvlanteam.LAN.Nmap_Scans (Target, Hosts_Total, Hosts_Up, Scan_Server, Scan_Command, Scan_Start, Scan_Finish, Elapsed, Latest, Tags) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)""", (target, len(hosts), sum(1 for h in hosts if h['up']), "扫描服务器", args.command, args.start, args.finish, args.elapsed, 1, args.tags) ) # 获取刚插入的Scan_ID(IDENTITY列用SCOPE_IDENTITY()) cursor.execute("SELECT SCOPE_IDENTITY()") scan_id = cursor.fetchone()[0] # 第二步:循环处理每个IP的Hosts和Ports for host in hosts: # 先处理Hosts表 cursor.execute( "UPDATE nvlanteam.LAN.Nmap_Scans_Hosts SET Latest = 0 WHERE Host_Address = %s AND Tags = %s", (host['address'], args.tags) ) cursor.execute( """INSERT INTO nvlanteam.LAN.Nmap_Scans_Hosts (Scan_ID, Host_Address, Host_Name, Region, Country, Location, Division, Subnet_Name, CIDR, Reason, RTT, Latest, Tags) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)""", (scan_id, host['address'], host['name'], host['region'], host['country'], host['location'], host['division'], host['subnet'], host['cidr'], host['reason'], host['rtt'], 1, args.tags) ) # 再处理Ports表 # 先更新Port为NULL的行 cursor.execute( "UPDATE nvlanteam.LAN.Nmap_Scans_Ports SET Latest = 0 WHERE Host_Address = %s AND Port IS NULL", (host['address'],) ) # 插入Port为NULL的记录(如果有的话) if host.get('port_null'): cursor.execute( """INSERT INTO nvlanteam.LAN.Nmap_Scans_Ports (Scan_ID, Host_Address, Port, Protocol, State, Reason, Reason_TTL, Latest) VALUES (%s, %s, NULL, %s, %s, %s, %s, %s)""", (scan_id, host['address'], host['port_null']['proto'], host['port_null']['state'], host['port_null']['reason'], host['port_null']['ttl'], 1) ) # 循环处理每个端口 for port in host['ports']: cursor.execute( "UPDATE nvlanteam.lan.Nmap_Scans_Ports SET Latest = 0 WHERE Host_Address = %s AND Port = %s", (host['address'], port['number']) ) cursor.execute( """INSERT INTO nvlanteam.LAN.Nmap_Scans_Ports (Scan_ID, Host_Address, Port, Protocol, State, Reason, Reason_TTL, Service_Name, Service_Product, Service_Version, Service_Extra_Info, Device_Type, Latest) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)""", (scan_id, host['address'], port['number'], port['proto'], port['state'], port['reason'], port['ttl'], port['service']['name'], port['service']['product'], port['service']['version'], port['service']['extra'], port['device_type'], 1) ) # 提交事务 conn.commit() print(f"扫描任务 {target} 执行成功") except Exception as e: # 发生异常回滚事务 conn.rollback() print(f"扫描任务 {target} 执行失败,已回滚:{str(e)}") raise e finally: cursor.close() conn.close()
4. 启用READ COMMITTED SNAPSHOT隔离级别
这个设置能让默认的READ COMMITTED隔离级别用快照读,减少读操作阻塞写操作的情况,间接降低死锁概率:
ALTER DATABASE 你的数据库名 SET READ_COMMITTED_SNAPSHOT ON;
5. 添加死锁重试机制
极端情况下还是可能出现死锁,所以给脚本加重试逻辑,捕获SQL Server的死锁错误(错误号1205)后重试整个事务:
MAX_RETRIES = 3 def run_scan_with_retry(target, hosts, args): retries = 0 while retries < MAX_RETRIES: try: run_scan_transaction(target, hosts, args) return except pymssql.OperationalError as e: # 检查是否是死锁错误 if e.args[0] == 1205: retries += 1 wait_time = 1 * retries # 指数退避,避免重复抢锁 print(f"死锁发生,{wait_time}秒后重试第{retries}次...") sleep(wait_time) else: raise e print(f"扫描任务 {target} 重试{MAX_RETRIES}次仍失败,放弃执行")
6. 捕获死锁日志定位问题
如果还是有死锁,开启SQL Server的跟踪标志1222,把死锁信息写入错误日志,这样就能看到具体是哪些资源被锁、哪些进程在等待:
DBCC TRACEON(1222, -1); -- 全局启用跟踪标志
之后你可以在SQL Server的错误日志里找到死锁的详细信息,针对性优化。
三、为什么SQL Server没自动终止死锁进程?
通常SQL Server会选一个“牺牲品”进程终止,但如果你的脚本没正确处理异常——比如捕获了异常但没关闭连接、没回滚事务,SQL Server可能没法正确终止该进程的事务。另外,死锁检测需要时间,如果你的事务持有锁的方式比较特殊,可能需要等一会儿才会触发。确保脚本在异常时正确回滚并关闭连接,SQL Server就能正常处理死锁了。
内容的提问来源于stack exchange,提问作者Sederfo

