Kafka集群:分区、副本和供应链事件流

Kafka 是高吞吐分布式消息系统,适合处理订单事件、库存变更、物流轨迹、仓储作业状态等业务流。学习 Kafka 集群时,不应该只关注配置项,更要理解 Topic、Partition、Replica、Consumer Group 如何共同保证吞吐、扩展性和可用性。

整体流程

Kafka 集群和供应链事件流

Kafka 集群核心概念

Kafka 的基本组件包括:

  1. Broker:Kafka 服务节点。
  2. Topic:消息主题,例如 order-created
  3. Partition:Topic 的分区,用于并行写入和消费。
  4. Replica:分区副本,用于容灾。
  5. Leader:当前负责读写的分区副本。
  6. Follower:跟随 Leader 同步数据。
  7. Consumer Group:消费者组,同一组内多个消费者共同消费分区。

分区是 Kafka 扩展吞吐的关键。一个 Topic 有多个分区,生产者可以并行写入,消费者组内的消费者也可以并行消费。

供应链事件流例子

订单创建成功后,订单服务可以发送事件:

1
2
3
4
5
6
7
{
"eventId": "evt-10001",
"eventType": "ORDER_CREATED",
"orderNo": "SO202606280001",
"warehouseId": 8,
"occurredAt": "2026-06-28T10:20:00"
}

库存服务消费事件后预占库存,仓储服务消费事件后生成出库任务,财务服务消费事件后创建应收记录:

1
2
3
4
Order Service -> order-created topic
Inventory Service -> reserve stock
WMS Service -> create outbound task
Finance Service -> create receivable

这种事件驱动方式能降低服务之间的同步耦合,但也要求下游具备幂等和补偿能力。

分区设计

如果同一订单的事件必须按顺序处理,应使用订单号作为 key:

1
2
3
ProducerRecord<String, OrderCreatedEvent> record =
new ProducerRecord<>("order-created", event.orderNo(), event);
kafkaTemplate.send(record);

相同 key 的消息会进入同一分区,从而在该分区内保持顺序。但这不代表整个 Topic 全局有序,Kafka 只保证单分区内有序。

分区数变更也会影响 key 到分区的映射。扩容期间,同一个 key 的新旧消息可能位于不同分区,因此对严格顺序敏感的业务需要设计版本化路由、停写迁移或在消费端按业务版本校验,不能把“使用相同 key”理解成永久顺序保证。

分区数量要结合吞吐和消费者并行度规划。分区太少会限制并发,分区太多会增加元数据、文件句柄和 rebalance 成本。

副本和可靠性

生产环境 Topic 应该配置多个副本:

1
2
3
4
5
kafka-topics.sh --create \
--topic order-created \
--partitions 12 \
--replication-factor 3 \
--bootstrap-server kafka-1:9092

生产者可靠性配置:

1
2
3
acks=all
enable.idempotence=true
retries=2147483647

acks=all 表示 Leader 等待 ISR 中满足要求的副本确认后才认为写入成功。Broker 端还应设置合适的 min.insync.replicas,例如副本数为 3 时设为 2,避免只剩一个同步副本时仍接受关键写入。开启幂等生产者后通常应保留充分的重试机会;具体默认值和兼容约束以当前 Kafka 客户端版本为准。生产者幂等只能约束单个生产者会话内的重试,业务消费者仍然要做幂等。

消费幂等

消息系统通常只能帮你降低丢消息和乱序风险,不能自动解决业务重复处理。消费者必须幂等。

1
2
3
4
5
6
7
8
9
10
@Transactional
public void handle(OrderCreatedEvent event) {
// event_id 上有唯一索引;插入成功才获得本次处理权。
if (!eventRepository.tryInsert(event.eventId())) {
return;
}

inventoryService.reserve(event.orderNo(), event.warehouseId());
eventRepository.markSucceeded(event.eventId());
}

“先查询是否存在、再插入”的写法存在并发竞态。更稳妥的方式是依赖数据库唯一约束原子地抢占事件处理权,并把事件记录、库存条件更新和状态变更放在同一事务中。若库存服务和事件表不在同一数据库,则需要 outbox、状态机或补偿任务处理跨系统一致性。

