【发布时间】:2019-12-29 09:50:55
【问题描述】:
我正在将一些 C# 代码转换为 scala 和 akka 流。
我的 c# 代码如下所示:
Task<Result1> GetPartialResult1Async(Request request) ...
Task<Result2> GetPartialResult2Async(Request request) ...
async Task<Result> GetResultAsync(Request request)
{
var result1 = await GetPartialResult1Async(request);
var result2 = await GetPartialResult2Async(request);
return new Result(request, result1, result2);
}
现在是 akka 流。我没有从 Request 到结果的 Task 的函数,而是从请求流向结果。
所以我已经有了以下两个流程:
val partialResult1Flow: Flow[Request, Result1, NotUsed] = ...
val partialResult2Flow: Flow[Request, Result2, NotUsed] = ...
但是我看不出如何将它们组合成一个完整的流程,因为在第一个流程上调用 via 我们会丢失原始请求,而通过在第二个流程上调用 via 我们会丢失第一个流程的结果。
所以我创建了一个看起来像这样的 WithState monad:
case class WithState[+TState, +TValue](value: TValue, state: TState) {
def map[TResult](func: TValue => TResult): WithState[TState, TResult] = {
WithState(func(value), state)
}
... bunch more helper functions go here
}
然后我将我的原始流程重写为如下所示:
def partialResult1Flow[TState]: Flow[WithState[TState, Request], WithState[TState, Result1]] = ...
def partialResult2Flow: Flow[WithState[TState, Request], WithState[TState, Result2]] = ...
并像这样使用它们:
val flow = Flow[Request]
.map(x => WithState(x, x))
.via(partialResult1Flow)
.map(x => WithState(x.state, (x.state, x.value))
.via(partialResult2Flow)
.map(x => Result(x.state._1, x.state._2, x.value))
现在这可行,但我当然不能保证如何使用 flow。所以我真的应该让它接受一个状态参数:
def flow[TState] = Flow[WithState[TState, Request]]
.map(x => WithState(x.value, (x.state, x.value)))
.via(partialResult1Flow)
.map(x => WithState(x.state._2, (x.state, x.value))
.via(partialResult2Flow)
.map(x => WithState(Result(x.state._1._2, x.state._2, x.value), x.state._1._1))
现在在这个阶段,我的代码变得非常难以阅读。我可以通过命名函数来清理它,并使用案例类而不是元组等,但基本上这里有很多附带的复杂性,这是很难避免的。
我错过了什么吗?这不是 Akka 流的好用例吗?有一些内置的方法吗?
【问题讨论】:
-
为什么不能使用 Future 而不是 Flow?它更简单,更容易组合。 Future 是 Scala 中的标准类,几乎每个 Scala 开发人员都知道如何使用它。
-
我的团队决定我们要使用 Akka 流,因为它有利于创建背压,我正在为我们的用例评估它。很可能这不是正确的解决方案,在这种情况下我们不会使用它。
标签: scala akka akka-stream