如何为OptionT[Future,_]显式指定ExecutionContext实现Akka SnapshotStore插件
问题背景
我正在编写自定义Akka SnapshotStore插件,当前需要实现如下方法:
def loadAsync(persistenceId: String, criteria: SnapshotSelectionCriteria): Future[Option[SelectedSnapshot]]
目前编写的实现代码如下:
import cats.data.OptionT import cats.implicits._ ... override def loadAsync( persistenceId: String, criteria: SnapshotSelectionCriteria ): Future[Option[SelectedSnapshot]] = { // 与原始插件实现逻辑一致 val metadata = snapshotMetadatas(persistenceId, criteria).sorted.takeRight(maxLoadAttempts) // 需要移除该全局隐式上下文导入! import scala.concurrent.ExecutionContext.Implicits.global val getSnapshotAndReportMetric = for { snapshot <- OptionT.fromOption[Future](getMaybeSnapshotFromMetadata(metadata)) _ <- OptionT.liftF(Future { observabilityService .recordMetric( LongCounterMetric(readVehicleSnapshotFromDiskCounter, SnapShotDirectoryScannerCommandOptions) ) }(directorySnapshotScanningDispatcher)) } yield snapshot getSnapshotAndReportMetric.value.recoverWith { // 若加载前已列出的旧快照被删除,则执行重试 case _: NoSuchFileException => loadAsync(persistenceId, criteria) }(streamDispatcher) }
供参考,getMaybeSnapshotFromMetadata方法的签名如下,我可以修改该方法签名为其添加Future包装,但最终返回的包装结果必须为Option[SelectedSnapshot]类型:
private def getMaybeSnapshotFromMetadata(metadata: Seq[SnapshotMetadata]): Option[SelectedSnapshot]
核心问题
上述代码可以正常编译,但编译通过的前提是导入了全局隐式ExecutionContext。我的目标是使用不同的显式执行上下文(不同的可配置dispatcher),但不知道在for推导式中,调用getMaybeSnapshotFromMetadata的第一行代码如何指定显式上下文:如果我使用OptionT.liftF包裹指定了streamDispatcher的Future调用,示例如下:
... snapshot <- OptionT.liftF( Future { getMaybeSnapshotFromMetadata(metadata)}(streamDispatcher))
会得到OptionT[Future, Option[SelectedSnapshot]]类型的嵌套结果,不符合预期。
备选实现
如果没有可行方案,我可以退而使用原生Future的andThen链式调用实现逻辑,备选实现代码如下:
Future { getMaybeSnapshotFromMetadata(metadata) }(streamDispatcher) .andThen(selectedSnapshot => { observabilityService .recordMetric( LongCounterMetric(readVehicleSnapshotFromDiskCounter, SnapShotDirectoryScannerCommandOptions) ) selectedSnapshot })(opentelemetryDispatcher) .recoverWith { // 若加载前已列出的旧快照被删除,则执行重试 case _: NoSuchFileException => loadAsync(persistenceId, criteria) }(streamDispatcher)
更新说明:我已确认for推导式中第二行的指标上报逻辑使用
liftF是可行方案,已更新对应代码块。
内容的提问来源于stack exchange,提问作者vasigorc
相关产品推荐
相关产品推荐

