JVM架构与类加载:订单服务从源码到运行

JVM 是 Java 程序的运行时基础。开发供应链系统时,我们通常关注订单、库存、仓储、物流这些业务模块,但每一次接口调用最终都会落到 JVM 的类加载、字节码执行、内存管理和垃圾回收上。掌握 JVM 架构的价值不是背概念,而是能解释线上现象:为什么服务启动慢、为什么类冲突、为什么热部署失败、为什么一次配置改动导致初始化异常。

整体流程图

JVM 类加载与执行流程

需要掌握的核心技能点

JVM 架构至少要掌握下面这些内容:

  1. JDK、JRE、JVM 的关系:JDK 提供编译、诊断和运行工具,JRE 提供运行环境,JVM 负责执行字节码。
  2. .java 到 .class 的流程:源码经过 javac 编译成字节码,JVM 再解释执行或通过 JIT 编译成本地机器码。
  3. 类加载机制:加载、验证、准备、解析、初始化。
  4. 类加载器体系:Bootstrap ClassLoader、Platform ClassLoader、Application ClassLoader、自定义 ClassLoader。
  5. 双亲委派模型:优先让父加载器加载类,避免核心类被篡改,也减少重复加载。
  6. 执行引擎:解释器、JIT 编译器、热点代码探测、方法内联、逃逸分析。
  7. 本地方法接口:Java 代码通过 JNI 调用操作系统或本地库能力。

供应链业务场景

假设有一个订单履约服务 OrderFulfillmentService,它负责把电商订单转成仓库出库单,并根据仓库能力选择履约策略:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
public interface FulfillmentPolicy {
String chooseWarehouse(String skuCode, String province);
}

public class DefaultFulfillmentPolicy implements FulfillmentPolicy {
static {
System.out.println("load warehouse routing rules");
}

@Override
public String chooseWarehouse(String skuCode, String province) {
if ("GD".equals(province)) {
return "SOUTH_WAREHOUSE";
}
return "CENTRAL_WAREHOUSE";
}
}

当业务代码第一次主动使用 DefaultFulfillmentPolicy 时,JVM 才会触发类初始化:

1
2
3
4
5
6
7
8
public class OrderFulfillmentService {
private final FulfillmentPolicy policy = new DefaultFulfillmentPolicy();

public String createOutboundOrder(String orderNo, String skuCode, String province) {
String warehouseCode = policy.chooseWarehouse(skuCode, province);
return orderNo + " -> " + warehouseCode;
}
}

这段代码背后发生了几件事:

  1. Application ClassLoader 找到 DefaultFulfillmentPolicy.class。
  2. JVM 验证字节码是否合法,避免非法访问栈、越界跳转等问题。
  3. JVM 为静态字段分配默认值,这一步叫准备。
  4. JVM 把符号引用解析成直接引用,例如方法、字段、类的真实内存入口。
  5. 执行 <clinit>,也就是静态代码块和静态变量赋值。

所以,供应链项目里如果把数据库连接、远程配置、缓存预热写进静态代码块,服务启动或首次访问时就可能出现类初始化失败。更合理的方式是把这些动作放到 Spring Bean 生命周期里,并做好失败重试和降级。

类加载冲突的典型问题

供应链系统经常集成 WMS、TMS、ERP、OMS 等外部系统,依赖包很容易变复杂。例如一个老的 WMS SDK 依赖 jackson 2.9,订单服务本身依赖 jackson 2.15,如果版本冲突,可能出现:

1
java.lang.NoSuchMethodError: com.fasterxml.jackson.databind.ObjectMapper.readerForUpdating

这不是编译期问题,而是运行期加载到的类版本和编译期预期不一致。排查时要关注:

1
2
3
mvn dependency:tree
javap -classpath target/classes com.example.OrderFulfillmentService
java -verbose:class -jar order-service.jar

-verbose:class 可以看到类从哪个 jar 加载。定位到冲突后,常见处理方式包括统一依赖版本、排除传递依赖、隔离插件 ClassLoader,或者把老 SDK 包装成独立适配服务。

双亲委派为什么重要

双亲委派的核心是:一个类加载器收到加载请求时,先委托父加载器尝试加载,父加载器加载不到时自己再加载。它的直接收益有两个:

  1. 安全:业务代码不能随便伪造 java.lang.String 这类核心类。
  2. 稳定:同一个基础类优先由上层加载器加载,减少重复定义带来的类型不一致。

