Java发布订阅案例

wen java案例 3

本文目录导读:

Java发布订阅案例

  1. 手写发布订阅模式
  2. JDK内置的Observer实现
  3. Spring事件驱动实现
  4. 高级特性扩展
  5. Maven依赖(Spring版本)

我将为您提供一个完整的Java发布订阅(Pub/Sub)模式案例,包含三种实现方式:手写实现、JDK内置实现和Spring实现。

手写发布订阅模式

基础接口定义

// 订阅者接口
public interface Subscriber {
    void onMessage(String topic, Object message);
}
// 发布者接口
public interface Publisher {
    void publish(String topic, Object message);
}
// 事件对象
public class Event {
    private String topic;
    private Object data;
    private long timestamp;
    public Event(String topic, Object data) {
        this.topic = topic;
        this.data = data;
        this.timestamp = System.currentTimeMillis();
    }
    // getters and setters
    public String getTopic() { return topic; }
    public Object getData() { return data; }
    public long getTimestamp() { return timestamp; }
}

事件总线实现

import java.util.*;
import java.util.concurrent.*;
public class EventBus {
    // 使用ConcurrentHashMap存储主题和订阅者列表
    private final ConcurrentMap<String, CopyOnWriteArrayList<Subscriber>> subscribers = new ConcurrentHashMap<>();
    // 线程池处理异步消息
    private final ExecutorService executor = Executors.newCachedThreadPool();
    // 订阅
    public void subscribe(String topic, Subscriber subscriber) {
        subscribers.computeIfAbsent(topic, k -> new CopyOnWriteArrayList<>())
                  .add(subscriber);
        System.out.println("订阅者订阅主题: " + topic);
    }
    // 取消订阅
    public void unsubscribe(String topic, Subscriber subscriber) {
        List<Subscriber> topicSubscribers = subscribers.get(topic);
        if (topicSubscribers != null) {
            topicSubscribers.remove(subscriber);
            System.out.println("订阅者取消订阅主题: " + topic);
        }
    }
    // 同步发布
    public void publishSync(String topic, Object message) {
        System.out.println("发布消息到主题 [" + topic + "]: " + message);
        List<Subscriber> topicSubscribers = subscribers.get(topic);
        if (topicSubscribers != null) {
            Event event = new Event(topic, message);
            topicSubscribers.forEach(subscriber -> 
                subscriber.onMessage(topic, event)
            );
        }
    }
    // 异步发布
    public void publishAsync(String topic, Object message) {
        System.out.println("异步发布消息到主题 [" + topic + "]: " + message);
        List<Subscriber> topicSubscribers = subscribers.get(topic);
        if (topicSubscribers != null) {
            Event event = new Event(topic, message);
            topicSubscribers.forEach(subscriber -> 
                executor.submit(() -> subscriber.onMessage(topic, event))
            );
        }
    }
    // 关闭线程池
    public void shutdown() {
        executor.shutdown();
    }
}

具体业务实现

// 新闻订阅者
public class NewsSubscriber implements Subscriber {
    private String name;
    public NewsSubscriber(String name) {
        this.name = name;
    }
    @Override
    public void onMessage(String topic, Object message) {
        Event event = (Event) message;
        System.out.println("[" + name + "] 收到新闻: " + event.getData() 
                          + " (时间: " + new java.text.SimpleDateFormat("HH:mm:ss")
                          .format(new java.util.Date(event.getTimestamp())) + ")");
    }
}
// 股票订阅者
public class StockSubscriber implements Subscriber {
    private String name;
    public StockSubscriber(String name) {
        this.name = name;
    }
    @Override
    public void onMessage(String topic, Object message) {
        Event event = (Event) message;
        System.out.println("[" + name + "] 收到股票消息: " + event.getData());
    }
}

测试代码