常见问题

第一,消费者堆积。要看是消费者处理慢、分区不够、下游数据库慢,还是单条消息失败反复重试。

第二,rebalance 频繁。消费者实例不稳定、处理时间过长或超时配置不合理都可能导致 rebalance。

第三,只关注 Kafka 不关注业务补偿。消息发送成功不代表业务最终成功,库存预占失败、仓储创建失败都需要补偿流程。

第四,把 Kafka 当 RPC。Kafka 适合异步事件,不适合需要立即返回强一致结果的流程。

小结

Kafka 集群的关键是分区提升吞吐、副本提升可用性、消费者组提升并行消费能力。供应链系统中,订单、库存、仓储、物流事件非常适合用 Kafka 解耦,但必须配套幂等、补偿、监控和告警,才能真正成为可靠的事件流基础设施。

快速排序:分治思想和供应链订单优先级排序

快速排序是一种典型的分治算法。它选择一个基准值,把数组拆分成“小于基准”和“大于基准”的两部分,再递归处理左右区间。平均时间复杂度是 O(n log n),在理解排序算法、递归和分治思想时非常重要。

整体流程

快速排序流程

基础思想

快速排序包含三步:

  1. 选择基准值 pivot。
  2. 分区:把小于 pivot 的元素放左边,大于 pivot 的元素放右边。
  3. 递归排序左右区间。

示例代码:

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
public void quickSort(int[] arr, int left, int right) {
if (left >= right) {
return;
}

int pivotIndex = partition(arr, left, right);
quickSort(arr, left, pivotIndex - 1);
quickSort(arr, pivotIndex + 1, right);
}

private int partition(int[] arr, int left, int right) {
int pivot = arr[right];
int storeIndex = left;

for (int i = left; i < right; i++) {
if (arr[i] <= pivot) {
swap(arr, storeIndex, i);
storeIndex++;
}
}
swap(arr, storeIndex, right);
return storeIndex;
}

private void swap(int[] arr, int i, int j) {
int temp = arr[i];
arr[i] = arr[j];
arr[j] = temp;
}

这段代码使用最后一个元素作为 pivot,便于理解,但生产实现通常会使用更稳健的 pivot 选择策略。

Pivot 选择和递归边界

快速排序的性能很依赖 pivot。如果每次 pivot 都把数组切得很不均匀,递归深度会接近 n,最坏时间复杂度会退化到 O(n^2)。例如数据已经按订单创建时间升序排列,而实现又总是选择最后一个元素作为 pivot,就容易出现这种问题。

一个常见优化是随机选择 pivot:

1
2
3
4
5
private int randomizedPartition(int[] arr, int left, int right) {
int pivotIndex = ThreadLocalRandom.current().nextInt(left, right + 1);
swap(arr, pivotIndex, right);
return partition(arr, left, right);
}

还有一种做法是“三数取中”,从左端、中间、右端选一个更接近中位数的值作为 pivot。它不能完全避免最坏情况,但能降低有序数据导致退化的概率。

递归边界也要写清楚:当 left >= right 时直接返回。否则空区间、单元素区间会继续递归,轻则浪费调用栈,重则造成栈溢出。

供应链业务例子

假设订单履约系统需要把待处理订单按综合优先级排序:

1
2
3
4
客户等级
订单时效
是否缺货风险
创建时间

如果订单量很大,排序应该优先下推到数据库或搜索引擎,让索引和分页机制发挥作用:

1
2
3
4
5
SELECT *
FROM scm_order
WHERE status = 'WAIT_FULFILL'
ORDER BY customer_level DESC, promise_time ASC, created_at ASC
LIMIT 100;

如果是在内存中对少量候选订单做二次排序,可以使用 Java 内置排序:

1
2
3
4
orders.sort(Comparator
.comparing(OrderCandidate::getCustomerLevel).reversed()
.thenComparing(OrderCandidate::getPromiseTime)
.thenComparing(OrderCandidate::getCreatedAt));

