- 示例工程
- 教程
【免费下载链接】java-design-patterns
Design patterns implemented in Java
导读
本文以 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); Map<Headers, String> getHeaders(); void setBody(String body); String getBody(); } public class SimpleMessage implements Message { private final Map<Headers, 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 Map<Headers, String> getHeaders() { return Collections.unmodifiableMap(headers); } @Override public void setBody(String body) { this.body = body; } @Override public String getBody() { return body; } }关键实现细节(对应源码 Message.java 与 SimpleMessage.java):
Headers枚举定义了消息头类型:DATE(时间戳)与SENDER(发送方名称);getHeaders()返回的是Collections.unmodifiableMap(headers),即不可修改视图,防止外部调用方篡改消息头;- 生产者在发送消息时会自动填充
DATE与SENDER两个头部,消息体body则由调用方传入。
毒丸的定义:POISON_PILL 常量
毒丸本身在Message接口中作为一个匿名内部类常量定义(见 Message.java):
Message POISON_PILL = new Message() { @Override public void addHeader(Headers header, String value) { throw poison(); } @Override public String getHeader(Headers header) { throw poison(); } @Override public Map<Headers, 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 BlockingQueue<Message> 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_PILL;InterruptedException被捕获后只记录日志、"允许线程退出",避免线程被异常打断时留下未清理状态。
仓库测试 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.java):
public 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-patterns
Design patterns implemented in Java
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考