Spring Boot集成Twilio WebSocket与Azure认知服务:音频回复重复问题排查求助
问题分析与解决方案
你的核心问题确实是多线程并发访问共享资源payloads导致的重复处理。Spring Boot的WebSocket处理是多线程的,handleTextMessage会被不同线程调用,而你当前的代码没有对payloads的操作做同步控制,这就会出现你怀疑的情况:在Azure处理的同时,其他线程仍在往payloads里添加新消息,导致清空操作前已经混入了新数据,最终重复播放。
具体代码问题拆解
- 非线程安全的集合操作:如果
payloads是普通的ArrayList,add()和size()判断不是原子操作。当线程A判断size>200开始处理时,线程B可能已经完成了add(),导致线程A处理的集合里包含了线程B刚加的元素,而清空payloads后,线程B的元素可能还没被处理(或者部分被处理),后续又会被重复处理。 - 全局
outputStream的污染:如果outputStream是全局变量,多线程处理时会导致不同批次的音频数据混在一起,进一步加剧重复播放的问题。 - 清空集合的时机问题:在
handleTextMessage里直接payloads.clear(),无法保证清空操作和其他线程的add()操作之间的原子性,会导致数据遗漏或重复。
修复方案
1. 使用线程安全的集合存储payload
把payloads换成ConcurrentLinkedQueue(线程安全的队列),避免并发修改异常,同时方便原子性地批量取出元素。
2. 重构消息处理逻辑,原子性批量处理
修改handleTextMessage,当队列元素达到阈值时,一次性取出所有元素进行处理,而不是判断size后再操作,彻底避免间隙问题。
3. 避免全局流,每次处理创建新的输出流
确保每个批次的音频处理都使用独立的ByteArrayOutputStream,防止不同批次的数据交叉污染。
修改后的代码示例
// 替换原有的payloads为线程安全的队列 private final ConcurrentLinkedQueue<String> payloads = new ConcurrentLinkedQueue<>(); // 如果sessions是普通Map,也建议换成线程安全的ConcurrentHashMap private final ConcurrentHashMap<WebSocketSession, YourAudioProcessor> sessions = new ConcurrentHashMap<>(); @Override public void handleTextMessage(WebSocketSession session, TextMessage message) throws IOException, InterruptedException { JsonNode request = jsonMapper.readTree(message.getPayload()); if (request.path("media").path("track").asText().equals("inbound")) { String base64EncodedAudio = request.path("media").path("payload").asText(); payloads.add(base64EncodedAudio); // 当元素数量达到200时,原子性取出所有元素批量处理 if (payloads.size() >= 200) { List<String> batch = new ArrayList<>(); String payload; // 一次性取出队列中所有元素,避免其他线程继续添加干扰 while ((payload = payloads.poll()) != null) { batch.add(payload); } if (!batch.isEmpty()) { String response = constructResponse(session, request, batch); session.sendMessage(new TextMessage(response)); } } } } private String constructResponse(WebSocketSession session, JsonNode request, List<String> batch) throws IOException { ObjectMapper objectMapper = new ObjectMapper(); String streamSid = request.path("streamSid").asText(); // 每个批次创建独立的输出流,避免数据污染 ByteArrayOutputStream outputStream = new ByteArrayOutputStream(); YourAudioProcessor processor = sessions.get(session); if (processor == null) { // 处理processor不存在的情况,比如抛出异常或初始化 throw new IllegalStateException("Audio processor not found for session"); } for (String payload : batch) { byte[] decoded = Base64.getDecoder().decode(payload); processor.pushData(decoded); byte[] processedBytes = processor.getBytes(); if (processedBytes != null && processedBytes.length > 0) { outputStream.write(processedBytes); } } byte[] encodedBytes = Base64.getEncoder().encode(outputStream.toByteArray()); return objectMapper.writeValueAsString(new OutBoundMessage("media", new Media(new String(encodedBytes)), streamSid)); }
额外注意事项
- 确认
YourAudioProcessor(也就是sessions.get(session)返回的对象)的线程安全性:如果同一个session的processor会被多个线程调用,也要确保pushData和getBytes方法是线程安全的,或者每个session对应独立的processor实例。 - 如果Azure的语音处理是异步的,可能需要等待处理完成后再获取
processedBytes,否则可能拿到不完整的数据。可以考虑用回调或Future来确保处理完成。
内容的提问来源于stack exchange,提问作者Kyozoku
相关产品推荐
相关产品推荐

