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

Scala Akka HTTP报响应实体100秒未订阅错误排查

Scala Akka HTTP 外部接口调用订阅超时问题排查

问题现象

  • 使用Postman请求目标外部服务接口,9秒内可获取完整返回结果,数据量约9000行JSON
  • 同接口通过Scala应用请求时,虽收到200 OK响应状态,但抛出如下警告日志:
[WARN] [06/09/2022 18:05:45.765] [default-akka.actor.default-dispatcher-9] [default/Pool(shared->http://ad-manager-api-production.ap-south-1.elasticbeanstalk.com:80)] [4 (WaitingForResponseEntitySubscription)] Response entity was not subscribed after 100 seconds. Make sure to read the response entity body or call `discardBytes()` on it. GET /admin/campaigns Empty -> 200 OK Chunked
  • 已尝试将response-entity-subscription-timeout配置参数调高到100秒,问题未修复。

核心实现代码

private val sendAndReceive = customSendAndReceive.getOrElse(HttpClientUtils.singleRequest)
...
def getActiveCampaigns: GetActiveCampaigns = () => {
    val request = HttpRequest(
      uri = s"$endpoint/admin/campaigns?status=PUBLISHED", // includes both PUBLISHED_READY and PUBLISHED_PAUSED
      method = HttpMethods.GET,
      headers = heathers
    )

    sendAndReceive(request).timed(getActiveCampaignsTimer).flatMap {
      case HttpResponse(StatusCodes.OK, _, entity, _) =>
        Unmarshal(entity).to[List[CampaignListDetailsDto]]
      case response@HttpResponse(_, _, _, _) =>
        response.discardEntityBytes()
        Future.failed(new RuntimeException(s"Ad manager service exception: $response"))
      case response =>
        log.error(s"Error calling ad manager service: $response")
        response.discardEntityBytes()
        Future.failed(new RuntimeException(s"Ad manager service exception: $response"))
    }
}
...
def getCampaignSpendData(getActiveCampaigns: GetActiveCampaigns, getCampaignTotalSpend: GetCampaignTotalSpend)(implicit ec: ExecutionContext): GetCampaignsSpendData = () => {
    getActiveCampaigns()
      .andThen {
        case Failure(t) => log.error("Failed to fetch ads from ad manager", t)
      }
      .flatMap {
        campaignList => Future.sequence(campaignList.map(campaign => budgetSpendPercentage(getCampaignTotalSpend)(campaign)))
      }
}

问题解答

1. 错误日志的具体含义

这个警告不代表网络层面未能拉取完整响应数据:

  • 日志中出现200 OK说明Akka HTTP已经成功和服务端建立连接、收到了完整的响应头,此时响应体是以分块(Chunked)流的形式暂存在底层,不会主动拉取。
  • Akka HTTP是基于响应式流实现的HTTP客户端,响应实体是懒加载的:收到响应头后,必须有代码主动消费(订阅)响应体,要么通过Unmarshal解析成业务对象,要么通过discardBytes()主动丢弃,否则底层不会拉取响应体数据。
  • 这个警告的触发逻辑是:Akka等了配置的100秒,始终没有等到任何代码订阅这个响应实体,为了避免连接泄漏,主动断开连接、丢弃未消费的响应体。
    你调高response-entity-subscription-timeout无效的核心原因是:问题不是等待订阅的时间太短,而是代码逻辑根本没有走到订阅响应体的路径,等再久也不会有订阅动作。

2. 排查与修复方案

按照优先级从高到低排查:

排查自定义timed算子的实现问题

这是最高概率的故障点:

  • 绝大多数自定义的请求计时算子如果实现不当,会在拿到响应头的瞬间就结束计时、返回结果,甚至会错误消费响应实体流、或者返回的对象不是原始携带可订阅实体的HttpResponse,导致后续的Unmarshal操作根本拿不到真实的响应实体,自然不会触发订阅。
  • 先临时注释掉.timed(getActiveCampaignsTimer)这一段,直接调用sendAndReceive(request).flatMap{...}测试,如果警告消失,就修正timed算子的实现:保证算子透传原始HttpResponse对象,不要提前消费实体,计时范围要覆盖「发送请求-拉取完整响应体-解析完成」的全流程,而不是仅统计拿到响应头的时间。

排查自定义sendAndReceive封装的问题

你当前用的是自定义sendAndReceive和默认实现的二选一逻辑,自定义封装很容易出现实体透传问题:

  • 比如在自定义逻辑里提前调用了toStrict、entity.dataBytes.runReduce等方法消费了响应实体,但没有把消费后生成的新实体向下传递。
  • 临时替换为Akka HTTP原生的Http().singleRequest(request)测试,如果问题消失,就检查自定义sendAndReceive的实体透传逻辑。

验证流物化是否正常

如果前两步排查后问题仍然存在,说明响应体流没有被正常物化启动,先在200 OK的处理分支加日志确认代码确实走到了Unmarshal逻辑,再将分支代码改为主动拉取全量响应体为严格实体再解析:

case HttpResponse(StatusCodes.OK, _, entity, _) =>
  // 30秒超时足够拉取9秒就能返回的接口数据
  entity.toStrict(30.seconds).flatMap { strictEntity =>
    Unmarshal(strictEntity).to[List[CampaignListDetailsDto]]
  }

toStrict会显式订阅响应流、拉取所有分块数据合并为内存中的完整实体,如果这一步能正常返回数据,说明之前的Unmarshal流程存在执行上下文问题——比如传入的ExecutionContext不是Akka流兼容的调度器,导致响应流根本没有被启动。

冗余代码清理

你代码中第三个匹配分支case response =>永远不会触发:Akka HTTP的singleRequest返回的Future成功值一定是HttpResponse类型,所有请求错误都会通过Future失败路径返回,这个分支属于冗余代码,可以删除避免逻辑误导。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 10:24:13