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

Scala中IOResult.wasSuccessful方法废弃后的替代实现咨询

解决IOResult.wasSuccessful废弃问题

根据Akka 2.6.0+的变更,IOResult.wasSuccessful方法已被废弃——因为IOResult.status现在始终是Success(Done),操作失败会直接触发Future进入Failure状态,不再通过IOResult返回错误。

针对你的代码,修改方式如下:

  1. 移除所有对sftpResult.wasSuccessful的条件判断,只要Future进入Success分支,就代表SFTP操作已成功。
  2. 原代码中!sftpResult.wasSuccessful分支的错误逻辑,现在会被Failure分支捕获,直接合并到该分支处理。

修改后的完整代码:

override def process(changedAfter: Option[ZonedDateTime], changedBefore: Option[ZonedDateTime])
                    (implicit fc: FlowContext): Future[(IOResult, Seq[Operation.Value])] = {
  val now = nowUtc
  val accountUpdatesFromDeletions: Source[AccountUpdate, NotUsed] =
    dbService
      .getDeleteActions(before = now)
      .map(deletionEventToAccountUpdate)

  val result = getAccountUpdates(changedAfter, changedBefore)
    .merge(getErrorUpdates)
    .merge(accountUpdatesFromDeletions)
    .mapAsync(4)(writer.writePSVString)
    .viaMat(creditBureauService.sendUpdate)(Keep.right)
    .mapAsync(4)(au =>
      for {
        _ <- dbService.performUpdate(au)
        _ <- performActionsDelete(now, au)
      } yield au.operation
    )
    .toMat(Sink.seq)(Keep.both)
    .withAttributes(ActorAttributes.supervisionStrategy(decider))
    .run()

  tupleFutureToFutureTuple(result) andThen {
    case Success((_, updateList)) =>
      val total = updateList.size
      val deleted = updateList.count(_ == Operation.DELETED)
      val updated = updateList.count(_ == Operation.UPDATED)
      val inserted = updateList.count(_ == Operation.INSERTED)
      log.info(s"SUCCESS! Uploaded $total accounts to Equifax.")
      log.info(s"There were $deleted deletions, " +
        s"$updated updates and " +
        s"$inserted insertions to the database")
      monitor.gauge("upload.process.batch.successful.total", total)
      monitor.gauge("upload.process.batch.successful.deleted", deleted)
      monitor.gauge("upload.process.batch.successful.updated", updated)
      monitor.gauge("upload.process.batch.successful.inserted", inserted)
    case Failure(e: Throwable) =>
      val sw = new StringWriter
      e.printStackTrace(new PrintWriter(sw))
      log.error(sw.toString)
      // 区分SFTP失败和其他失败的监控指标
      if (e.getMessage.contains("SFTP") || e.getClass.getSimpleName.contains("SFTP")) {
        monitor.gauge("upload.process.batch.sftp.failed", 1)
      } else {
        monitor.gauge("upload.process.batch.failed", 1)
      }
  }
}

关键修改说明

  • 移除冗余的wasSuccessful条件判断,简化分支逻辑。
  • 原SFTP失败的错误处理合并到Failure分支,通过异常信息区分SFTP相关失败和其他类型失败,保留原有监控指标的区分逻辑。
  • 不再使用sftpResult.getError(),因为操作失败时异常会直接通过Future抛出,无需从IOResult获取。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 21:35:18