0%

1. 引言与适用范围

本文是一份面向独立开发者和小型技术团队的上线方法指南,不是可以原样复制到任意环境的生产脚本。文中的命令、服务名、目录、镜像和网络策略都应结合实际项目核对后再执行。

适用场景是:已有一个能够在本地运行的 AI 辅助内容系统,希望以较低风险完成第一次真实生产闭环。目标不是一次性建成完整平台,而是尽快验证“生成内容是否有价值、人工审核是否可控、发布链路是否稳定”。

本文不适用于高并发、多租户、多区域部署、强合规行业或复杂多平台同步发布场景。这些场景需要更完整的权限、审计、容量、灾备和发布治理设计。

核心原则只有一句话:

先让产品在受控条件下产生真实价值,再用真实数据决定下一步投入。

2. 为什么受控 MVP 应优先于复杂发布治理

没有真实运行数据时,团队很难判断真正的瓶颈在哪里。提前投入多环境流水线、灰度发布、复杂编排和多平台分发,可能会让工程看起来更完整,却不能证明用户是否愿意使用生成结果。

MVP 阶段应优先建立四条边界:

  1. 生产可用优先于功能齐全:至少有一个稳定入口可以完成一次真实任务。
  2. 人工审核优先于自动发布:未经明确批准的 AI 内容不得进入发布队列。
  3. 单平台优先于多平台:先把一个发布目标做稳定,再扩展格式适配和认证集成。
  4. 可恢复优先于高自动化:数据库、内容仓库和备份必须有明确恢复路径。

这并不意味着忽略工程质量,而是把工程投入集中到最影响业务闭环的风险上:数据安全、生成质量、人工门禁、发布幂等性和失败可见性。

3. 执行前提与环境假设

一台 2 核 4 GB 云服务器可以作为单用户、低并发、AI 请求串行执行场景的起始规格,但它不是容量承诺。实际资源需求应根据 JVM 内存、数据库规模、AI 响应时间、并发任务数和内容构建过程测量。

生产环境至少需要:

  • 已安装并验证可用的 Docker Engine 与 Docker Compose v2;
  • 已固定并经过测试的 MySQL、Redis、JDK、Spring Boot 和 Caddy 版本;
  • 一个可以构建或拉取的应用镜像;
  • 可持久化的 MySQL、Redis 和应用数据卷;
  • 一个可写且可推送的 Hexo Git 仓库;
  • 一个用于公网健康检查或站点访问的域名;
  • 一个不直接暴露到公网的管理入口,例如 SSH 隧道、TAT 或私有网络。

不要只依据“版本足够新”判断兼容性。更可靠的方式是固定精确镜像版本或 digest,并用同一组版本完成部署、备份恢复和回滚演练。

4. 用 Docker Compose 建立最小可靠生产闭环

最小生产拓扑可以包含四类服务:

  • Spring Boot 应用;
  • MySQL;
  • Redis;
  • Caddy 或其他 HTTPS 入口。

应用、数据库和缓存可以位于同一个 Compose 项目网络中,但 MySQL 与 Redis 不应直接映射到公网。管理接口也不应因“方便调试”而公开暴露;可以只将应用端口绑定到 127.0.0.1,再通过 SSH 隧道访问后台。

敏感信息应与普通配置分离。数据库密码、Redis 密码和 AI API Key 不应写进镜像、代码仓库或可被普通用户读取的配置文件。优先使用 Compose secrets 或只读 secret 文件,让容器通过 /run/secrets/... 读取。宿主机上的源文件应限制权限,并明确哪些服务可以挂载对应 secret。

MySQL 客户端密码不要直接出现在命令行参数中。自动化脚本可使用受保护的 option file、login path 或挂载的 secret 文件,并确保日志不会输出凭证。不要把宿主机路径直接假设为容器内路径,必须显式挂载或通过 Compose secret 注入。

上线前至少验证:

  • 四个服务处于 running,有健康检查的服务处于 healthy
  • 应用只能通过预期入口访问;
  • MySQL、Redis 没有公网端口;
  • 管理页面只能经受控通道访问;
  • 应用数据卷、数据库卷和 Hexo 工作区均可读写;
  • 自动生成和自动发布默认关闭。

5. 只升级应用容器的发布策略

生产升级应使用已经构建并验证的不可变镜像,推荐固定到镜像 digest。不要在正式服务器上临时修改代码后直接构建一个无法追踪的镜像。

一次应用升级的基本顺序可以是:

  1. 记录当前应用镜像、容器 ID、数据库版本和关键数据校验值;
  2. 验证 Compose 配置;
  3. 拉取目标不可变镜像;
  4. 只重新创建应用服务;
  5. 验证应用健康、迁移版本和核心数据;
  6. 保留原镜像与回滚证据。

只更新应用服务时,可使用类似下面的命令,但服务名和环境文件必须以项目实际配置为准:

1
2
docker compose pull app
docker compose up -d --no-deps app

--no-deps 的作用是避免因应用更新而启动或重新创建依赖服务。不要无条件执行 docker compose down,更不能使用 down -v 或删除现有数据库卷。

数据库备份与迁移

涉及数据库变更前,应生成最终备份并验证其可恢复性。对以 InnoDB 为主的数据库,mysqldump --single-transaction --quick 可以在不长时间阻塞业务写入的情况下取得事务一致性快照;但执行期间的 DDL,以及 MyISAM、MEMORY 等非事务表,不具备同等一致性保证。

备份文件存在不等于备份有效。至少应在隔离环境中完成一次恢复演练:

  1. 创建新的空数据卷;
  2. 使用兼容的 MySQL 镜像启动临时实例;
  3. 有限时等待数据库就绪;
  4. 导入备份;
  5. 核对迁移版本、关键表行数和核心业务记录;
  6. 只清理本次演练创建的临时资源。