学习快速排序的意义,不是让业务代码手写排序,而是理解“分而治之”的思路。比如订单分仓、波次拆分、库存重算,都可以把大任务拆成多个小区间并行处理。

分治思想在履约任务里的应用

分治思想在供应链系统中很常见。假设一天有几十万张待履约订单,如果直接由一个任务串行处理,会遇到处理时间长、失败重试成本高、单点压力大的问题。更合理的做法是先拆分:

1
按仓库拆分 -> 按承诺发货日期拆分 -> 按波次拆分 -> 每个分片独立处理

拆分以后,每个任务只处理一个相对小的订单集合。失败时可以只重试某个仓库、某个波次,而不是重跑整批数据。

在 Java 里可以用线程池并行处理这些分片:

1
2
3
for (OrderShard shard : shards) {
executor.submit(() -> fulfillmentService.process(shard));
}

这里和快速排序类似:先把大问题切小,再分别处理,最后合并结果。区别是排序算法合并的是有序区间,业务系统合并的是处理状态、异常结果和监控指标。

复杂度和风险

快速排序特点:

  1. 平均时间复杂度:O(n log n)
  2. 最坏时间复杂度:O(n^2),通常发生在 pivot 选择很差且数据分区极不均衡时。
  3. 空间复杂度:平均 O(log n),来自递归调用栈。
  4. 稳定性:通常不稳定。

为了降低最坏情况风险,可以使用随机 pivot、三数取中,或者在小区间切换为插入排序。JDK 内部排序实现已经处理了很多细节,业务代码一般不要重复造轮子。

小结

快速排序的核心价值是分治。供应链系统中,排序只是表象,更重要的是理解如何拆分大任务、减少无效扫描、控制单次处理的数据量。真正落地时,要优先使用数据库排序、索引、分页和 JDK 排序工具。

冒泡排序:从基础实现到供应链任务排序

冒泡排序是一种基础排序算法。它通过相邻元素两两比较,把较大的元素逐步交换到数组末尾。实际项目里很少直接使用冒泡排序处理大数据量,但它适合用来理解排序的基本思想、稳定性、时间复杂度和提前终止优化。

整体流程

冒泡排序流程

基础实现

假设有一组仓库拣货任务,需要按优先级从小到大排序:

1
int[] priorities = {8, 3, 2, 6, 7, 9};

冒泡排序实现:

1
2
3
4
5
6
7
8
9
10
11
public void bubbleSort(int[] arr) {
for (int i = 0; i < arr.length - 1; i++) {
for (int j = 0; j < arr.length - i - 1; j++) {
if (arr[j] > arr[j + 1]) {
int temp = arr[j];
arr[j] = arr[j + 1];
arr[j + 1] = temp;
}
}
}
}

每一轮结束后,当前未排序区间中的最大值会被交换到末尾。因此内层循环的边界是 arr.length - i - 1

提前终止优化

如果某一轮没有发生交换,说明数组已经有序,可以提前结束:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
public void bubbleSortWithBreak(int[] arr) {
for (int i = 0; i < arr.length - 1; i++) {
boolean swapped = false;
for (int j = 0; j < arr.length - i - 1; j++) {
if (arr[j] > arr[j + 1]) {
int temp = arr[j];
arr[j] = arr[j + 1];
arr[j + 1] = temp;
swapped = true;
}
}
if (!swapped) {
break;
}
}
}

这个优化对“基本有序”的数据有明显收益。例如仓库任务列表已经按创建时间大致排好,只是少量紧急任务插入,提前终止可以减少无效比较。

边界条件与测试

排序算法最容易出错的地方不是核心循环,而是边界条件。至少要覆盖空数组、单元素数组、已经有序、完全逆序、存在重复值这几类输入:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
public static void main(String[] args) {
int[][] cases = {
{},
{1},
{1, 2, 3},
{5, 4, 3, 2, 1},
{3, 1, 3, 2}
};

for (int[] item : cases) {
bubbleSortWithBreak(item);
System.out.println(Arrays.toString(item));
}
}

