Java并发工具类:采购对账和库存汇总如何并行协作

Java 并发工具类解决的是线程之间的协作问题。供应链系统里,很多流程不是简单加锁,而是多个任务并行执行后汇总结果,或者限制同时访问某个下游系统的并发量。常用工具包括 CountDownLatchSemaphoreCompletableFuture

工具选择流程

Java 并发工具选择流程

CountDownLatch:等待多个任务完成

采购对账时,需要同时加载三类数据:

  • 采购入库单。
  • 供应商发票。
  • 付款记录。

三类数据查询互不依赖,可以并行查询,最后汇总差异。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
public ReconcileResult reconcile(long supplierId, LocalDate month) throws InterruptedException {
CountDownLatch latch = new CountDownLatch(3);

AtomicReference<List<Receipt>> receiptsRef = new AtomicReference<>();
AtomicReference<List<Invoice>> invoicesRef = new AtomicReference<>();
AtomicReference<List<Payment>> paymentsRef = new AtomicReference<>();

executor.execute(() -> {
try {
receiptsRef.set(receiptService.query(supplierId, month));
} finally {
latch.countDown();
}
});

executor.execute(() -> {
try {
invoicesRef.set(invoiceService.query(supplierId, month));
} finally {
latch.countDown();
}
});

executor.execute(() -> {
try {
paymentsRef.set(paymentService.query(supplierId, month));
} finally {
latch.countDown();
}
});

latch.await();
return reconcileEngine.compare(receiptsRef.get(), invoicesRef.get(), paymentsRef.get());
}

这里的收益是缩短对账等待时间。原来三类数据串行查询,现在可以并行加载。

Semaphore:限制并发访问下游

物流轨迹同步可能调用承运商 API。承运商接口有 QPS 限制,不能因为系统里有大量线程就无限请求。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
public class CarrierTrackClient {
private final Semaphore semaphore = new Semaphore(20);

public TrackInfo queryTrack(String trackingNo) {
boolean acquired = false;
try {
acquired = semaphore.tryAcquire(1, 2, TimeUnit.SECONDS);
if (!acquired) {
throw new BizException("承运商接口繁忙,请稍后重试");
}
return doQuery(trackingNo);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new BizException("查询被中断");
} finally {
if (acquired) {
semaphore.release();
}
}
}
}

Semaphore 控制的是并发许可数。这里最多允许 20 个线程同时访问承运商接口,保护下游系统,也保护自己。

CompletableFuture:异步编排更清晰

CompletableFuture 适合多个异步任务组合。订单详情页需要同时展示订单基础信息、库存状态、物流轨迹、应收金额:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
public OrderDetail detail(long orderId) {
CompletableFuture<Order> orderFuture =
CompletableFuture.supplyAsync(() -> orderService.get(orderId), executor);

CompletableFuture<InventoryView> inventoryFuture =
CompletableFuture.supplyAsync(() -> inventoryService.viewByOrder(orderId), executor);

CompletableFuture<TrackInfo> trackFuture =
CompletableFuture.supplyAsync(() -> trackService.queryByOrder(orderId), executor);

CompletableFuture<Receivable> receivableFuture =
CompletableFuture.supplyAsync(() -> financeService.receivable(orderId), executor);

CompletableFuture.allOf(orderFuture, inventoryFuture, trackFuture, receivableFuture).join();

return new OrderDetail(
orderFuture.join(),
inventoryFuture.join(),
trackFuture.join(),
receivableFuture.join()
);
}

注意要传入业务线程池,不要默认依赖 ForkJoinPool.commonPool(),否则不同业务会混用同一个公共线程池,排查困难。

异常处理

异步任务一定要处理异常:

1
2
3
4
CompletableFuture<TrackInfo> trackFuture =
CompletableFuture
.supplyAsync(() -> trackService.queryByOrder(orderId), executor)
.exceptionally(e -> TrackInfo.empty("物流轨迹暂不可用"));

供应链系统的详情页通常允许部分信息降级,比如物流轨迹临时失败不应该导致整个订单详情不可用。但结算、扣库存这类核心流程不能随意吞异常。

选择建议