但有些场景会打破或绕开双亲委派,例如 JDBC SPI、应用服务器隔离、插件化系统。供应链中如果要让不同仓库客户使用不同的计费插件,可以自定义 ClassLoader 隔离插件依赖,但要明确边界:插件可以依赖公共接口,不应该反向依赖主应用内部实现。

实战建议

JVM 架构和类加载要能落到排查能力上:

  1. 看到 ClassNotFoundException,优先判断运行时 classpath 是否缺包。
  2. 看到 NoClassDefFoundError,除了缺包,还要判断类初始化是否失败过。
  3. 看到 NoSuchMethodError 或 NoSuchFieldError,优先怀疑依赖版本冲突。
  4. 看到启动阶段变慢,检查静态初始化、Spring 扫描范围、反射和代理生成。
  5. 设计插件系统时,先定义稳定接口,再考虑 ClassLoader 隔离。

JVM 不是脱离业务的底层知识。对于订单履约、库存同步、仓储调度这类高并发服务,类加载决定了服务启动和依赖边界,执行引擎决定了热点路径性能,诊断工具决定了线上问题能否快速收敛。

Kafka生产者与消费者原理:分区、确认与消费进度

为什么要理解消息链路

Kafka 这种组件,刚开始学的时候很容易记成几个名词:Producer、Consumer、Topic、Partition、Broker、Offset。真正用到项目里以后才会发现,Kafka 的难点不在于“会不会发消息”,而在于吞吐量、顺序性、可靠性和消费进度之间的取舍。

在供应链、订单、库存、ERP 同步这类系统里,Kafka 很常见。比如订单创建后通知库存系统预占库存,采购单状态变化后通知财务系统生成应付记录,仓库出库后通知报表系统刷新数据。这些业务都有一个共同点:消息量可能很大,但业务又不能随便丢。

所以理解 Kafka,不能只看 API,要看完整链路:生产者怎么把消息写进去,Broker 怎么存,消费者怎么拉取,offset 怎么提交,失败时怎么恢复。

Kafka生产者消费者工作流程

生产、存储与消费的核心机制

Kafka 的核心设计可以概括成三句话。

第一,生产者不是一条一条傻发,而是会把消息按 topic 和 partition 组织起来,经过序列化、分区选择、批量缓存后再发送给 Broker。

第二,Broker 不是把消息存在一个普通队列里,而是把消息追加到分区日志中。分区是 Kafka 并行能力的基础,日志追加是它高吞吐的基础。

第三,消费者不是等 Broker 推消息,而是主动 poll 拉取。消费者属于某个 consumer group,同一个 group 里一个分区同一时刻通常只会分给一个消费者处理,处理进度通过 offset 记录。

Kafka 的吞吐量来自批量、顺序写、分区并行和零拷贝等机制;可靠性来自副本、ack、幂等、事务和 offset 提交策略。项目里真正要做的,是根据业务重要性选择合适参数,而不是一味追求“最快”。

供应链事件处理示例

先看生产者。生产者发送一条消息时,大致会经历这些步骤:

  1. 业务代码构造消息,比如订单号、业务类型、变更时间。
  2. 序列化,把对象转成字节。
  3. 选择分区,如果指定 key,通常会根据 key hash 到固定 partition。
  4. 放入本地缓冲区,按 batch 组织。
  5. Sender 线程把 batch 发给对应 Broker。
  6. Broker 追加日志并根据 ack 策略返回结果。

如果要增大生产吞吐量,常见方向有几个。

batch.size 可以调大,让更多消息合并成一个批次。批量越充分,网络请求次数越少。

linger.ms 可以适当增加,让生产者多等几毫秒凑批次。它会牺牲一点延迟,换更高吞吐。

compression.type 可以使用 lz4、snappy 或 zstd,减少网络传输和磁盘占用。消息体较大时效果明显。

分区数要足够。一个 topic 如果只有一个 partition,再多消费者也无法在同一个 consumer group 内并行消费这个 topic。

生产者可以设置 acks=all、开启幂等 enable.idempotence=true,并合理配置重试。这样可以在 Broker 短暂失败时自动恢复,同时避免重试造成重复写入。

再看消费者。消费者是通过 poll 拉取消息,处理后提交 offset。这里最关键的是 offset 提交时机。

如果先提交 offset 再处理业务,消费者宕机后这批消息可能永远不会再处理,容易丢消息。