数据库迁移应由 Flyway 或 Liquibase 等工具管理。破坏性迁移必须单独评估,必要时安排停写窗口。不要假设切回旧应用镜像会自动还原数据库,也不要假设迁移工具一定具备可靠的反向迁移能力。

6. AI 生成、人工审核、单平台发布与验收

一个安全的真实闭环不应只写成抽象的 draft → approved → published,而应与系统实际状态和门禁一致。

以当前系统的实现为例,流程如下。

第一步:创建并执行生成任务

人工提交一个生成任务。任务状态从 PENDING 进入 RUNNING,最终可能是:

  • SUCCESS
  • FAILED
  • SKIPPED
  • CANCELLED

生成过程中应记录实际 Provider、模型、阶段、尝试次数、Token 使用量和失败代码,但不得记录 API Key 或完整敏感 Prompt。

第二步:生成待审核草稿

任务成功后创建草稿。此时草稿生命周期为 AI_GENERATED,独立审核字段 reviewStatusPENDING

这个阶段只说明正文已经生成并落盘,不代表已经具备发布条件,也不应自动产生发布任务。

第三步:生成平台变体

当前实现要求审核批准前存在完整的平台变体。因此即使首轮只计划发布 Hexo,也需要先生成系统要求的 CSDN、WECHAT 和 HEXO 三个平台版本。

这一步的目的不是同时发布三个平台,而是形成可供人工逐项检查的确定版本。平台正文发生修改时,应更新内容摘要并使旧审核结果失效,避免发布未经查看的新版本。

第四步:人工审阅

审核人员至少检查:

  • 事实和技术命令是否准确;
  • 是否包含内部路径、账号、密钥或真实服务器信息;
  • 标题、摘要和正文是否一致;
  • Markdown 与目标平台格式是否兼容;
  • AI 质量审查是否存在 blocking finding;
  • 各平台版本号是否与当前页面快照一致。

发现问题时,应先修改规范化正文,再重新生成平台变体并重新审阅。

第五步:明确批准

批准成功后,reviewStatus 变为 APPROVED,草稿生命周期进入 REVIEWED

当前版本能够记录审核状态和内容版本一致性,但不能据此声称已经持久化“审批人身份”和“审批时间”。如果业务需要责任追踪,应把审批主体、时间、来源 IP 或会话信息作为后续审计能力单独实现。

第六步:只创建 Hexo 发布任务

批准后只为 Hexo 创建发布任务,不调用会为全部平台建单的批量接口。

发布任务状态可能包括:

  • PENDING
  • RUNNING
  • SUCCESS
  • RETRYABLE_FAILED
  • NEED_MANUAL
  • FAILED
  • CANCELLED

任务领取必须具备幂等性,已经运行或成功的任务不能因重复请求再次发布。

第七步:执行 Hexo Git 发布

当前 Hexo 发布并不是调用内容平台 API。发布器会把 Hexo 版本写入仓库的文章目录,校验工作区和分支,创建 Git commit,并在启用推送时推送到指定远端分支。

因此验收证据应包括:

  • 新文章文件存在;
  • Git 工作区干净;
  • 本地与远端分支指向预期提交;
  • 发布任务为 SUCCESS
  • 记录了外部访问 URL 或可验证的文章路径。

第八步:公网验收

最后从公网访问实际文章,确认:

  • 页面返回成功;
  • 标题与正文正确;
  • 中文和代码块没有乱码;
  • 没有泄露管理入口或内部信息;
  • 失败重试和发布日志保持可查询。

在这一流程中,任何发布任务都必须出现在人工批准之后。

7. 风险与回滚方案

应用升级后无法启动

先判断问题是否只涉及应用镜像。如果数据库结构仍与旧版本兼容,可以恢复上一不可变镜像并只重新创建应用容器。

如果无法确认旧应用是否兼容升级后的数据库,不要直接让旧镜像连接当前数据卷。应保留现有卷不删除,用新的空卷和已验证备份恢复旧版本环境,再执行数据核对。

数据库迁移失败

停止继续写入,保留现场证据。根据已验证的恢复方案,从迁移前备份恢复数据库,并回退应用镜像。不要把“切回旧镜像”描述为数据库回滚。

Caddy 配置错误

先使用配置校验命令确认语法。回滚时恢复上一版配置并执行受控 reload;只有 reload 不可用或进程异常时才考虑重启入口服务。不要用入口配置问题作为重建整个 Compose 项目的理由。

Hexo 发布失败

不要盲目重复推送。先检查:

  • 工作区是否干净;
  • 当前分支是否正确;
  • 远端是否可访问;
  • 同一 slug 是否已经存在不同内容;
  • 本地提交是否已产生但远端推送失败。

根据失败类型选择自动重试、人工处理或终止,不得覆盖远端已有的不同文章。

服务器或凭证故障

数据库备份应保存到独立存储位置,并设置恢复演练和保留策略。凭证泄露时应立即轮换,而不是只删除泄露文件。管理后台应保持非公网暴露。

8. 上线后用真实数据决定迭代方向

上线后应持续记录最少的一组指标:

  • AI 生成成功率;
  • 各生成阶段耗时和 Token 消耗;
  • 人工审核通过率;
  • 从生成到批准的耗时;
  • 发布成功率;
  • 重试与人工处理比例;
  • 从批准到公网可见的耗时。

