
示例工程教程【免费下载链接】java-design-patternsDesign patterns implemented in Java项目地址https://gitcode.com/GitHub_Trending/ja/java-design-patterns点击查看免费下载导读本文以 java-design-patterns 仓库中的 Poison Pill毒丸模式为对象讲解如何通过一条特殊的毒丸消息在生产-消费者消息队列中实现优雅、可控的线程停机。读完本文你将掌握毒丸模式的适用场景、仓库中 Message / SimpleMessageQueue / Producer / Consumer 的完整实现细节、App 入口的完整运行示例以及该模式的优缺点与真实世界的应用如 Akka。模式概览Poison Pill毒丸是一个预先定义好的、众所周知的特殊数据元素它的作用是为一套独立运行的分布式消费流程提供优雅的graceful关闭方式。在消息交换语境下毒丸就是一条已知的、特殊的消息结构消费者一旦读到它就知道消息交换已经结束从而安全地退出消费循环。该模式在仓库中的位置模式文档poison-pill/README.md源码目录poison-pill/src/main/java/com/iluwatar/poison/pill/测试目录poison-pill/src/test/java/com/iluwatar/poison/pill/从仓库中的 App 源码注释可以看出这个模式的定位是终止生产者-消费者模式的方案之一由生产者负责通知消费者消息交换已经结束并拒绝后续任何消息消费者收到毒丸后会停止从队列中读取消息。文档还特别提醒必须保证毒丸是消费者从队列中读到的最后一条消息如果你使用了带优先级的队列这一点会变得棘手。现实世界的类比一家商店打烊时店员在门口挂上已打烊的牌子。牌子并不把店内正在购物的顾客赶出去而是告诉新顾客不再接待店员继续服务店内剩余顾客等他们买完东西再锁门关灯。毒丸消息的作用与此完全一致——它告诉消费者不再接收新任务但允许消费者把手头队列中剩余的任务处理完再优雅退出。核心组件Message 与 SimpleMessage消息结构由接口Message与实现类SimpleMessage组成。public interface Message { enum Headers { DATE, SENDER } void addHeader(Headers header, String value); String getHeader(Headers header); MapHeaders, String getHeaders(); void setBody(String body); String getBody(); } public class SimpleMessage implements Message { private final MapHeaders, String headers new HashMap(); private String body; Override public void addHeader(Headers header, String value) { headers.put(header, value); } Override public String getHeader(Headers header) { return headers.get(header); } Override public MapHeaders, String getHeaders() { return Collections.unmodifiableMap(headers); } Override public void setBody(String body) { this.body body; } Override public String getBody() { return body; } }关键实现细节对应源码 Message.java 与 SimpleMessage.javaHeaders枚举定义了消息头类型DATE时间戳与SENDER发送方名称getHeaders()返回的是Collections.unmodifiableMap(headers)即不可修改视图防止外部调用方篡改消息头生产者在发送消息时会自动填充DATE与SENDER两个头部消息体body则由调用方传入。毒丸的定义POISON_PILL 常量毒丸本身在Message接口中作为一个匿名内部类常量定义见 Message.javaMessage POISON_PILL new Message() { Override public void addHeader(Headers header, String value) { throw poison(); } Override public String getHeader(Headers header) { throw poison(); } Override public MapHeaders, String getHeaders() { throw poison(); } Override public void setBody(String body) { throw poison(); } Override public String getBody() { throw poison(); } private RuntimeException poison() { return new UnsupportedOperationException(Poison); } };设计要点毒丸是一个共享的单例标记对象。App 源码注释指出简单场景下毒丸可以只是一个null引用但持有唯一的、独立的共享对象标记命名为 Poison 或 Poison Pill更清晰、更具自描述性毒丸的所有消息操作方法都会抛出UnsupportedOperationException(Poison)从行为层面杜绝了对毒丸进行读写操作的可能——它不携带任何真实消息内容只是一个终止信号测试 PoisonMessageTest.java 对上述五个方法逐一断言其抛出UnsupportedOperationException印证了这一设计。消息队列抽象MqPublishPoint / MqSubscribePoint / MessageQueue队列层被拆分为两个单一职责接口与一个组合接口public interface MqPublishPoint { void put(Message msg) throws InterruptedException; } public interface MqSubscribePoint { Message take() throws InterruptedException; } public interface MessageQueue extends MqPublishPoint, MqSubscribePoint { }SimpleMessageQueue同时实现这三个接口内部封装了 JDK 的阻塞队列public class SimpleMessageQueue implements MessageQueue { private final BlockingQueueMessage queue; public SimpleMessageQueue(int bound) { queue new ArrayBlockingQueue(bound); } Override public void put(Message msg) throws InterruptedException { queue.put(msg); } Override public Message take() throws InterruptedException { return queue.take(); } }说明对应源码 SimpleMessageQueue.java构造函数参数bound指定ArrayBlockingQueue的有界容量。示例中new SimpleMessageQueue(10000)即队列最多容纳 10000 条消息put/take均为阻塞语义队列满时put阻塞等待空间队列空时take阻塞等待消息天然满足生产者-消费者模型的需求接口拆分的好处Producer只依赖MqPublishPoint发布方视角Consumer只依赖MqSubscribePoint订阅方视角通过接口隔离降低耦合。生产者 Producer发送消息与投递毒丸public class Producer { private final MqPublishPoint queue; private final String name; private boolean isStopped; public Producer(String name, MqPublishPoint queue) { this.name name; this.queue queue; this.isStopped false; } public void send(String body) { if (isStopped) { throw new IllegalStateException(String.format( Producer %s was stopped and fail to deliver requested message [%s]., body, name)); } var msg new SimpleMessage(); msg.addHeader(Headers.DATE, new Date().toString()); msg.addHeader(Headers.SENDER, name); msg.setBody(body); try { queue.put(msg); } catch (InterruptedException e) { // allow thread to exit LOGGER.error(Exception caught., e); } } public void stop() { isStopped true; try { queue.put(Message.POISON_PILL); } catch (InterruptedException e) { // allow thread to exit LOGGER.error(Exception caught., e); } } }对应源码 Producer.java要点send首先检查isStopped标志生产者一旦停止就不再接受任何新消息此时调用send会抛出IllegalStateException消息体与生产者名称会拼入异常信息stop()是关闭流程的核心先将isStopped置为true再向队列投递Message.POISON_PILLInterruptedException被捕获后只记录日志、允许线程退出避免线程被异常打断时留下未清理状态。仓库测试 ProducerTest.java 验证了两点send(Hello!)后发出的消息 SENDER 头等于producer、DATE 头非空、body 等于Hello!stop()后publishPoint.put(eq(Message.POISON_PILL))被调用且再次send抛出IllegalStateException。消费者 Consumer识别毒丸并优雅退出public class Consumer { private final MqSubscribePoint queue; private final String name; public Consumer(String name, MqSubscribePoint queue) { this.name name; this.queue queue; } public void consume() { while (true) { try { var msg queue.take(); if (Message.POISON_PILL.equals(msg)) { LOGGER.info(Consumer {} receive request to terminate., name); break; } var sender msg.getHeader(Headers.SENDER); var body msg.getBody(); LOGGER.info(Message [{}] from [{}] received by [{}], body, sender, name); } catch (InterruptedException e) { // allow thread to exit LOGGER.error(Exception caught., e); return; } } } }对应源码 Consumer.java要点consume()运行在无限循环中不断take消息关键判断Message.POISON_PILL.equals(msg)一旦取到毒丸打印终止日志并break退出循环不会再去处理队列中残留的任何后续消息——这正对应毒丸必须是最后一条被读取的消息的前提约束若take抛出InterruptedException则直接return结束消费线程。测试 ConsumerTest.java 构造了一个包含两条真实消息、一条POISON_PILL、以及一条迟到消息的队列断言消费者处理完前两条消息后遇到毒丸即打印Consumer NSA receive request to terminate.并退出——毒丸之后的消息不会被消费验证了停机信号只对排在其前的消息生效。完整示例App 入口与运行输出仓库提供了可直接运行的完整示例App.javapublic static void main(String[] args) { var queue new SimpleMessageQueue(10000); final var producer new Producer(PRODUCER_1, queue); final var consumer new Consumer(CONSUMER_1, queue); new Thread(consumer::consume).start(); new Thread(() - { producer.send(hand shake); producer.send(some very important information); producer.send(bye!); producer.stop(); }).start(); }运行流程创建容量为 10000 的有界队列创建名为PRODUCER_1的生产者与名为CONSUMER_1的消费者消费者线程先启动并阻塞在queue.take()等待消息生产者线程依次发送hand shake、some very important information、bye!三条消息最后调用stop()投递毒丸消费者按顺序消费三条真实消息后读到毒丸打印终止日志并退出。程序输出时序信息因运行环境而异Message [hand shake] from [PRODUCER_1] received by [CONSUMER_1] Message [some very important information] from [PRODUCER_1] received by [CONSUMER_1] Message [bye!] from [PRODUCER_1] received by [CONSUMER_1] Consumer CONSUMER_1 receive request to terminate.在仓库中你还可以通过 Maven 直接运行验证mvn -pl poison-pill compile exec:java或运行该模块的单元测试验证各组件行为AppTest、ConsumerTest、ProducerTest、PoisonMessageTest、SimpleMessageTest见 poison-pill/src/test/java/com/iluwatar/poison/pill/。适用场景当满足以下条件时适合使用毒丸模式对应 poison-pill/README.md 与英文版文档的适用性说明需要从一个线程/进程向另一个线程/进程发送终止信号——这是毒丸模式最核心的诉求系统处于多线程环境要求健壮的容错能力与消费者无缝停机fault tolerance and seamless consumer shutdown典型的生产者-消费者场景中需要告知消费者消息处理已结束需要保证消费者在处理完队列中剩余消息之后才关闭而不是被外部强制打断。优点与权衡优点简化消费者停机流程停机逻辑被收拢为识别一条特殊消息无需复杂的线程协作原语保证排队的任务处理完成消费者在遇到毒丸前会继续消费队列中的既有消息任务不会丢失停机逻辑与主处理逻辑解耦生产者和消费者只关心消息本身停机信号也是消息的一种职责边界清晰。权衡与注意点消费者必须主动检查毒丸每次循环多一次equals比较存在少量运行时开销毒丸识别失败将导致无限阻塞如果消费者因实现问题如没有正确比较、或队列带优先级导致毒丸被插到真实消息之后无法识别毒丸take()会一直阻塞线程无法退出。因此必须保证毒丸是消费者读到的最后一条消息带优先级队列时尤其需要注意若同一队列有多个消费者需要约定好毒丸的数量与投递方式每个消费者一条否则部分消费者可能收不到停机信号。真实世界中的应用Akka Actor 框架akka.actor.PoisonPill是毒丸模式在 Actor 系统中的经典实现——向 Actor 发送 PoisonPill 消息即可使其优雅停止处理后续消息并终止Java ExecutorService 停机通过提交一个特殊任务作为停机信号通知线程池不再接受新任务各类消息系统在队列处理流程中使用一条约定好的特殊消息标识队列处理结束。与相关设计模式的关系Producer-Consumer毒丸模式常与生产者-消费者模式搭配使用负责其中的通信与消费者停机环节消息队列Message Queue消息队列类实现经常借助毒丸标识队列处理流程的终止Observer可用于在停机事件发生时通知订阅者。总结毒丸模式把停机这件跨线程的棘手事抽象成了一条特殊的、众所周知的、不携带业务内容的消息生产者投递它消费者识别它并优雅退出。在 java-design-patterns 仓库的 poison-pill 模块中你可以看到一个最小但完整的 Java 实现——从Message/POISON_PILL常量、有界阻塞队列SimpleMessageQueue到Producer.stop()与Consumer.consume()的配合再到测试用例对毒丸之后的消息不被消费这一语义的验证。它的核心价值在于既保证了剩余任务处理完成又让停机逻辑变得简单、可读、可测试代价则是要求消费端严格遵守毒丸必须是最后一条消息的约定。赞分享示例工程教程【免费下载链接】java-design-patternsDesign patterns implemented in Java项目地址https://gitcode.com/GitHub_Trending/ja/java-design-patterns点击查看免费下载相关推荐fabio 平滑下线指南深入解析 proxy.deregistergraceperiod 优雅停机配置fabio 平滑下线指南深入解析 proxy.deregistergraceperiod 优雅停机配置 导读 本文聚焦 fabio 的优雅停机机制中一个关键配后端API网关微服务WPF UI通知系统SnackbarService消息队列设计WPF UI通知系统SnackbarService消息队列设计 引言告别通知风暴构建优雅的WPF消息系统 你是否还在为WPF应用中的通知管理烦恼重叠显示UI组件桌面应用创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考