如果先处理业务再提交 offset,宕机后可能重复消费,但至少消息不会丢。大多数核心业务更接受“重复但可幂等”,而不是“直接丢”。

所以在订单、库存这类场景里,我更倾向于手动提交 offset:

1
2
3
4
poll 消息
执行业务处理
写入业务库或幂等表
处理成功后 commit offset

为了防止重复消费,业务侧要做幂等。比如用消息唯一 ID 建一张消费记录表,或者让订单状态流转本身具备幂等判断:已经处理过的状态不再重复扣减库存。

可靠性与顺序性风险

第一个坑,是只调大分区数,不看消费者处理能力。分区数增加能提高并行度,但也会带来更多文件句柄、更多 leader 选举和更复杂的再均衡。分区不是越多越好。

第二个坑,是为了吞吐把 acks 调成 0。这样生产者发出去就不管了,速度很快,但 Broker 是否收到并不确定。日志、埋点可以这么考虑,订单状态、库存变更就不应该这么随意。

第三个坑,是自动提交 offset。自动提交很方便,但它提交的是消费进度,不是业务成功。只要业务处理和 offset 提交之间没有绑定,就要接受消息丢失或重复的风险。

第四个坑,是忽略 rebalance。消费者数量变化、心跳超时、poll 时间过长,都可能触发再均衡。处理单条消息耗时很长时,要注意 max.poll.interval.ms 和批量大小,避免消费者被踢出 group。

第五个坑,是把 Kafka 当数据库。Kafka 适合做日志流和消息流,不适合承担复杂查询。业务状态仍然要落到数据库、缓存或搜索系统中。

总结

Kafka 的生产者负责高效、可靠地把消息写入分区日志;消费者负责按分区拉取消息、处理业务并提交 offset。吞吐量靠批量、压缩、分区和顺序写;可靠性靠副本、ack、幂等、事务和手动提交。

在真实项目里,我会按业务重要性分层:日志类消息可以优先吞吐,核心业务消息优先可靠;允许重复,但不能无声丢失。只要这条原则清楚,Kafka 参数就不会乱调。

MySQL锁机制详解:从行锁到Next-Key Lock

SQL 里面的锁,表面上看是数据库为了防止并发冲突做的限制,实际本质是数据库在并发读写之间做秩序管理。只要系统里存在多个事务同时读写同一批数据,就一定会遇到锁。锁设计得好,系统可以同时保证数据正确和较高吞吐;锁用得不好,轻则接口变慢,重则死锁、阻塞、库存扣错、订单状态错乱。

很多开发者第一次接触锁,是因为线上出现了 Lock wait timeout exceeded 或 Deadlock found。但真正理解锁,不能只背“共享锁、排他锁、行锁、表锁”这些名称,而要看清楚三个问题:锁保护的对象是什么、锁之间是否兼容、锁在事务什么时候加上和释放。

SQL 锁类型和事务执行流程

为什么数据库需要锁

数据库锁要解决的核心问题是并发一致性。

假设供应链系统里有一条库存记录:

1
2
3
sku_id = 1001
warehouse_id = 8
available_qty = 10

现在两个订单同时提交,每个订单都要锁定 8 件库存。如果两个事务都先读到 available_qty = 10,然后都认为库存足够,再分别把库存更新为 2,就会出现超卖。数据库必须让这两个更新按某种顺序执行,或者让其中一个事务发现条件已经不满足。

锁就是这个顺序的基础。它告诉数据库:某个事务正在读或写某个资源,其他事务能不能同时读、能不能同时写、要不要等待。

共享锁和排他锁

共享锁也叫 S 锁,英文是 Shared Lock。它表示当前事务要读取数据,并且希望读取期间数据不要被别人修改。

排他锁也叫 X 锁,英文是 Exclusive Lock。它表示当前事务要修改数据,其他事务不能同时修改,也通常不能再加共享锁读取同一行。

它们的兼容关系可以简单理解为:

已有锁 新申请共享锁 新申请排他锁
共享锁 兼容 不兼容
排他锁 不兼容 不兼容

共享锁之间兼容,是因为多个事务同时读同一行数据不会破坏数据。排他锁和任何锁都不兼容,是因为写操作必须独占资源。

在 MySQL InnoDB 中,可以通过下面的方式显式加锁:

1
2
3
SELECT * FROM inventory
WHERE sku_id = 1001
LOCK IN SHARE MODE;

