之前在某数据平台做调度核心改造,每天晚上几十个离线任务抢着去写同一张表,Redis 分布式锁改了一轮又一轮,连接池加了又加,还是会在凌晨的 GC 尖峰里翻车。后来我把锁迁移到了 Kafka 上——没用 Redis,也没用 ZooKeeper,只靠着 Kafka 消费组的分区分配协议,就稳定跑完了整个季度的大数据任务调度。这篇文章就从原理、实现、踩坑、适用边界四个层面,把 Kafka 分布式锁这件事讲透。
1. 为什么我放弃了 Redis 锁,转投 Kafka 阵营
1.1 大数据任务调度里的"抢活干"问题
先说场景。大数据平台里最怕的一件事就是重复计算:两个 Spark 任务同时启动,都去覆盖同一张 Hive 表,后写的一方把先写的数据冲掉,第二天报表数据全是乱的。为了避免这种事,我们通常会给关键任务加分布式锁,谁拿到锁谁执行,其他人排队等待。
这类锁的特点是:锁粒度特别粗、持有时间特别长、竞争频率并不高。一个离线任务可能跑几分钟甚至几十分钟,一天也就抢一次。它不像电商场景里抢库存那样毫秒级高频操作,而是典型的"低频长持"模型。也正是这个特点,决定了后面的技术选型路径。
1.2 Redis 锁在大数据集群上的三道坎
最开始我们用的是 Redis 分布式锁,SET NX PX一条命令把锁设上,配合 Lua 脚本释放,这套方案在常规业务系统里很成熟。但放在大数据集群上,问题就出来了。
第一道坎是连接集中爆发。离线任务调度时刻是固定的,夜里零点整,几十个 Worker 节点同时尝试加锁。这些节点跟 Kafka、HDFS、YARN 都有长连接,现在又多了 Redis 连接池,瞬间上百个连接涌过去,Redis 所在的网络带宽和 CPU 直接被打满,锁获取的成功率反而更低了。
第二道坎是 JVM 停顿导致的误判。我们的调度节点跑在 Java 里,Full GC 一停就是几秒。Redis 锁有 TTL,业务线程还在 GC 里憋着,锁已经过期了,别的 Worker 拿到锁进来干活,两个任务就同时跑起来了。延长 TTL 能缓解,但延长了又会导致节点崩溃后锁要等很久才能被释放,这是分布式锁老生常谈的矛盾。
第三道坎是可用性问题。Redis 主从切换的情况下,SET NX的原子性没法保证跨节点;如果用 RedLock,等于要维护至少三套 Redis 实例,专门为一个"每天抢一次"的大数据任务加锁,运维成本完全不成比例。
1.3 为什么 Kafka 反而是个更顺手的选项
大数据集群里几乎一定有 Kafka,因为它本来就承担日志采集、消息管道、实时计算这些职责。既然 Kafka 已经存在,那我为什么不多用一个现成组件的能力,而要专门再去维持一套 Redis?
Kafka 的消费组协议里藏着一个天然的互斥机制:同一个消费组里的多个消费者订阅同一个 topic,同一个分区只会被分配给其中一个消费者实例。这个"一个分区只能被一个消费者消费"的语义,本身就是一把分布式锁——分区就是锁资源,消费者就是持锁者。
后面我会详细拆这套机制,但从选型思路上先给个结论:如果你的锁是秒级甚至分钟级的长持锁、竞争不频繁、集群里本来就有 Kafka,那么用 Kafka 做分布式锁是完全靠谱的。
2. Kafka 能当锁用,靠的是消费组协议而不是消息队列
2.1 把分区想象成"只有一个座位的房间"
Kafka 的底层存储单元是 partition 分区。Producer 发消息时按 key 哈希决定进哪个分区,Consumer 消费时则是按分区来分配。这里不扯太深,你只需要记住一个关键点:同一个消费组内,任何一个分区在同一时刻只会被一个消费者实例持有。
换句话说,如果一个 topic 只有一个分区,那这个分区就是"只有一个座位的房间":N 个消费者挤在同一个房间里,最终只有一个能坐到座位上。坐下的就是拿到锁的,站着的就是在等待的。
这个模型是不是很像锁?太像了。座位只有一个,谁坐下就是谁的,这个人退场(关闭消费者),Kafka 协调器会重新分配座位,站在旁边的某个人就能坐下。
2.2 Consumer Group 协议的完整加锁流程
很多同学对 Kafka 消费组的认知停留在"它能负载均衡",但从不知道它内部是怎么协商的。这里补一下完整流程,因为这正是 Kafka 分布式锁的原理基础。
一个消费者实例启动后,会经历这么几步:
- 向 Broker 发送 FindCoordinator 请求,找到消费组对应的协调器节点;
- 协调器返回当前消费组成员列表,消费者发起 JoinGroup 请求;
- 消费组内部选出一个 Leader 消费者,由它执行分区分配策略;
- 分配结果通过 SyncGroup 请求下发给所有成员;
- 之后每个成员持续发送 Heartbeat 心跳,保持 Group 内身份存活。
关键在第 3 步到第 4 步:如果是单分区 topic,RangeAssignor 或者 RoundRobinAssignor 里任意一种分配策略,最终都会把这个唯一的分区分配给一个成员,其他成员分配结果为空。
所以"抢锁"在 Kafka 世界里就是"加入一个消费组并等待分区分配完成"。拿到分区分配,就是加锁成功;没拿到,就是加锁失败等待重试。
2.3 锁的持有、释放与自动转移
持锁的逻辑很简单:拿到的消费者实例一直保持心跳,持续 poll 这个分区里的消息,业务代码在"持锁期间"执行。但真正让这个方案有价值的是释放和转移机制。
第一是主动释放。持锁的消费者调用consumer.close()时,会向协调器发送 LeaveGroup 请求,Broker 感知到成员退出后,触发 rebalance 重新分配分区,其他等待者就有机会拿到这把锁。
第二是故障转移。如果持锁的节点崩溃,没有来得及发 LeaveGroup,协调器会等session.timeout.ms超时后强制把该成员踢出,随后触发 rebalance,把分区分配给别人。这个"崩溃自动转移"是 Redis 锁里最麻烦的场景,在 Kafka 里反而成了默认行为。
这就带来一个重要结论:用 Kafka 做分布式锁,不需要自己写续期、也不需要自己处理持锁者宕机后的锁释放,消费组协议把这两件事全做了。我用了三年 Redis 锁,期间写过的续期守护线程、看门狗脚本,到 Kafka 方案下一条都没用上。
3. 从分区规划到报废机制:完整实现 Kafka 分布式锁
3.1 前置条件:Topic 规划与众核式命名
实现 Kafka 分布式锁之前,先要规划好 Topic。这里有一条铁律:一个锁资源必须对应一个单分区的 topic。为什么后面会在踩坑部分单独讲,总之先记住,一个锁一个独立 topic,分区数设为 1,副本数按集群容错设置(生产环境建议 3)。
创建命令:
kafka-topics.sh --create \ --bootstrap-server kafka01:9092,kafka02:9092,kafka03:9092 \ --topic lock:task:payment-reconcile \ --partitions 1 \ --replication-factor 3命名方式我用的是lock:+ 业务域 + 资源 ID,比如lock:task:payment-reconcile、lock:task:user-profile-build。为什么这么设计?因为每个锁独立的 topic,可以让 Kafka 的分区分配互不干扰;同时 topic 名字本身也能在 Kafka Manager、命令行工具里一眼看出是哪把锁。
3.2 加锁流程的 Java 骨架
下面这段代码是我在实际项目里用过的核心逻辑的简化版。完整代码还包含配置中心下发、监控埋点、线程池管理,但骨架就是这几行:
public class KafkaDistributedLock implements AutoCloseable { private final String brokers; private final String topic; private final String lockGroupId; private volatile KafkaConsumer<String, String> consumer; private volatile boolean locked; public KafkaDistributedLock(String brokers, String topic, String lockGroupId) { this.brokers = brokers; this.topic = topic; this.lockGroupId = lockGroupId; } public boolean tryLock(Duration waitTime) { Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, brokers); props.put(ConsumerConfig.GROUP_ID_CONFIG, lockGroupId); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, RangeAssignor.class.getName()); props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "15000"); props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, "3000"); props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, "600000"); consumer = new KafkaConsumer<>(props); consumer.subscribe(Collections.singletonList(topic)); long deadline = System.currentTimeMillis() + waitTime.toMillis(); while (System.currentTimeMillis() < deadline) { consumer.poll(Duration.ofMillis(500)); if (!consumer.assignment().isEmpty()) { locked = true; return true; } } consumer.close(); consumer = null; return false; } public boolean isHeld() { return locked; } @Override public void close() { if (consumer != null) { consumer.close(); consumer = null; locked = false; } } }核心就一句话:调用poll()等待消费组完成分区分配,然后检查consumer.assignment()是否为空。不为空,说明这个实例拿到了唯一的分区,也就是拿到了锁。
优化点:最好把${lockGroupId}和${topic}设计成同一个字符串体系,比如lock:task:payment-reconcile既是 topic 名,也是 group.id。这样能避免后面坑一里说的"多个锁共用一个消费组"的问题。
3.3 关键参数清单与推荐值
Kafka 消费者的参数很多,但真正决定锁行为的就是下面这几个。我整理了一张表,方便直接抄作业:
| 参数 | 推荐值 | 说明 |
|---|---|---|
session.timeout.ms | 15000 | Broker 判定消费者下线的时间,相当于 Redis 锁的 TTL,但它是基于心跳的,不需要业务续期 |
heartbeat.interval.ms | 3000 | 心跳间隔,建议是 session.timeout 的 1/5 左右,留足重试空间 |
max.poll.interval.ms | 600000 | 两次 poll 之间允许的最大间隔,业务在锁内跑超过这个时间会被踢出 |
enable.auto.commit | false | 锁场景不需要提交位点,关掉自动提交避免无意义的 offset 写入 |
auto.offset.reset | earliest | 只影响消息消费起点,干扰锁逻辑,按经验统一 earliest 最省心 |
partition.assignment.strategy | RangeAssignor | 单分区场景下任意策略结果都一样,搞清楚原理后不用纠结 |
我重点说说session.timeout.ms和max.poll.interval.ms的区别,这是最容易弄混的两兄弟。
session.timeout.ms是心跳层面的判死标准——消费者多久不发心跳就会被踢出。max.poll.interval.ms是业务处理层面的判死标准——消费者多久不调用 poll 方法就会被踢出。如果业务在锁内跑了 10 分钟没 poll,即使心跳一直正常,协调器也会认为消费者卡死了,强制触发 rebalance。所以锁的持有时间上限在 Kafka 里同样存在,只不过它不是锁超时,而是max.poll.interval.ms。
3.4 业务线程与锁的配合方式
实际使用的时候,锁的加锁和业务执行通常不在同一个线程。比如一个 Worker 节点启动后,定时任务触发 tryLock,但业务逻辑可能需要用线程池并发执行。
我的做法是:把 Kafka 消费者放进一个单独的线程循环里,持续调用poll()维持消费组身份;拿到 assignment 之后再把业务任务提交给另一个线程。这样poll()一直有线程在调用,max.poll.interval.ms不会超时,锁的持有时间上限就只取决于session.timeout.ms了,而心跳是不会因为业务阻塞而停掉的。
需要注意一点:这个方案里的锁超时完全由 Kafka 协调器判断,不存在"锁还在手上但时间到了"这种 Redis 式问题。只要有消费者实例在正常发心跳,分区分配就始终归它。真正的持锁者退出只有三种情况:主动 close、进程崩溃、心跳超时。
4. 实测踩过的四个坑:锁丢失、羊群效应、延迟失控、分区数翻车
4.1 坑一:同一个 group.id 被不同业务复用,直接把锁压死
踩坑现象:上线第一个星期,发现某个任务 A 一直拿不到锁,但任务 B 却长时间持锁。日志里看 Kafka 消费组,成员列表有两个,分区始终只分配给了 B。代码检查半天,最后发现问题出在 group.id 上:两个任务的锁都用了同一个 group.id,而同一个消费组下,单分区只能分配给一个消费者。
在 Kafka 里,消费组是锁的边界标识。不同的锁必须用不同的消费组,就像不同的门锁必须配不同的钥匙。如果把所有锁都挂在同一个消费组下,那么无论多少个 topic,这些 topic 的所有分区都会被当成"同一组资源"来分配,互斥逻辑全线崩溃。
解决方式也很简单:group.id 一定要带上锁资源的唯一标识,我用的是和 topic 名完全一致的字面量,从源头杜绝复用。
4.2 坑二:Topic 分区数大于 1,互斥被静默破坏
踩坑现象:压测时发现四个 Worker 同时"拿到锁",四个任务并行跑了起来,数据被覆盖。排查时先怀疑代码逻辑,后来用命令看了下 topic 元数据:
kafka-topics.sh --describe --bootstrap-server kafka01:9092 --topic lock:task:user-profile-build结果 Topic 分区数是 4。原因可能是之前创建 topic 时用了默认分区数配置。单分区 topic 被多消费者抢,是"座位只有一个";而 4 个分区的 topic 被 4 个消费者抢,Kafka 的负载均衡会自然地给每个人分配一个分区,于是每个人都以为自己拿到了锁。
这就是我在前面强调"一个锁必须对应单分区 topic"的原因。这个坑最隐蔽,因为代码逻辑没有任何问题,assignment()不为空就是加锁成功,但分区是 4 个,互斥完全失效。修复方式就是删掉 topic 重新建单分区版本的,或者干脆新建一个分区数为 1 的新 topic 迁移过去。
4.3 坑三:业务处理超时后被踢出,锁在手里"蒸发"
踩坑现象:某个数据补跑任务,业务逻辑跑了 20 分钟,结果跑到第 12 分钟时另一个 worker 也开始跑同一个任务了。日志显示第一个消费者从消费组里被移除了。
根因就是max.poll.interval.ms。我们的业务逻辑虽然不长,但中间有一次外部接口等待了超过 10 分钟,业务线程阻塞的同时 poll 也停了,协调器等不到 poll 调用,判定消费者"处理不过来",强制触发 rebalance,把锁分给了别人。
解决方案分两层:第一,把max.poll.interval.ms调大到业务最大执行时间以上,比如 30 分钟;第二,更稳妥的做法是单独开一个线程只做 poll 维持身份,业务逻辑放到别的线程里执行,两者互不干扰。我在 3.4 节说的就是这个结构。
4.4 坑四:大量实例同时抢锁,rebalance 羊群效应
踩坑现象:白天手动触发了一次全量任务,几十个 Worker 同时尝试抢锁。结果不光抢锁的延迟从秒级暴涨到 30 秒以上,而且抢锁期间整个集群的消息消费延迟也飙高,监控里出现明显的数据堆积。
原因在于 Kafka 的 rebalance 机制:只要有成员加入或退出,消费组就触发一次全量重分配。几十个消费者同时挤进一个消费组,协调器要处理海量的 JoinGroup 请求,还要重新计算分区分配,整个流程被放大。
解决思路有两个。一是错峰:给每个 Worker 的抢锁动作加上随机延迟,比如 30 到 90 秒之间随机,避免集中在同一秒发起。二是减少抢锁失败后的重启流:抢锁失败后不要立即 close 消费者然后重试,因为每一次 close 和重新 join 都会触发 rebalance。正确姿势是消费者一直保持存活,周期性地 poll,直到某次 rebalance 把分区分配给自己。
4.5 顺带说下 Kafka 消息延迟高在这里是个什么信号
很多同学用 Kafka 过程中会遇到"消息延迟高"的警报,在这里我要额外提一句:在使用 Kafka 分布式锁的消费组里,如果发现该消费组的消费延迟飙升,往往不是 Broker 的问题,而是 rebalance 或者 poll 阻塞引起的消费停滞。
排查路径建议按下述顺序走:
- 看消费组状态:
kafka-consumer-groups.sh --describe --group lock:task:xxx --members,确认成员数是否远大于预期; - 看是否有消费者频繁加入退出:
kafka-consumer-groups.sh --describe --group lock:task:xxx看 GROUP 状态频繁从 Stable 变 PreparingRebalance; - 检查业务线程是否长期阻塞导致 poll 调用停摆。
这三点确认完,基本就能定位延迟的根因了。
5. 用 Kafka 做锁的边界:哪些场景合适,哪些会害了你
5.1 分布式锁三套主流方案的对照
做了整个项目之后,我对三种方案的定位有了更清醒的认识,这里直接给张对比表:
| 维度 | Redis 锁 | ZooKeeper 锁 | Kafka 消费组锁 |
|---|---|---|---|
| 获取延迟 | 毫秒级 | 毫秒级 | 百毫秒到秒级(rebalance 耗时) |
| 锁粒度 | 任意字符串 | 任意路径 | 一个 topic 一个锁,天然粗粒度 |
| 超时机制 | 需要业务续期(TTL) | 会话过期自动释放 | 心跳/session 超时自动释放 |
| 可重入 | 基本靠自研 | 需要自研 | 不支持,每把锁一次持有 |
| 公平性 | 不保证 | 临时顺序节点天然 FIFO | 不保证,靠消费组重新分配 |
| 运维依赖 | 需要独立 Redis 集群 | 需要 ZooKeeper 集群 | 复用已有 Kafka 集群 |
| 崩溃转移 | 需要等 TTL 过期 | 会话超时后自动 | 协调器判死后 rebalance |
一句话总结:Redis 锁适合高频短锁、ZK 锁适合需要公平队列的场景、Kafka 锁适合大数据批处理里低频长持的粗粒度互斥。
5.2 哪些场景千万别用 Kafka 做锁
不是所有场景都适合拿 Kafka 当分布式锁用。下面这几种情况我建议你还是老老实实用 Redis:
第一,高频短锁。比如秒杀扣库存、Redis 缓存更新这类毫秒级操作,Kafka 加锁光 rebalance 就需要几百毫秒甚至数秒,完全不可接受。
第二,需要公平锁的场景。Kafka 的分区分配不保证先到先得,谁拿到完全看 rebalance 时机。
第三,锁的持有时间特别不均匀。如果你既要保护 50 毫秒的操作,又要保护 5 分钟的操作,用同一套 Kafka 锁会导致参数配置顾此失彼,Redis 锁配 TTL + 续期反而更灵活。
第四,Kafka 集群本身已经成为瓶颈的场景。Kafka 客户端大量加入退出消费组产生 rebalance 风暴,会影响集群上其他核心链路的消息消费。
5.3 我的最终使用建议
经过这个项目,我个人的体会是:分布式锁没有银弹,方案的价值取决于它解决的问题域。大数据场景里,"抢锁"的机会本来就少,锁持有时间本来就长,集群里本来就有 Kafka,这三个条件叠加在一起时,Kafka 锁的优势远比 Redis 明显。
如果真要上生产,建议先从非核心任务灰度一个月,把消费组监控、rebalance 频率指标都接好,再逐步铺到关键链路。这期间重点盯两个指标:锁获取耗时 p99,以及持锁期间的 rebalance 次数。这两个数字正常,基本就稳了。