本文目录导读:

我将为您提供一个完整的Java优雅停机案例,包含多种实现方案。
基于Spring Boot的优雅停机
Spring Boot 2.3+ 内置支持
# application.yml
server:
shutdown: graceful # 开启优雅停机
spring:
lifecycle:
timeout-per-shutdown-phase: 30s # 停机等待时间
Spring Boot 完整示例
@SpringBootApplication
public class GracefulShutdownApplication {
public static void main(String[] args) {
SpringApplication app = new SpringApplication(GracefulShutdownApplication.class);
app.run(args);
}
}
// 异步任务处理器
@Component
public class AsyncTaskProcessor {
private static final Logger log = LoggerFactory.getLogger(AsyncTaskProcessor.class);
private final ExecutorService executorService = Executors.newFixedThreadPool(10);
private volatile boolean running = true;
@PreDestroy
public void shutdown() {
log.info("开始停止异步任务处理器...");
running = false;
executorService.shutdown();
try {
if (!executorService.awaitTermination(30, TimeUnit.SECONDS)) {
executorService.shutdownNow();
if (!executorService.awaitTermination(30, TimeUnit.SECONDS)) {
log.error("线程池未能正常终止");
}
}
} catch (InterruptedException e) {
executorService.shutdownNow();
Thread.currentThread().interrupt();
}
log.info("异步任务处理器已停止");
}
public void submitTask(Runnable task) {
if (running) {
executorService.submit(task);
}
}
}
// 优雅停机控制器
@RestController
public class ShutdownController {
private final ApplicationContext context;
public ShutdownController(ApplicationContext context) {
this.context = context;
}
@PostMapping("/shutdown")
public String shutdown() {
// 触发优雅停机
Thread thread = new Thread(() -> {
try {
// 等待当前请求处理完成
Thread.sleep(3000);
// 关闭Spring容器
((ConfigurableApplicationContext) context).close();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
thread.start();
return "Shutdown initiated";
}
}
基于Java原生的优雅停机
使用ShutdownHook
public class GracefulShutdownExample {
private static final Logger log = LoggerFactory.getLogger(GracefulShutdownExample.class);
private final ExecutorService executorService = Executors.newFixedThreadPool(10);
private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(2);
private volatile boolean running = true;
public static void main(String[] args) {
GracefulShutdownExample app = new GracefulShutdownExample();
app.start();
// 注册ShutdownHook
Runtime.getRuntime().addShutdownHook(new Thread(app::shutdown, "Shutdown-Hook-Thread"));
log.info("应用启动成功,按Ctrl+C停止");
}
public void start() {
// 启动定时任务
scheduler.scheduleAtFixedRate(() -> {
if (running) {
log.info("执行定时任务...");
}
}, 0, 5, TimeUnit.SECONDS);
// 模拟接收请求
executorService.submit(() -> {
while (running) {
try {
Thread.sleep(1000);
log.info("处理请求...");
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
});
}
public void shutdown() {
log.info("开始优雅停机...");
// 1. 停止接收新任务
running = false;
// 2. 关闭定时任务
scheduler.shutdown();
// 3. 等待任务完成
executorService.shutdown();
try {
if (!executorService.awaitTermination(30, TimeUnit.SECONDS)) {
log.warn("线程池未在30秒内完成,强制关闭");
executorService.shutdownNow();
if (!executorService.awaitTermination(10, TimeUnit.SECONDS)) {
log.error("线程池未能正常终止");
}
}
} catch (InterruptedException e) {
log.error("停机过程被中断", e);
Thread.currentThread().interrupt();
}
// 4. 关闭资源
closeResources();
log.info("优雅停机完成");
}
private void closeResources() {
// 关闭数据库连接、文件句柄等资源
log.info("释放系统资源");
}
}
基于Netty的优雅停机
public class NettyGracefulShutdown {
private static final Logger log = LoggerFactory.getLogger(NettyGracefulShutdown.class);
private EventLoopGroup bossGroup;
private EventLoopGroup workerGroup;
private Channel serverChannel;
public void start() throws InterruptedException {
bossGroup = new NioEventLoopGroup(1);
workerGroup = new NioEventLoopGroup();
try {
ServerBootstrap bootstrap = new ServerBootstrap();
bootstrap.group(bossGroup, workerGroup)
.channel(NioServerSocketChannel.class)
.childHandler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) {
ch.pipeline().addLast(new SimpleChannelInboundHandler<Object>() {
@Override
protected void channelRead0(ChannelHandlerContext ctx, Object msg) {
log.info("收到消息: {}", msg);
// 处理消息
ctx.writeAndFlush("响应");
}
});
}
});
ChannelFuture future = bootstrap.bind(8080).sync();
serverChannel = future.channel();
log.info("Netty服务启动成功");
// 注册优雅停机钩子
Runtime.getRuntime().addShutdownHook(new Thread(this::shutdown));
// 等待服务器关闭
serverChannel.closeFuture().sync();
} finally {
shutdown();
}
}
public void shutdown() {
log.info("Netty开始优雅停机...");
if (serverChannel != null && serverChannel.isOpen()) {
// 停止接受新连接
serverChannel.close().syncUninterruptibly();
}
// 优雅关闭事件循环组
if (bossGroup != null) {
bossGroup.shutdownGracefully(5, 15, TimeUnit.SECONDS).syncUninterruptibly();
}
if (workerGroup != null) {
workerGroup.shutdownGracefully(5, 15, TimeUnit.SECONDS).syncUninterruptibly();
}
log.info("Netty优雅停机完成");
}
public static void main(String[] args) throws InterruptedException {
new NettyGracefulShutdown().start();
}
}
基于Kafka消费者的优雅停机
@Component
public class KafkaGracefulShutdown {
private static final Logger log = LoggerFactory.getLogger(KafkaGracefulShutdown.class);
@Autowired
private ConsumerFactory<String, String> consumerFactory;
private final Map<String, KafkaConsumer<String, String>> consumers = new ConcurrentHashMap<>();
private final AtomicBoolean running = new AtomicBoolean(true);
@PostConstruct
public void startConsumers() {
// 启动多个消费者线程
for (int i = 0; i < 3; i++) {
Thread consumerThread = new Thread(() -> consumeMessage("consumer-" + i));
consumerThread.setName("kafka-consumer-" + i);
consumerThread.start();
}
}
private void consumeMessage(String consumerName) {
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "my-group");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Collections.singletonList("my-topic"));
consumers.put(consumerName, consumer);
try {
while (running.get()) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
for (ConsumerRecord<String, String> record : records) {
String message = record.value();
log.info("消费者 {} 处理消息: {}", consumerName, message);
// 处理业务逻辑
try {
if (!processMessage(message)) {
// 处理失败,等待下次重试
log.warn("消息处理失败: {}", message);
} else {
// 手动提交offset
consumer.commitSync();
}
} catch (Exception e) {
log.error("处理消息异常", e);
}
}
}
} catch (WakeupException e) {
log.info("消费者 {} 被唤醒,准备停止", consumerName);
} finally {
consumer.close();
log.info("消费者 {} 已关闭", consumerName);
}
}
private boolean processMessage(String message) {
// 业务处理逻辑
return true;
}
@PreDestroy
public void shutdown() {
log.info("Kafka消费者开始优雅停机...");
// 停止消费
running.set(false);
// 唤醒所有消费者
consumers.forEach((name, consumer) -> {
consumer.wakeup();
});
// 等待消费者线程结束
try {
Thread.sleep(10000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
// 清理
consumers.clear();
log.info("Kafka消费者优雅停机完成");
}
}
综合实践案例 - 完整的企业级实现
@SpringBootApplication
public class EnterpriseGracefulShutdown {
private static final Logger log = LoggerFactory.getLogger(EnterpriseGracefulShutdown.class);
@Autowired
private HealthCheckController healthCheckController;
@Autowired
private RequestTrackingService requestTrackingService;
public static void main(String[] args) {
// 设置JVM参数
System.setProperty("server.shutdown", "graceful");
System.setProperty("spring.lifecycle.timeout-per-shutdown-phase", "30s");
SpringApplication app = new SpringApplication(EnterpriseGracefulShutdown.class);
app.addListeners(new AppShutdownListener());
app.run(args);
// 打印启动信息
log.info("企业应用启动成功");
}
}
// 请求跟踪服务
@Component
public class RequestTrackingService {
private final Map<Long, Long> activeRequests = new ConcurrentHashMap<>();
private final AtomicLong requestCounter = new AtomicLong(0);
private volatile boolean acceptingRequests = true;
public long beginRequest() {
long requestId = requestCounter.incrementAndGet();
activeRequests.put(requestId, System.currentTimeMillis());
return requestId;
}
public void endRequest(long requestId) {
activeRequests.remove(requestId);
}
public int getActiveRequestCount() {
return activeRequests.size();
}
public void stopAcceptingRequests() {
acceptingRequests = false;
log.info("停止接收新请求");
}
public void resumeAcceptingRequests() {
acceptingRequests = true;
log.info("恢复接收新请求");
}
public boolean canAcceptRequests() {
return acceptingRequests;
}
public void waitForRequestsToComplete(long timeout, TimeUnit unit) {
long timeoutMillis = unit.toMillis(timeout);
long startTime = System.currentTimeMillis();
while (!activeRequests.isEmpty()) {
if (System.currentTimeMillis() - startTime > timeoutMillis) {
log.warn("等待请求完成超时,仍有 {} 个活跃请求", activeRequests.size());
break;
}
try {
Thread.sleep(100);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
log.info("所有请求已完成处理,共等待 {} ms", System.currentTimeMillis() - startTime);
}
}
// 健康检查控制器
@RestController
public class HealthCheckController {
private static final Logger log = LoggerFactory.getLogger(HealthCheckController.class);
@Autowired
private RequestTrackingService requestTrackingService;
@GetMapping("/health")
public Map<String, Object> health() {
Map<String, Object> status = new HashMap<>();
status.put("status", "UP");
status.put("activeRequests", requestTrackingService.getActiveRequestCount());
status.put("acceptingRequests", requestTrackingService.canAcceptRequests());
return status;
}
@PostMapping("/maintenance")
public String toggleMaintenance() {
if (requestTrackingService.canAcceptRequests()) {
requestTrackingService.stopAcceptingRequests();
return "Maintenance mode enabled, waiting for active requests to complete...";
} else {
requestTrackingService.resumeAcceptingRequests();
return "Maintenance mode disabled";
}
}
}
// 应用监听器
@Component
public class AppShutdownListener implements ApplicationListener<ContextClosedEvent> {
private static final Logger log = LoggerFactory.getLogger(AppShutdownListener.class);
@Autowired
private RequestTrackingService requestTrackingService;
@Autowired
private ExecutorService requestExecutor;
@Override
public void onApplicationEvent(ContextClosedEvent event) {
log.info("应用关闭事件触发,执行优雅停机...");
// 1. 停止接收新请求
requestTrackingService.stopAcceptingRequests();
// 2. 等待正在处理的请求完成
requestTrackingService.waitForRequestsToComplete(30, TimeUnit.SECONDS);
// 3. 关闭业务线程池
requestExecutor.shutdown();
try {
if (!requestExecutor.awaitTermination(10, TimeUnit.SECONDS)) {
requestExecutor.shutdownNow();
log.warn("业务线程池强制关闭");
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
log.error("线程池关闭中断", e);
}
// 4. 关闭其他资源
closeDatabaseConnections();
closeCacheConnections();
log.info("优雅停机流程完成");
}
private void closeDatabaseConnections() {
// 关闭数据库连接池
log.info("关闭数据库连接池");
}
private void closeCacheConnections() {
// 关闭Redis等缓存连接
log.info("关闭缓存连接");
}
}
测试与验证
@SpringBootTest
public class GracefulShutdownTest {
private static final Logger log = LoggerFactory.getLogger(GracefulShutdownTest.class);
@Test
public void testGracefulShutdown() throws Exception {
// 获取应用上下文
ConfigurableApplicationContext context = SpringApplication.run(TestConfig.class);
// 模拟发送请求
TestRequestSender sender = new TestRequestSender();
sender.sendRequestsAsynchronously();
// 触发停机
Thread shutdownThread = new Thread(() -> {
try {
Thread.sleep(2000); // 等待请求开始处理
log.info("触发优雅停机...");
context.close();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
shutdownThread.start();
// 等待停机完成
shutdownThread.join(30000);
// 验证:请求是否处理完成
assertTrue("请求应全部完成", sender.allRequestsCompleted());
log.info("测试通过:优雅停机成功完成");
}
@SpringBootConfiguration
static class TestConfig {
@Bean
public ExecutorService requestExecutor() {
return Executors.newFixedThreadPool(5);
}
@Bean
public RequestTrackingService requestTrackingService() {
return new RequestTrackingService();
}
}
}
关键要点
- 分层停机:从业务层到基础设施层,按顺序关闭
- 超时控制:设置合理的超时时间,避免无限等待
- 资源释放:确保数据库连接、线程池等资源被正确关闭
- 状态检查:通过健康检查API监控停机进度
- 信号处理:正确捕获SIGTERM、Ctrl+C等信号
- 日志记录:详细记录停机过程,便于问题排查
这个案例涵盖了Java优雅停机的核心实现方式和最佳实践,可以根据具体需求选用不同的方案。