或者:

1
2
3
SELECT * FROM inventory
WHERE sku_id = 1001
FOR UPDATE;

LOCK IN SHARE MODE 倾向于加共享锁,FOR UPDATE 会对命中的记录加排他锁。业务里更常见的是更新语句自动加排他锁:

1
2
3
4
UPDATE inventory
SET available_qty = available_qty - 8
WHERE sku_id = 1001
AND available_qty >= 8;

这条 SQL 在更新命中的记录时,会自动对相关记录加排他锁。

表锁和行锁

表锁保护的是整张表。一个事务锁住表以后,其他事务对这张表的读写可能都会受到影响。表锁粒度大,管理简单,但并发能力弱。

行锁保护的是某一行或某个索引范围。行锁粒度小,并发能力强,但实现复杂,也更容易出现死锁。

举个例子,ERP 系统里有一张 order 表。如果系统对整张订单表加表锁,那么一个用户修改订单时,其他用户可能连其他订单也无法修改。这对高并发系统非常不友好。

如果使用行锁,一个用户修改订单 A001,另一个用户修改订单 A002,两者互不影响。只有两个事务同时修改同一张订单时,才需要等待。

InnoDB 支持行级锁,但有一个非常重要的前提:行锁通常是加在索引上的。如果查询条件没有命中索引,数据库可能扫描大量记录,锁范围也会扩大,甚至表现得像锁了很多行。

例如:

1
2
3
UPDATE order_info
SET status = 'CLOSED'
WHERE order_no = 'SO20220917001';

如果 order_no 有唯一索引,InnoDB 可以精准锁住这一行。如果 order_no 没有索引,数据库需要扫描全表判断哪些行满足条件,锁冲突风险就会明显增加。

意向锁

意向锁是很多人容易忽略的一类锁。它不是直接锁某一行业务数据,而是用来协调表锁和行锁。

假设事务 A 已经对订单表中的某一行加了排他行锁。此时事务 B 想对整张订单表加表级排他锁。数据库必须知道表里是否已经有行锁,否则就要扫描整张表逐行检查,成本很高。

意向锁就是一个提示:某个事务打算在这张表里的某些行上加锁。

常见意向锁有:

  • IS,意向共享锁,表示事务准备在某些行上加共享锁。
  • IX,意向排他锁,表示事务准备在某些行上加排他锁。

当事务要给某行加共享锁时,会先在表上加 IS 锁。当事务要给某行加排他锁时,会先在表上加 IX 锁。

意向锁的价值是让表级锁判断冲突更快。它像是在表门口挂了一个牌子:里面已经有人在某些行上操作,整表加锁前先看看是否兼容。

记录锁

记录锁是 InnoDB 最容易理解的行锁,它锁住的是索引上的一条记录。

例如库存表有唯一索引:

1
UNIQUE KEY uk_sku_warehouse (sku_id, warehouse_id)

执行:

1
2
3
4
5
SELECT *
FROM inventory
WHERE sku_id = 1001
AND warehouse_id = 8
FOR UPDATE;

如果命中一条记录,InnoDB 会对这条索引记录加排他记录锁。其他事务再想更新同一条库存记录,就必须等待当前事务提交或回滚。

记录锁适合解决“同一行业务数据不能被并发修改”的问题,例如订单状态流转、库存数量变更、账户余额扣减。

间隙锁

间隙锁锁住的不是已经存在的记录,而是索引记录之间的空隙。它的目的主要是防止幻读。

假设库存预警表里已有预警阈值:

1
threshold: 10, 20, 50

一个事务执行范围查询:

1
2
3
4
SELECT *
FROM stock_warning_rule
WHERE threshold BETWEEN 10 AND 50
FOR UPDATE;

如果数据库只锁住 10、20、50 这几条已经存在的记录,另一个事务仍然可以插入 threshold = 30 的新记录。第一个事务再次查询时,就会发现多了一条之前不存在的数据,这就是幻读。

间隙锁会锁住索引范围中的空隙,让其他事务不能在这个范围里插入新记录。它牺牲了一部分并发能力,换取范围查询的一致性。

需要注意,间隙锁依赖索引范围。如果 SQL 没有合适索引,锁范围可能比预期大很多。

Next-Key Lock

Next-Key Lock 可以理解为记录锁加间隙锁。它既锁住已经存在的索引记录,也锁住记录前后的范围。