在供应链任务排序里,重复值尤其常见。比如两个拣货任务优先级相同,或者两个订单承诺发货时间一致。冒泡排序在只使用 > 交换时是稳定的,相同优先级任务不会改变原有顺序。这一点可以帮助我们理解“稳定排序”为什么对业务有意义:同优先级时保留创建顺序,能减少人工理解成本。

供应链业务例子

假设 WMS 中有少量待处理任务,需要在内存中按优先级排序:

1
2
3
4
5
public class PickTask {
private String taskNo;
private int priority;
private LocalDateTime createdAt;
}

如果任务数量只有十几个,冒泡排序可以作为教学示例理解排序过程。但在真实系统中,不建议自己手写冒泡排序处理任务列表,应该优先使用 JDK 提供的排序:

1
2
3
tasks.sort(Comparator
.comparingInt(PickTask::getPriority)
.thenComparing(PickTask::getCreatedAt));

业务代码更应该关注排序规则是否正确:优先级、创建时间、仓库、波次、客户等级等字段的优先顺序。

生产系统里怎么选择排序方式

生产系统通常不应该把大量数据查到 JVM 里再排序。供应链系统的数据规模很容易放大:订单可能是几十万级,库存流水可能是千万级,仓库任务在促销期间也会快速增长。

更常见的选择是:

  1. 数据库排序:适合按索引字段分页查询,例如 created_atprioritystatus
  2. 搜索引擎排序:适合复杂筛选和全文检索,例如订单号、客户名称、商品名称组合查询。
  3. JDK 内置排序:适合已经筛选出的小批量候选集,例如 100 条待分配任务做二次排序。
  4. 手写排序:主要用于学习、面试和解释算法过程,不建议作为业务主实现。

也就是说,冒泡排序的工程价值在于训练基本功,而不是替代成熟排序能力。真正写业务代码时,要先判断数据量、排序字段、分页方式和索引条件。

复杂度和稳定性

冒泡排序特点:

  1. 最好时间复杂度:O(n),前提是加了提前终止且数据已经有序。
  2. 平均和最坏时间复杂度:O(n^2)
  3. 空间复杂度:O(1)
  4. 稳定性:稳定。相等元素不会因为 > 比较而交换相对顺序。

稳定性在业务里有意义。例如两个拣货任务优先级相同,稳定排序可以保留原来的创建顺序。

小结

冒泡排序适合理解排序思想,不适合大规模生产数据。供应链系统里的订单、库存、任务、流水通常数量很大,生产代码应该使用数据库排序、索引排序或 JDK 内置排序。学习冒泡排序的重点,是理解相邻比较、边界缩小、提前终止和稳定性。

递归:从基础概念到供应链BOM和组织树处理

递归是一种函数直接或间接调用自身的编程方式。它适合处理天然具有层级结构的问题,例如树、目录、菜单、组织架构、BOM 物料清单。递归代码通常简洁,但如果缺少终止条件或层级过深,容易造成死循环或栈溢出。

整体流程

递归处理流程

递归的两个条件

写递归必须明确两个条件:

  1. 终止条件:什么时候停止继续调用。
  2. 递推关系:当前问题如何拆成更小的同类问题。

最经典的阶乘示例:

1
2
3
4
5
6
public int factorial(int n) {
if (n <= 1) {
return 1;
}
return n * factorial(n - 1);
}

n <= 1 是终止条件,factorial(n - 1) 是递推关系。

供应链例子:BOM 物料树

制造和供应链系统中经常有 BOM。一个成品由多个半成品组成,半成品又由原材料组成,这就是典型树结构。

1
2
3
4
5
public class BomNode {
private String materialCode;
private int quantity;
private List<BomNode> children;
}

计算一个成品需要多少原材料,可以递归遍历:

