当系统里同时跑着多个Agent去抓取数据、更新数据库、调用外部接口时,并发冲突几乎是绕不开的问题。典型表现包括:两个Agent同时更新同一条记录导致后写入的数据覆盖了先写入的,同一个任务被重复领取执行,或者外部接口被并发调用触发了限流。这类问题的根源在于多个执行单元在无协调的情况下竞争同一份共享资源。解决思路主要有两条:一是用锁机制让竞争方排队,二是用队列把并发操作转换成串行操作。本文结合实际代码详细分析这两种方案的原理、实现和取舍。

一、并发冲突的根源与典型场景
要解决问题,先要弄清楚冲突是怎么产生的。Agent系统中的并发冲突本质上是一种竞态条件(Race Condition):多个执行流在访问共享资源时,最终结果取决于它们的执行顺序,而这个顺序是不确定的。
举个例子,一个数据采集系统部署了5个Agent,负责更新同一个数据库中的商品价格表。Agent A读取到价格为100,准备更新为120;几乎同一时刻,Agent B读取到价格也为100,准备更新为150。如果A先写入、B后写入,最终价格是150,B的操作看似成功了,但A中间做的促销计算逻辑就丢失了。这就是典型的“丢失更新”问题,而且这类错误往往不会报错,只是数据悄悄变错,排查起来非常困难。
除了数据层面的冲突,还有任务层面的冲突。多个Agent从任务表中领取待处理任务时,如果领取逻辑没有做原子控制,同一个任务可能被两个Agent同时取走并执行,造成重复调用外部API、重复扣费等实际损失。此外,缓存失效时的缓存击穿、定时任务在多实例部署下的重复触发,也都是并发冲突的常见变体。
分析这些场景可以归纳出一个共性:只要存在“先检查再操作”这类非原子的复合动作,就可能被并发插入打乱。所以解决思路也很明确,要么把复合动作变成原子的(加锁),要么干脆取消竞争(排队)。
二、锁机制:从单机锁到分布式锁
1. 悲观锁与乐观锁的选择
单进程内的并发可以用语言自带的互斥量解决,比如Python的threading.Lock或Go的sync.Mutex。但Agent系统通常是多进程、多机器部署的,进程内锁不起作用,这时需要数据库锁或分布式锁。
数据库层面的悲观锁可以直接用SELECT ... FOR UPDATE,在事务提交前锁住目标行,其他事务只能等待。这种方式实现简单、强一致,适合冲突频繁的场景,缺点是持有锁的时间较长,吞吐量会下降。乐观锁则不真正加锁,而是在更新时校验版本号:
-- 乐观锁更新,affected_rows为0说明发生了并发修改 UPDATE product SET price = 150, version = version + 1 WHERE id = 1001 AND version = 3;
乐观锁适合读多写少的场景,冲突时需要业务层重试。如果Agent更新频率很高,大量重试反而会拖垮系统,此时悲观锁更合适。
2. Redis分布式锁的实现
跨机器的Agent协调最常用的是Redis分布式锁。核心命令是SET key value NX EX seconds,NX保证只有key不存在时才能设置成功,EX设置过期时间防止死锁。用Python实现一个带重试的锁获取逻辑:
import redis
import time
import uuid
class AgentDistributedLock:
def __init__(self, redis_client, lock_name, timeout=10):
self.redis = redis_client
self.lock_key = f"agent:lock:{lock_name}"
self.lock_id = str(uuid.uuid4()) # 唯一标识,防止误删别人的锁
self.timeout = timeout
def acquire(self, retries=3, interval=0.2):
for _ in range(retries):
# NX保证互斥,EX设置过期时间防止死锁
if self.redis.set(self.lock_key, self.lock_id, nx=True, ex=self.timeout):
return True
time.sleep(interval)
return False
def release(self):
# 用Lua脚本保证"判断+删除"的原子性
script = """
if redis.call('get', KEYS[1]) == ARGV[1] then
return redis.call('del', KEYS[1])
else
return 0
end
"""
return self.redis.eval(script, 1, self.lock_key, self.lock_id)
这里有两个容易踩坑的细节。第一,锁的value必须是唯一标识(比如UUID),释放锁时先比对再删除,否则Agent A因网络延迟阻塞后锁过期,Agent B拿到锁,A恢复后直接删除就会误删B的锁。第二,判断和删除必须放在Lua脚本里原子执行,否则两步之间仍然存在被插队的窗口。
如果对锁的可靠性要求极高,比如涉及资金操作,建议使用Redisson或者etcd、ZooKeeper实现的分布式锁,它们提供了看门狗机制自动续期,避免业务没执行完锁先过期的尴尬局面。
三、队列方案:把并发竞争变成顺序执行
1. 为什么队列能消除冲突
锁的思路是“竞争时排队”,队列的思路更彻底——“根本不让你们同时到达”。把针对同一资源的操作封装成消息发到队列,由单一消费者按顺序处理,天然就没有并发冲突。对于Agent系统,队列还有一个额外好处:削峰填谷。当大量Agent同时产生任务时,队列起到缓冲作用,下游按自己的处理能力匀速消费。
常见做法有两种:一是同一个key的操作路由到同一个分区,比如用商品ID做分区键,保证同一商品的操作严格有序;二是全局单队列串行,实现简单但吞吐量受限。Kafka、RocketMQ支持按key分区,RabbitMQ可以用一致性哈希交换机实现类似效果,Redis的List或Stream则适合轻量级场景。
用Redis Stream实现按key分组的消费逻辑示例:
import redis
r = redis.Redis()
# 生产端:Agent把任务投递到Stream
def submit_task(product_id, action, payload):
r.xadd("agent:tasks", {
"product_id": product_id,
"action": action,
"payload": payload
})
# 消费端:用消费者组保证每个消息只被处理一次
def consume(consumer_name):
while True:
# 阻塞读取,组内消息不会重复投递给其他消费者
entries = r.xreadgroup(
group_name="agent-group",
consumername=consumer_name,
streams={"agent:tasks": ">"},
count=10,
block=5000
)
for stream, messages in entries:
for msg_id, fields in messages:
handle(fields)
r.xack("agent:tasks", "agent-group", msg_id)
2. 幂等设计不可缺席
队列能保证顺序,但不能保证消息只被投递一次。网络抖动、消费者重启都会导致消息重发,所以消费逻辑必须幂等。常用手段包括:给每个任务生成唯一任务ID,处理前先查Redis或数据库的已处理标记;或者用INSERT ... ON DUPLICATE KEY这类数据库特性天然去重。
def handle(fields):
task_id = fields[b"task_id"].decode()
# setnx为原子操作,已存在说明处理过,直接跳过
if not r.set(f"agent:done:{task_id}", 1, nx=True, ex=86400):
return # 幂等拦截
# 正常执行业务逻辑
do_business(fields)
四、锁与队列的组合使用建议
实际工程中,锁和队列往往不是二选一,而是分层配合。一个可参考的架构是:Agent产生任务后先写入队列做削峰和缓冲,消费端从队列取任务,对涉及热点资源的操作在处理前再加短时分布式锁,锁粒度尽量小(比如锁单个资源ID而不是全局锁),锁内只做最小临界区的操作,耗时长的网络调用放到锁外。
选型上可以遵循几个原则:冲突少且能接受重试的,用乐观锁最轻量;冲突激烈、临界区短的,用Redis分布式锁;需要严格顺序、高吞吐的,用分区的消息队列;对一致性要求极高的关键链路,优先考虑数据库悲观锁或etcd锁,并在写入侧做好对账兜底。另外无论用哪种方案,都要设计好超时和降级逻辑——锁拿不到时是阻塞等待、快速失败还是放弃任务,应该在编码前就想清楚,否则线上出现锁饥饿时系统会大面积卡住。
最后提醒一点,加锁和排队都会降低并发度,不要为了“安全”给所有操作都套上锁。先梳理清楚哪些资源真正共享、哪些操作会产生冲突,只对必要的临界区做保护,这样才能在正确性和性能之间找到平衡点。