不同工具适合不同问题:

工具 解决的问题 供应链例子
CountDownLatch 等待多个任务全部完成 采购对账同时加载入库单、发票、付款记录
Semaphore 控制同时访问数量 限制承运商轨迹接口并发
CompletableFuture 异步编排和结果聚合 订单详情并行查询库存、物流、财务
BlockingQueue 生产者消费者 仓储波次任务排队处理

选择工具时先描述线程之间的关系:是等待、限流、结果组合,还是任务排队。关系清楚了,工具选择通常就清楚了。不要为了使用高级 API 而把简单流程写复杂。

小结

并发工具类的价值是让线程协作更清晰。CountDownLatch 适合等待多个并行任务完成;Semaphore 适合限制下游并发;CompletableFuture 适合异步任务编排。供应链系统使用这些工具时,必须区分查询类流程和交易类流程:查询可以并行和降级,交易必须保证状态一致、异常可追踪。

Java线程池实战:订单履约任务如何稳定提吞吐

线程池是 Java 多线程在业务系统中最常用的落地方式。它解决两个问题:复用线程,减少创建销毁成本;限制并发,避免请求无限堆积压垮系统。供应链系统里的订单履约、库存同步、物流轨迹拉取、报表生成,都应该用线程池管理并发。

履约线程池流程

订单履约线程池流程

不要无限创建线程

错误示例:

1
2
3
for (Long orderId : orderIds) {
new Thread(() -> fulfillmentService.fulfill(orderId)).start();
}

如果一次批处理有 5000 个订单,这段代码会创建 5000 个线程。线程本身占内存,调度也有成本,下游数据库、WMS、TMS 都可能被打爆。

应该使用 ThreadPoolExecutor

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
ThreadPoolExecutor executor = new ThreadPoolExecutor(
8,
16,
60,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(1000),
new ThreadFactory() {
private final AtomicInteger index = new AtomicInteger();

@Override
public Thread newThread(Runnable r) {
return new Thread(r, "fulfillment-worker-" + index.incrementAndGet());
}
},
new ThreadPoolExecutor.CallerRunsPolicy()
);

这个线程池最多 16 个工作线程,队列最多缓存 1000 个任务。超过能力后,CallerRunsPolicy 会让提交任务的线程自己执行,形成反压。

七个核心参数

ThreadPoolExecutor 的关键参数包括:

  • corePoolSize:核心线程数。
  • maximumPoolSize:最大线程数。
  • keepAliveTime:非核心线程空闲多久回收。
  • unit:时间单位。
  • workQueue:任务队列。
  • threadFactory:线程工厂。
  • handler:拒绝策略。

供应链系统里,线程数不是越大越好。订单履约通常包含数据库、库存服务、仓储服务、物流服务调用,属于 IO 密集型,可以适度提高线程数。但如果下游 TMS 每秒只允许 100 次请求,线程池再大也只会制造超时。

订单履约 demo

批量履约订单时,可以这样提交任务:

1
2
3
4
5
6
7
8
9
10
11
public void fulfillBatch(List<Long> orderIds) {
for (Long orderId : orderIds) {
executor.execute(() -> {
try {
fulfillmentService.fulfill(orderId);
} catch (Exception e) {
fulfillmentLogService.recordFailed(orderId, e.getMessage());
}
});
}
}

任务内部要保证幂等:

1
2
3
4
5
6
7
8
9
10
11
@Transactional
public void fulfill(Long orderId) {
int affected = orderMapper.markFulfilling(orderId);
if (affected != 1) {
return;
}

inventoryService.reserve(orderId);
warehouseTaskService.createPickTask(orderId);
orderMapper.markFulfilled(orderId);
}

对应状态更新:

1
2
3
4
UPDATE scm_sales_order
SET status = 'FULFILLING'
WHERE id = #{orderId}
AND status = 'WAIT_FULFILL';

线程池负责并发执行,数据库状态机负责避免重复履约。

队列选择