1
2
3
4
5
6
7
8
9
10
11
12
public void collectMaterial(BomNode node, int multiplier, Map<String, Integer> result) {
int requiredQty = node.getQuantity() * multiplier;

if (node.getChildren() == null || node.getChildren().isEmpty()) {
result.merge(node.getMaterialCode(), requiredQty, Integer::sum);
return;
}

for (BomNode child : node.getChildren()) {
collectMaterial(child, requiredQty, result);
}
}

这段代码表达的是:如果当前节点已经是叶子物料,就汇总数量;否则继续处理子节点。

防止循环依赖

业务数据不一定天然正确。BOM 里如果出现 A 包含 B,B 又包含 A,递归就会无限执行。生产代码必须加访问路径检查:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
public void collectMaterial(BomNode node,
int multiplier,
Map<String, Integer> result,
Set<String> path) {
if (!path.add(node.getMaterialCode())) {
throw new BizException("BOM存在循环依赖: " + node.getMaterialCode());
}

int requiredQty = node.getQuantity() * multiplier;
if (node.getChildren() == null || node.getChildren().isEmpty()) {
result.merge(node.getMaterialCode(), requiredQty, Integer::sum);
} else {
for (BomNode child : node.getChildren()) {
collectMaterial(child, requiredQty, result, path);
}
}

path.remove(node.getMaterialCode());
}

path 保存当前递归路径,不是全局已访问集合。这样既能识别当前链路上的循环,又不会误伤其他分支复用同一个物料的正常情况。

递归的风险

递归常见风险有三类:

  1. 没有终止条件,导致无限递归。
  2. 层级过深,导致 StackOverflowError
  3. 重复计算,导致性能指数级下降。

如果层级可能非常深,可以改成显式栈。显式栈不仅要保存节点,还要保存从父节点累积下来的数量:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
record BomFrame(BomNode node, long multiplier) {}

Map<String, Long> result = new HashMap<>();
Deque<BomFrame> stack = new ArrayDeque<>();
stack.push(new BomFrame(root, 1L));
while (!stack.isEmpty()) {
BomFrame frame = stack.pop();
BomNode current = frame.node();
long requiredQty = Math.multiplyExact(current.getQuantity(), frame.multiplier());

if (current.getChildren() == null || current.getChildren().isEmpty()) {
result.merge(current.getMaterialCode(), requiredQty, Math::addExact);
continue;
}

for (BomNode child : current.getChildren()) {
stack.push(new BomFrame(child, requiredQty));
}
}

示例使用 Math.multiplyExactMath.addExact 主动暴露数量溢出。真实 BOM 还应校验单位换算、替代料、生效日期和损耗率;涉及小数数量时应使用 BigDecimal,不能用整数示例直接承载生产计算。

如果存在大量重复子问题,可以使用缓存或动态规划。

小结

递归适合表达层级结构。供应链系统里的 BOM、仓库库区库位树、组织权限树、菜单树都适合用递归建模。生产代码中必须补上终止条件、循环依赖检查、最大深度限制和异常处理,否则递归很容易从优雅实现变成线上风险。

MyBatis与Spring整合:事务边界和库存扣减

MyBatis 和 Spring 整合后,Mapper 的创建、事务管理、数据源配置都交给 Spring 容器统一管理。业务代码不需要手动创建 SqlSession,也不应该在 Service 里手动提交事务。对于供应链系统来说,这一点非常关键,因为库存扣减、订单状态流转、库存流水写入必须处在清晰的事务边界内。

整体流程

MyBatis 与 Spring 事务流程

Spring 整合 MyBatis 的核心

整合后通常有三层:

  1. Controller:接收请求,做参数校验和返回结果。
  2. Service:编排业务流程,定义事务边界。
  3. Mapper:执行 SQL,不写业务判断。

Spring Boot 项目里常见依赖是:

1
2
3
4
5
<dependency>
<groupId>org.mybatis.spring.boot</groupId>
<artifactId>mybatis-spring-boot-starter</artifactId>
<version>3.0.3</version>
</dependency>

配置示例:

1
2
3
4
5
6
7
8
9
spring:
datasource:
url: jdbc:mysql://localhost:3306/scm?useUnicode=true&characterEncoding=utf8
username: scm_user
password: your_password
mybatis:
mapper-locations: classpath*:mapper/**/*.xml
configuration:
map-underscore-to-camel-case: true

