4.6 发布订阅模式


文档摘要

4.6 发布/订阅模式 第四章:Zookeeper 应用场景 - 4.6 发布/订阅模式 在分布式系统中,组件之间的解耦和异步通信至关重要。发布/订阅模式(Publish/Subscribe Pattern,简称 Pub/Sub)作为一种强大的消息通信模式,允许消息的发送者(发布者)和接收者(订阅者)在不知道彼此的情况下进行通信。这种模式极大地降低了系统组件之间的耦合度,提高了系统的灵活性和可扩展性。 Zookeeper,作为一个分布式协调服务,虽然其核心功能并非消息队列,但它凭借其独特的数据模型、watcher机制和可靠性,成为了构建分布式发布/订阅系统的理想基础设施。本章节将深入探讨如何利用Zookeeper实现发布/订阅模式,并结合代码实践进行详细讲解。 4.6.

4.6 发布/订阅模式

第四章:Zookeeper 应用场景 - 4.6 发布/订阅模式

在分布式系统中,组件之间的解耦和异步通信至关重要。发布/订阅模式(Publish/Subscribe Pattern,简称 Pub/Sub)作为一种强大的消息通信模式,允许消息的发送者(发布者)和接收者(订阅者)在不知道彼此的情况下进行通信。这种模式极大地降低了系统组件之间的耦合度,提高了系统的灵活性和可扩展性。

Zookeeper,作为一个分布式协调服务,虽然其核心功能并非消息队列,但它凭借其独特的数据模型、watcher机制和可靠性,成为了构建分布式发布/订阅系统的理想基础设施。本章节将深入探讨如何利用Zookeeper实现发布/订阅模式,并结合代码实践进行详细讲解。

4.6.1 发布/订阅模式概述

发布/订阅模式是一种消息范式,其中消息的发送者(发布者)不会直接将消息发送给特定的接收者(订阅者),而是将消息发布到主题(Topic)频道(Channel)。订阅者可以订阅一个或多个主题,系统会将发布到这些主题的消息推送给所有订阅者。

关键角色:

  • 发布者(Publisher): 负责创建消息,并将消息发布到指定的主题。发布者无需知道订阅者的存在或数量。

  • 订阅者(Subscriber): 负责订阅感兴趣的主题,并接收发布到这些主题的消息。订阅者无需知道发布者的信息。

  • 主题(Topic)/频道(Channel): 消息的逻辑通道。发布者将消息发布到主题,订阅者订阅主题以接收消息。

优势:

  • 解耦性: 发布者和订阅者之间完全解耦,彼此无需感知对方的存在,降低了系统组件之间的依赖性。

  • 异步性: 发布者发布消息后无需等待订阅者的响应,提高了系统的响应速度和吞吐量。

  • 灵活性: 可以动态地添加或删除发布者和订阅者,系统具有良好的可扩展性和灵活性。

  • 广播性: 一个消息可以被多个订阅者接收,实现消息的广播。

应用场景:

发布/订阅模式广泛应用于各种分布式系统中,例如:

  • 配置中心更新通知: 当配置中心的数据发生变化时,通过发布/订阅模式通知所有订阅了该配置的应用程序。

  • 实时消息推送: 例如,社交媒体的实时消息推送、新闻资讯的实时更新等。

  • 事件驱动架构: 系统中的各个组件通过发布和订阅事件进行通信,构建事件驱动的架构。

  • 日志收集系统: 各个应用服务器将日志发布到统一的主题,日志收集系统订阅该主题进行日志收集和分析。

4.6.2 基于Zookeeper实现发布/订阅模式的原理

Zookeeper本身并没有内置的消息队列功能,但其强大的数据模型和watcher机制使其能够有效地模拟发布/订阅模式。在Zookeeper中,我们可以利用节点(ZNode)来模拟主题,利用watcher机制来实现订阅和通知。

核心思想:

  1. 主题表示: 在Zookeeper中,使用一个持久节点(Persistent ZNode)来代表一个主题。例如,可以创建一个名为 /topics 的根节点,然后在 /topics 下为每个主题创建子节点,如 /topics/topicA/topics/topicB 等。

  2. 发布消息: 发布者要发布消息到某个主题时,它会在该主题节点下创建一个临时顺序节点(Ephemeral Sequential ZNode)。节点的数据部分存储消息内容。使用临时顺序节点的原因是:

    • 顺序性: 保证消息的发布顺序。

    • 临时性: 如果发布者崩溃,临时节点会自动删除,避免残留无用的消息节点。

  3. 订阅主题: 订阅者要订阅某个主题时,它需要对该主题节点注册Watcher。Watcher 监听主题节点的子节点变化事件(Child Node Change Event)

  4. 消息通知: 当发布者在主题节点下创建新的临时顺序节点(发布消息)时,Zookeeper会触发主题节点的子节点变化事件。所有订阅了该主题的订阅者都会收到Watcher通知。

  5. 获取消息: 订阅者收到Watcher通知后,会重新获取主题节点的子节点列表。新增的子节点就是新发布的消息节点。订阅者可以获取这些节点的数据,即消息内容。

  6. 消息确认与清理 (可选): 订阅者处理完消息后,可以选择删除对应的消息节点(临时顺序节点),以清理Zookeeper上的消息。但通常情况下,由于是临时节点,发布者断开连接后会自动删除,如果需要持久化消息,则需要更复杂的设计。