常见队列:

  • ArrayBlockingQueue:有界数组队列,容量固定,适合明确限流。
  • LinkedBlockingQueue:链表队列,可有界也可无界;业务中必须设置容量。
  • SynchronousQueue:不存储任务,直接移交线程,适合快速扩容线程的场景。
  • PriorityBlockingQueue:优先级队列,适合高优先级订单先处理。

供应链系统建议优先使用有界队列。无界队列在高峰期会隐藏问题,直到内存被耗尽。

监控线程池

线程池必须监控:

1
2
3
4
5
6
7
8
public ThreadPoolStats stats() {
return new ThreadPoolStats(
executor.getPoolSize(),
executor.getActiveCount(),
executor.getQueue().size(),
executor.getCompletedTaskCount()
);
}

核心指标:

  • 活跃线程数是否长期接近最大线程数。
  • 队列长度是否持续增长。
  • 拒绝任务是否出现。
  • 单任务耗时是否变长。

如果队列持续增长,不要只加线程。要确认瓶颈是数据库、外部接口、锁等待还是代码慢。

容量估算

线程池参数要从业务容量倒推,而不是凭经验随手写。以订单履约为例:

  1. 单个履约任务平均耗时 200ms。
  2. 其中 150ms 在等待数据库、库存服务、WMS 或 TMS。
  3. 目标吞吐是每秒 80 单。

粗略估算并发数:

1
2
3
并发数 = 目标吞吐 * 平均耗时
= 80 * 0.2
= 16

所以核心线程数可以从 16 附近开始压测,再结合 CPU、数据库连接池、WMS 限流和错误率调整。队列大小也要有业务含义,例如最多允许堆积 1000 单,超过后让调用方降级、延迟重试或进入消息队列,而不是无限排队。

线程池调优的目标是稳定吞吐,不是把线程数调到最大。高峰期如果队列持续上涨,说明系统消费能力不足,需要扩容、拆批、限流或优化单任务耗时。

小结

线程池的作用是稳定地控制并发,而不是无上限地提高并发。供应链系统中,线程池适合订单履约、库存同步、物流轨迹拉取等后台任务。设计时要使用有界队列、明确拒绝策略、处理任务异常、保证业务幂等,并持续监控线程池运行状态。

Java内存模型与锁:库存同步中的可见性和互斥

Java 多线程的核心问题可以归纳为三类:原子性、可见性、有序性。Java 内存模型,也就是 JMM,定义了线程之间如何看见彼此的写入,以及哪些同步动作能建立 happens-before 关系。供应链系统里,库存同步、价格缓存、仓库配置刷新都离不开这些基础。

可见性和互斥流程

JMM 可见性和锁边界流程

可见性问题

库存同步任务通常有一个停止标志:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
public class InventorySyncJob implements Runnable {
private boolean running = true;

public void stop() {
running = false;
}

@Override
public void run() {
while (running) {
syncOnce();
}
}
}

管理线程调用 stop() 后,工作线程不一定马上看到 running = false。这就是可见性问题。

可以使用 volatile

1
2
3
4
5
6
7
8
9
10
11
12
13
14
public class InventorySyncJob implements Runnable {
private volatile boolean running = true;

public void stop() {
running = false;
}

@Override
public void run() {
while (running) {
syncOnce();
}
}
}

volatile 保证写入对其他线程可见,并限制相关指令重排。它适合任务开关、配置引用、状态标志。

volatile 不保证复合操作原子性

下面这个计数器是错误的:

1
2
3
4
5
6
7
public class SyncCounter {
private volatile int successCount;

public void success() {
successCount++;
}
}

successCount++ 包含读取、加一、写回三个步骤。多个线程同时执行会丢失更新。正确方式是使用原子类:

1
2
3
4
5
6
7
public class SyncCounter {
private final AtomicInteger successCount = new AtomicInteger();

public void success() {
successCount.incrementAndGet();
}
}

如果是高并发指标统计,可以用 LongAdder

1
2
3
4
5
6
7
8
9
10
11
public class SyncCounter {
private final LongAdder successCount = new LongAdder();

public void success() {
successCount.increment();
}

public long value() {
return successCount.sum();
}
}

synchronized 解决互斥

如果多个线程要修改同一个本地库存快照 Map,需要互斥控制:

