Scala如何让Future[Iterable]对应的POST请求按顺序串行执行
解决代码
import scala.concurrent.Future // 注意请替换为你业务中实际使用的ExecutionContext,不要直接用全局默认 import scala.concurrent.ExecutionContext.Implicits.global val totalProcessResult: Future[Unit] = preparePartitioningDataset(dataset, "sdp_id").flatMap { partitions => // 先将Iterator转为List,固定所有分区值,避免重复遍历问题 val partitionList = partitions.toList // 用foldLeft拼接Future链实现串行执行 partitionList.foldLeft(Future.successful(())) { (prevCompleted, currentPartition) => // 等待前一个分区的全流程执行完成,再启动当前分区的逻辑 prevCompleted.flatMap { _ => val whereStatement = s"SDP_ID = '$currentPartition'" // 执行当前分区的查询 datasetService.getFullDatasetResultIterable( dataset = dataset, format = format._1, limit = none[Int], where = whereStatement.some ).flatMap { queryResult => // 查询完成后执行通知,等通知返回后再进入下一个分区的处理 notifyPartitionedDataset( bearerToken = bearerToken, endpoint = endpoint, dataset = queryResult ) } } } }
实现原理
- 核心是通过
foldLeft将所有分区的处理逻辑拼接为一条链式Future,初始值为已经完成的空FutureFuture.successful(()) - 每一步的查询、通知逻辑都包裹在
flatMap内部,只有前一步Future执行完成后才会创建当前步骤的Future,从根源上避免了并行启动 - 最终返回的
totalProcessResult是整个串行流程的总Future,可根据业务需要添加成功/异常回调处理整体执行结果
注意事项
- 若分区数量极大,转List会占用少量内存,属于串行执行的正常取舍
- 不要在
flatMap外部提前创建getFullDatasetResultIterable的Future实例,否则Future会立刻启动执行,仍会出现并行问题 - 默认逻辑下如果某一步查询或通知失败,整个流程会直接终止,后续任务不再执行;如果需要失败后继续执行下一个分区,可在对应Future后添加
recover逻辑处理异常
内容的提问来源于stack exchange,提问作者Tizianoreica
相关产品推荐
相关产品推荐

