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
相关产品推荐
相关产品推荐

