当 finally 跑得比消费者慢:异步状态更新里的并发回退陷阱

异步状态更新里的并发回退陷阱

最近有一个数据对接的需求,背景是这样的:该项目用来收集其它系统的已有的监控数据,使用一些聚类算法或机器学习算法试试能不能对未发生问题的设备进行预测,预测其生命周期、健康度等。数据获取无非两种形式,要么 pull 要么 push,基于项目目前的结构,最终决定提供 rest api 由其它平台 push 给我们分钟级的监控数据。

设计这样的系统需要考虑什么?

  1. 性能:tomcat 默认 200 线程,如果对方大并发来推,一方面系统可能承受不住,另一方面对方的线程迟迟拿不到响应,可能会导致服务崩溃。
  2. 可靠性:对方的数据在到来之后,需要经过清洗、转换、拼 sql、写时序库,每一步都可能出问题。网络可能有波动,连接池也可能等待超时,系统要能应对这些非正常情况。

消息队列在这种场景下太合适了,tomcat 线程收到监控数据后推送到消息队列里,消费者根据自身情况慢慢消费,如果消费失败还可以重试,线程在消息发送到 mq 后就能返回,也不会拖垮对方平台。

当然事情要是这么简单也不会有这篇文章了~

问题1: 监控数据每一条都可能是几百个指标项,如果直接发数据到 mq 消息体太大,mq 的磁盘占用和吞吐都会有影响。

问题2: 收到对方的数据就发 mq,如果发 mq 失败了怎么办(网络波动、mq 宕机),该给对方返回什么信息?如何补推?

基于此,我完善了系统的设计:

  1. 收到请求体后直接存关系数据库,状态置为 PENDING;
  2. 提交任务到线程池,然后直接返回成功;任务内容为:
    1. 把监控数据在关系库的主键发送到 MQ;
    2. 根据发送结果,更新状态为 MQ_SEND 或 MQ_SEND_FAILED;
  3. 消费者消费消息时,根据 id 找到对应的监控数据,进行反序列化、转换、拼 sql、写时序库的逻辑;
  4. 根据写时序库的结果,更新状态为 PARSED 或 PARSED_FAILED;
  5. 在 xxl-job 中添加任务,定时处理发送 MQ 失败的记录;

整体逻辑就是这样,当然其中也包含了很多细节,如反序列化失败、写时序库失败等。项目使用了 RabbitMQ、SpringCloudStream。

服务器的内存资源还可以,但 cpu 资源不太行,而想要分钟级入库,当然需要并发执行了。众所周知 jdk21 的两大特性:虚拟线程与ZGC。虚拟线程作为用户级线程,数量可以轻松上千;ZGC 的垃圾回收毫秒级停顿,处理大堆比较好用。项目就充分的使用了虚拟线程池来提高处理能力,因为写时序库、发 MQ 都是 IO型任务。

  1. 接收数据,不做处理直接入库,提交发送 mq 的任务到虚拟线程池
1
2
3
4
5
6
7
8
9
10
11
public void acceptTsData(String body) {
long nextId = uidSegmentHolder.nextId();
SyncRecord record = new SyncRecord();
record.setData(body)
.setStatus(SyncStatus.PENDING.code())
.setRetryCount(0)
.setId(nextId)
.setDeleted(GeneralConsts.NOT_DELETED);
syncRecordDao.save(record);
generalExecutor.submit(() -> sendToMq(nextId));
}
  1. 发送消息到 mq 并更新数据库状态
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
private void sendToMq(Long id) {
String code = SyncStatus.MQ_SENT_FAILED.code();
String msg = "";
boolean success;
try {
success = tsDataProducer.produce(id);
if (success) {
code = SyncStatus.MQ_SENT.code();
msg = "发送消息到MQ成功";
} else {
log.error("发送解析消息到 MQ 失败 id={}", id);
msg = "发送消息到MQ失败,程序无异常";
}
} catch (Exception e) {
log.error("发送解析消息到 MQ 失败 id={}, err={}", id, e.getMessage(), e);
msg = "发送消息到MQ失败,有异常";
} finally {
syncTransactionService.updateStatus(id, code, msg);
}
}
  1. 消费者消费数据
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
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
public Consumer<Long> tsData() {
return messageId -> {
String body = syncRecordDao.readDataOnly(messageId);
if (StrUtil.isBlank(body)) {
log.error("id={}的记录请求体为空,无需解析", messageId);
syncTransactionService.updateStatus(messageId, SyncStatus.PARSED.code(), "请求体为空,无需解析");
return;
}

TsDataRequest dataRequest;
try {
dataRequest = JSON.parseObject(body, TsDataRequest.class);
} catch (Exception e) {
log.error("id={}的记录请求体转 json 失败!", messageId);
syncTransactionService.updateStatus(messageId, SyncStatus.PARSE_FAILED.code(), "转json失败");
return;
}

TsDataParseResult parseResult = tsDataParser.parse(dataRequest);
List<String> sqls = parseResult.sqls();
if (sqls.isEmpty()) {
log.error("id={}的记录请求体中无有效数据,没有可对应的指标项!", messageId);
syncTransactionService.updateStatus(messageId, SyncStatus.PARSED.code(), "请求体无有效数据");
return;
}

AtomicBoolean flag = new AtomicBoolean(true);
// 每条 SQL 一个虚拟线程
List<CompletableFuture<Void>> futures = new ArrayList<>(sqls.size());
for (String sql : sqls) {
futures.add(CompletableFuture.runAsync(() -> {
try {
tsJdbcTemplate.update(sql);
} catch (Throwable t) {
log.error("处理[messageId:{},sql:{}]时出现异常", messageId, sql, t);
flag.set(false);
}
}, tsSqlExecutor));
}
// 聚合等待
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).exceptionally(ex -> {
log.error("等待异步任务时发生异常", ex);
flag.set(false); // 如果 join 异常,也标记为失败
return null;
}).thenRun(() -> {
// 无论成功、失败还是异常,都会执行这里
try {
if (flag.get()) {
syncTransactionService.updateStatus(messageId, SyncStatus.PARSED.code(), "");
} else {
syncTransactionService.updateStatus(messageId, SyncStatus.PARSE_FAILED.code(),
"sql 插入时序库时失败!");
}
} catch (Exception e) {
log.error("最终更新状态失败, messageId:{}", messageId, e);
}
}).join();
};