流程图 (graph TD):

流程详解:

  1. 发布者发布消息: 发布者连接Zookeeper,在指定主题节点(例如 /topics/topicA)下创建一个临时顺序节点,节点数据为消息内容。Zookeeper会自动为节点添加顺序编号,例如 /topics/topicA/msg-000001/topics/topicA/msg-000002 等。

  2. 订阅者订阅主题: 订阅者连接Zookeeper,对指定主题节点(例如 /topics/topicA)注册一个子节点变化的Watcher。

  3. Zookeeper事件通知: 当发布者创建新的消息节点后,Zookeeper服务器检测到主题节点的子节点发生变化,触发子节点变化事件。

  4. Watcher回调: 注册在主题节点上的所有Watcher会被触发,Zookeeper客户端会将事件通知发送给订阅者。

  5. 订阅者处理通知并获取消息: 订阅者收到Watcher通知后,知道主题有新消息发布。它会重新向Zookeeper获取主题节点的子节点列表。

  6. 获取消息内容: 订阅者遍历子节点列表,获取新增的消息节点,并读取节点数据,即消息内容。

  7. 消息处理: 订阅者根据消息内容进行相应的业务处理。

  8. 重新注册Watcher: 为了持续接收后续的消息,订阅者在处理完当前消息后,需要重新注册Watcher,以便监听下一次的子节点变化事件。这是一个关键步骤,因为Watcher是一次性的,触发一次后需要重新注册才能继续监听。

4.6.3 代码实践:基于Zookeeper实现发布/订阅模式 (Java)

以下是一个使用Java Zookeeper客户端 Curator 实现发布/订阅模式的示例代码。为了简化示例,我们只关注核心的发布和订阅逻辑。

1. 添加 Curator 依赖 (Maven)

<dependency> <groupId>org.apache.curator</groupId> <artifactId>curator-framework</artifactId> <version>5.2.0</version> <!-- 使用最新稳定版本 --> </dependency> <dependency> <groupId>org.apache.curator</groupId> <artifactId>curator-recipes</artifactId> <version>5.2.0</version> <!-- 使用最新稳定版本 --> </dependency>

2. 定义常量和工具类