在 InnoDB 的可重复读隔离级别下,范围查询加锁时经常会使用 Next-Key Lock 来防止幻读。

例如:

1
2
3
4
5
SELECT *
FROM purchase_order
WHERE supplier_id = 88
AND amount BETWEEN 10000 AND 50000
FOR UPDATE;

如果 supplier_id, amount 上有联合索引,InnoDB 会锁定这个索引范围内的记录和间隙。其他事务不能随便插入符合这个范围的新采购单。

Next-Key Lock 的好处是一致性强,坏处是容易让范围更新、范围查询变得更容易互相阻塞。因此业务 SQL 要尽量让范围条件走合适索引,避免锁住过大的范围。

乐观锁

乐观锁不是数据库引擎内部固定的一种锁,而是一种并发控制思想。它假设冲突不常发生,所以不提前阻塞别人,而是在提交更新时检查数据有没有被别人改过。

最常见实现是版本号字段:

1
2
3
4
5
6
7
UPDATE inventory
SET available_qty = available_qty - 8,
version = version + 1
WHERE sku_id = 1001
AND warehouse_id = 8
AND available_qty >= 8
AND version = 12;

如果更新影响行数为 1,说明版本没变,扣减成功。如果影响行数为 0,说明数据已经被别人改过,当前事务需要重试或提示失败。

乐观锁适合读多写少、冲突概率低的场景。例如商品资料编辑、客户档案修改、配置项更新。它的优点是不会长时间阻塞,缺点是冲突发生时需要业务处理重试和失败提示。

悲观锁

悲观锁也是一种思想。它假设冲突很可能发生,所以在操作前先把数据锁住。

典型写法是:

1
2
3
4
5
SELECT *
FROM inventory
WHERE sku_id = 1001
AND warehouse_id = 8
FOR UPDATE;

当前事务拿到锁以后,再计算库存、写入订单、更新库存。其他事务想修改同一条库存记录,就必须等待。

悲观锁适合写冲突高、数据不能错的场景,例如库存扣减、余额扣减、核心单据状态流转。缺点也明显:事务时间越长,等待越多,吞吐越低。

所以悲观锁一定要控制事务范围。不要在持有锁期间调用外部接口、发送 MQ、请求第三方系统,也不要在事务里做复杂计算。

元数据锁

元数据锁也叫 MDL,Metadata Lock。它保护的是表结构,而不是具体业务行。

当一个事务正在查询或修改某张表时,数据库会持有这张表的元数据锁,防止另一个会话同时修改表结构。否则就可能出现一个事务读表读到一半,另一个事务把字段删了。

常见问题是:一个长事务一直不提交,导致 ALTER TABLE 等 DDL 操作被阻塞;DDL 又反过来阻塞后续普通查询,最后形成一串等待。

例如:

1
2
3
BEGIN;
SELECT * FROM order_info WHERE id = 1;
-- 长时间不提交

此时另一个会话执行:

1
ALTER TABLE order_info ADD COLUMN source_type varchar(32);

DDL 可能会等待前面的事务释放 MDL。后续新的查询又可能排在 DDL 后面,导致业务接口突然大面积变慢。

线上做 DDL 时,必须关注长事务和元数据锁等待。大表变更最好使用在线 DDL 工具或低峰期执行。

自增锁

自增锁用于处理自增主键分配。多个事务同时插入数据时,数据库要保证自增 ID 不重复。

InnoDB 对自增锁做过很多优化,不同配置下表现不同。简单理解,普通插入通常可以较快分配自增值;批量插入、INSERT ... SELECT 这类语句可能持有自增相关锁更久。

业务上不建议依赖自增 ID 的连续性。事务回滚、插入失败、并发插入都可能造成 ID 跳号。自增 ID 的目标是唯一和大体递增,不是绝对连续。

死锁是怎么发生的

死锁是两个或多个事务互相等待对方释放锁。

例如供应链系统里同时更新订单和库存:

事务 A:

1
2
1. 锁订单 O1001
2. 再锁库存 S1001

事务 B:

1
2
1. 锁库存 S1001
2. 再锁订单 O1001

事务 A 拿到了订单锁,等待库存锁。事务 B 拿到了库存锁,等待订单锁。双方都不释放,就形成死锁。

数据库通常会检测死锁,并主动回滚其中一个事务。业务系统看到的就是死锁异常。