这些指标可以帮助团队区分不同问题:

  • 生成失败率高:先检查 Provider 稳定性、超时和重试;
  • blocking finding 多:先优化 Prompt、模型参数和质量规则;
  • 审核修改量大:先改善内容结构和平台变体;
  • 发布失败率高:先修复 Git 工作区、凭证、远端和幂等逻辑;
  • 人工审核耗时长:再考虑更好的差异对比、批量操作和审批体验。

只有当人工闭环稳定后,才适合逐步启用定时生成、更多平台和更高自动化。

9. 上线检查清单与常见误区

上线检查清单

  • 应用、MySQL、Redis 和入口服务健康;
  • 应用使用固定镜像版本或 digest;
  • MySQL、Redis 和管理后台未直接暴露公网;
  • 公网只开放业务真正需要的入口,管理通过 TAT、SSH 隧道或私有网络完成;
  • 自动生成保持关闭,直到人工闭环稳定;
  • 数据库备份已完成恢复演练;
  • 应用升级不会重建 MySQL、Redis 或删除卷;
  • Hexo 仓库分支正确、工作区干净、远端可推送;
  • 至少完成一次“生成—审核—批准—Hexo 发布—公网验收”;
  • 失败状态、日志和重试入口可查询。

常见误区

  1. 过早引入 Kubernetes,而单机 Compose 尚未形成稳定业务闭环;
  2. 把 AI 生成成功误认为内容已经可以发布;
  3. 未生成并审阅平台变体就直接批准;
  4. 调用批量发布接口,意外为多个平台创建任务;
  5. 应用升级时顺带重建数据库和缓存;
  6. 备份从未做过恢复演练;
  7. 在命令行、日志、镜像或仓库中泄露凭证;
  8. 把管理后台直接暴露到公网;
  9. 为了自动化而启用定时生成,却没有先证明人工流程稳定。

受控 MVP 的成功标准不是“架构看起来足够复杂”,而是系统能够在明确边界内稳定完成一次真实业务闭环,并且失败时知道如何停止、定位和恢复。

Welcome to Hexo! This is your very first post. Check documentation for more info. If you get any problems when using Hexo, you can find the answer in troubleshooting or you can ask me on GitHub.

Quick Start

Create a new post

1
$ hexo new "My New Post"

More info: Writing

Run server

1
$ hexo server

More info: Server

Generate static files

1
$ hexo generate

More info: Generating

Deploy to remote sites

1
$ hexo deploy

More info: Deployment

服务部署完成不等于稳定运行。上线后的前30分钟往往是故障高发期,掌握健康检查与日志分析的组合拳,能将平均故障恢复时间(MTTR)从小时级压缩到分钟级。本文围绕6个高频故障场景,给出可直接落地的排查路径。

一、健康检查机制搭建:第一道防线

健康检查不是可选项,而是服务存活的基本探针。没有它,负载均衡会将流量打到已崩溃的实例上。

1.1 Spring Boot Actuator 基础配置

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
# application.yml
management:
endpoints:
web:
exposure:
include: health,info,metrics,prometheus
endpoint:
health:
show-details: always
probes:
enabled: true
health:
livenessstate:
enabled: true
readinessstate:
enabled: true

1.2 K8s 探针配置

1
2
3
4
5
6
7
8
9
10
11
12
13
14
livenessProbe:
httpGet:
path: /actuator/health/liveness
port: 8080
initialDelaySeconds: 30
periodSeconds: 10
failureThreshold: 3
readinessProbe:
httpGet:
path: /actuator/health/readiness
port: 8080
initialDelaySeconds: 20
periodSeconds: 5
failureThreshold: 3

执行要点liveness 失败会触发容器重启,readiness 失败只会摘除流量。初始延迟必须大于服务启动时间,否则会陷入重启死循环。

二、启动失败:端口冲突与Bean初始化异常

部署后服务直接无法启动,日志是最快的信息源。

2.1 端口占用

日志特征:

1
Caused by: java.net.BindException: Address already in use

排查命令:

1
2
3
4
5
6
# 查看端口占用
ss -tlnp | grep 8080
# 或
lsof -i :8080
# 杀掉占用进程
kill -9 <PID>

2.2 Bean创建循环依赖

日志特征:

1
BeanCurrentlyInCreationException: Error creating bean with name 'userService'

处理步骤

  1. 在日志中搜索 BeanCurrentlyInCreationException,定位Bean名称
  2. 检查是否新增了带 @PostConstruct 的初始化逻辑
  3. 使用 @Lazy 打破循环依赖,或重构为事件驱动模式

2.3 配置缺失导致启动失败

1
2
# 快速过滤启动失败原因
java -jar app.jar 2>&1 | grep -E "(ERROR|Caused by|BindingFailure)"

三、OOM与内存泄漏:从日志到堆转储

OOM不会总是立即崩溃,但日志中会留下明确的死亡轨迹。

3.1 常见OOM类型与日志关键词

OOM类型 日志关键词 根因方向
Java heap space java.lang.OutOfMemoryError: Java heap space 大对象/内存泄漏
GC overhead GC overhead limit exceeded Full GC频繁且回收率低
Metaspace OutOfMemoryError: Metaspace 动态类加载过多
Direct memory OutOfMemoryError: Direct buffer memory Netty/NIO堆外内存

3.2 自动堆转储配置

1
2
3
4
# JVM启动参数
-XX:+HeapDumpOnOutOfMemoryError
-XX:HeapDumpPath=/data/dumps/heapdump.hprof
-XX:OnOutOfMemoryError="sh /data/scripts/alert.sh %p"

3.3 在线诊断(不重启服务)

1
2
3
4
5
6
7
8
# 获取堆内存直方图
jmap -histo:live <PID> | head -30