1
2
3
4
5
6
7
8
9
10
11
public class InventorySnapshotCache {
private final Map<String, Integer> cache = new HashMap<>();

public synchronized void put(String key, Integer qty) {
cache.put(key, qty);
}

public synchronized Integer get(String key) {
return cache.get(key);
}
}

synchronized 修饰实例方法时,锁对象是 this。同一个对象上的同步方法互斥。

更推荐使用私有锁对象,避免外部代码锁住当前实例:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
public class InventorySnapshotCache {
private final Object lock = new Object();
private final Map<String, Integer> cache = new HashMap<>();

public void put(String key, Integer qty) {
synchronized (lock) {
cache.put(key, qty);
}
}

public Integer get(String key) {
synchronized (lock) {
return cache.get(key);
}
}
}

锁的边界

供应链系统一般是多实例部署。synchronized 只能保护当前 JVM 内存,不能保护数据库里的库存余额。如果两个应用实例同时扣同一条库存,Java 本地锁没有任何作用。

库存扣减必须落到数据库条件更新:

1
2
3
4
5
6
UPDATE scm_inventory
SET available_qty = available_qty - #{qty},
locked_qty = locked_qty + #{qty}
WHERE warehouse_id = #{warehouseId}
AND sku_id = #{skuId}
AND available_qty >= #{qty};

Java 锁适合保护本地缓存、内存队列、对象状态;数据库事务和行锁负责保护最终业务数据。

happens-before 怎么理解

JMM 里最实用的判断工具是 happens-before。它不是描述时间先后,而是描述一个线程的写入是否对另一个线程可见。

常见规则包括:

  1. 对一个 volatile 变量的写,happens-before 后续对这个变量的读。
  2. 对一个锁的解锁,happens-before 后续对同一把锁的加锁。
  3. 线程 start() 之前的操作,happens-before 新线程中的操作。
  4. 一个线程中的操作,按程序顺序 happens-before 后续操作。

例如库存规则缓存刷新时,可以使用整体替换引用:

1
2
3
4
5
6
7
8
9
10
11
public class InventoryRuleHolder {
private volatile InventoryRuleSnapshot snapshot = InventoryRuleSnapshot.empty();

public void refresh() {
snapshot = inventoryRuleRepository.loadSnapshot();
}

public InventoryRuleSnapshot current() {
return snapshot;
}
}

这里 volatile 保证查询线程能看到新的快照引用。但快照对象本身应该是不可变的,否则引用可见不代表内部集合并发修改安全。

小结

JMM 是理解 Java 多线程的基础。volatile 解决可见性和有序性,不解决复合操作原子性;synchronized 解决单 JVM 内共享状态的互斥;数据库锁解决跨实例的业务数据一致性。供应链系统里要明确每把锁保护的对象,不能用本地锁替代数据库并发控制。

Java多线程基础:供应链系统为什么需要并发处理

Java 多线程不是为了把代码写复杂,而是为了在合适的业务场景里提高吞吐、降低等待时间、充分利用 CPU 和 IO 资源。供应链系统里有大量典型场景:订单创建后要校验库存、校验客户信用、计算运费、匹配促销、写操作日志;仓库出库时要生成波次、分配库位、通知 WMS、刷新库存快照。这些动作有些可以并行,有些必须串行,多线程的价值就是把这两类工作区分清楚。

并发决策流程

Java 多线程决策流程

线程和进程

进程是操作系统资源分配的基本单位,线程是 CPU 调度执行的基本单位。一个 Java 应用进程里可以有很多线程,比如 Tomcat 请求线程、业务线程池线程、GC 线程、定时任务线程。

在供应链系统中,一个订单接口请求通常由一个 Web 容器线程处理。如果接口里所有操作都串行执行,请求耗时会被每一步累加。如果某些步骤没有依赖关系,就可以并行执行。

线程生命周期

Java 线程常见状态包括:

  • NEW:线程对象已创建但未启动。
  • RUNNABLE:可运行,等待 CPU 调度或正在执行。
  • BLOCKED:等待进入 synchronized 临界区。
  • WAITING:无限期等待其他线程通知。
  • TIMED_WAITING:限时等待,比如 sleep、带超时的 wait
  • TERMINATED:线程执行结束。

