如何用OkHttp接收服务端推送的无限分块数据
问题描述
我有一个HTTP服务器,客户端发送请求后,服务器会保持连接并持续推送无限的分块字符串。我知道现在用WebSocket更合适,但这是旧项目,没法修改服务端代码。
服务端代码(server.kt):
// server.kt package com.example.long_http import io.vertx.core.AbstractVerticle import io.vertx.core.Promise import io.vertx.core.Vertx class MainVerticle : AbstractVerticle() { override fun start(startPromise: Promise<Void>) { vertx .createHttpServer() .requestHandler { req -> var i = 0 req.response().setChunked(true).putHeader("Content-Type", "text/plain") val timer = vertx.setPeriodic(2000) { req.response().write("hello ${System.currentTimeMillis()}") println("write ${System.currentTimeMillis()}") } req.response().closeHandler { vertx.cancelTimer(timer) println("close") } } .listen(8888) { http -> if (http.succeeded()) { startPromise.complete() println("HTTP server started on port 8888") } else { startPromise.fail(http.cause()); } } } } fun main() { Vertx.vertx().deployVerticle(MainVerticle()) }
我尝试用OkHttp接收这些分块字符串,但没能成功。
客户端代码(client.kt):
// client.kt package com.example.long_http import okhttp3.* import java.io.IOException fun main() { val client = OkHttpClient() val request = Request.Builder().url("http://localhost:8888").build() client.newCall(request).enqueue(handler()) } class handler : Callback { override fun onFailure(call: Call, e: IOException) { e.printStackTrace() } override fun onResponse(call: Call, response: Response) { println("onResponse") val stream = response.body!!.byteStream().bufferedReader() while (true) { var line = stream.readLine() println(line) } } }
解决方案
问题出在客户端的读取逻辑:服务端推送的分块内容没有换行符,而readLine()会一直等待换行符或流关闭,导致无法实时获取推送内容。
修正方案1:字节流读取
改用逐字节读取可用内容的方式,无需依赖换行符:
// 修改后的client.kt package com.example.long_http import okhttp3.* import java.io.IOException import java.io.InputStream fun main() { val client = OkHttpClient() val request = Request.Builder().url("http://localhost:8888").build() client.newCall(request).enqueue(object : Callback { override fun onFailure(call: Call, e: IOException) { e.printStackTrace() } override fun onResponse(call: Call, response: Response) { println("onResponse") val inputStream: InputStream = response.body!!.byteStream() val buffer = ByteArray(1024) var bytesRead: Int try { while (inputStream.read(buffer).also { bytesRead = it } != -1) { val content = String(buffer, 0, bytesRead) print(content) // 服务端无换行,用print避免多余换行 } } catch (e: IOException) { e.printStackTrace() } finally { response.body?.close() } } }) // 保持主线程存活,避免程序直接退出 Thread.currentThread().join() }
修正方案2:用OkHttp Source读取分块
利用OkHttp的source()直接读取原始分块数据,更贴合HTTP分块传输的特性:
// 修改后的client.kt package com.example.long_http import okhttp3.* import java.io.IOException fun main() { val client = OkHttpClient() val request = Request.Builder().url("http://localhost:8888").build() client.newCall(request).enqueue(object : Callback { override fun onFailure(call: Call, e: IOException) { e.printStackTrace() } override fun onResponse(call: Call, response: Response) { println("onResponse") val source = response.body!!.source() try { while (!source.exhausted()) { val chunk = source.readUtf8() println("收到推送: $chunk") } } catch (e: IOException) { e.printStackTrace() } finally { response.body?.close() } } }) // 保持主线程存活 Thread.currentThread().join() }
额外注意点
- 原客户端单独定义
handler类无必要,改用匿名内部类更简洁。 - 必须保持主线程存活,否则OkHttp的异步回调线程结束后程序会直接退出。
内容的提问来源于stack exchange,提问作者Heo
相关产品推荐
相关产品推荐