# 手动触发堆转储
jmap -dump:format=b,file=/tmp/heap.hprof <PID>

# 查看GC情况
jstat -gcutil <PID> 1000 10

可执行建议:使用 Arthas 在线诊断工具,无需重启即可分析:

1
2
3
4
5
6
# 启动Arthas
java -jar arthas-boot.jar <PID>
# 查看大对象
heapdump --live /tmp/dump.hprof
# 监控方法执行耗时
trace com.example.ServiceImpl methodName

四、数据库连接池耗尽:健康检查的隐藏信号

健康检查返回 DOWN 但服务进程存活,这是典型的资源耗尽场景。

4.1 健康检查暴露的数据库状态

1
2
3
4
5
6
7
8
9
10
11
12
13
{
"status": "DOWN",
"components": {
"db": {
"status": "DOWN",
"details": {
"database": "MySQL",
"validationQuery": "SELECT 1",
"error": "HikariPool-1 - Connection is not available"
}
}
}
}

4.2 日志中的连接池告警

1
2
WARN  HikariPool-1 - Thread starvation or clock leap detected
WARN HikariPool-1 - Failed to obtain connection in 30000ms

4.3 排查路径

1
2
3
4
5
6
7
8
9
# 1. 查看数据库当前连接数
SHOW PROCESSLIST;
SHOW STATUS LIKE 'Threads_connected';

# 2. 检查连接池配置
grep -A5 hikari application.yml

# 3. 查找未关闭的连接(日志中的泄漏线索)
grep -C3 "leak" app.log

连接池调优参考

1
2
3
4
5
6
spring:
datasource:
hikari:
maximum-pool-size: 20
connection-timeout: 3000
leak-detection-threshold: 60000

leak-detection-threshold 开启后,连接超过60秒未归还会打印泄漏堆栈,这是定位慢查询和未关连接的利器。

五、依赖服务超时与熔断:链路日志追踪

下游服务超时是级联故障的起点,日志中的时间线是还原现场的关键。

5.1 超时日志特征

1
2
ERROR FeignClient - Read timed out executing POST http://order-service/api/orders
WARN Hystrix - Command timeout, fallback triggered

5.2 使用Trace ID串联调用链

1
2
3
4
5
6
7
8
9
10
11
12
// MDC注入TraceId
@Component
public class TraceIdFilter implements Filter {
@Override
public void doFilter(ServletRequest req, ServletResponse res, FilterChain chain) {
String traceId = ((HttpServletRequest) req).getHeader("X-Trace-Id");
if (traceId == null) traceId = UUID.randomUUID().toString().replace("-", "");
MDC.put("traceId", traceId);
chain.doFilter(req, res);
MDC.clear();
}
}

日志格式配置:

1
<pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%X{traceId}] %-5level %logger{36} - %msg%n</pattern>

5.3 快速定位超时根因

1
2
3
4
5
# 按TraceId串联全链路日志
grep "a1b2c3d4" /data/logs/*.log | sort -k1

# 统计各服务耗时分布
grep "duration" app.log | awk -F'duration=' '{print $2}' | sort -n | tail -20

熔断配置建议

1
2
3
4
5
6
7
8
9
resilience4j:
circuitbreaker:
instances:
orderService:
failure-rate-threshold: 50
slow-call-rate-threshold: 60
slow-call-duration-threshold: 2s
wait-duration-in-open-state: 30s
sliding-window-size: 20

六、日志治理与告警自动化

散落的日志毫无价值,结构化 + 告警才是运维闭环。

6.1 结构化日志输出

1
2
3
4
5
6
<!-- logback-spring.xml -->
<appender name="JSON" class="ch.qos.logback.core.ConsoleAppender">
<encoder class="net.logstash.logback.encoder.LogstashEncoder">
<customFields>{"app":"order-service","env":"prod"}</customFields>
</encoder>
</appender>

输出示例:

1
{"@timestamp":"2024-01-15T10:30:00.123Z","level":"ERROR","logger":"c.e.OrderService","message":"order creation failed","traceId":"a1b2c3d4","app":"order-service"}

6.2 关键告警规则

告警项 日志关键词 触发条件 告警级别
服务不可用 Started.*failed 1次/5分钟 P0
OOM预兆 Full GC >5次/分钟 P1
连接池耗尽 Connection is not available >3次/分钟 P1
下游超时 Read timed out >10次/分钟 P2
熔断触发 circuit breaker opened 任意1次 P1

6.3 ELK + Alertmanager 联动

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
# Elasticsearch查询:最近5分钟ERROR日志按服务聚合
curl -X GET "localhost:9200/app-logs-*/_search" -H 'Content-Type: application/json' -d'
{
"query": {
"bool": {
"must": [
{"match": {"level": "ERROR"}},
{"range": {"@timestamp": {"gte": "now-5m"}}}
]
}
},
"aggs": {
"by_app": {"terms": {"field": "app.keyword"}}
}
}
'

落地检查清单

  • 健康检查端点已暴露且K8s探针已配置
  • OOM自动堆转储参数已添加
  • 日志包含TraceId且全链路贯通
  • ERROR级别日志已接入告警
  • 连接池泄漏检测已开启
  • 熔断器慢调用阈值已配置

健康检查告诉你”是否坏了”,日志告诉你”哪里坏了”。两者配合使用,配合结构化的日志输出和自动化告警,才能在故障发生的黄金5分钟内完成定位和恢复。

本文是对原有中文文章的生产加固修订,继续沿用原 slug 与原 URL。修订后的标题、正文语言和技术范围保持一致。本文基线为 JDK 21 + Spring Boot 3.2 及以上版本;涉及 JDK 24 的内容会单独标注,避免把后续版本能力误写为 JDK 21 已具备。

