Java Redis发布订阅案例如何实操

wen java案例 29

本文目录导读:

Java Redis发布订阅案例如何实操

  1. 环境准备
  2. 实现订阅者(Subscriber)
  3. 实现发布者(Publisher)
  4. 完整的测试案例
  5. 高级应用:对象消息传输
  6. 运行说明
  7. 优化建议

我来为你展示一个完整的Java Redis发布订阅实操案例。

环境准备

Maven依赖

<dependency>
    <groupId>redis.clients</groupId>
    <artifactId>jedis</artifactId>
    <version>4.4.3</version>
</dependency>

实现订阅者(Subscriber)

消息监听器

import redis.clients.jedis.JedisPubSub;
public class RedisMsgSubscriber extends JedisPubSub {
    @Override
    public void onMessage(String channel, String message) {
        System.out.println("收到消息 - 频道: " + channel + ", 消息: " + message);
    }
    @Override
    public void onPMessage(String pattern, String channel, String message) {
        System.out.println("收到模式匹配消息 - 模式: " + pattern + 
                          ", 频道: " + channel + ", 消息: " + message);
    }
    @Override
    public void onSubscribe(String channel, int subscribedChannels) {
        System.out.println("订阅频道: " + channel + ", 当前订阅数: " + subscribedChannels);
    }
    @Override
    public void onUnsubscribe(String channel, int subscribedChannels) {
        System.out.println("取消订阅频道: " + channel + ", 当前订阅数: " + subscribedChannels);
    }
    @Override
    public void onPSubscribe(String pattern, int subscribedChannels) {
        System.out.println("模式订阅: " + pattern + ", 当前订阅数: " + subscribedChannels);
    }
    @Override
    public void onPUnsubscribe(String pattern, int subscribedChannels) {
        System.out.println("取消模式订阅: " + pattern + ", 当前订阅数: " + subscribedChannels);
    }
}

订阅服务

import redis.clients.jedis.Jedis;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
public class RedisSubscriberService {
    private static final String REDIS_HOST = "localhost";
    private static final int REDIS_PORT = 6379;
    public void subscribeChannel(String channel) {
        ExecutorService executor = Executors.newSingleThreadExecutor();
        executor.submit(() -> {
            try (Jedis jedis = new Jedis(REDIS_HOST, REDIS_PORT)) {
                System.out.println("开始订阅频道: " + channel);
                jedis.subscribe(new RedisMsgSubscriber(), channel);
            } catch (Exception e) {
                e.printStackTrace();
            }
        });
    }
    public void subscribePattern(String pattern) {
        ExecutorService executor = Executors.newSingleThreadExecutor();
        executor.submit(() -> {
            try (Jedis jedis = new Jedis(REDIS_HOST, REDIS_PORT)) {
                System.out.println("开始模式订阅: " + pattern);
                jedis.psubscribe(new RedisMsgSubscriber(), pattern);
            } catch (Exception e) {
                e.printStackTrace();
            }
        });
    }
}

实现发布者(Publisher)