public class PubSubTest {
    public static void main(String[] args) throws InterruptedException {
        EventBus eventBus = new EventBus();
        // 创建订阅者
        NewsSubscriber newsSub1 = new NewsSubscriber("News-订阅者1");
        NewsSubscriber newsSub2 = new NewsSubscriber("News-订阅者2");
        StockSubscriber stockSub1 = new StockSubscriber("Stock-订阅者1");
        // 订阅
        eventBus.subscribe("news", newsSub1);
        eventBus.subscribe("news", newsSub2);
        eventBus.subscribe("stock", stockSub1);
        // 发布消息
        System.out.println("=== 同步发布新闻 ===");
        eventBus.publishSync("news", "Java 21 发布了新特性!");
        System.out.println("\n=== 异步发布股票 ===");
        eventBus.publishAsync("stock", "股票代码:600001 上涨5%");
        // 异步等待
        Thread.sleep(1000);
        // 取消订阅
        System.out.println("\n=== 取消订阅 ===");
        eventBus.unsubscribe("news", newsSub2);
        eventBus.publishSync("news", "第二条新闻消息");
        // 关闭
        eventBus.shutdown();
    }
}

JDK内置的Observer实现

import java.util.Observable;
import java.util.Observer;
// 新闻发布者(继承Observable)
public class NewsPublisher extends Observable {
    private String news;
    public void publishNews(String news) {
        this.news = news;
        setChanged();
        notifyObservers(news);
    }
}
// 新闻订阅者(实现Observer)
public class NewsObserver implements Observer {
    private String name;
    public NewsObserver(String name) {
        this.name = name;
    }
    @Override
    public void update(Observable o, Object arg) {
        System.out.println("[" + name + "] 收到新闻: " + arg);
    }
}
// JDK实现测试
public class JdkPubSubTest {
    public static void main(String[] args) {
        NewsPublisher publisher = new NewsPublisher();
        NewsObserver observer1 = new NewsObserver("观察者1");
        NewsObserver observer2 = new NewsObserver("观察者2");
        publisher.addObserver(observer1);
        publisher.addObserver(observer2);
        publisher.publishNews("JDK自带的观察者模式");
    }
}

Spring事件驱动实现

定义事件

import org.springframework.context.ApplicationEvent;
// 自定义事件
public class MessageEvent extends ApplicationEvent {
    private String topic;
    private String message;
    public MessageEvent(Object source, String topic, String message) {
        super(source);
        this.topic = topic;
        this.message = message;
    }
    public String getTopic() { return topic; }
    public String getMessage() { return message; }
}

事件监听器

import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Component;
@Component
public class MessageListener {
    @EventListener
    public void handleMessageEvent(MessageEvent event) {
        System.out.println("监听器收到消息 - 主题: " + event.getTopic() 
                          + ", 内容: " + event.getMessage());
    }
    // 条件监听器
    @EventListener(condition = "#event.topic == 'important'")
    public void handleImportantMessage(MessageEvent event) {
        System.out.println("重要消息处理: " + event.getMessage());
    }
}

事件发布者

import org.springframework.context.ApplicationEventPublisher;
import org.springframework.context.ApplicationEventPublisherAware;
import org.springframework.stereotype.Component;
@Component
public class MessagePublisher implements ApplicationEventPublisherAware {
    private ApplicationEventPublisher publisher;
    @Override
    public void setApplicationEventPublisher(ApplicationEventPublisher publisher) {
        this.publisher = publisher;
    }
    public void publishMessage(String topic, String message) {
        MessageEvent event = new MessageEvent(this, topic, message);
        publisher.publishEvent(event);
    }
}

Spring配置和启动

import org.springframework.context.annotation.AnnotationConfigApplicationContext;
import org.springframework.context.annotation.ComponentScan;
import org.springframework.context.annotation.Configuration;
@Configuration
@ComponentScan("com.example.pubsub")
public class SpringPubSubApplication {
    public static void main(String[] args) {
        AnnotationConfigApplicationContext context = 
            new AnnotationConfigApplicationContext(SpringPubSubApplication.class);
        MessagePublisher publisher = context.getBean(MessagePublisher.class);
        // 发布普通消息
        publisher.publishMessage("normal", "普通消息示例");
        // 发布重要消息
        publisher.publishMessage("important", "重要消息示例");
        context.close();
    }
}

高级特性扩展

带过滤器的消息总线

