告别Thread.sleep:用分区串行与版本号保证并发写库顺序 📅 发布时间:2026/8/30 3:31:49 👁 浏览次数: 某个版本的批量导入模块上线后业务方连续反馈数据被覆盖。我检查了日志同一用户的两条记录明明先提交的那条最后生效的却是后提交的旧数据。代码逻辑没有明显问题唯一可疑的是当时那个“I got a bad idea”我在每个线程处理完任务后加了一句Thread.sleep(200)以为让后写的线程晚一点执行就能避开并发写入冲突。这个坏主意在本地跑了几轮都正常上线后却暴露出它的脆弱也让我重新理解了“保证顺序”应该靠什么实现。下面就从当时的业务场景开始复盘把坏主意替换成真正可靠的分区串行方案并补充数据库层兜底、并发验证、常见问题和线上排查路径。如果你也写过“先等一下再写库”“延迟几秒应该能避开”这类代码这篇文章可以帮你把思路理顺。1. 先从那个坏主意讲起用 sleep 错开线程执行时间1.1 业务场景同用户记录必须按提交顺序生效当时的场景是批量导入用户资料修改记录。导入文件里每个用户会有多条记录例如用户1001先来一条seq1的内容后面又来一条seq2的内容。业务要求最终数据库里必须是seq2的结果因为seq2代表用户后续的修改晚于seq1。如果并发执行时seq1的任务比seq2晚写库数据库里最终就会留下旧数据。也就是说这个场景有一个严格约束不同用户之间可以并行处理互不影响。同一用户的多条记录必须按照seq从小到大依次生效。为了提升导入吞吐代码用了固定大小的线程池把每条记录作为独立任务提交。问题就出在“同一个用户的多条记录被提交到了不同线程”。1.2 坏主意第一版代码长什么样第一版代码大致是这样ExecutorService pool Executors.newFixedThreadPool(8); for (ImportItem item : importItems) { pool.submit(() - { try { // 模拟解析、校验耗时 Thread.sleep(50); // 坏主意用固定 sleep 让后续任务晚点执行 Thread.sleep(200); importMapper.updateByUserId(item.getUserId(), item.getContent()); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }); }这个写法的逻辑是每条任务处理完成后都强制等待 200 毫秒这样线程之间就会被错开。开发者当时的想法是先提交的任务先开始执行后提交的任务晚点执行数据库里自然就是先提交的先生效。这个逻辑最大的问题是把“线程开始执行的时间先后”和“数据库最终生效的顺序”画了等号。真实系统里两者没有任何确定关系。1.3 本地测试为什么能通过本地跑这个坏主意时数据量通常只有几十条线程池 8 个线程并发度不高本机数据库写入也很快往往刚好是“先提交的先写完”。如果测试数据里同一用户只有一条记录更不可能暴露顺序问题。可以对比一下本地环境和生产环境的差异因素本地环境生产环境数据量几十条到几百条几万条甚至更多线程并发度8 个线程偶发竞争多个应用实例、连接池并发写库耗时毫秒级波动小受磁盘 IO、网络、大事务影响任务耗时稳定性比较稳定GC、IO 抖动会放大耗时验证方式看最后一条数据需要检查终态和日志本地测试通过不代表方案正确只是没有触发那个概率事件。并发问题里“偶发”“复现不出来”本来就是最危险的情况因为它会让人误以为代码没有问题。2. 为什么 sleep 无法承担“保证顺序”这个责任2.1 sleep 的语义和线程调度机制先明确一点Thread.sleep(200)的含义是让当前线程进入TIMED_WAITING状态至少等待 200 毫秒。它不等于“200 毫秒后一定立刻恢复执行”更不等于“恢复后一定能抢到锁、一定先写库”。线程恢复运行取决于操作系统的抢占式调度以及当前 CPU 核心是否空闲。如果系统负载高一个线程从 sleep 中被唤醒到真正获得执行权可能还要再等几十毫秒甚至更久。这个时间完全不受业务代码控制。Thread.sleep(200); // 下面这行代码什么时候执行不由 200 决定而由调度器决定 importMapper.updateByUserId(...);所以固定 sleep 的时间本质上是在赌“所有环境的行为都一样”。本地机器和生产机器的 CPU 数量、负载、数据库配置完全不同这个赌没有任何理论支撑。2.2 sleep 不解决“写库生效顺序”一条导入记录从提交到最终生效至少经过三个环节业务提交顺序导入文件里的先后顺序。线程执行顺序任务被 worker 线程取到的顺序。数据库生效顺序UPDATE真正提交并让后续读到结果的顺序。Thread.sleep最多只能影响第 2 环的一部分对第 1 环和第 3 环没有任何约束力。业务方真正关心的是最终数据库里的结果也就是第 3 环。即使线程 A 先执行UPDATE线程 B 后执行UPDATE在事务隔离、锁等待、批量提交等场景下A 的提交也不一定比 B 更早对数据库生效。靠 sleep 去控制一个自己控制不了的环节结果自然是偶发错乱。2.3 固定时间与真实耗时之间的鸿沟sleep(200)本身隐含了一个假设200 毫秒足够让上一个线程把数据库写完。但这个假设在真实环境里很难成立。数据库写入耗时可能超过 200 毫秒的场景包括数据库连接池繁忙获取连接就需要几百毫秒。目标表数据量大更新语句走了慢查询。有触发器、复杂索引或大字段写入放大明显。应用服务器发生 Full GC线程被暂停。数据库主从延迟虽然主库生效但从库读不到。只要有一次任务的耗时超过 200 毫秒顺序就乱了整个方案就失去了意义。2.4 使用 sleep 还会带来连锁问题除了不能保证顺序固定 sleep 还会引入新的问题吞吐被硬编码拖慢。每条任务至少等 200 毫秒导入一万条记录最少也要 2000 秒。如果在synchronized代码块里调用 sleep会放大临界区让其他线程长时间阻塞。中断处理容易被忽略。线程池关闭时还在 sleep 的任务会收到中断信号必须正确处理InterruptedException。无法扩展。调整线程数、更换机器后固定 sleep 值全部变成魔法数字。核心结论时间只适合用来等待耗时资源比如网络 IO、磁盘 IO、锁释放时间不适合用来“制造顺序”。顺序必须依靠确定性的数据结构、锁或者版本号来保证。3. 正确方向一按业务键分区分区内串行执行3.1 核心设计同 key 必须落在同一处理链路上保证同一用户顺序的正确做法不是让线程等一会而是从根源上让同一用户的多条记录永远只属于同一条处理链路。这个设计可以叫“按 key 分区串行化”也可以叫 sharded executor。具体做法是用业务标识如userId计算分片索引。根据分片索引选择一个单线程执行器。同一个userId永远提交到同一个单线程执行器。单线程执行器内部使用 FIFO 队列任务按提交顺序依次执行。这样不同用户的任务仍然可以并行而同一用户的任务天然串行既保证了顺序又不会把整体吞吐拖成单线程。3.2 用分区执行器实现下面是一个最小可用的OrderedExecutorimport java.util.concurrent.*; import java.util.concurrent.atomic.AtomicInteger; public class OrderedExecutor { private final ExecutorService[] shards; public OrderedExecutor(int shardCount, int queueCapacity) { this.shards new ExecutorService[shardCount]; for (int i 0; i shardCount; i) { shards[i] new ThreadPoolExecutor( 1, 1, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue(queueCapacity), new NamedThreadFactory(order-shard- i), new ThreadPoolExecutor.AbortPolicy() ); } } public void submit(String key, Runnable task) { int index (key.hashCode() Integer.MAX_VALUE) % shards.length; shards[index].execute(task); } public void shutdown() { for (ExecutorService shard : shards) { shard.shutdown(); } } private static final class NamedThreadFactory implements ThreadFactory { private final String prefix; private final AtomicInteger seq new AtomicInteger(); NamedThreadFactory(String prefix) { this.prefix prefix; } Override public Thread newThread(Runnable r) { Thread t new Thread(r, prefix - seq.incrementAndGet()); t.setDaemon(false); return t; } } }使用方式OrderedExecutor executor new OrderedExecutor(8, 1000); for (ImportItem item : importItems) { executor.submit(item.getUserId(), () - { importMapper.updateByUserId(item.getUserId(), item.getContent()); }); } executor.shutdown();关键点在于new ThreadPoolExecutor(1, 1, ...)。核心线程数和最大线程数都是 1队列使用 FIFO 的LinkedBlockingQueue这样同一分片内的任务严格按照提交顺序执行。3.3 关键细节key 必须稳定分区逻辑完全依赖key.hashCode()所以 key 的选择必须稳定且唯一。推荐用业务主键作为 key例如用户 ID、订单 ID。不要在任务执行过程中改变 key不要让 key 包含时间戳、随机数这种会变化的内容。如果 key 为 null所有 null 都会进入同一个分片造成热点。提交任务前要先做空值处理。另外hashCode()本身是允许不同 key 落在同一分片的这不会影响正确性只会影响并行度。例如用户 A 和用户 B 的 hashCode 恰好落在同一个分片那么这两个用户的任务会串行执行吞吐降低一点但数据不会错。3.4 与直接加 synchronized 的区别有人会问为什么不直接对userId加锁synchronized (userId.intern()) { importMapper.updateByUserId(...); }这个思路有两个问题锁只能保证同一时间只有一个线程进入不能保证先进入的线程一定先执行完。如果线程 A 抢到锁后因为 GC 暂停线程 B 可能会先完成写库。String.intern()驻留字符串会占用堆内存大量用户 ID 场景下容易造成内存压力。单线程执行器虽然本质上也是一种“锁”但它比synchronized多了一重 FIFO 队列约束。队列保证了任务在“等待进入执行”这个阶段就排好了顺序这是普通锁做不到的。4. 正确方向二数据库层用版本号做兜底4.1 为什么分区串行化之后还需要兜底分区串行化已经解决了“同一用户并发写库”的问题但线上还会有其他路径进入数据任务失败后重试重试任务不小心走了别的线程。同一个导入批次被重复提交。人工修复数据时直接执行了 SQL。历史任务和当前任务在清理后同时到达。这些场景都需要数据库层有一个“最后防线”保证即使服务端顺序出问题旧数据也不能覆盖新数据。4.2 用 version 字段实现乐观锁在用户资料表上增加版本号字段CREATE TABLE user_profile ( user_id VARCHAR(64) PRIMARY KEY, nickname VARCHAR(128) NOT NULL, version BIGINT NOT NULL DEFAULT 0, updated_at DATETIME NOT NULL );更新时不直接按user_id更新而是带上version条件UPDATE user_profile SET nickname #{content}, version version 1, updated_at NOW() WHERE user_id #{userId} AND version #{expectedVersion};执行后返回影响行数。返回 1更新成功版本号已经加一。返回 0说明当前数据库里的版本号已经不是读取时的版本号说明有其他任务先修改了数据。Java 侧逻辑UserProfile current userProfileMapper.selectByUserId(userId); int rows userProfileMapper.updateWithVersion( userId, content, current.getVersion() ); if (rows 0) { // 版本冲突说明有更晚的数据写入了 // 根据业务决定是忽略还是重新读取后比较 seq }4.3 配合业务 seq 判断是否覆盖如果导入记录本身带有seq字段也可以持久化到数据库里作为“已生效的最新 seq”UPDATE user_profile SET nickname #{content}, last_seq #{seq}, version version 1, updated_at NOW() WHERE user_id #{userId} AND last_seq #{seq};last_seq #{seq}这个条件保证了只有 seq 更大的记录才能覆盖旧数据即使重试任务乱序到达也不会回退。4.4 幂等设计导入任务一般还要考虑重复执行的问题。可以设计一张导入任务明细表CREATE TABLE import_item ( id BIGINT AUTO_INCREMENT PRIMARY KEY, batch_id VARCHAR(64) NOT NULL, row_no INT NOT NULL, user_id VARCHAR(64) NOT NULL, content VARCHAR(255) NOT NULL, status TINYINT NOT NULL DEFAULT 0, UNIQUE KEY uk_batch_row (batch_id, row_no) );(batch_id, row_no)唯一索引可以保证同一条导入记录只生效一次。重复提交时插入或更新可以根据唯一键幂等处理。分区串行和版本号不是二选一而是两层防线。分区串行负责在应用内避免并发写版本号负责在数据层阻止过期数据覆盖。生产环境建议两层都做。5. 用一个小工程跑通替代方案5.1 工程结构和依赖用 Java 17 和 Maven代码只依赖 JDK 自带的java.util.concurrent不需要额外框架。src/main/java/com/example/seq/ ├── OrderedExecutor.java ├── ImportItem.java ├── ImportService.java └── ConcurrentOrderTest.java导入流程的入口类public class ImportService { private final OrderedExecutor executor; public ImportService(OrderedExecutor executor) { this.executor executor; } public void process(ListImportItem items) { for (ImportItem item : items) { executor.submit(item.getUserId(), () - { // 这里执行真实写库逻辑 applyItem(item); }); } } private void applyItem(ImportItem item) { // 实际项目里读取当前版本比较 seq执行带 version 的 UPDATE } }5.2 模拟并发导入并验证顺序下面代码构造了 20 个用户每个用户 50 条记录记录带seq然后打乱顺序提交。最终校验每个用户是否按 seq 递增生效。import java.util.*; import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; public class ConcurrentOrderTest { public static void main(String[] args) throws Exception { int userCount 20; int perUser 50; // 记录每个用户当前生效的 seq MapString, AtomicInteger lastSeqMap new ConcurrentHashMap(); AtomicInteger orderError new AtomicInteger(); ListImportItem items new ArrayList(); for (int u 0; u userCount; u) { String userId user- u; for (int seq 0; seq perUser; seq) { items.add(new ImportItem(userId, seq)); } } Collections.shuffle(items); OrderedExecutor executor new OrderedExecutor(8, 1000); CountDownLatch latch new CountDownLatch(items.size()); for (ImportItem item : items) { executor.submit(item.getUserId(), () - { AtomicInteger current lastSeqMap.computeIfAbsent( item.getUserId(), k - new AtomicInteger(-1) ); int old current.get(); if (item.getSeq() old) { orderError.incrementAndGet(); } else { while (!current.compareAndSet(old, item.getSeq())) { old current.get(); if (item.getSeq() old) { orderError.incrementAndGet(); break; } } } latch.countDown(); }); } if (!latch.await(30, TimeUnit.SECONDS)) { System.out.println(任务未在限定时间内完成); return; } executor.shutdown(); System.out.println(乱序次数: orderError.get()); long finishCount lastSeqMap.values().stream() .filter(v - v.get() perUser - 1) .count(); System.out.println(最终状态正确的用户数: finishCount / userCount); } static class ImportItem { final String userId; final int seq; ImportItem(String userId, int seq) { this.userId userId; this.seq seq; } String getUserId() { return userId; } int getSeq() { return seq; } } }5.3 运行结果和判定标准正常输出乱序次数: 0 最终状态正确的用户数: 20 / 20这个测试的关键是“打乱顺序 多线程并发”。如果 sleep 方案跑这个测试乱序次数会明显大于 0分区串行方案则每轮都应该是 0。需要说明的是这里验证的是“应用层任务是否按序生效”真实项目还要验证数据库最终状态也就是在applyItem执行完再查一次库确认实际写入的内容符合预期。6. 常见问题与排查路径6.1 替换方案后顺序仍然错乱问题现象常见原因检查方式解决方案同一用户数据仍被旧记录覆盖key 不稳定同一用户落到了不同分片打印key.hashCode()和分片索引确认同一用户一定落在同一分片使用稳定业务 ID 作为 key禁止动态拼接 key任务内部又开了子线程异步写库执行器只保证任务开始时不乱不保证子线程不乱检查任务代码里是否有new Thread、线程池提交统一通过OrderedExecutor.submit写库恢复重试时绕过了执行器失败任务直接使用别的线程池重跑查看重试日志里的线程名和队列重试必须走同一个submit(key, task)链路关闭执行器后还有任务提交调用方未等待执行器结束就继续写入数据检查shutdown()后是否还有写库调用使用awaitTermination等待队列清空OrderedExecutor只能约束经过它提交的任务。如果某条代码路径绕过它直接写库顺序保护就失效了。6.2 某个分片队列堆积现象是大部分分片空闲某个分片队列深度持续增长。可能的原因是多个热门用户 ID 哈希碰撞到了同一个分片。某个用户的任务数量远大于其他用户。单条任务耗时过长拖慢了整个分片。排查方式给队列增加监控定期记录每个分片的getQueue().size()。在任务提交时打点统计 key 对应的任务数量。解决方案增加分片数量例如从 8 提到 32。对超大任务集合按批次拆分限制单个 key 的待处理数量。如果单个 key 的任务量极端可以考虑退化为“该 key 单独专用分片”。6.3 提交任务时抛 RejectedExecutionException线程池队列满了默认的AbortPolicy会直接抛出异常并丢弃任务。在导入场景下这会造成数据丢失。解决方式加大队列容量例如从 1000 提到 10000但要注意内存占用。使用CallerRunsPolicy队列满时由提交线程自己执行任务避免丢任务。更稳妥的做法是把失败任务写入待重试表稍后重新提交。使用CallerRunsPolicy时要注意提交线程执行任务会占用调用方时间如果调用方是同步入口可能会有阻塞风险。6.4 数据没有错乱但导入速度比预期慢同一个用户的多条记录必须串行这是业务决定的不是优化能绕开的。如果同一用户记录很多整体吞吐自然会受影响。可以优化的方向把“解析、校验”等不涉及顺序的步骤提前并行执行只把“写库”放到分片执行器里。分片数量根据 CPU 和数据库连接池容量调整不要盲目加大。批量化写库例如同一用户的多条记录合并成一条批量 SQL减少网络往返。7. 那些容易混进来的同类坏主意7.1 用 sleep 等待异步接口返回值类似的坏主意还有这种写法FutureString future asyncClient.submit(request); Thread.sleep(1000); String result future.get(); // 实际上 get 本身就会阻塞等待Future.get()本来就支持阻塞等待先 sleep 再 get 只会增加不必要的等待时间。更糟糕的是如果异步任务执行了 2 秒sleep 1 秒后get()仍会继续阻塞如果异步任务只用了 100 毫秒sleep 1 秒就是白白浪费。正确做法是直接get(timeout)或者使用CompletableFuture的回调方法CompletableFutureString future CompletableFuture.supplyAsync(() - requestRemote()); String result future.get(10, TimeUnit.SECONDS);7.2 用固定 sleep 模拟重试退避失败重试时固定 sleep 几秒再重试也不能有效应对持续故障。如果下游服务已经过载固定间隔重试只会加剧压力。更合理的退避策略是指数退避加随机抖动long base 1000L; int attempt 0; while (attempt maxAttempts) { try { return callRemote(); } catch (Exception e) { attempt; long interval Math.min(30000L, base Math.min(attempt, 5)); interval ThreadLocalRandom.current().nextLong(0, 500); Thread.sleep(interval); } }退避的用途是“等故障恢复”顺序保证的用途是“让任务按既定顺序生效”两者不能互相替代。7.3 用固定线程数解决所有并发问题Executors.newFixedThreadPool(8)本身没有任何顺序保证能力。它只是创建了一个线程池任务在队列里的顺序由提交顺序决定但任务执行完成后对共享数据的影响顺序完全取决于线程调度的实际结果。线程池解决的是资源复用和并发度控制它不解决“多个并发操作之间如何排序”的问题。把线程池当作顺序保证工具是并发编程里最常见的误解之一。7.4 用“本地跑一次能过”作为并发验证如果验证方式只是启动服务、导入一遍数据、看看结果这是远远不够的。并发问题具有概率性必须在不同并发度、不同数据分布下重复验证。建议在 CI 里加入压力验证构造同一个 key 的多条记录打乱顺序提交。强制线程池使用接近生产的并发度。循环执行数百次只要出现一次乱序就判定失败。校验最终数据库状态而不是只看运行日志。8. 复盘坏主意出现时应该怎么处理8.1 把“感觉可行”转成可验证的假设当脑子里冒出“加个 sleep 应该能行”的想法时先把它写成一条技术假设“只要后执行的线程晚 200 毫秒写库就不会覆盖先前线程的数据”。然后追问三个问题200 毫秒这个值从哪里来如果先执行的线程耗时超过 200 毫秒会发生什么为什么不用队列、锁、版本号这些确定性手段这些问题通常只需要一分钟就能想清楚但能避免上线后花两天排查数据错乱。8.2 最省力的验证方法并发测试打乱顺序任何涉及顺序的业务都应该有一个对应的并发顺序测试。测试的核心不是“能跑通”而是“乱序时能发现”。测试要点同一业务 key 的多个任务seq 从 0 递增。提交时对任务列表做Collections.shuffle。使用CountDownLatch等待所有任务结束。校验每个 key 最终生效的 seq 是否为最大值。重复执行多轮只要有一轮乱序就修改方案。如果坏主意版本能过这种测试那大概率不是逻辑真的对了而是测试没有覆盖到关键路径。8.3 上线前检查清单生产环境引入新的并发方案之前建议把下面几项过一遍检查项要求同一业务 key 是否必然串行通过分区执行器或锁保证是否依赖固定 sleep 或固定延时必须改成确定性机制关键写库是否有 version 或 seq 条件不能无条件覆盖任务失败重试是否走同一条链路重试不能绕过顺序控制是否有并发顺序测试CI 中必须存在并稳定通过队列是否有监控分片队列深度、任务耗时、异常率执行器关闭时是否等待任务完成使用awaitTermination防止提前退出数据库是否有唯一键或乐观锁兜底防止重复任务和人工修数导致覆盖8.4 这个错误带给我的真正教训时间只能用来等待耗时资源不能用来伪造顺序。顺序问题的本质是“谁先提交、谁先生效”这必须靠数据结构、锁、版本号这些确定性的手段来解决。以后再听到“加个 sleep 就好”的时候第一反应应该是把它翻译成另一句话这个方案不能确定顺序。然后回到队列、分区、版本号上面把“不确定”变成“确定”。并发代码里最昂贵的往往是“看起来能跑”的代码。它们不会在测试环境报错只会在生产环境选择一个最不合适的时间点暴露问题。