排查供应链批处理卡顿时,线程状态很关键。大量 BLOCKED 说明锁竞争严重,大量 WAITING 可能是队列无任务或等待条件未满足,大量 RUNNABLE 但 CPU 很高可能是计算任务过重或死循环。

创建线程的方式

最原始的方式是直接创建 Thread

1
2
3
4
Thread thread = new Thread(() -> {
System.out.println("同步库存快照");
});
thread.start();

这种方式适合理解原理,不适合业务系统大量使用。真实项目应该使用线程池,避免频繁创建和销毁线程。

Runnable 没有返回值,适合只执行动作:

1
2
Runnable task = () -> inventorySyncService.syncWarehouse(8L);
new Thread(task).start();

Callable 有返回值,可以配合 Future

1
2
3
4
Callable<Integer> task = () -> inventoryQueryService.countAvailableSku(8L);
FutureTask<Integer> future = new FutureTask<>(task);
new Thread(future).start();
Integer count = future.get();

但业务里更常用线程池提交任务。

供应链订单校验的并行 demo

创建销售订单时,下面三个校验互不依赖:

  • 查询库存是否足够。
  • 查询客户信用额度。
  • 查询承运商是否覆盖收货地址。

串行写法:

1
2
3
4
5
6
public OrderCheckResult checkOrder(OrderCreateCommand command) {
InventoryCheckResult inventory = inventoryService.check(command.items());
CreditCheckResult credit = creditService.check(command.customerId(), command.amount());
CarrierCheckResult carrier = carrierService.check(command.address());
return OrderCheckResult.merge(inventory, credit, carrier);
}

如果三个调用分别耗时 80ms、60ms、50ms,总耗时至少接近 190ms。可以用线程池并行:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
public OrderCheckResult checkOrder(OrderCreateCommand command) throws Exception {
Future<InventoryCheckResult> inventoryFuture =
executor.submit(() -> inventoryService.check(command.items()));

Future<CreditCheckResult> creditFuture =
executor.submit(() -> creditService.check(command.customerId(), command.amount()));

Future<CarrierCheckResult> carrierFuture =
executor.submit(() -> carrierService.check(command.address()));

return OrderCheckResult.merge(
inventoryFuture.get(),
creditFuture.get(),
carrierFuture.get()
);
}

并行后,请求耗时接近最慢的那个任务,而不是三个任务之和。这里的收益来自 IO 等待重叠,而不是 CPU 被神奇地加速。

多线程不适合什么

多线程不是默认选项。以下场景要谨慎:

  • 任务之间有严格顺序,比如先扣库存再生成库存流水。
  • 数据共享复杂,容易产生竞态。
  • 单个任务非常小,线程切换成本超过收益。
  • 下游系统已经限流,并发只会把压力打爆。
  • 业务要求强事务一致性,不能拆成并发步骤。

比如订单扣库存不能简单把每个 SKU 丢给不同线程同时扣。如果一个订单内多个 SKU 需要同一个事务保证全部成功或全部失败,并行拆分反而会增加补偿复杂度。

落地判断标准

判断一个供应链流程是否适合多线程,可以按下面的顺序确认:

  1. 是否存在互不依赖的子任务。订单详情聚合、库存报表查询、物流轨迹批量拉取通常适合。
  2. 是否主要等待 IO。数据库、RPC、HTTP 调用占主导时,并发更容易产生收益。
  3. 是否会修改同一份业务事实。库存余额、订单状态、付款状态这类数据必须先保证一致性。
  4. 下游是否有容量限制。承运商接口、WMS 接口、数据库连接池都需要限流。
  5. 失败结果能否被处理。查询类任务可以降级,交易类任务必须明确回滚或补偿。

多线程优化前后要记录耗时、吞吐、错误率、队列长度和下游压力。只看单次请求变快是不够的,系统还要在高峰流量下保持稳定。

小结