定时任务的逻辑就不放了,无非是查询一段时间前 状态为发送MQ失败 且 重试次数<最大重试次数 的记录,一段时间前是避免和其它线程重复执行;最大重试次数是避免无限重试。

整体流程顺畅,关键流程异步并行处理,主线程迅速返回,消息队列解耦合。然后就出现了一个折磨我好几天的问题:

每隔一段时间就会出现状态是 MQ_SENT 即发送到 MQ 成功但一直没有被更新为消费成功或消费失败的记录,于是我在消费者的方法第一行打印了 收到id={}的消息,最后一行打印了 id={}的消息处理完成。

然而并没有什么用,我甚至看到了 id 为某个值的记录打印出了这两行日志但状态依然是 MQ_SENT 的,而且 MQ 中消息无积压,死信队列里也没数据。

1
2
2026-07-28 08:46:06.658  WARN [TID:N/A] 1 --- c.e.bhws.tsdb.mq.TsDataMqConsumer        : 收到了 id=214806855055587554 的消费消息
2026-07-28 08:46:06.753 WARN [TID:N/A] 1 --- c.e.bhws.tsdb.mq.TsDataMqConsumer : id=214806855055587554 的消费消息处理完成,处理结果:true,影响行数:1
  • 如果是事务回滚,那么应该所有的记录都是 MQ_SENT,为什么只有少量是 MQ_SENT?
  • 如果是消费者根本没有收到消息,那么为什么会有处理日志?
  • 如果是发送MQ失败,那么为什么状态是 MQ_SENT?

我最开始认为是发送MQ失败,由于异步等原因返回了成功,但随后打印出的消费者日志否定了这个可能。

随后我猜测是不是消费者配置的原因,消费失败了,但MQ的配置是拉取了消息就算消费成功,于是事务回滚+消息队列无积压,然而绝大多数数据又是正常的,也没办法解释。

最后发现事情的真相是:消费者消费完数据立刻更新状态为 PARSED,而代码段2的第18行,finnally 更新状态为 MQ_SENT 执行可能会在消费者更新之后,导致数据被重新更新为 MQ_SENT。

监控数据对接流程

既然根因是”发送 MQ 成功后,消费者比 finally 更快更新状态”,核心是需要给状态更新加上单调性约束——允许向前推进,禁止回退。一个简单的解决方案就是,在发送到MQ后的更新 sql 里,加上 where status not in (PARSED, PARSED_FAILED)。这也像是常说的分布式协同里的内容,当多个请求同时到达数据库时,如何决定到底谁先谁后?

顺便一提,虚拟线程的并发果然强大,该服务 pod 只分配了 1 个 cpu,数据处理起来几乎没有任何压力,但是注意如果用户线程数量不进行任何限制,有可能会等待数据库连接池超时。


当 finally 跑得比消费者慢:异步状态更新里的并发回退陷阱
https://zhuwenjie0716.github.io/2026/07/28/当-finally-跑得比消费者慢:异步状态更新里的并发回退陷阱/
作者
Wenjie Zhu
发布于
2026年7月28日
许可协议