ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

使用HttpURLConnection调用SSE采坑记录

使用HttpURLConnection调用SSE采坑记录 背景本系统为客服系统对接了后端智能辅助系统当消费者发文字时问题本系统会将文字发给辅助系统辅助系统会返回相关回答供客服参考。过程如下消费者发送问题文字给客服消费者问题文字推送给客服上屏客服前端调用本系统接口将问题发过去本系统将调用透传给智能辅助系统并将结果回传给前端3、4两步均采用的SSE本系统调用智能辅助时纯透传没有额外逻辑。本来我们系统基于Spring的RestTemplate写了一套类似feign的框架简化调用后端。因为RestTemplate不能很好的处理SSE(RestTemplate返回时会关闭连接所以一次调用必须读取完整的响应那就没办法实现流式效果)因此本次调用不能使用之前的框架但是我们又不想引入新的开源组件因此决定使用JDK自带的HttpURLConnection进行远程调用。备注:流式效果是指像deepseek那样一个字一个字的返回给前端上屏所以后台只要获取到内容就得写回给前端而不是像通常的接口后台获取到完整的结果后一次写回前端。# 代码有简化publicvoidquireSopAnswer(Paramparam,HttpServletResponseresponse){setHead(response);HttpURLConnectionhttpURLConnectiongetHttpURLConnection(param);InputStreaminputStreamhttpURLConnection.getInputStream();OutputStreamoutputStreamresponse.getOutputStream()while(true){byte[]tmpnewbyte[8192];intreadinputStream.read(tmp);outputStream.write(tmp,0,read);outputStream.flush();if(read0){break;}}}该代码测试的时候发现流式效果不是很明显感觉所有内容都是挤在一起上屏的。后面参考隔壁项目组的写法能很好的实现流式的效果。# 代码有简化publicvoidquireSopAnswer(HttpServletRequestrequest,HttpServletResponseresponse){setHead(response);AsyncContextasyncContextrequest.startAsync(request,response);getWebClient().post().retrieve().bodyToFlux(newParameterizedTypeReferenceServerSentEventString(){}).subscribe(eventSourcet-receiveAnswerAndSend(eventSourcet,asyncContext),err-LOGGER.error(err),asyncContext::complete);}privatevoidreceiveAnswerAndSend(ServerSentEventStringeventSourcet,AsyncContextasyncContext){HttpServletResponseresp(HttpServletResponse)asyncContext.getResponse();StringBuildersbnewStringBuilder();sb.append(event: ).append(eventSourcet.event()).append(\n);sb.append(data: ).append(eventSourcet.data()).append(\n);sb.append(\n);resp.getWriter().write(sb.toString());resp.getWriter().flush();}既然代码不符合预期那就要搞清楚根因是啥。定位过程因为已经有正确的写法效果好的写法为了快速确认原因因此首先通过对比差异针对差异点进行论证分析。如果此方法仍未找到合理的原因那就得根据第一性原理先找到前端为什么一大段文字一起上屏然后逐步追踪线索找到最终根因。分析一先快速定界确认我们后台收到后端的数据是不是本来就是挤在一起的增加日志每次read的时候打印read的字节数以及内容。确认是不是后台自己一次将内容都挤在一次返回的导致流式效果不明显。结果从日志来看确实是流式输出而不是所有内容挤在一起发送的。而且我们也看了后端的代码他们也确实是一个一个的event独立返回的并不会缓冲等着一起返回。2025-11-1515:59:42,300SseClient:200-------readed number-----2025-11-1515:59:42,301SseClient:201-282025-11-1515:59:42,302SseClient:203-------Date-----2025-11-1515:59:42,303SseClient:204-event:requestId[N]data:{\conte2025-11-1515:59:42,304SseClient:200-------readed number-----2025-11-1515:59:42,306SseClient:201-1272025-11-1515:59:42,307SseClient:203-------Date-----2025-11-1515:59:42,308SseClient:204-nt\:\373BEB7200DDB2BABF0AC14DE6E6383F9F847C1FC85C422A\,\conversionId\:\a8bef8a85f1ed27fc3fff22de4747a412106\,\fallback\:null}[N][N]2025-11-1515:59:42,323SseClient:200-------readed number-----2025-11-1515:59:42,324SseClient:201-282025-11-1515:59:42,325SseClient:203-------Date-----2025-11-1515:59:42,327SseClient:204-event:handle_thoughts[N]data:{2025-11-1515:59:42,328SseClient:200-------readed number-----2025-11-1515:59:42,329SseClient:201-2002025-11-1515:59:42,330SseClient:203-------Date-----2025-11-1515:59:42,331SseClient:204-\content\:\1\,\conversionId\:\a8bef8a85f1ed27fc3fff22de4747a412106\,\fallback\:null}[N][N]event:handle_thoughts[N]data:{\content\:\.\,\conversionId\:\a8bef8a85f1ed27fc3fff22de4747a412106\,\fallback\:null}[N][N]2025-11-1515:59:42,351SseClient:200-------readed number-----2025-11-1515:59:42,353SseClient:201-282025-11-1515:59:42,354SseClient:203-------Date-----2025-11-1515:59:42,355SseClient:204-event:handle_thoughts[N]data:{2025-11-1515:59:42,356SseClient:200-------readed number-----2025-11-1515:59:42,357SseClient:201-2052025-11-1515:59:42,358SseClient:203-------Date-----2025-11-1515:59:42,359SseClient:204-\content\:\ 发\,\conversionId\:\a8bef8a85f1ed27fc3fff22de4747a412106\,\fallback\:null}[N][N]event:handle_thoughts[N]data:{\content\:\送\,\conversionId\:\a8bef8a85f1ed27fc3fff22de4747a412106\,\fallback\:null}[N][N]2025-11-1515:59:42,361SseClient:200-------readed number-----2025-11-1515:59:42,362SseClient:201-282025-11-1515:59:42,363SseClient:203-------Date-----2025-11-1515:59:42,364SseClient:204-event:handle_thoughts[N]data:{2025-11-1515:59:42,365SseClient:200-------readed number-----2025-11-1515:59:42,366SseClient:201-862025-11-1515:59:42,367SseClient:203-------Date-----2025-11-1515:59:42,368SseClient:204-\content\:\H\,\conversionId\:\a8bef8a85f1ed27fc3fff22de4747a412106\,\fallback\:null}[N][N]......分析二分析两种写法的差异通过修改第一种写法保证2者行为一致并测试效果来确认是那个差异导致的问题个人看到两者明显的差异是第二种写法一次返回一个完整的event。而写法一通过日志可以看到一个完整的event被拆开了分好几次写回前端。首先先分析了产生该现象的原因然后在修改第一种写法保证它也是一次返回完整的event然后测试效果。第一种写法为什么本系统收到后端的event被拆成好几部分而不是一次读取到完整的event。先确认是不是后端问题即后端持续流式输出的时候本身就不是一次返回完整的event。为了确认该问题首先找后台要到了他们的代码进行分析他们代码确实是一次写回完整的event。然后我们也通过抓包确认通过中间件后该行为并未改变我们并不是直接调用后台服务器中间还经过了几个中间件的转发。通过抓包也确认我们系统每次接受到的报文也确实是一个完整的event。根据我掌握的知识来看底层获取的tcp报文在应用层read的时候不会被拆分。因为底层的报文都是完整的放到缓冲区供应用层read除非应用层传的接收数组比缓存区当前已缓冲的数据要小一次读取不完导致截断。本系统并不属于该场景因此该行为应该是HttpURLConnection自己的行为因此需要对源码进行分析。另外从报文里面也发现了其他的疑点响应采用的是trunked编码(Transfer-Encoding: chunked)该编码每个分片前面都有长度字段代表本次响应内容的长度如上述截图中的9b “e4”。如果HttpURLConnection不解析HTTP协议即对底层的字节流完全不做处理那么前端收到的报文会无端多出代表chunked编码长度的字段而导致内容语法错误。为此分析源码确认了以上2点疑点。通过阅读代码发现HttpURLConnection会解析HTTP协议如果是Transfer-Encoding: chunked他会按照该编码格式解析数据并过滤掉chunked编码中多余的块长度字段。上面的日志也证明了这一点。sun.net.www.protocol.http.HttpURLConnection#getInputStreamsun.net.www.protocol.http.HttpURLConnection#getInputStream0sun.net.www.http.HttpClient#parseHTTPsun.net.www.http.HttpClient#parseHTTPHeader另外也通过阅读ChunkedInputStream的public synchronized int read(byte b[], int off, int len)发现为啥底层tcp一帧报文HttpURLConnection需要read多次而不是一次就read就返回全部按说数据在缓冲区都是ready的不存在读阻塞一说。publicChunkedInputStream(InputStreamin,HttpClienthc,MessageHeaderresponses)throwsIOException{...# 构造函数state 被初始化为STATE_AWAITING_CHUNK_HEADERstateSTATE_AWAITING_CHUNK_HEADER;}publicsynchronizedintread(byteb[],intoff,intlen)throwsIOException{......# 省略了无关的代码intavailchunkCount-chunkPos;if(avail0){if(stateSTATE_READING_CHUNK){returnfastRead(b,off,len);}# state初始状态为STATE_AWAITING_CHUNK_HEADER代码会走这里 availreadAhead(true);if(avail0){return-1;/* EOF */}}......returncnt;}privateintreadAhead(booleanallowBlocking)throwsIOException{......# 省略无关代码if(allowBlocking){# 走的这个分支returnreadAheadBlocking();}else{returnreadAheadNonBlocking();}}privateintreadAheadBlocking()throwsIOException{# 代码有省略do{....../* * We must read into the raw buffer so make sure there is space * available. We use a size of 32 to avoid too much chunk data * being read into the raw buffer. */# 重点在这里第一次读时只会只读32字节方便先解析出chunk的大小然后在逐步读取剩下的大小 # 因为此时还不知道单个chunk的大小buff设多大合适并不知道所以使用32byte先尝试读取一下ensureRawAvailable(32);intnread;try{nreadin.read(rawData,rawCount,rawData.length-rawCount);}......rawCountnread;processRaw();}while(chunkCount0);/* * Return the number of chunked bytes available to read */returnchunkCount-chunkPos;}从上面可以看到应用层调用read的时候最终会调用到ChunkedInputStream的read但是ChunkedInputStream初次调用时只会尝试读取32byte弄清楚chunk的大小后在去设置buff以读取chunk剩余的内容。这个与我们日志也能映射上每次读取完整的event时第一次读取总是拿到的是28个字节的内容因为chunked编码的块长度字段会占用4字节(2字节空格换行符)。试验抹除差异后是否有效果publicvoidquireSopAnswer(SopAgentReqsopAgentReq,HttpServletResponseresponse){setHead(response);try(OutputStreamoutputStreamresponse.getOutputStream();BufferedReaderbrnewBufferedReader(newInputStreamReader(getSseInputStream(para),StandardCharsets.UTF_8))){Stringline;StringBuilderbuildernewStringBuilder();while(true){linebr.readLine();if(linenull){break;}# 为了快速验证这里硬编码了请忽略这里的不优雅 builder.append(line).append(\n);# 读取event行 builder.append(line).append(\n);# 读取data行 builder.append(line).append(\n);# 读取空行 outputStream.write(builder.toString().getBytes(StandardCharsets.UTF_8));outputStream.flush();}}}验证发现仍然没有效果看来并不是该差异导致的问题。想来也是SSE是业界标准的协议只要按照标准格式返回数据客户端应该是能好好处理这种场景的netty中俗称的粘包。此路不通我观察到第二种写法用的servlet的startAsync。经验证也不是该差异导致的问题。此次通过差异点确认问题所在已基本走进死胡同了因此整理思路从源头出发逐步分析。分析三分析前端为什么会整段文字一起上屏为了排除前端库的问题(可能前端处理粘包问题处理的不好导致内容没办法及时丢给应用层代码最终文字挤在了一起或者换个说法应用层瞬间拿到了大段的文字)经对event-source-plus组件的初步分析排除此问题。因此自然而然的想确认前端是如何接收到数据的。通过F12观察到很多event是同一时间接收到的。难道是网速太快导致数据最终到达前端时挤在了一起于是后台每次write的时候我都sleep 100ms进行降速在观察效果。最终观察到的效果并无区别前端还是大段文字一起上屏而且从F12看到还是大量event在同一时间被接受到。此时我感觉像是每次flush并没生效而是多次write后缓存区满了后一次写回前端。对于此疑问我做了两件事进行确认同时打开浏览器F12及后台日志监控一边观察前端F12中接收的数据一边观察日志打印。结果write的日志不断在打印但是前端F12没有显示有数据接收然后突然前端一下接收大量数据此时我已经基本确认flush没有生效。使用tcpdump监控网络报文确认后台确实一次批量写回数据的。如果flush不生效tomcat会等缓冲区满了后写回数据默认的缓冲区大小是8k这个与tcpdump监控到的网络报文一致。servlet的flush函数一定会强制将结果写回的此时未生效我第一想到的就是框架对HttpServletResponse做了包装并重写了OutputStream的flush方法因此打断点确认。我们使用的CXF框架因此代码里面拿到的是CXF包装后的HttpServletResponseFilterOutputStream为ServletOutputStreamFilter其flush方法为空至此所有疑问得以解决。publicclassHttpServletResponseFilterextendsHttpServletResponseWrapper{# 代码有省略OverridepublicServletOutputStreamgetOutputStream()throwsIOException{returnnewServletOutputStreamFilter(super.getOutputStream(),m);}}试验将最开始的代码搬到原生servlet写的接口中进行验证。验证结果ok。
返回列表