import redis.clients.jedis.Jedis;
public class RedisPublisher {
    private static final String REDIS_HOST = "localhost";
    private static final int REDIS_PORT = 6379;
    public void publish(String channel, String message) {
        try (Jedis jedis = new Jedis(REDIS_HOST, REDIS_PORT)) {
            Long result = jedis.publish(channel, message);
            System.out.println("发布消息完成 - 频道: " + channel + 
                              ", 消息: " + message + 
                              ", 接收者数量: " + result);
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
    public void publishMultipleMessages(String channel, List<String> messages) {
        try (Jedis jedis = new Jedis(REDIS_HOST, REDIS_PORT)) {
            for (String message : messages) {
                Long result = jedis.publish(channel, message);
                System.out.println("发布消息 - 频道: " + channel + 
                                  ", 消息: " + message + 
                                  ", 接收者数量: " + result);
                Thread.sleep(1000); // 模拟消息间隔
            }
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

完整的测试案例

import java.util.Arrays;
import java.util.List;
import java.util.concurrent.TimeUnit;
public class PubSubDemo {
    public static void main(String[] args) throws InterruptedException {
        // 1. 创建订阅服务
        RedisSubscriberService subscriberService = new RedisSubscriberService();
        // 2. 创建发布者
        RedisPublisher publisher = new RedisPublisher();
        // 3. 订阅具体频道
        System.out.println("=== 订阅具体频道 ===");
        subscriberService.subscribeChannel("news");
        subscriberService.subscribeChannel("sports");
        Thread.sleep(1000); // 等待订阅完成
        // 4. 订阅模式(匹配所有以.news结尾的频道)
        System.out.println("\n=== 订阅模式 ===");
        subscriberService.subscribePattern("*.news");
        Thread.sleep(1000);
        // 5. 发布消息到具体频道
        System.out.println("\n=== 发布消息到news频道 ===");
        publisher.publish("news", "今天天气很好");
        publisher.publish("news", "发布了新版Java");
        Thread.sleep(500);
        System.out.println("\n=== 发布消息到sports频道 ===");
        publisher.publish("sports", "NBA总决赛即将开始");
        Thread.sleep(500);
        // 6. 批量发布消息
        System.out.println("\n=== 批量发布消息 ===");
        List<String> messages = Arrays.asList(
            "消息1: Hello",
            "消息2: World", 
            "消息3: Redis Pub/Sub"
        );
        publisher.publishMultipleMessages("news", messages);
        // 7. 测试模式匹配
        System.out.println("\n=== 测试模式匹配 ===");
        publisher.publish("tech.news", "AI技术发展");
        publisher.publish("sports.news", "世界杯");
        // 8. 等待消息处理完成
        TimeUnit.SECONDS.sleep(2);
        System.out.println("\nDemo完成!");
    }
}

高级应用:对象消息传输

import com.fasterxml.jackson.databind.ObjectMapper;
public class AdvancedPubSubDemo {
    private static final ObjectMapper objectMapper = new ObjectMapper();
    // 消息对象
    static class OrderMessage {
        private String orderId;
        private String userId;
        private double amount;
        private long timestamp;
        // getters and setters
        public String getOrderId() { return orderId; }
        public void setOrderId(String orderId) { this.orderId = orderId; }
        public String getUserId() { return userId; }
        public void setUserId(String userId) { this.userId = userId; }
        public double getAmount() { return amount; }
        public void setAmount(double amount) { this.amount = amount; }
        public long getTimestamp() { return timestamp; }
        public void setTimestamp(long timestamp) { this.timestamp = timestamp; }
        @Override
        public String toString() {
            return "OrderMessage{" +
                    "orderId='" + orderId + '\'' +
                    ", userId='" + userId + '\'' +
                    ", amount=" + amount +
                    ", timestamp=" + timestamp +
                    '}';
        }
    }
    // 发布订单消息
    public static void publishOrder(String channel, OrderMessage order) {
        try {
            String jsonMessage = objectMapper.writeValueAsString(order);
            try (Jedis jedis = new Jedis("localhost", 6379)) {
                jedis.publish(channel, jsonMessage);
                System.out.println("发布订单消息: " + jsonMessage);
            }
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
    // 高级订阅者(JSON消息处理)
    static class AdvancedSubscriber extends JedisPubSub {
        @Override
        public void onMessage(String channel, String message) {
            try {
                if (message.startsWith("{")) {  // JSON消息
                    OrderMessage order = objectMapper.readValue(message, OrderMessage.class);
                    System.out.println("接收到订单: " + order);
                    // 处理订单逻辑
                    processOrder(order);
                } else {
                    System.out.println("普通文本消息: " + message);
                }
            } catch (Exception e) {
                System.err.println("消息处理失败: " + e.getMessage());
            }
        }
        private void processOrder(OrderMessage order) {
            // 模拟订单处理
            System.out.println("处理订单: " + order.getOrderId() + 
                             ", 用户: " + order.getUserId() + 
                             ", 金额: " + order.getAmount());
        }
    }
    public static void main(String[] args) throws InterruptedException {
        // 启动高级订阅者
        ExecutorService executor = Executors.newSingleThreadExecutor();
        executor.submit(() -> {
            try (Jedis jedis = new Jedis("localhost", 6379)) {
                jedis.subscribe(new AdvancedSubscriber(), "orders");
            }
        });
        Thread.sleep(1000);
        // 发布订单
        OrderMessage order1 = new OrderMessage();
        order1.setOrderId("ORD001");
        order1.setUserId("USER001");
        order1.setAmount(99.99);
        order1.setTimestamp(System.currentTimeMillis());
        publishOrder("orders", order1);
        Thread.sleep(500);
        OrderMessage order2 = new OrderMessage();
        order2.setOrderId("ORD002");
        order2.setUserId("USER002");
        order2.setAmount(199.99);
        order2.setTimestamp(System.currentTimeMillis());
        publishOrder("orders", order2);
        Thread.sleep(2000);
        executor.shutdown();
    }
}

运行说明

启动Redis

# 本地启动Redis服务
redis-server
# 或使用Docker
docker run -d -p 6379:6379 redis

注意事项

  1. Redis版本:确保使用Redis 2.0以上版本
  2. 网络连接:默认连接localhost:6379,根据实际情况修改
  3. 线程安全:Jedis实例不是线程安全的,每个线程使用独立的连接
  4. 连接池:生产环境建议使用连接池
  5. 消息持久化:Pub/Sub不会持久化消息,订阅者离线会丢失消息

优化建议

// 使用连接池
import redis.clients.jedis.JedisPool;
import redis.clients.jedis.JedisPoolConfig;
public class RedisPoolConfig {
    private static final String REDIS_HOST = "localhost";
    private static final int REDIS_PORT = 6379;
    public static JedisPool createPool() {
        JedisPoolConfig config = new JedisPoolConfig();
        config.setMaxTotal(10);
        config.setMaxIdle(5);
        config.setMinIdle(2);
        config.setTestOnBorrow(true);
        config.setTestOnReturn(true);
        return new JedisPool(config, REDIS_HOST, REDIS_PORT);
    }
}

这个完整的案例涵盖了Redis发布订阅的基础用法和高级应用,你可以根据实际需求进行调整和扩展。

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