本文目录导读:

我来为你展示一个完整的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
注意事项
- Redis版本:确保使用Redis 2.0以上版本
- 网络连接:默认连接localhost:6379,根据实际情况修改
- 线程安全:Jedis实例不是线程安全的,每个线程使用独立的连接
- 连接池:生产环境建议使用连接池
- 消息持久化: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发布订阅的基础用法和高级应用,你可以根据实际需求进行调整和扩展。