1. 先明确版本边界

虚拟线程已经在 JDK 21 正式交付。JEP 444 的状态为 Delivered,Release 为 21,因此在 JDK 21 中使用 Thread.ofVirtual()Thread.startVirtualThread(...)Executors.newVirtualThreadPerTaskExecutor() 不需要启用预览特性。

Spring Boot 从 3.2 开始提供虚拟线程集成。在 Java 21 或更高版本上,可以通过下面的配置启用:

1
2
3
4
spring:
threads:
virtual:
enabled: true

如果应用主要依赖 @Scheduled 等后台任务维持进程存活,还应考虑:

1
2
3
spring:
main:
keep-alive: true

原因是虚拟线程属于 daemon thread;当 JVM 中只剩 daemon thread 时,进程可以退出。这个配置不是所有 Web 应用都必需,但对没有其他非 daemon thread 的任务型应用尤其重要。

版本差异必须明确区分:

能力 实际版本 生产含义
虚拟线程正式交付 JDK 21 / JEP 444 可在正式生产基线使用,无需 --enable-preview
Spring Boot 虚拟线程开关 Spring Boot 3.2 使用 spring.threads.virtual.enabled=true
synchronized 场景大幅消除 pinning JDK 24 / JEP 491 不是 JDK 21 能力;JDK 21 仍需关注 monitor pinning
CompletableFuture.cancel(true) 中断运行任务 不支持 mayInterruptIfRunningCompletableFuture 没有效果

2. 虚拟线程适合什么,不适合什么

虚拟线程主要解决的是“线程因阻塞 I/O 长时间占用平台线程”的扩展性问题,例如:

  • JDBC 查询;
  • 同步 HTTP 调用;
  • 文件读写;
  • 阻塞式消息或 RPC 客户端;
  • 一个请求或批次项对应一条清晰同步调用链的业务。

它不会让 CPU 密集计算自动变快。大量 JSON 计算、压缩、加密、图像处理或复杂规则计算仍受 CPU 核数限制。虚拟线程也不会扩大数据库连接池、HTTP 连接池、文件句柄和下游限流额度。

因此生产设计应遵循两个原则:

  1. 虚拟线程负责承载阻塞任务。
  2. Semaphore、连接池和下游限流负责控制稀缺资源。

不要为了“控制虚拟线程数量”重新建立一个固定大小的虚拟线程池。对于数据库连接、第三方接口并发等有限资源,应使用 Semaphore 或资源池本身进行准入控制。

3. 配置后必须验证任务真的运行在虚拟线程上

配置存在并不等于所有代码路径都会切换为虚拟线程。Spring Boot 管理的异步执行器和调度器可以使用虚拟线程,但应用自行创建的 ThreadPoolExecutor、第三方框架内部执行器或显式指定的平台线程工厂不会自动改变。

最直接的运行时验证是:

1
2
3
4
5
Thread current = Thread.currentThread();
System.out.printf(
"thread=%s virtual=%s%n",
current.getName(),
current.isVirtual());

验收标准应检查 Thread.isVirtual(),不要依赖线程名称猜测。线程名称可以因 Spring Boot 版本、线程工厂或业务配置而变化。

对于批量任务,还可以生成包含虚拟线程的线程转储:

1
jcmd <pid> Thread.dump_to_file -format=json /tmp/threads.json

JSON 线程转储适合程序化分析,可检查目标任务栈中的 virtual: true。传统的 Thread.print 也有诊断价值,但在大量虚拟线程场景下,Thread.dump_to_file 更适合完整观察和工具处理。

4. 用显式准入控制保护数据库和外部接口

下面的完整示例只依赖 JDK 21,包含:

  • newVirtualThreadPerTaskExecutor()
  • Semaphore 下游并发上限;
  • 准入等待超时;
  • 单项执行 deadline;
  • Future.cancel(true) 的协作式取消请求;
  • transient/permanent failure 分类;
  • 批次结果聚合;
  • 可预测的 executor 关闭流程。

该示例已使用以下命令独立编译并运行:

1
2
javac --release 21 VirtualThreadBatchProcessor.java
java VirtualThreadBatchProcessor
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
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
import java.util.concurrent.CancellationException;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.Semaphore;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.logging.Level;
import java.util.logging.Logger;

