Java 如何高并发批量处理支付结算相关的业务

发布时间:2026/9/27 5:00:52
Java 如何高并发批量处理支付结算相关的业务 结算场景的并发难的从来不是「把循环改成多线程」而是在百万级订单的吞吐之下同时守住资金安全的三条底线不重复、不遗漏、可对账。本文从任务拆分、线程模型、幂等与状态机到热点行合并、事务边界与对账闭环把整条工程链路拆开讲清楚。目录结算是个什么样的并发问题总体形态结算流水线并发底座线程池与虚拟线程分片策略吞吐与正确性一起定幂等与状态机资金安全的最后防线数据访问层流式读、批量写、热点合并事务边界与一致性闭环实战组装分片调度器与进度账本踩坑清单与上线检查表01 · 结算是个什么样的并发问题先把它和「普通高并发」区分开。交易链路的高并发是横向的请求从网关涌进来无状态服务直接扩容就能扛。结算链路的并发是纵向的一批几十万到几百万条已成交的订单要在有限的时间窗内比如 T1 日终后的两小时全部算清、入账、出款、出报表。前者考验接入层后者考验的是一整条批处理流水线。结算业务有三个业务属性直接决定了技术方案长什么样资金安全高于一切。多结一次是资损少结一次是客诉和监管问题。任何并发优化都必须让位于「不重复、不遗漏」。天然幂等压力。批任务会重试、MQ 会重投、机器会重启同一条结算指令被执行两次是常态而不是意外方案必须默认它会发生。强审计。每一分钱的每一步流转都要留痕出错要能定位到具体哪一条、卡在哪个状态。所以结算批处理的正确目标函数是在时间窗内跑完的前提下把一致性做扎实再去优化吞吐。顺序反了就是在资损的地基上盖楼。02 · 总体形态结算流水线无论日终结算、批量代付还是分账工程形态高度一致可以抽象成五段┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │ 取数切片 │ → │ 任务分片 │ → │ 并发执行 │ → │ 幂等回写 │ → │ 账务落账 │ │ 水印分页 │ │ 按结算账户 │ │ 线程池 / │ │ 状态机CAS │ │ 分户/总账 │ │ │ │ 分组 │ │ 分片隔离 │ │ Outbox │ │ │ └──────────┘ └──────────┘ └────┬─────┘ └──────────┘ └────┬─────┘ │ shard × N │ ▼ ▼ ┌─────────────────────────────────────────────────────────────────────┐ │ 对账闭环批次 totals 校验笔数 金额→ 渠道账单交叉核对 │ │ → 差异挂起表 → 重试 / 人工 │ │ final safety net — runs after every batch, every day │ └─────────────────────────────────────────────────────────────────────┘图 1 · 结算流水线。并发段为第三段对账层横跨全链路是独立于主流程的最后防线。这个形状里有两个容易忽视的设计决策「取数」和「执行」解耦取数用游标或水印向前推进执行侧通过有界队列承接两层速度不一致时由队列和反压机制吸收而不是让取数线程无脑灌。对账不是主流程的一部分——它是旁路的、每天必跑的独立闭环主流程的任何 bug 都应该能被对账兜住并暴露出来。下面各节逐段展开。03 · 并发底座线程池与虚拟线程3.1 批处理线程池core max队列必须有意地小在线服务怕雪崩倾向「小核心 大队列 弹性扩容」批任务要的是稳定的并行度——并发度直接对应数据库连接占用和分片隔离来回伸缩反而破坏可预期性。所以核心参数取core max跳过 JDK 线程池「队列满了才扩容」的陷阱// SettlementPools.java · 批处理线程池 public final class SettlementPools { ​ /** 结算批池与在线交易池物理隔离DB 密集型任务 */ public static ThreadPoolExecutor newBatchPool(int workers, int queueSize) { AtomicInteger seq new AtomicInteger(1); return new ThreadPoolExecutor( workers, workers, // core max并行度恒定 0L, TimeUnit.MILLISECONDS, new ArrayBlockingQueue(queueSize), // 有界队列 反压的容量 r - { Thread t new Thread(r, settle-worker- seq.getAndIncrement()); t.setDaemon(false); // 日终批要能阻止 JVM 直接退出 return t; }, new ThreadPoolExecutor.CallerRunsPolicy() // 队列满 → 提交线程自己跑 ); } }三个细节值得展开CallerRunsPolicy是批处理里最优雅的反压队列满时提交任务的线程也就是取数线程自己执行一个分片取数速度自动降到消费速度不需要任何额外的限流代码。线程要命名。出款卡住时jstack里一眼看到settle-worker-7卡在哪个账户和看到一堆pool-3-thread-7排查效率完全不同。优雅停机是资金安全要求。日终任务收到 kill 信号时要shutdown()awaitTermination()等在途分片收尾否则一批「回写了一半」的订单要靠补偿任务猜状态。3.2 线程数怎么定先找到真正的瓶颈结算任务绝大多数是 IO 密集——线程大部分时间在等数据库、等下游接口。经验公式CPU 核数 × (1 等待/计算比)只在上游无限速时成立。现实中批处理的并行度上限通常不是 CPU而是数据库连接池和热点行的锁100 个线程抢 20 个连接80 个线程纯排队除了上下文切换什么都没换来。合理做法是让线程数 ≤ 连接池留给批任务的配额再压测验证。任务类型实际瓶颈并行度定法计算型汇率换算、费率引擎CPU≈ 核数或核数 − 1DB 密集状态回写、账务落账连接池 / 行锁≤ 批任务专用连接配额热点行数再封顶下游 RPC 密集渠道代付出款渠道限流对齐渠道 TPS 配额超发只会换来 4293.3 虚拟线程JDK 21 之后的新选项结算任务「大量阻塞、少量计算」的画像正好是虚拟线程的主场。每个分片一个虚拟线程阻塞在数据库 IO 上几乎无成本代码模型也从「回调编排」回到朴实的同步顺序写法// 虚拟线程版调度 · JDK 21 try (var executor Executors.newVirtualThreadPerTaskExecutor()) { ListFutureShardResult futures shards.stream() .map(shard - executor.submit(() - settleShard(shard))) .toList(); // ... 汇总 futures }但有一条必须写进团队规范虚拟线程解除了线程数限制没有解除数据库连接数限制。一万个虚拟线程抢一个 50 连接的池排队发生在getConnection()上吞吐纹丝不动还可能把await撑爆连接等待超时。正确姿势是虚拟线程 Semaphore或固定大小的ExecutorService包 DB 段双层限流并发度由信号量说了算虚拟线程只负责把「等连接」变得便宜。⚠️ Pitfall—parallelStream()不要出现在结算代码里。它共享全局的ForkJoinPool.commonPool()容量与 CPU 绑定一旦流内出现 DB 调用阻塞会饿死同一 JVM 里所有其他 parallelStream 使用方没有隔离、没有命名、没有反压异常还会被吞进流的聚合处。它的位置在内存中的纯计算不在资金链路。04 · 分片策略吞吐与正确性一起定分片不是「把 List 均匀切成 N 块」这么简单。分片键决定了哪些操作会被并发执行、哪些会被串行化——这既是性能问题更是正确性问题。结算场景的黄金法则同一结算账户或同一商户、同一出款主体的订单串行不同账户之间并发。原因有二其一同账户并行更新余额会产生行锁竞争甚至按不同顺序加锁导致死锁其二账户维度的约束比如「待结算余额不能扣成负数」只有在串行视角下才容易校验。按账户分组后锁天然无冲突UPDATE的 CAS 也有清晰的先后语义。// 按账户分片 提交线程池 MapString, ListSettlementOrder shards orders.stream() .collect(Collectors.groupingBy(SettlementOrder::getAccountNo)); ​ ListCompletableFutureShardResult futures shards.entrySet().stream() .map(e - CompletableFuture .supplyAsync(() - settleShard(e.getKey(), e.getValue()), pool) .orTimeout(15, TimeUnit.MINUTES) // 长尾分片保护 .exceptionally(ex - ShardResult.failed(e.getKey(), ex))) // 失败隔离 .toList(); ​ CompletableFuture.allOf(futures.toArray(CompletableFuture[]::new)).join(); BatchReport report BatchReport.of( futures.stream().map(CompletableFuture::join).toList());两个工程上必须处理的现实问题大象商户。头部一两个商户可能占掉订单量的两三成按账户分组后它们成了超长分片把整体时长拖成长尾。解法是把大商户降级到独立泳道单独的执行队列、更细的二级分片按子单时间窗甚至提前到日间增量结算别让它和几十万小商户挤同一条跑道。分片粒度 vs 单事务大小。一个分片太大单事务内回写行数膨胀undo log 和锁持有时间都失控太小则调度开销占比上升。实践中按「一个分片 一到两个批次回写每批几百行」来定粒度而不是一刀切的固定条数。05 · 幂等与状态机资金安全的最后防线前面说过重试是结算链路的常态。任务重跑、MQ 至少一次投递、机器重启后补偿——同一订单的结算逻辑被执行多次是必然。幂等要做三层且要明确哪层是「优化」、哪层是「保证」第一层数据库唯一约束保证。结算单号、出款流水号上的唯一索引是任何并发 bug 都穿透不过去的物理防线。所有上游校验都可能被绕过唯一索引不会。第二层状态机 CAS保证。每一步流转都带着「期望前置状态」去更新影响行数为 0 就说明状态已被别人推进直接幂等退出-- 状态机 CAS · 只有预期状态才能被改写 UPDATE settlement_order SET status SUCCESS, finish_time NOW(), version version 1 WHERE order_no ? AND status PROCESSING; -- ← 期望前置状态// Java 侧updated 0 → 已被处理成功/回滚/被其他消费者抢走→ 幂等返回 int updated jdbc.update(sql, orderNo); if (updated 0) { log.info([settle] idempotent skip, orderNo{}, orderNo); return; }第三层分布式去重优化。RedisSETNX之类的前置过滤可以省掉大量无效落库但它只能当加速器不能当正确性依据——Redis 故障、驱逐、主从切换的瞬间去重就会漏。架构评审时看到「幂等靠 Redis」基本可以直接判不合格。状态机的价值不止幂等。它把「一条结算单走到哪了」变成可查询、可对账的事实卡在PROCESSING超过阈值就是告警补偿任务按状态精确续跑而不是靠日志猜。状态集合要小而严格INIT → PROCESSING → SUCCESS / FAILED / CANCELED并且只进不退——任何「从 SUCCESS 改回 PROCESSING」的需求出现时正确动作是生成一张反向新单而不是改历史。Note— 乐观锁的version字段和状态机 CAS 可以合一WHERE status ?本身就是一次 CAS。区分场景——单行状态流转用状态条件就够跨表的多行一致性交给事务边界见第 07 节。06 · 数据访问层流式读、批量写、热点合并6.1 读别用 OFFSET 翻页LIMIT 100000, 1000会让数据库先扫过十万行再丢弃。批处理取数只有两种正确姿势ID 水印分页WHERE id ? ORDER BY id LIMIT 1000每批记下末尾 ID。扫过的每一段都有据可查任务中断后从水印续跑天然支持断点恢复。游标流式读JDBCfetchSizeMySQL 需useCursorFetchtrue或fetchSizeInteger.MIN_VALUE的流式模式MyBatis 的CursorT是它的封装。适合一次全表扫的场景注意流式期间连接被独占读完必须关闭。// ID 水印推进 · 可断点续跑的取数 long cursor watermarkDao.get(jobName); // 从上次断点开始 ListSettlementOrder page; while (!(page mapper.selectAfter(batchDay, cursor, 1000)).isEmpty()) { dispatcher.submit(page); // 内部按账户分片后进线程池 cursor page.get(page.size() - 1).getId(); watermarkDao.advance(jobName, cursor); // 水印与批内回写解耦见 §07 }6.2 写批量 JDBC且确认批量真的生效逐行UPDATE是批处理最常见的隐形杀手——每行一次网络往返 一次 binlog 事件。改用jdbc.batchUpdate()后还要确认两件事MySQL 连接串上要有rewriteBatchedStatementstrue否则驱动还是逐条发批量大小落在几百到一千的区间再大收益递减且单包超限。这两点都验证过批量写才有数量级的差异。6.3 热点行合并把 N 次扣减变 1 次有些分片键选法按渠道、按日期会让同一账户的几十笔订单散落在多个并发分片里对同一行settle_account反复UPDATE ... pending_out pending_out - ?行锁排队并行度形同虚设。解法是内存合并、单次落账// 同账户金额合并后单次扣减 // 分片内先按账户聚合N 笔 → 1 笔净额 ListObject[] batch shard.stream() .collect(Collectors.groupingBy( SettlementOrder::getAccountNo, Collectors.summingLong(SettlementOrder::getAmount))) .entrySet().stream() .map(e - new Object[]{e.getValue(), e.getKey()}) .toList(); ​ jdbc.batchUpdate( UPDATE settle_account SET pending_out pending_out - ? WHERE account_no ?, batch); // 每账户 1 次行锁替代原来的 N 次更进一步如果热点账户平台主账户、大商户的频次已经高到单行更新都成为瓶颈正确的架构演进是账务改追加式流水表只做INSERT无锁竞争余额改成异步物化的汇总视图。结算系统的成熟度很大程度上就体现在从「直接改余额」走到「流水 物化余额」的那一步。07 · 事务边界与一致性闭环7.1 反模式一个事务包住整个批「整个批次一个事务要么全成要么全败」听起来很稳实际上是灾难百万行的 undo log、分钟级的锁持有、从库延迟雪崩、回滚时间比执行时间还长。批处理的一致性单位是分片或单笔不是整批。整批的一致性靠「状态机 断点续跑 对账」在更高一层拼出来而不是靠一个大事务硬扛。7.2 模式短事务 Outbox每一笔或每个小批次的结算需要同时完成两个动作改结算单状态、记一条账务消息。跨服务时不能靠「先提交事务再发 MQ」——中间宕机就是丢单。标准解法是本地消息表Outbox账务消息和状态变更写在同一个本地事务里投递交给独立的中继任务。// Outbox · 状态回写与消息落库同事务 Transactional void settleOne(SettlementOrder o) { int updated jdbc.update( UPDATE settlement_order SET status PROCESSING, version version 1 WHERE order_no ? AND status INIT , o.getOrderNo()); if (updated 0) return; // 幂等退出 ​ jdbc.update( INSERT INTO outbox(topic, biz_no, payload, status) VALUES (LEDGER_ENTRY, ?, ?, NEW) , o.getOrderNo(), Jsons.write(o)); // 与上面的 UPDATE 原子提交 } ​ // 中继任务独立调度扫 outbox 中 status NEW 的行 → 发 MQ → 标记 SENT // 消费端账务服务按 biz_no 幂等 —— MQ 是 at-least-once重复投递必须被消化7.3 对账接受一切都会出错然后每天证明它没出错上面所有机制叠加后仍然存在出错窗口时钟、人、灰度期间的旧代码、下游的异步丢失。成熟结算系统的态度是假设主链路一定有漏用对账做最后闭环批内自校每批结束比对「取数侧笔数/金额合计」与「回写成功笔数/金额合计」不等即挂起告警。这是 O(1) 成本的强断言。渠道对账日终拉渠道账单文件与本地出款流水逐笔核对差异进「待查差异表」配自动重试掉单补拉与人工处理流程。总分核对分户账余额合计 总账科目余额防的是账务侧的逻辑性错误。关于分布式事务TCC / Saga结算内部链路通常用不上它们——本地事务 Outbox 对账已经覆盖了绝大多数一致性需求而且每一步都可人工接手。TCC 留给真正跨实时的资金操作如联机充值冻结不要为了「看起来严谨」把日终批做成 Saga 编排复杂度会吃掉你。08 · 实战组装分片调度器与进度账本把前面的构件组装起来核心是一个分片调度器对外只暴露「给我一批订单」内部完成分片、并发、失败隔离与进度汇总。配套一个进程内的进度账本让「跑到哪了、还要多久、哪片卡住」随时可答——日终批的可观测性不是锦上添花是运维的救命稻草。// ShardDispatcher · 分片调度骨架 public class ShardDispatcher { ​ private final ThreadPoolExecutor pool; private final ProgressTracker progress; // AtomicLongdone / failed / sum ​ public BatchReport submit(String batchDay, ListSettlementOrder orders) { MapString, ListSettlementOrder shards orders.stream() .collect(Collectors.groupingBy(SettlementOrder::getAccountNo)); ​ progress.beginBatch(batchDay, orders.size(), orders.stream().mapToLong(SettlementOrder::getAmount).sum()); ​ ListCompletableFutureShardResult futures shards.entrySet().stream() .map(e - CompletableFuture .supplyAsync(() - settleShard(e.getKey(), e.getValue()), pool) .orTimeout(15, TimeUnit.MINUTES) .exceptionally(ex - ShardResult.failed(e.getKey(), ex))) .toList(); ​ CompletableFuture.allOf(futures.toArray(CompletableFuture[]::new)).join(); return progress.finishBatch( futures.stream().map(CompletableFuture::join).toList()); } ​ private ShardResult settleShard(String accountNo, ListSettlementOrder shard) { for (SettlementOrder o : shard) { try { settleService.settleOne(o); // §05 CAS §07 Outbox短事务 } catch (DataAccessException e) { retryTemplate.execute(ctx - { // 分片内有限重试 settleService.settleOne(o); return null; }); progress.markFailed(o); // 重试仍失败 → 差异表不阻塞后续 } } return ShardResult.ok(accountNo, shard.size()); } }几个组装期的关键取舍失败隔离的粒度到分片不到整批。一个分片超时或异常记入结果由对账和补偿任务处理绝不让join()抛异常打断其他几百个正常分片。重试是有限的。分片内对单笔做 2–3 次退避重试仍失败就落差异表。无限重试的批任务等于把故障藏起来。进度定期外化。ProgressTracker除了内存计数应周期性把快照写入任务表含水印位置这样无论是查询进度还是断点续跑都不依赖进程存活。09 · 踩坑清单与上线检查表最后把散在各节的坑集中成一张表。左边是在真实系统里见过或亲手写过的问题右边是对应的正确姿势。踩坑正确姿势用parallelStream()跑结算独立线程池 / 虚拟线程命名、隔离、可反压Executors.newFixedThreadPool无界队列手工ThreadPoolExecutor有界队列 CallerRunsPolicy线程数拍脑袋设成 200对齐连接池配额与渠道限流压测定数幂等只靠 RedisSETNX唯一索引 状态机 CAS 是保证Redis 只做加速LIMIT offset, size深翻页取数ID 水印分页或游标流式读支持断点续跑批量写没开rewriteBatchedStatements开启并验证批量真生效batch size 数百级热点账户逐笔扣减行锁排队分片内合并净额一次落账高频户改流水 物化余额大事务包整批一致性单位 分片/单笔整批靠状态机 对账先提交事务再发 MQ宕机丢单Outbox 同事务落消息中继投递消费端幂等没有对账靠用户客诉发现少结批内 totals 断言 渠道账单核对 总分核对每日必跑kill -9停批状态悬空优雅停机等在途分片重启按状态 水印续跑大象商户拖成长尾无感知长尾检测最慢分片 / 中位耗时 大户独立泳道上线前的最后一问永远是这三个这个批重跑一遍结果会不会变幂等中断在任意一步能不能从断点续跑可恢复跑完之后有没有一个独立机制能证明它没错可对账三问皆有答案并发度只是调参问题三问有一个含糊先补防线再谈吞吐。

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询