减少死锁的关键方法是:

  • 多表更新保持固定顺序。
  • 批量更新时按主键排序。
  • 事务尽量短。
  • 查询条件命中索引,减少锁范围。
  • 避免在事务中做远程调用。
  • 捕获死锁异常,对幂等操作做有限重试。

怎么排查锁等待

线上出现锁等待时,不要只盯着慢 SQL。要看谁在等锁,谁持有锁,事务已经执行了多久。

MySQL 里常用的排查方向包括:

1
SHOW PROCESSLIST;

查看当前连接状态。

1
SHOW ENGINE INNODB STATUS;

查看最近死锁、锁等待、事务信息。

在 MySQL 8 中,也可以通过 performance_schema 和 sys 库查看锁等待关系。

排查时重点看:

  • 哪个事务持有锁。
  • 持锁事务执行了多久。
  • 等待的 SQL 是什么。
  • 是否存在未提交长事务。
  • SQL 是否走了索引。
  • 是否有 DDL 和普通业务 SQL 互相阻塞。

实际开发里的用锁建议

第一,能用一条原子 SQL 解决的,不要拆成先查再改。

库存扣减推荐写成:

1
2
3
4
5
6
UPDATE inventory
SET available_qty = available_qty - 8,
locked_qty = locked_qty + 8
WHERE sku_id = 1001
AND warehouse_id = 8
AND available_qty >= 8;

然后根据影响行数判断是否成功。

第二,加锁查询必须有合适索引。

FOR UPDATE 不是魔法。如果条件没有索引,锁范围会扩大,性能和并发都会出问题。

第三,事务里只放必须保持一致的操作。

订单创建、库存锁定、订单主表写入可以放在事务里;短信通知、日志上报、消息推送应该放到事务提交后。

第四,统一更新顺序。

如果业务规定先锁订单,再锁库存,再锁财务单据,那么所有代码都要遵守这个顺序。不要一个接口先锁订单,另一个接口先锁库存。

第五,锁冲突高的热点数据要做业务拆分。

例如某个爆款 SKU 的库存行成为热点,可以按仓库、批次、库存桶拆分,降低单行竞争。

总结

SQL 里的锁不是孤立概念,而是一套并发控制体系。共享锁和排他锁决定读写是否兼容,表锁和行锁决定锁粒度,意向锁协调表级和行级锁,记录锁保护已有记录,间隙锁和 Next-Key Lock 保护索引范围,乐观锁和悲观锁是两种业务并发控制思想,元数据锁保护表结构,自增锁保证自增值分配。

真正写业务代码时,最重要的不是记住所有锁名,而是控制三个东西:索引、事务范围、更新顺序。索引决定锁得准不准,事务范围决定锁持有多久,更新顺序决定是否容易死锁。把这三点做好,绝大多数 SQL 锁问题都会少很多。

Java并发集合与排查:仓储任务并发安全怎么落地

Java 多线程最终要落到数据结构和排查能力上。供应链系统里,仓储任务调度、库存缓存、波次队列、接口指标都会用到并发集合和原子类。选择正确的数据结构,可以减少手写锁;具备排查能力,才能在线上出现卡顿时定位问题。

数据结构和排查流程

并发集合和线程排查流程

ConcurrentHashMap:本地任务状态缓存

仓储系统可能需要缓存正在处理的上架任务,避免同一个任务在当前实例内重复提交:

1
2
3
4
5
6
7
8
9
10
11
public class PutawayTaskRegistry {
private final ConcurrentHashMap<Long, PutawayTask> runningTasks = new ConcurrentHashMap<>();

public boolean register(PutawayTask task) {
return runningTasks.putIfAbsent(task.id(), task) == null;
}

public void unregister(Long taskId) {
runningTasks.remove(taskId);
}
}

putIfAbsent 是原子操作。多个线程同时注册同一个任务,只有一个会成功。

注意,这只保护当前 JVM。多实例部署时,仍然要靠数据库状态条件防重:

1
2
3
4
UPDATE scm_putaway_task
SET status = 'RUNNING'
WHERE id = #{taskId}
AND status = 'WAITING';

BlockingQueue:生产者消费者

仓库波次任务可以用阻塞队列实现生产者消费者:

1
2
3
4
5
6
7
8
9
10
11
public class WaveDispatchQueue {
private final BlockingQueue<WaveTask> queue = new ArrayBlockingQueue<>(1000);

public void submit(WaveTask task) throws InterruptedException {
queue.put(task);
}

public WaveTask take() throws InterruptedException {
return queue.take();
}
}

