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

如何用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()
}

额外注意点

  1. 原客户端单独定义handler类无必要,改用匿名内部类更简洁。
  2. 必须保持主线程存活,否则OkHttp的异步回调线程结束后程序会直接退出。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 13:50:45