【问题标题】:How to group incoming events from infinite stream?如何对来自无限流的传入事件进行分组?
【发布时间】:2023-03-04 22:52:01
【问题描述】:

我有无限的事件流:

(timestamp, session_uid, traffic)

...
(1448089943, session-1, 10)
(1448089944, session-1, 20)
(1448089945, session-2, 50)
(1448089946, session-1, 30)
(1448089947, session-2, 10)
(1448089948, session-3, 10)
...

我想按 session_uid 对这些事件进行分组并计算每个会话的流量总和。

我写了一个akka-streams 流,它适用于有限流使用groupBy(我的代码基于食谱中的this 示例)。但是对于无限流,它不会起作用,因为groupBy 函数应该处理所有传入的流,只有在这之后才准备好返回结果。

我认为我应该使用超时实现分组,即如果我没有收到指定 stream_uid 的事件超过 5 分钟,我应该返回这个 session_uid 的分组事件。但是如何实现它只使用akka-streams

【问题讨论】:

  • 我正在使用akka-streams(如标签中所述)。
  • 不是时间限制,更新次数可以吗?例如。每 10,000 次更新产生一个分组。
  • 如果你有无限的蒸汽,groupby 怎么能“处理所有传入的流,只有在这之后才准备好返回结果”?如果流真的是无限的,这将永远不会发生。
  • @RamonJ.RomeroyVigil 任何会话总是在无限流中具有第一个和最后一个事件。事件之间的间隔不能超过 N 分钟,但不能超过 M 分钟(M > N,即 M=5 分钟,如上所述)。在不同的情况下,流可以在此时间间隔内处理 N 或 N 百万个事件,因此具有固定批量大小的解决方案是不可接受的。

标签: scala akka-stream


【解决方案1】:

我想出了一个有点gnarly 的解决方案,但我认为它可以完成工作。

基本思想是使用Source的keepAlive方法作为触发完成的定时器。

但要做到这一点,我们首先必须对数据进行一点抽象。计时器将需要从原始 Source 发送触发器或另一个元组值,因此:

sealed trait Data

object TimerTrigger extends Data
case class Value(tstamp : Long, session_uid : String, traffic : Int) extends Data

然后将我们的元组源转换为值源。我们仍将使用groupBy 进行类似于您的有限流案例的分组:

val originalSource : Source[(Long, String, Int), Unit] = ???

type IDGroup = (String, Source[Value, Unit]) //uid -> Source of Values for uid

val groupedDataSource : Source[IDGroup, Unit] = 
  originalSource.map(t => Value(t._1, t._2, t._3))
                .groupBy(_.session_uid)

棘手的部分是处理只是元组的分组:(String, Source[Value,Unit])。如果时间已经过去,我们需要计时器来通知我们,所以我们需要另一个抽象来知道我们是否仍在计算,或者我们是否由于超时而完成了计算:

sealed trait Sum {
  val sum : Int
}
case class StillComputing(val sum : Int) extends Sum
case class ComputedSum(val sum : Int) extends Sum

val zeroSum : Sum = StillComputing(0)

现在我们可以耗尽每个组的源头。如果值的来源在timeOut 之后没有产生任何东西,keepAlive 将发送一个TimerTrigger。然后,keepAlive 中的 Data 与 TimerTrigger 或来自原始 Source 的新值进行模式匹配:

val evaluateSum : ((Sum , Data)) => Sum = {
  case (runningSum, data) => { 
    data match {
      case TimerTrigger => ComputedSum(runningSum.sum)
      case v : Value    => StillComputing(runningSum.sum + v.traffic)
    }
  }
}//end val evaluateSum

type SumResult = (String, Future[Int]) // uid -> Future of traffic sum for uid

def handleGroup(timeOut : FiniteDuration)(idGroup : IDGroup) : SumResult = 
  idGroup._1 -> idGroup._2.keepAlive(timeOut, () => TimerTrigger)
                          .scan(zeroSum)(evaluateSum)
                          .collect {case c : ComputedSum => c.sum}
                          .runWith(Sink.head)

该集合应用于仅匹配完成总和的部分函数,​​因此只有在计时器触发后才能到达 Sink。

然后我们将这个处理程序应用到每个出现的分组:

val timeOut = FiniteDuration(5, MINUTES)

val sumSource : Source[SumResult, Unit] = 
  groupedDataSource map handleGroup(timeOut)

我们现在有一个 (String,Future[Int]) 的 Source,它是 session_uid 和该 ID 的流量总和的 Future。

就像我说的,复杂但符合要求。另外,我不完全确定如果一个 uid 已经分组并已超时,但随后会出现具有相同 uid 的新值会发生什么。

【讨论】:

  • 非常感谢!您的回答给了我几个有用的想法,但在我的任务中它有几个问题。例如它在groupBy 函数中有一个problem with consumption memory
  • @maxd 任何解决方案都会有潜在的内存问题,因为活动 uid 的数量可能会变得很大,因此运行总和的数量可能会变得很大。但不客气,祝黑客愉快。
【解决方案2】:

这似乎是Source.groupedWithin 的用例:

def groupedWithin(n: Int, d: FiniteDuration): Source[List[Out], Mat]

“将这个流分成在一个时间窗口内接收到的元素组,或者受给定元素数量的限制,无论先发生什么。”

Here's the link to the docs

【讨论】:

  • 这也是不可接受的,因为会话可以分成几个块。
【解决方案3】:

也许你可以通过演员简单地实现它

case class SessionCount(name: String)

class Hello private() extends Actor {
  var sessionMap = Map[String, Int]()

  override def receive: Receive = {
    case (_, session: String, _) =>
      sessionMap = sessionMap + (session -> (sessionMap.getOrElse(session, 0) + 1))

    case SessionCount(name: String) => sender() ! sessionMap.get(name).getOrElse(0)
  }
}


object Hello {
  private val actor = ActorSystem.apply().actorOf(Props(new Hello))
  private implicit val timeOver = Timeout(10, TimeUnit.SECONDS)
  type Value = (String, String, String)

  def add(value: Value) = actor ! value

  def count(name:String) = (actor ? SessionCount(name )).mapTo[Int]
}

【讨论】:

  • 我知道可以解决我的任务只使用 Akka 演员,但我认为 akka-streams 也应该为我的任务提供解决方案。我想找到它。仅供参考:您的代码示例计算会话数,但不是每个会话的总流量。此外,它没有正确处理无限流。
  • 怎么用words.forearch(e=>Hello.add(e))
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多