Scala STTP异步客户端超时问题及Flink流性能优化咨询
问题描述
出现异常:
java.util.concurrent.TimeoutException: Request timeout to ******* after 60000 ms.
该异常导致Flink流中的一个算子运行极慢,进而造成Kafka消息大量积压。
想请教以下问题:
- 超时发生时异步客户端压力陡增的根本原因是什么?是否与异步后端有关?
- 我们具备请求失败重试机制,若将超时时间从60秒缩短至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
相关产品推荐
相关产品推荐