生产者生成波次任务,消费者线程处理任务:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
public void startWorkers() {
for (int i = 0; i < 8; i++) {
executor.execute(() -> {
while (!Thread.currentThread().isInterrupted()) {
try {
WaveTask task = queue.take();
waveService.dispatch(task);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
});
}
}

阻塞队列的好处是自然支持等待和唤醒,不需要自己写 wait/notify。

CopyOnWriteArrayList:读多写少配置

仓库作业规则通常读多写少,比如上架策略规则:

1
2
3
4
5
6
7
8
9
10
11
12
public class PutawayRuleHolder {
private final CopyOnWriteArrayList<PutawayRule> rules = new CopyOnWriteArrayList<>();

public List<PutawayRule> currentRules() {
return rules;
}

public void reload(List<PutawayRule> latest) {
rules.clear();
rules.addAll(latest);
}
}

CopyOnWriteArrayList 写入时复制数组,读取时不加锁。它适合规则、配置、监听器列表,不适合高频写入的数据。

CAS 和 AtomicReference

本地策略快照可以用 AtomicReference 原子替换:

1
2
3
4
5
6
7
8
9
10
11
12
public class RoutingPolicyCache {
private final AtomicReference<RoutingPolicy> policyRef =
new AtomicReference<>(RoutingPolicy.defaultPolicy());

public RoutingPolicy current() {
return policyRef.get();
}

public void refresh(RoutingPolicy latest) {
policyRef.set(latest);
}
}

如果要防止旧版本覆盖新版本:

1
2
3
4
5
6
7
8
9
10
11
public boolean refreshIfNewer(RoutingPolicy latest) {
while (true) {
RoutingPolicy current = policyRef.get();
if (latest.version() <= current.version()) {
return false;
}
if (policyRef.compareAndSet(current, latest)) {
return true;
}
}
}

这适合单 JVM 内的配置引用。业务库存余额不应该用本地 CAS 保存,因为库存需要跨实例一致、事务回滚和审计。

死锁排查

Java 死锁常见于多个线程以不同顺序获取锁。例如两个仓储任务同时锁两个库位:

1
2
3
4
5
synchronized (binA) {
synchronized (binB) {
transfer();
}
}

另一个线程反过来:

1
2
3
4
5
synchronized (binB) {
synchronized (binA) {
transfer();
}
}

解决办法是固定加锁顺序:

1
2
3
4
5
6
7
8
9
10
public void transfer(Bin from, Bin to) {
Bin first = from.id() < to.id() ? from : to;
Bin second = from.id() < to.id() ? to : from;

synchronized (first) {
synchronized (second) {
doTransfer(from, to);
}
}
}

线上排查可以用:

1
jstack <pid>

或者 Arthas:

1
thread -b

重点看线程是否大量 BLOCKED,以及堆栈里等待的是哪把锁。

线上排查顺序

仓储任务卡顿时,可以按这个顺序排查:

  1. 看线程池:活跃线程数、队列长度、拒绝次数是否异常。
  2. 看线程状态:jstack 或 Arthas thread 查看是否大量 BLOCKED、WAITING。
  3. 看锁对象:定位堆栈里等待的是 Java 对象锁、数据库锁,还是队列阻塞。
  4. 看业务状态:仓储任务是否卡在同一个波次、库位、SKU 或外部 WMS 调用。
  5. 看数据一致性:确认任务状态机是否允许重复执行、失败重试和人工恢复。

排查时不要只看到 ConcurrentHashMap 就认为线程安全已经解决。并发集合只能保证集合自身操作安全,不能自动保证“任务状态 + 数据库记录 + 外部调用”的业务一致性。

小结

并发集合和原子类能减少手写锁,让并发代码更清晰。供应链系统里,ConcurrentHashMap 适合本地任务注册,BlockingQueue 适合生产者消费者,CopyOnWriteArrayList 适合读多写少规则,AtomicReference 适合配置快照替换。核心业务数据仍然要用数据库状态机和事务保护。线上卡顿时,要能用 jstack、Arthas 定位线程阻塞和死锁。

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

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

工具选择流程

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 适合异步任务编排。供应链系统使用这些工具时,必须区分查询类流程和交易类流程:查询可以并行和降级,交易必须保证状态一致、异常可追踪。