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

Scala STTP异步客户端超时问题及Flink流性能优化咨询

问题描述

出现异常:

java.util.concurrent.TimeoutException: Request timeout to ******* after 60000 ms.

该异常导致Flink流中的一个算子运行极慢,进而造成Kafka消息大量积压。

想请教以下问题:

  1. 超时发生时异步客户端压力陡增的根本原因是什么?是否与异步后端有关?
  2. 我们具备请求失败重试机制,若将超时时间从60秒缩短至3秒,是否能让请求更快进入重试流程,不再占用后端资源?这么做的弊端(如大量请求进入重试)和收益分别是什么?
  3. 有没有更优的异常处理方案?

附相关代码实现:

import java.sql.Timestamp
import java.text.SimpleDateFormat
import java.time.Instant
import java.util.{Calendar, UUID}
import io.circe._

import scala.util.Try
import org.apache.flink.util.concurrent.Executors
import org.apache.flink.streaming.api.scala.async.{AsyncFunction, ResultFuture}
import org.slf4j.{Logger, LoggerFactory}
import sttp.client._
import sttp.client.asynchttpclient.future.AsyncHttpClientFutureBackend
import sttp.model.Uri

import scala.concurrent.{ExecutionContext, Future}
import scala.util.{Failure, Success}

class MySink(prop: Prop)
    extends AsyncFunction[Record, CallDetails] {
  @transient lazy val directExecutor: ExecutionContext =
    ExecutionContext.fromExecutor(Executors.directExecutor())

  implicit lazy val backend = AsyncHttpClientFutureBackend()
  val uri = Uri.apply(prop.Uri())


  override def asyncInvoke(
    record: BatchRecord,
    resultFuture: ResultFuture[CallDetails]
  ): Unit = {

    try {

      val batchUUID = com.datastax.driver.core.utils.UUIDs.timeBased()
      val timeStamp = getTimestampAndDate
      val token = record.sessionDetails

      val future = sendData(
        token.token.tokenString,
        record
      )

      future.onComplete {
        case Success(response) => {
          resultFuture.complete(Iterable(recordCall(response, record, batchUUID, time)))
        }
        case Failure(exception) =>
          resultFuture.complete(Iterable(recordCall(exception, record, batchUUID, time)))
      }(directExecutor)
    } catch {
      case e: Exception =>
        resultFuture.completeExceptionally(e)
    }

  }

  def sendData(
    token: String,
    record: Json,
    testUri: Option[Uri] = None
  ): Future[Response[Either[String, String]]] = {

    val jsonFormatted = record.toString()

    val dataRequest = basicRequest
      .contentType("application/json")
      .auth.bearer(token)
      .body(
        s"""$jsonFormatted"""
      )
      .post((testUri.getOrElse(uri)))

    dataRequest.send()
  }
问题解答

1. 超时引发异步客户端压力陡增的根因

超时发生时压力陡增的核心原因主要有两点:

  • 资源池耗尽:你使用的AsyncHttpClientFutureBackend默认维护连接池和线程池处理异步请求。当目标服务响应变慢,大量请求处于等待超时的状态,会持续占用连接池中的连接和线程池中的线程,新请求无法获取资源只能排队,直接导致客户端侧压力堆积。
  • 请求堆积连锁反应:Flink的AsyncFunction默认有并发限制(若未通过AsyncDataStream.unorderedWait/orderedWait配置并发数),超时请求会占用并发槽位,后续请求无法及时处理,进一步加剧算子处理缓慢;而Kafka消息积压又会反过来推高请求量,形成恶性循环。

这确实和异步后端直接相关——如果异步客户端的资源池(连接、线程)配置与业务吞吐量不匹配,遇到服务超时就会快速耗尽资源,引发压力陡增。

2. 缩短超时时间的收益与弊端

收益

  • 快速释放资源:超时时间从60秒缩至3秒,等待超时的请求会更快失败,及时释放连接、线程等资源,避免资源被长期占用,让新请求能更快获取资源处理。
  • 加快重试节奏:若重试基于失败触发,更快的超时能让失败请求更早进入重试流程,在目标服务恢复时可更快完成数据投递。

弊端

  • 重试风暴风险:如果目标服务只是短暂抖动,大量请求在3秒内超时会触发批量重试,进一步加大目标服务压力,可能导致服务彻底不可用,形成“超时→重试→更超时”的恶性循环。
  • 无效重试增加:若目标服务处于长期故障,缩短超时会导致更多无效重试,浪费客户端资源和网络带宽。
  • 数据重复风险:如果目标服务实际已处理请求但响应超时,重试会导致数据重复投递,需要业务侧具备幂等处理能力。

是否可行需结合以下因素判断:目标服务的故障特性(短暂抖动/长期故障)、重试机制是否有退避策略(如指数退避)、业务是否能容忍重复数据。

3. 更优的异常处理方案

结合你的代码和场景,推荐以下优化方向:

(1)配置异步客户端资源池与超时

显式配置AsyncHttpClientFutureBackend的资源参数,避免默认配置不满足吞吐量需求:

implicit lazy val backend = AsyncHttpClientFutureBackend.usingConfig(
  AsyncHttpClientConfig.Builder()
    .setMaxConnections(200) // 根据业务吞吐量调整
    .setMaxConnectionsPerHost(100)
    .setConnectionTimeout(3000) // 连接超时
    .setRequestTimeout(3000) // 请求超时
    .build()
)

同时在sttp请求中显式设置超时,覆盖后端默认配置:

val dataRequest = basicRequest
  .timeout(3000, 3000) // 连接超时、请求超时
  .contentType("application/json")
  // ... 其他请求配置

(2)优化Flink异步算子的并发与重试

  • 配置AsyncDataStream的并发数,匹配异步客户端的资源池大小:
    AsyncDataStream.unorderedWait(
      inputStream,
      new MySink(prop),
      5000, // 算子超时时间
      TimeUnit.MILLISECONDS,
      100 // 并发数,对应客户端连接池大小
    )
    
  • 给重试机制增加指数退避策略,避免重试风暴:比如第一次重试间隔1秒,第二次2秒,第四次8秒,最多重试3次。可通过封装scala.concurrent.Future实现,或借助Flink的RetryStrategy。

(3)异常分类处理与降级

  • 区分超时异常和其他异常(如网络错误、权限错误):对于权限错误这类无法通过重试解决的异常,直接标记失败并记录,不进入重试流程。
  • 增加降级逻辑:当目标服务超时率超过阈值时,暂时将数据写入备用存储(如Kafka死信队列),待服务恢复后再批量处理,避免堵塞主流程。

(4)代码层面小优化

  • 你的directExecutor使用了Executors.directExecutor(),会在当前线程执行回调,可能阻塞Flink算子线程,建议改用Flink提供的ExecutionContext或单独的线程池处理回调。
  • 避免每次sendData都将record转为字符串,可提前在asyncInvoke中处理,减少重复操作。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 13:45:24