/**
* JDK 21-only example: virtual-thread-per-task execution with explicit admission
* control, per-item deadlines, cooperative cancellation, and typed aggregation.
*/
public final class VirtualThreadBatchProcessor implements AutoCloseable {

private static final Logger LOGGER =
Logger.getLogger(VirtualThreadBatchProcessor.class.getName());

public record BatchItem(long id, String payload) {
public BatchItem {
Objects.requireNonNull(payload, "payload");
}
}

public enum FailureKind {
OVERLOADED,
TIMEOUT,
INTERRUPTED,
TRANSIENT,
PERMANENT,
UNEXPECTED
}

public sealed interface BatchResult permits Success, Failure {
long itemId();
}

public record Success(long itemId, String value) implements BatchResult {
public Success {
Objects.requireNonNull(value, "value");
}
}

public record Failure(
long itemId,
FailureKind kind,
String message,
boolean retryable
) implements BatchResult {
public Failure {
Objects.requireNonNull(kind, "kind");
Objects.requireNonNull(message, "message");
}
}

@FunctionalInterface
public interface BlockingOperation {
String execute(BatchItem item) throws Exception;
}

public static final class TransientBatchException extends Exception {
public TransientBatchException(String message) {
super(message);
}
}

public static final class PermanentBatchException extends Exception {
public PermanentBatchException(String message) {
super(message);
}
}

private record Submitted(
BatchItem item,
Future<BatchResult> future,
long deadlineNanos
) {
}

private final Semaphore downstreamPermits;
private final Duration admissionTimeout;
private final Duration taskTimeout;
private final BlockingOperation operation;
private final ExecutorService executor;

public VirtualThreadBatchProcessor(
int maxConcurrentDownstreamCalls,
Duration admissionTimeout,
Duration taskTimeout,
BlockingOperation operation
) {
if (maxConcurrentDownstreamCalls <= 0) {
throw new IllegalArgumentException(
"maxConcurrentDownstreamCalls must be positive");
}
this.admissionTimeout = requirePositive(admissionTimeout, "admissionTimeout");
this.taskTimeout = requirePositive(taskTimeout, "taskTimeout");
this.operation = Objects.requireNonNull(operation, "operation");
this.downstreamPermits = new Semaphore(maxConcurrentDownstreamCalls);
this.executor = Executors.newVirtualThreadPerTaskExecutor();
}

/**
* Processes a bounded input page concurrently. Callers should page very large
* datasets rather than materializing an unbounded number of Future instances.
*/
public List<BatchResult> process(List<BatchItem> items)
throws InterruptedException {
Objects.requireNonNull(items, "items");

List<Submitted> submitted = new ArrayList<>(items.size());
for (BatchItem item : items) {
Objects.requireNonNull(item, "items must not contain null");
long deadline = System.nanoTime() + taskTimeout.toNanos();
Future<BatchResult> future = executor.submit(() -> executeBounded(item));
submitted.add(new Submitted(item, future, deadline));
}

List<BatchResult> results = new ArrayList<>(submitted.size());
try {
for (Submitted task : submitted) {
results.add(await(task));
}
return List.copyOf(results);
} catch (InterruptedException interrupted) {
cancelAll(submitted);
Thread.currentThread().interrupt();
throw interrupted;
}
}

private BatchResult executeBounded(BatchItem item) {
boolean acquired = false;
try {
acquired = downstreamPermits.tryAcquire(
admissionTimeout.toNanos(), TimeUnit.NANOSECONDS);
if (!acquired) {
return new Failure(
item.id(),
FailureKind.OVERLOADED,
"downstream admission timeout",
true);
}

String value = operation.execute(item);
return new Success(item.id(), value);
} catch (InterruptedException interrupted) {
Thread.currentThread().interrupt();
return new Failure(
item.id(),
FailureKind.INTERRUPTED,
"task observed interruption",
true);
} catch (TransientBatchException transientFailure) {
return new Failure(
item.id(),
FailureKind.TRANSIENT,
safeMessage(transientFailure),
true);
} catch (PermanentBatchException permanentFailure) {
return new Failure(
item.id(),
FailureKind.PERMANENT,
safeMessage(permanentFailure),
false);
} catch (Exception unexpected) {
LOGGER.log(Level.WARNING,
"Unexpected batch failure for item " + item.id(), unexpected);
return new Failure(
item.id(),
FailureKind.UNEXPECTED,
unexpected.getClass().getSimpleName() + ": " + safeMessage(unexpected),
false);
} finally {
if (acquired) {
downstreamPermits.release();
}
}
}

private BatchResult await(Submitted task) throws InterruptedException {
long remainingNanos = task.deadlineNanos() - System.nanoTime();
if (remainingNanos <= 0) {
return timeout(task, "deadline elapsed before result collection");
}

try {
return task.future().get(remainingNanos, TimeUnit.NANOSECONDS);
} catch (TimeoutException timeout) {
return timeout(task, "per-item execution deadline exceeded");
} catch (CancellationException cancelled) {
return new Failure(
task.item().id(),
FailureKind.INTERRUPTED,
"task was cancelled",
true);
} catch (ExecutionException executionFailure) {
Throwable cause = executionFailure.getCause();
String type = cause == null
? executionFailure.getClass().getSimpleName()
: cause.getClass().getSimpleName();
String message = cause == null
? safeMessage(executionFailure)
: safeMessage(cause);
return new Failure(
task.item().id(),
FailureKind.UNEXPECTED,
type + ": " + message,
false);
}
}

private BatchResult timeout(Submitted task, String reason) {
// Future.cancel(true) requests interruption when the implementation knows
// the runner. It does not prove that the blocking library stopped.
boolean cancellationRequested = task.future().cancel(true);
return new Failure(
task.item().id(),
FailureKind.TIMEOUT,
reason + "; cancellationRequested=" + cancellationRequested,
true);
}

private static void cancelAll(List<Submitted> submitted) {
for (Submitted task : submitted) {
task.future().cancel(true);
}
}

private static Duration requirePositive(Duration value, String name) {
Objects.requireNonNull(value, name);
if (value.isZero() || value.isNegative()) {
throw new IllegalArgumentException(name + " must be positive");
}
return value;
}

private static String safeMessage(Throwable throwable) {
String message = throwable.getMessage();
return message == null || message.isBlank() ? "no detail" : message;
}

@Override
public void close() {
executor.shutdown();
try {
if (!executor.awaitTermination(30, TimeUnit.SECONDS)) {
List<Runnable> notStarted = executor.shutdownNow();
LOGGER.warning(() -> "Forced executor shutdown; notStarted="
+ notStarted.size());
if (!executor.awaitTermination(30, TimeUnit.SECONDS)) {
LOGGER.severe("Virtual-thread executor did not terminate");
}
}
} catch (InterruptedException interrupted) {
executor.shutdownNow();
Thread.currentThread().interrupt();
}
}

public static void main(String[] args) throws Exception {
try (VirtualThreadBatchProcessor processor =
new VirtualThreadBatchProcessor(
4,
Duration.ofSeconds(1),
Duration.ofSeconds(2),
item -> {
Thread.sleep(50);
return item.payload().toUpperCase();
})) {
List<BatchResult> results = processor.process(List.of(
new BatchItem(1, "alpha"),
new BatchItem(2, "beta"),
new BatchItem(3, "gamma")));
results.forEach(System.out::println);
}
}
}

