You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Boot中基于Redis实现跨服务阻塞式消息获取的问题

基于Redis的跨服务阻塞式通信实现问题排查与修复

需求背景

  • Service A(Spring Boot):需暴露/api-publish和/api-consume两个端点,前者向Redis主题topic1发布消息,后者需阻塞等待Service B处理完成后,从topic2收到通知再返回响应
  • Service B(ExpressJS):监听topic1的消息,处理完成后向topic2发送"已处理完成"的通知

现有代码核心问题

  1. 端点路径不匹配:当前控制器提供的是/api/v1/public/test,不符合需求的/api-consume,且缺少/api-publish端点
  2. RedisTemplate泛型不明确:配置的是通配符泛型RedisTemplate<?,?>,但注入时用RedisTemplate<String, Object>,会引发类型转换警告
  3. 消息串扰风险:RedisMessageSubscriber是单例Bean,所有请求共用同一个BlockingQueue,会导致不同请求的处理结果互相干扰
  4. 无请求标识关联:未在消息中携带唯一请求ID,Service B无法针对性返回处理结果,订阅者无法区分消息归属

修复后的完整实现

1. Redis配置类(修正泛型与命名)

@Configuration
public class RedisConfig {

    @Value("${spring.data.redis.host}")
    private String redisHost;

    @Value("${spring.data.redis.port}")
    private int redisPort;

    @Bean
    public RedisConnectionFactory lettuceConnectionFactory() {
        return new LettuceConnectionFactory(new RedisStandaloneConfiguration(redisHost, redisPort));
    }

    // 明确泛型为<String, String>,匹配字符串消息场景
    @Bean
    public RedisTemplate<String, String> redisTemplate() {
        RedisTemplate<String, String> template = new RedisTemplate<>();
        template.setConnectionFactory(lettuceConnectionFactory());
        template.setHashKeySerializer(new StringRedisSerializer());
        template.setKeySerializer(new StringRedisSerializer());
        template.setHashValueSerializer(new StringRedisSerializer());
        template.setValueSerializer(new StringRedisSerializer());
        return template;
    }

    @Bean
    MessageListenerAdapter messageListener() {
        return new MessageListenerAdapter(new RedisMessageSubscriber());
    }

    @Bean
    ChannelTopic topic2() {
        return new ChannelTopic("topic2");
    }

    @Bean
    RedisMessageListenerContainer redisContainer() {
        RedisMessageListenerContainer container = new RedisMessageListenerContainer();
        container.setConnectionFactory(lettuceConnectionFactory());
        container.addMessageListener(messageListener(), topic2());
        return container;
    }
}

2. 消息订阅服务(解决并发串扰)

通过ConcurrentHashMap维护请求ID与专属阻塞队列的映射,确保每个请求只接收自己的处理结果:

@Service
public class RedisMessageSubscriber implements MessageListener {
    // 键:请求唯一ID,值:对应请求的阻塞队列
    private final ConcurrentHashMap<String, BlockingQueue<String>> requestQueueMap = new ConcurrentHashMap<>();

    @Override
    public void onMessage(Message message, byte[] pattern) {
        String receivedMsg = new String(message.getBody());
        // 按约定格式拆分请求ID和处理结果
        String[] msgParts = receivedMsg.split(":", 2);
        if (msgParts.length == 2) {
            String requestId = msgParts[0];
            String result = msgParts[1];
            BlockingQueue<String> targetQueue = requestQueueMap.get(requestId);
            if (targetQueue != null) {
                targetQueue.offer(result);
                requestQueueMap.remove(requestId); // 清理已完成的请求映射,避免内存泄漏
            }
        }
    }

    // 根据请求ID等待处理结果
    public String waitForResult(String requestId, long timeout, TimeUnit unit) throws InterruptedException {
        BlockingQueue<String> queue = new LinkedBlockingQueue<>();
        requestQueueMap.put(requestId, queue);
        try {
            String result = queue.poll(timeout, unit);
            if (result == null) {
                requestQueueMap.remove(requestId); // 超时后清理无效映射
            }
            return result;
        } catch (InterruptedException e) {
            requestQueueMap.remove(requestId); // 中断后清理映射
            Thread.currentThread().interrupt();
            throw e;
        }
    }

    // 生成唯一请求ID
    public String generateRequestId() {
        return UUID.randomUUID().toString();
    }
}

3. 控制器(实现需求端点)

@RestController
public class MessageController {
    private final RedisTemplate<String, String> redisTemplate;
    private final RedisMessageSubscriber subscriber;

    // 构造注入替代@Autowired,符合Spring最佳实践
    public MessageController(RedisTemplate<String, String> redisTemplate, RedisMessageSubscriber subscriber) {
        this.redisTemplate = redisTemplate;
        this.subscriber = subscriber;
    }

    @PostMapping("/api-publish")
    public ResponseEntity<String> publishMessage(@RequestBody String content) {
        String requestId = subscriber.generateRequestId();
        // 消息格式:requestId:业务内容,让Service B处理后携带相同ID返回
        String message = requestId + ":" + content;
        redisTemplate.convertAndSend("topic1", message);
        return ResponseEntity.ok("消息已发布,请求ID:" + requestId);
    }

    @GetMapping("/api-consume")
    public ResponseEntity<String> consumeResult(@RequestParam String requestId) {
        try {
            String result = subscriber.waitForResult(requestId, 5, TimeUnit.SECONDS);
            if (result != null) {
                return ResponseEntity.ok("收到Service B处理结果:" + result);
            } else {
                return ResponseEntity.status(HttpStatus.REQUEST_TIMEOUT).body("等待处理结果超时");
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            return ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR).body("获取处理结果失败");
        }
    }
}

4. Service B的ExpressJS核心实现

const redis = require('redis');
// 创建订阅者和发布者客户端
const subscriber = redis.createClient({ host: '你的Redis地址', port: 6379 });
const publisher = redis.createClient({ host: '你的Redis地址', port: 6379 });

subscriber.subscribe('topic1');

// 监听topic1消息并处理
subscriber.on('message', (channel, message) => {
    const [requestId, content] = message.split(':', 2);
    // 模拟业务处理逻辑
    console.log(`处理业务内容:${content}`);
    // 处理完成后向topic2发送带请求ID的通知
    publisher.publish('topic2', `${requestId}:处理完成,内容:${content}`);
});

关键修正总结

  1. 对齐需求的端点路径,补全/api-publish实现
  2. 引入请求ID机制,解决多请求场景下的消息串扰问题
  3. 明确RedisTemplate泛型,消除类型转换风险
  4. 优化订阅者的队列管理,避免内存泄漏

内容的提问来源于stack exchange,提问作者ginbarca

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.20 21:14:57