public interface MessageFilter {
    boolean accept(String topic, Object message);
}
public class FilteredEventBus extends EventBus {
    private final Map<String, List<MessageFilter>> filters = new HashMap<>();
    public void addFilter(String topic, MessageFilter filter) {
        filters.computeIfAbsent(topic, k -> new CopyOnWriteArrayList<>())
               .add(filter);
    }
    @Override
    public void publishSync(String topic, Object message) {
        List<MessageFilter> topicFilters = filters.get(topic);
        if (topicFilters != null && !topicFilters.isEmpty()) {
            boolean shouldPublish = topicFilters.stream()
                .allMatch(filter -> filter.accept(topic, message));
            if (shouldPublish) {
                super.publishSync(topic, message);
            }
        } else {
            super.publishSync(topic, message);
        }
    }
}

支持通配符订阅

public class WildcardEventBus extends EventBus {
    private final Map<Subscriber, List<String>> subscriberPatterns = new ConcurrentHashMap<>();
    public void subscribePattern(String pattern, Subscriber subscriber) {
        subscriberPatterns.computeIfAbsent(subscriber, k -> new CopyOnWriteArrayList<>())
                         .add(pattern);
    }
    @Override
    public void publishSync(String topic, Object message) {
        super.publishSync(topic, message);
        // 检查通配符订阅者
        subscriberPatterns.forEach((subscriber, patterns) -> {
            boolean match = patterns.stream().anyMatch(pattern -> 
                topic.matches(pattern.replace("*", ".*")));
            if (match) {
                subscriber.onMessage(topic, message);
            }
        });
    }
}

完整测试演示

public class CompletePubSubDemo {
    public static void main(String[] args) throws InterruptedException {
        System.out.println("=== 发布订阅模式完整演示 ===\n");
        // 使用带过滤功能的EventBus
        FilteredEventBus eventBus = new FilteredEventBus();
        // 创建订阅者
        Subscriber logSubscriber = (topic, message) -> {
            if (message instanceof Event) {
                Event event = (Event) message;
                System.out.println("[" + new java.text.SimpleDateFormat("HH:mm:ss")
                    .format(new java.util.Date(event.getTimestamp())) + "] "
                    + topic + ": " + event.getData());
            }
        };
        Subscriber criticalSubscriber = (topic, message) -> {
            if (message instanceof Event) {
                Event event = (Event) message;
                System.out.println("🔴 紧急订阅者处理: [" + topic + "] " + event.getData());
            }
        };
        // 订阅
        eventBus.subscribe("user.login", logSubscriber);
        eventBus.subscribe("user.important", criticalSubscriber);
        // 添加过滤器
        eventBus.addFilter("user.login", (topic, message) -> {
            if (message instanceof Event) {
                Event event = (Event) message;
                return !String.valueOf(event.getData()).contains("BANNED");
            }
            return true;
        });
        // 发布消息
        System.out.println("--- 发布普通消息 ---");
        eventBus.publishSync("user.login", "用户 admin 登录成功");
        System.out.println("\n--- 发布被过滤的消息 ---");
        eventBus.publishSync("user.login", "用户 banned_user BANNED");
        System.out.println("\n--- 发布紧急消息 ---");
        eventBus.publishSync("user.important", "检测到异常访问!");
        // 异步演示
        System.out.println("\n--- 异步发布演示 ---");
        eventBus.publishAsync("user.login", "用户 async_user 异步登录");
        Thread.sleep(1000);
        // 关闭
        eventBus.shutdown();
    }
}

Maven依赖(Spring版本)

<dependencies>
    <!-- Spring Context -->
    <dependency>
        <groupId>org.springframework</groupId>
        <artifactId>spring-context</artifactId>
        <version>5.3.30</version>
    </dependency>
</dependencies>

这个案例展示了:

  1. 手写实现 - 完整理解发布订阅的原理
  2. JDK内置实现 - 使用Observer模式
  3. Spring实现 - 企业级事件驱动开发
  4. 高级特性 - 过滤器、通配符等扩展

您可以根据实际需求选择合适的实现方式,手写实现适合学习和轻量级场景,Spring实现适合企业级应用。

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