4.1 为什么在任务内部获取 Semaphore

如果先获取 permit,再提交任务,而任务在真正开始前被取消,任务中的 finally 不会执行,容易造成 permit 泄漏。示例让任务启动后再执行 tryAcquire,无论成功、异常还是中断,都由同一条 finally 路径释放 permit。

这会允许一部分虚拟线程等待 permit,但虚拟线程本身成本较低;真正受保护的是数据库连接、HTTP 并发等稀缺资源。对于百万级输入,不应一次性创建百万个 Future,而应在上游分页读取或分块提交。

4.2 并发上限如何确定

maxConcurrentDownstreamCalls 不应凭感觉设置。常见约束包括:

  • 数据库连接池最大连接数;
  • 服务内其他请求必须保留的连接;
  • 外部 API 的并发或 QPS 限额;
  • 单请求内存占用;
  • 下游 P95/P99 延迟;
  • 超时后资源释放速度。

例如连接池最大 30,在线请求通常占用 15,批任务就不应直接把并发设为 30。应为在线流量、连接回收抖动和管理操作保留余量,再通过压测确定批任务上限。

5. 正确理解超时与取消

5.1 CompletableFuture.cancel(true) 不会中断运行任务

CompletableFuture 官方文档明确说明:mayInterruptIfRunning 在该实现中没有效果,因为 CompletableFuture 不使用中断控制处理过程。

因此下面的说法是错误的:

1
completableFuture.cancel(true); // 错误理解:一定会中断底层虚拟线程

它可以把 CompletableFuture 标记为取消,并使依赖阶段异常完成,但不能据此断言底层 JDBC、HTTP 或文件 I/O 已停止。

5.2 ExecutorService.submit 返回的 Future 可以请求中断,但仍不是强制终止

Future.cancel(true) 的语义是:当实现知道执行线程时,尝试中断正在运行的任务。这里有三个重要限定:

  1. 是“attempt”,不是强制杀死线程;
  2. 阻塞库必须正确响应中断;
  3. 返回 true 表示取消请求被接受,不等于底层资源已经释放。

示例使用 ExecutorService.submit(...) 返回的 Future,超时后调用 cancel(true),同时把结果记录为 TIMEOUT。这比使用 CompletableFuture.cancel(true) 宣称“已中断”更准确,但仍必须配置资源自身的超时。

5.3 生产超时必须下沉到资源层

线程级 deadline 只是最后一道控制,真正可靠的超时通常来自具体客户端:

  • JDBC:连接超时、socket 超时、事务超时、查询超时;
  • HTTP:connect timeout、request timeout、read timeout;
  • Redis/RPC:连接和命令超时;
  • 文件或对象存储:客户端请求超时;
  • 业务层:幂等键、状态机和补偿任务。

不要假设“发出 interrupt”就必然回收数据库连接。超时压测必须同时观察连接池 active、idle、pending 和 leak 指标。

6. 失败聚合:先分类,再决定重试

批量任务不能把所有异常都归为“失败后重试”。至少应区分:

类型 示例 默认策略
OVERLOADED 获取下游 permit 超时 延迟重试或回队列
TIMEOUT 超过单项 deadline 确认资源已释放后重试
INTERRUPTED 任务或批次被取消 根据任务状态恢复
TRANSIENT 短时网络故障、可恢复下游错误 有界退避重试
PERMANENT 参数非法、业务状态冲突 不自动重试,进入人工或数据修复
UNEXPECTED 未分类代码异常 告警并人工分析,默认不无限重试

重试必须满足:

  • 操作具备幂等性;
  • 有最大次数;
  • 有退避和抖动;
  • 保存最后错误和下一次执行时间;
  • 永久失败不会反复占用队列;
  • 批次结果能关联到具体 item ID。

结果聚合不应只返回一个 boolean。示例中的 SuccessFailure record 可以直接用于统计成功数、失败类别、重试候选和人工处理清单。

7. JFR 验证虚拟线程与 pinning

JDK Flight Recorder 提供以下虚拟线程事件:

  • jdk.VirtualThreadStart:虚拟线程开始,默认关闭;
  • jdk.VirtualThreadEnd:虚拟线程结束,默认关闭;
  • jdk.VirtualThreadPinned:虚拟线程 pinning 超过阈值,JDK 21 默认启用,默认阈值为 20 ms;
  • jdk.VirtualThreadSubmitFailed:虚拟线程启动或 unpark 提交失败,默认启用。

可以在压测期间启动一段有界 JFR 记录:

1
jcmd <pid> JFR.start   name=virtual-thread-check   settings=profile   duration=60s   filename=/tmp/virtual-thread-check.jfr

记录结束后打印相关事件:

1
jfr print   --events jdk.VirtualThreadStart,jdk.VirtualThreadEnd,jdk.VirtualThreadPinned,jdk.VirtualThreadSubmitFailed   /tmp/virtual-thread-check.jfr

如果需要稳定采集 Start/End,应通过 JDK Mission Control 或自定义 JFR 配置显式开启,因为它们默认关闭。不能因为输出中没有 Start/End 就认定应用没有虚拟线程。

7.1 JDK 21 与 JDK 24 的 pinning 差异

