Java优雅停机案例

wen java案例 2

本文目录导读:

Java优雅停机案例

  1. 基于Spring Boot的优雅停机
  2. 基于Java原生的优雅停机
  3. 基于Netty的优雅停机
  4. 基于Kafka消费者的优雅停机
  5. 综合实践案例 - 完整的企业级实现
  6. 测试与验证
  7. 关键要点

我将为您提供一个完整的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();
        }
    }
}

关键要点

  1. 分层停机:从业务层到基础设施层,按顺序关闭
  2. 超时控制:设置合理的超时时间,避免无限等待
  3. 资源释放:确保数据库连接、线程池等资源被正确关闭
  4. 状态检查:通过健康检查API监控停机进度
  5. 信号处理:正确捕获SIGTERM、Ctrl+C等信号
  6. 日志记录:详细记录停机过程,便于问题排查

这个案例涵盖了Java优雅停机的核心实现方式和最佳实践,可以根据具体需求选用不同的方案。

抱歉,评论功能暂时关闭!