Mapper 扫描:

1
2
3
4
5
6
7
@SpringBootApplication
@MapperScan("com.example.scm.**.mapper")
public class ScmApplication {
public static void main(String[] args) {
SpringApplication.run(ScmApplication.class, args);
}
}

供应链例子:库存预占事务

订单创建时,需要把可用库存转成锁定库存,并写入库存流水:

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
@Service
public class StockReserveService {
private final InventoryMapper inventoryMapper;
private final InventoryLogMapper inventoryLogMapper;

public StockReserveService(InventoryMapper inventoryMapper,
InventoryLogMapper inventoryLogMapper) {
this.inventoryMapper = inventoryMapper;
this.inventoryLogMapper = inventoryLogMapper;
}

@Transactional(rollbackFor = Exception.class)
public void reserve(ReserveStockCommand command) {
int affectedRows = inventoryMapper.reserve(
command.warehouseId(),
command.skuId(),
command.quantity()
);
if (affectedRows != 1) {
throw new BizException("库存不足,无法预占");
}

inventoryLogMapper.insertReserveLog(command);
}
}

Mapper SQL:

1
2
3
4
5
6
7
8
9
<update id="reserve">
UPDATE scm_inventory
SET available_qty = available_qty - #{quantity},
locked_qty = locked_qty + #{quantity},
updated_at = NOW()
WHERE warehouse_id = #{warehouseId}
AND sku_id = #{skuId}
AND available_qty >= #{quantity}
</update>

这条 SQL 同时完成判断和扣减,避免先查库存再更新导致并发超卖。@Transactional 保证库存表更新和库存流水写入要么同时成功,要么同时回滚。

事务边界应该放在哪里

事务应该放在 Service 层,而不是 Controller 或 Mapper 层。

原因是:

  1. Controller 负责协议和参数,不应该持有数据库事务。
  2. Mapper 只执行单条或少量 SQL,不知道完整业务流程。
  3. Service 能表达业务一致性边界,例如“预占库存 + 写流水”必须同事务。

错误写法是把远程调用放在事务中:

1
2
3
4
5
6
@Transactional
public void reserveAndNotify(ReserveStockCommand command) {
inventoryMapper.reserve(command.warehouseId(), command.skuId(), command.quantity());
wmsClient.notifyReserve(command.orderNo());
inventoryLogMapper.insertReserveLog(command);
}

远程调用慢或失败时,数据库锁会被长时间持有,容易引发锁等待。更稳妥的方式是事务内只做本地数据变更,事务提交后再发事件:

1
2
3
4
5
6
7
8
9
@Transactional(rollbackFor = Exception.class)
public void reserve(ReserveStockCommand command) {
int affectedRows = inventoryMapper.reserve(command.warehouseId(), command.skuId(), command.quantity());
if (affectedRows != 1) {
throw new BizException("库存不足,无法预占");
}
inventoryLogMapper.insertReserveLog(command);
eventPublisher.publish(new StockReservedEvent(command.orderNo()));
}

如果要确保事件和数据库事务一致,可以使用事务消息、outbox 表或 @TransactionalEventListener

常见问题

第一,自调用导致事务不生效。Spring 声明式事务通常通过代理实现,同一个类内部方法直接调用可能绕过代理。

第二,异常被吞掉导致事务提交。如果捕获异常后不再抛出,Spring 可能认为方法正常结束。

第三,事务范围太大。事务里不要做 HTTP 调用、文件上传、大量循环查询。

第四,只依赖本地锁保护数据库数据。多实例部署后,本地 synchronized 只能锁住当前 JVM,不能保护数据库里的库存。

小结

MyBatis 与 Spring 整合的重点不是配置本身,而是把 SQL、Mapper、Service 事务边界组织清楚。供应链系统里最容易出问题的是库存、订单状态和流水一致性,写代码时要明确哪些操作必须同事务,哪些操作应该事务后异步处理。