在 JDK 21 中,虚拟线程执行阻塞操作时,如果位于 synchronized 方法或代码块内,或者进入 native/foreign function,可能被 pin 到 carrier thread。频繁且长时间的 pinning 会降低扩展性。

JEP 491 在 JDK 24 交付后,JVM 可以让位于 synchronized 与 monitor 相关等待中的虚拟线程释放 carrier,消除了绝大多数此类 pinning。这个改进不能反向写成 JDK 21 已具备。

因此:

  • JDK 21:重点检查长时间阻塞是否发生在 synchronized 或 native 调用内;
  • JDK 24 及以后:monitor pinning 大幅改善,但仍应通过 JFR 观察实际运行情况;
  • 升级 JDK 前后都应进行同一套压力测试,不应只依据版本号推断性能。

8. Spring Boot 生产落地注意事项

8.1 线程池参数可能不再生效

Spring Boot 官方文档指出,启用虚拟线程后,用于配置传统线程池的属性不再产生原来的限制效果,因为虚拟线程由 JVM 全局的平台线程调度器承载,而不是专用固定线程池。

因此不能继续依赖 core-sizemax-size 等参数限制数据库并发。下游并发限制应迁移到 Semaphore、连接池和客户端限流配置。

8.2 不要混淆 Web 请求线程和自建批处理线程

spring.threads.virtual.enabled=true 影响 Spring Boot 自动配置管理的执行路径,但本文示例显式创建自己的 virtual-thread-per-task executor。两种方式可以共存,不过必须明确:

  • 哪个组件创建线程;
  • 谁负责生命周期;
  • 谁负责准入和超时;
  • 应用关闭时谁取消未完成任务。

自建 processor 应注册为 Spring Bean,并在 Bean 销毁时调用 close();不要每处理一条数据就创建和销毁 executor。

8.3 事务边界保持在线程内部

如果每个批次项需要独立数据库事务,应让事务方法在线程执行入口内调用。不要在提交任务的外层线程开启事务后,期待事务上下文自动传播到虚拟线程。

同时需要避免把不可线程安全的 request-scoped 对象、可变集合或数据库会话在多个任务间共享。

9. 上线前验证清单

配置与版本

  • 运行时为 JDK 21 或更高版本;
  • Spring Boot 为 3.2 或更高版本;
  • spring.threads.virtual.enabled=true 已在目标 profile 生效;
  • 任务型应用已评估 spring.main.keep-alive=true
  • 没有把 JEP 491 能力错误归入 JDK 21。

线程与并发

  • 业务日志或测试断言证明 Thread.isVirtual()==true
  • 下游并发由 Semaphore 或资源池明确限制;
  • 大数据集按页或分块提交,没有一次性创建无界 Future;
  • 批次重叠执行受到锁、租约或调度状态控制。

超时与取消

  • 没有宣称 CompletableFuture.cancel(true) 会中断运行任务;
  • Future.cancel(true) 仅被描述为中断尝试;
  • JDBC、HTTP 和 RPC 客户端均配置自身超时;
  • 超时压力测试后连接池能够回落到稳定水位;
  • 中断处理会恢复调用线程 interrupt flag。

失败与恢复

  • transient/permanent failure 分类清晰;
  • 重试有最大次数、退避和幂等保护;
  • 单项失败不会让整个批次状态丢失;
  • 永久失败进入人工或数据修复队列;
  • 任务取消、应用重启和重复执行都有恢复策略。

可观测性

  • 有批次总量、成功、失败、超时、拒绝和重试指标;
  • 有 in-flight 下游调用数;
  • 有数据库连接池与 HTTP 连接池指标;
  • 在压测中采集并检查 jdk.VirtualThreadPinned
  • 需要时使用自定义 JFR 配置开启 Start/End 事件。

10. 回滚与故障恢复

虚拟线程上线应保持可逆:

  1. 保留关闭虚拟线程配置的旧 profile;
  2. 保留平台线程时期的吞吐、延迟和连接池基线;
  3. 配置变更与代码改造分阶段发布;
  4. 如果出现连接池耗尽、下游限流、pinning 激增或进程无法退出,先停止新批次准入;
  5. 等待或取消在途任务,确认资源释放;
  6. 回滚 spring.threads.virtual.enabled 或应用版本;
  7. 使用任务状态与幂等键恢复未完成项,而不是直接重复整个批次。

回滚目标不是简单地“把线程改回去”,而是保证未完成任务可识别、重复执行不会产生副作用、下游资源可以恢复到稳定状态。

11. 结论

虚拟线程让同步阻塞代码可以用更直接的 thread-per-task 模型获得更高并发,但它不是无界并发开关。生产加固的核心是:

  • 用准确的 JDK 和 Spring Boot 版本边界描述能力;
  • Thread.isVirtual()、线程转储和 JFR 验证真实运行状态;
  • 用 Semaphore 和资源池保护下游;
  • 区分 CompletableFuture.cancel(true) 与普通 Future.cancel(true)
  • 把超时下沉到 JDBC、HTTP 等资源层;
  • 对失败进行结构化分类和有界恢复;
  • 用 JFR 验证 pinning,而不是凭配置推测。

这样才能把“启用虚拟线程”从一项配置变化,变成可验证、可观测、可回滚的生产能力。

官方参考资料

  1. OpenJDK JEP 444: Virtual Threads
  2. OpenJDK JEP 491: Synchronize Virtual Threads without Pinning
  3. Java SE 21 CompletableFuture API
  4. Java SE 21 Future API
  5. Java SE 21 Virtual Threads Guide
  6. Spring Boot 3.2 Reference: Virtual Threads