Java 多线程首先是一种工程手段,不是语法技巧。供应链系统使用多线程,核心收益是提升吞吐、缩短 IO 密集型流程耗时、让后台批处理更高效。使用前必须判断任务是否独立、是否共享状态、是否需要事务一致性。能并行的校验和查询可以并行,涉及最终库存和单据状态的数据修改必须保持清晰的事务边界。

Redis入门:数据结构、缓存和供应链业务场景

Redis 是高性能内存数据存储,常用于缓存、分布式锁、计数器、排行榜、队列和限流。学习 Redis 不应该停留在安装命令,而要理解它的数据结构适合解决什么业务问题,以及缓存一致性、击穿、穿透、雪崩这些生产风险。

整体流程

Redis 数据结构和供应链场景

Redis 适合做什么

Redis 的核心优势是低延迟和丰富的数据结构。常用结构包括:

  1. String:缓存单个值、计数器、分布式锁 value。
  2. Hash:缓存对象字段,例如 SKU 基础信息。
  3. List:简单队列,但复杂消息场景更推荐消息队列。
  4. Set:去重集合,例如活动参与用户。
  5. ZSet:带分数排序,例如热销商品排行。
  6. Stream:Redis 5 引入的消息流,适合轻量事件流。

供应链系统里,Redis 常见用途是缓存库存展示值、仓库路由规则、承运商配置、SKU 基础资料,以及做接口限流和幂等控制。

缓存 SKU 信息

SKU 基础资料读多写少,适合缓存:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
public SkuInfo getSkuInfo(String skuCode) {
String key = "sku:info:" + skuCode;
SkuInfo cached = redisTemplate.opsForValue().get(key);
if (cached != null) {
return cached;
}

SkuInfo skuInfo = skuRepository.findByCode(skuCode);
if (skuInfo == null) {
redisTemplate.opsForValue().set(key, SkuInfo.empty(), Duration.ofMinutes(5));
return null;
}

redisTemplate.opsForValue().set(key, skuInfo, Duration.ofHours(6));
return skuInfo;
}

这里要注意两点:

  1. 空值也可以短时间缓存,防止缓存穿透。
  2. 缓存必须设置过期时间,避免长期脏数据。

缓存库存要谨慎

库存是强一致性敏感数据。展示库存可以缓存,但下单扣减不能只依赖 Redis 缓存值。

更稳妥的做法是:

1
2
3
下单预占 -> 数据库条件更新保证一致性
库存展示 -> Redis 缓存提升查询性能
库存变更 -> 发布事件刷新或删除缓存

示例:

1
2
3
4
5
public void refreshInventoryCache(long warehouseId, long skuId) {
InventoryRecord record = inventoryRepository.find(warehouseId, skuId);
String key = "inventory:view:" + warehouseId + ":" + skuId;
redisTemplate.opsForValue().set(key, record.toView(), Duration.ofMinutes(10));
}

缓存可以提升读性能,但数据库仍然是库存一致性的最终来源。

分布式锁基本原则

Redis 可以实现分布式锁,但要满足几个条件:

1
2
3
4
5
Boolean ok = redisTemplate.opsForValue().setIfAbsent(
"lock:order:" + orderNo,
requestId,
Duration.ofSeconds(10)
);

释放锁时必须校验 value,避免删除别人的锁。生产环境建议使用成熟客户端,例如 Redisson,并结合业务幂等设计,而不是只依赖锁。

缓存风险

缓存穿透:请求不存在的数据,绕过缓存打到数据库。处理方式是参数校验、空值缓存、布隆过滤器。

缓存击穿:热点 key 过期,大量请求同时访问数据库。处理方式是互斥重建、逻辑过期、热点 key 不短过期。

缓存雪崩:大量 key 同时过期。处理方式是过期时间加随机值、分批预热、多级缓存。

缓存脏读:数据库已变更,缓存未刷新。处理方式是先更新数据库,再删除缓存,必要时通过消息进行二次删除或异步刷新。

小结

Redis 入门要围绕数据结构和业务场景学习。供应链系统里,缓存能显著提升 SKU、仓库、路由、库存展示等读性能,但订单扣减、库存预占、财务结算这类核心一致性流程不能只依赖缓存。Redis 是加速器,不应该成为一致性的唯一防线。