import org.apache.curator.framework.CuratorFramework; import org.apache.curator.framework.CuratorFrameworkFactory; import org.apache.curator.retry.ExponentialBackoffRetry; public class ZookeeperUtils { private static final String CONNECT_STRING = "your_zookeeper_address:2181"; // 替换为你的Zookeeper地址 private static final int SESSION_TIMEOUT_MS = 5000; private static final int CONNECTION_TIMEOUT_MS = 5000; private static final int RETRY_BASE_SLEEP_TIME_MS = 1000; private static final int MAX_RETRIES = 3; public static CuratorFramework createCuratorFramework() { ExponentialBackoffRetry retryPolicy = new ExponentialBackoffRetry(RETRY_BASE_SLEEP_TIME_MS, MAX_RETRIES); return CuratorFrameworkFactory.newClient(CONNECT_STRING, SESSION_TIMEOUT_MS, CONNECTION_TIMEOUT_MS, retryPolicy); } }

3. 发布者 (Publisher) 代码

import org.apache.curator.framework.CuratorFramework; import org.apache.curator.framework.api.transaction.CuratorTransactionResult; import org.apache.zookeeper.CreateMode; public class Publisher { private final CuratorFramework client; private final String topicPath; public Publisher(String topic, CuratorFramework client) { this.topicPath = "/topics/" + topic; this.client = client; try { // 确保主题节点存在 if (client.checkExists().forPath(topicPath) == null) { client.create().creatingParentsIfNeeded().forPath(topicPath); } } catch (Exception e) { e.printStackTrace(); } } public void publish(String message) throws Exception { String messagePath = topicPath + "/msg-"; String createdPath = client.create() .withMode(CreateMode.EPHEMERAL_SEQUENTIAL) .forPath(messagePath, message.getBytes()); System.out.println("Published message to path: " + createdPath + ", message: " + message); } public static void main(String[] args) throws Exception { CuratorFramework client = ZookeeperUtils.createCuratorFramework(); client.start(); String topic = "myTopic"; Publisher publisher = new Publisher(topic, client); for (int i = 0; i < 5; i++) { String message = "Message-" + i; publisher.publish(message); Thread.sleep(1000); // 模拟发布消息的间隔 } client.close(); } }

4. 订阅者 (Subscriber) 代码

import org.apache.curator.framework.CuratorFramework; import org.apache.curator.framework.recipes.cache.ChildData; import org.apache.curator.framework.recipes.cache.NodeCache; import org.apache.curator.framework.recipes.cache.NodeCacheListener; import org.apache.curator.framework.recipes.cache.PathChildrenCache; import org.apache.curator.framework.recipes.cache.PathChildrenCacheEvent; import org.apache.curator.framework.recipes.cache.PathChildrenCacheListener; import java.nio.charset.StandardCharsets; import java.util.List; public class Subscriber { private final CuratorFramework client; private final String topicPath; public Subscriber(String topic, CuratorFramework client) { this.topicPath = "/topics/" + topic; this.client = client; } public void subscribe() throws Exception { PathChildrenCache cache = new PathChildrenCache(client, topicPath, true); // true 表示缓存节点数据 cache.getListenable().addListener(new PathChildrenCacheListener() { @Override public void childEvent(CuratorFramework client, PathChildrenCacheEvent event) throws Exception { switch (event.getType()) { case CHILD_ADDED: ChildData data = event.getData(); if (data != null) { String message = new String(data.getData(), StandardCharsets.UTF_8); String messagePath = data.getPath(); System.out.println("Received message from path: " + messagePath + ", message: " + message); // 在这里处理接收到的消息 } break; case CHILD_REMOVED: // 可以处理消息节点被删除的事件,例如消息确认机制 break; case CHILD_UPDATED: // 如果需要支持消息更新,可以处理 break; default: break; } } }); cache.start(); System.out.println("Subscribed to topic: " + topicPath + ", waiting for messages..."); } public static void main(String[] args) throws Exception { CuratorFramework client = ZookeeperUtils.createCuratorFramework(); client.start(); String topic = "myTopic"; Subscriber subscriber = new Subscriber(topic, client); subscriber.subscribe(); // 保持订阅者程序运行,持续接收消息 System.in.read(); cache.close(); // 关闭 PathChildrenCache client.close(); } }

代码详解:

  • ZookeeperUtils: 工具类,用于创建 CuratorFramework 客户端实例,配置 Zookeeper 连接信息和重试策略。你需要将 your_zookeeper_address:2181 替换为你的实际 Zookeeper 服务器地址。

  • Publisher:

    • Publisher(String topic, CuratorFramework client) 构造函数:

      • 初始化 topicPath/topics/topic名称

      • 检查主题节点是否存在,如果不存在则创建持久节点作为主题。

    • publish(String message) 方法:

      • 构建消息节点路径 topicPath + "/msg-"

      • 使用 client.create().withMode(CreateMode.EPHEMERAL_SEQUENTIAL).forPath(messagePath, message.getBytes()) 创建临时顺序节点,并将消息内容写入节点数据。

      • 打印发布成功的消息路径和内容。

    • main 方法:

      • 创建 CuratorFramework 客户端。

      • 创建 Publisher 实例,指定主题名称。

      • 循环发布 5 条消息,每条消息间隔 1 秒。

      • 关闭 CuratorFramework 客户端。

  • Subscriber:

    • Subscriber(String topic, CuratorFramework client) 构造函数:

      • 初始化 topicPath/topics/topic名称
    • subscribe() 方法:

      • 创建 PathChildrenCache 实例,用于监听主题节点的子节点变化。PathChildrenCache 是 Curator 提供的用于缓存和监听子节点变化的工具,非常方便。

      • cache.getListenable().addListener(...) 添加 PathChildrenCacheListener 监听器,处理子节点事件。

      • childEvent 方法中,根据事件类型 event.getType() 进行处理:

        • CHILD_ADDED: 当有新的子节点添加时(发布者发布消息),获取子节点数据 event.getData(),提取消息内容并打印。在这里可以编写消息处理的业务逻辑。

        • CHILD_REMOVED, CHILD_UPDATED: 可以根据需要处理节点删除和更新事件,例如实现消息确认机制或消息更新通知。

      • cache.start() 启动 PathChildrenCache,开始监听。

      • 打印订阅成功的提示信息,并进入等待状态,持续接收消息。

    • main 方法:

      • 创建 CuratorFramework 客户端。

      • 创建 Subscriber 实例,指定主题名称。

      • 调用 subscriber.subscribe() 启动订阅。

      • 使用 System.in.read() 阻塞主线程,保持订阅者程序运行,直到用户手动停止。

      • 最后关闭 PathChildrenCache 和 CuratorFramework 客户端。

运行示例:

  1. 确保你的 Zookeeper 服务正常运行。

  2. 修改 ZookeeperUtils.CONNECT_STRING 为你的 Zookeeper 地址。

  3. 编译并分别运行 Publisher.javaSubscriber.java

  4. 先运行 Subscriber.java,它会订阅 myTopic 主题并等待消息。

  5. 再运行 Publisher.java,它会向 myTopic 主题发布 5 条消息。

  6. 你将在 Subscriber 的控制台看到接收到的消息。

4.6.4 深入探讨与注意事项

1. 消息持久化与可靠性:

  • Zookeeper 本身不适合作为消息队列进行大规模的消息持久化存储。Zookeeper 的设计目标是分布式协调,而不是高吞吐量的消息传递。

  • 临时顺序节点的生命周期与发布者的会话绑定。如果发布者断开连接,临时节点会自动删除,消息会丢失。

  • 如果需要消息持久化和更高的可靠性,应该考虑使用专门的消息队列系统,如 Kafka、RabbitMQ 等。

  • Zookeeper 在发布/订阅模式中更适合用于轻量级的通知和协调,例如配置变更通知、服务状态更新等。

2. 消息顺序性:

  • 使用临时顺序节点可以保证同一个发布者发布的消息在同一个主题内是顺序的。因为 Zookeeper 会为每个临时顺序节点分配递增的顺序号。

  • 不同发布者发布的消息,或者不同主题的消息,无法保证全局顺序。

  • 如果对消息顺序有严格要求,需要根据具体业务场景进行更精细的设计,例如使用更复杂的排序算法或引入全局顺序ID生成器。

3. 扩展性与性能:

  • Zookeeper 集群本身具有良好的扩展性和可靠性。

  • 但当主题数量非常多,或者消息发布频率非常高时,Zookeeper 的性能可能会成为瓶颈。

  • 大量的 Watcher 会增加 Zookeeper 服务器的压力。

  • 对于高吞吐量的消息发布/订阅场景,仍然建议使用专门的消息队列系统。

  • 在 Zookeeper 中实现发布/订阅模式时,应该控制主题的数量和消息的频率,避免对 Zookeeper 服务器造成过大的压力。

4. Watcher 的特性:

  • 一次性触发: Zookeeper 的 Watcher 是一次性触发的。一旦 Watcher 被触发,它就会失效,需要重新注册才能继续监听。在订阅者代码中,PathChildrenCache 已经帮我们处理了 Watcher 的自动重新注册,简化了开发。

  • 顺序性保证: Zookeeper 保证事件通知的顺序性,即事件通知的顺序与事件发生的顺序一致。

  • 可靠性保证: Zookeeper 保证事件通知的可靠性,即事件通知不会丢失,除非 Zookeeper 集群发生故障。

5. 错误处理与重试:

  • 在发布者和订阅者代码中,需要处理 Zookeeper 连接异常、节点操作异常等。

  • 可以使用 Curator 提供的重试策略(如 ExponentialBackoffRetry)来提高系统的健壮性。

  • 对于订阅者,需要考虑处理消息失败的情况,例如消息处理逻辑出错、网络异常等。可以引入消息重试机制或死信队列等。

6. 安全性:

  • 如果 Zookeeper 集群部署在不安全的环境中,需要考虑安全性问题。

  • 可以使用 Zookeeper 提供的 ACL (Access Control List) 机制来控制对 Zookeeper 节点的访问权限,保护主题数据和防止未授权的发布和订阅。

4.6.5 总结

Zookeeper 提供了一种轻量级、可靠的机制来实现发布/订阅模式。虽然它并非专门的消息队列系统,但在分布式系统中,对于配置中心更新通知、服务状态变更广播等低吞吐量、高可靠性的场景,使用 Zookeeper 实现发布/订阅模式是一个简单有效的方案。

通过本章节的学习,我们了解了基于 Zookeeper 实现发布/订阅模式的原理、代码实践以及一些重要的注意事项。希望这些内容能够帮助你更好地理解和应用 Zookeeper,构建更加健壮和灵活的分布式系统。

关键要点回顾:

  • 使用持久节点作为主题,临时顺序节点作为消息。

  • 利用 Watcher 机制监听主题节点的子节点变化事件。

  • 订阅者收到通知后,重新获取子节点列表并处理新消息。

  • Watcher 是一次性的,需要重新注册才能持续监听。

  • Zookeeper 适合轻量级的通知和协调场景,不适合高吞吐量的消息队列场景。

  • 需要关注消息持久化、顺序性、扩展性、错误处理和安全性等问题。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U