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

基于相同ID合并两个Akka Stream流的实现方案问询

嘿,针对你这个基于ID合并两个流、要求先到先匹配且保留未匹配元素的需求,我整理了Akka Stream、Spark Structured Streaming和Flink三个主流框架的实现方案,完全贴合你的要求👇


核心需求回顾

  • 两个流(Foo/Bar)按相同id合并为MergeResult
  • 同ID重复元素采用先到先匹配策略(比如第一个Foo(1)匹配第一个Bar(1))
  • 未匹配元素暂存内存(有状态),等待后续匹配
  • 源元素有序,需严格遵循顺序逻辑

Akka Stream 实现方案

Akka Stream没有直接内置这种带优先级的有状态Join算子,咱们可以手动用statefulMapConcat维护缓存,逻辑非常直观:

import akka.stream.scaladsl.{Source, Merge, Flow}
import scala.collection.mutable.{Queue, Map}

// 先把两个流合并为Either类型,统一处理
val mergedStream = Merge[Either[Foo, Bar]](2)
  .add(sourceOne.map(Left(_)), sourceTwo.map(Right(_)))

// 自定义有状态Flow,维护未匹配元素的队列(先进先出保证先到先匹配)
val mergeFlow = Flow[Either[Foo, Bar]].statefulMapConcat { () =>
  // 每个ID对应一个未匹配Foo的队列
  val fooCache: Map[Int, Queue[String]] = Map.empty
  // 每个ID对应一个未匹配Bar的队列
  val barCache: Map[Int, Queue[String]] = Map.empty

  element => {
    element match {
      case Left(Foo(id, fooVal)) =>
        barCache.get(id) match {
          case Some(barQueue) if barQueue.nonEmpty =>
            // 找到匹配的Bar,输出结果并移除队列头部的元素
            val barVal = barQueue.dequeue()
            if (barQueue.isEmpty) barCache.remove(id)
            List(MergeResult(id, fooVal, barVal))
          case _ =>
            // 无匹配Bar,加入Foo缓存队列
            fooCache.getOrElseUpdate(id, Queue.empty).enqueue(fooVal)
            Nil // 暂时不输出
        }
      case Right(Bar(id, barVal)) =>
        fooCache.get(id) match {
          case Some(fooQueue) if fooQueue.nonEmpty =>
            // 找到匹配的Foo,输出结果并移除队列头部的元素
            val fooVal = fooQueue.dequeue()
            if (fooQueue.isEmpty) fooCache.remove(id)
            List(MergeResult(id, fooVal, barVal))
          case _ =>
            // 无匹配Foo,加入Bar缓存队列
            barCache.getOrElseUpdate(id, Queue.empty).enqueue(barVal)
            Nil // 暂时不输出
        }
    }
  }
}

// 最终结果流
val resultStream = mergedStream.via(mergeFlow)

说明:用Map+Queue的组合保证每个ID的未匹配元素按顺序存储,完全符合先到先匹配的要求,未匹配元素会一直留在内存中(可自行添加超时清理逻辑)。


Spark Structured Streaming 实现方案

Spark的mapGroupsWithState可以自定义状态逻辑,完美适配你的需求,支持无限期保留未匹配元素:

import org.apache.spark.sql.{SparkSession, Dataset}
import org.apache.spark.sql.streaming.{GroupState, GroupStateTimeout}

// 统一输入类型,方便合并处理
sealed trait Event
case class FooEvent(id: Int, value: String) extends Event
case class BarEvent(id: Int, value: String) extends Event

// 状态类型:存储每个ID的未匹配Foo/Bar队列
case class MatchState(fooQueue: List[String], barQueue: List[String])

val spark = SparkSession.builder().appName("StreamMerge").getOrCreate()
import spark.implicits._

// 转换原始流为统一Event类型
val fooStream: Dataset[FooEvent] = sourceOne.toDF().as[Foo].map(f => FooEvent(f.id, f.value))
val barStream: Dataset[BarEvent] = sourceTwo.toDF().as[Bar].map(b => BarEvent(b.id, b.value))

// 合并流并按ID分组
val groupedEvents = fooStream.union(barStream).keyBy(_.id)

// 自定义状态处理逻辑
val resultStream = groupedEvents.mapGroupsWithState[MatchState, MergeResult](
  GroupStateTimeout.NoTimeout // 永不超时,保留未匹配元素
) { (id: Int, events: Iterator[Event], state: GroupState[MatchState]) =>
  val currentState = state.getOption.getOrElse(MatchState(Nil, Nil))
  var newFooQueue = currentState.fooQueue
  var newBarQueue = currentState.barQueue
  val output = scala.collection.mutable.ListBuffer[MergeResult]()

  events.foreach {
    case FooEvent(_, fooVal) =>
      newBarQueue match {
        case barVal :: rest =>
          // 匹配到第一个Bar,输出结果
          output += MergeResult(id, fooVal, barVal)
          newBarQueue = rest
        case Nil =>
          // 加入Foo队列末尾,保证顺序
          newFooQueue = newFooQueue :+ fooVal
      }
    case BarEvent(_, barVal) =>
      newFooQueue match {
        case fooVal :: rest =>
          // 匹配到第一个Foo,输出结果
          output += MergeResult(id, fooVal, barVal)
          newFooQueue = rest
        case Nil =>
          // 加入Bar队列末尾,保证顺序
          newBarQueue = newBarQueue :+ barVal
      }
  }

  // 更新状态:如果还有未匹配元素则保留,否则清空
  if (newFooQueue.nonEmpty || newBarQueue.nonEmpty) {
    state.update(MatchState(newFooQueue, newBarQueue))
  } else {
    state.remove()
  }

  output.toIterator
}

// 输出结果到控制台(可替换为其他Sink)
resultStream.writeStream.format("console").start().awaitTermination()

说明:通过MatchState维护每个ID的未匹配队列,每次处理事件时优先匹配队列头部的元素,严格遵循先到先匹配规则。


Flink的CoProcessFunction是处理双源流有状态匹配的神器,允许咱们分别处理两个流的元素,并为每个Key维护独立状态:

import org.apache.flink.streaming.api.scala._
import org.apache.flink.streaming.api.functions.co.CoProcessFunction
import org.apache.flink.util.Collector
import org.apache.flink.api.common.state.{ListState, ListStateDescriptor}

val env = StreamExecutionEnvironment.getExecutionEnvironment

// 将原始流转换为KeyedStream(按ID分组)
val fooKeyedStream = env.fromCollection(sourceOne.toList).keyBy(_.id)
val barKeyedStream = env.fromCollection(sourceTwo.toList).keyBy(_.id)

// 自定义CoProcessFunction,处理两个流的元素
class MergeCoProcessFunction extends CoProcessFunction[Foo, Bar, MergeResult] {
  // 存储未匹配的Foo(有序ListState保证先到先匹配)
  lazy val fooState: ListState[String] = getRuntimeContext.getListState(
    new ListStateDescriptor[String]("foo-cache", classOf[String])
  )
  // 存储未匹配的Bar
  lazy val barState: ListState[String] = getRuntimeContext.getListState(
    new ListStateDescriptor[String]("bar-cache", classOf[String])
  )

  // 处理Foo流的元素
  override def processElement1(foo: Foo, ctx: CoProcessFunction[Foo, Bar, MergeResult]#Context, out: Collector[MergeResult]): Unit = {
    val barIter = barState.get().iterator()
    if (barIter.hasNext) {
      // 匹配到第一个Bar,输出结果
      val barVal = barIter.next()
      out.collect(MergeResult(foo.id, foo.value, barVal))
      // 移除已匹配的Bar,更新状态
      val remainingBars = barState.get().filter(_ != barVal).toList
      barState.update(remainingBars)
    } else {
      // 无匹配Bar,加入Foo状态
      fooState.add(foo.value)
    }
  }

  // 处理Bar流的元素
  override def processElement2(bar: Bar, ctx: CoProcessFunction[Foo, Bar, MergeResult]#Context, out: Collector[MergeResult]): Unit = {
    val fooIter = fooState.get().iterator()
    if (fooIter.hasNext) {
      // 匹配到第一个Foo,输出结果
      val fooVal = fooIter.next()
      out.collect(MergeResult(bar.id, fooVal, bar.value))
      // 移除已匹配的Foo,更新状态
      val remainingFoos = fooState.get().filter(_ != fooVal).toList
      fooState.update(remainingFoos)
    } else {
      // 无匹配Foo,加入Bar状态
      barState.add(bar.value)
    }
  }
}

// 连接两个Keyed流并应用自定义处理逻辑
val resultStream = fooKeyedStream.connect(barKeyedStream).process(new MergeCoProcessFunction())

// 输出结果
resultStream.print()
env.execute("Merge Streams")

说明:Flink的ListState是有序的,保证先到的元素先被匹配,未匹配元素会一直保存在状态中(可通过TimerService添加超时清理逻辑)。


内容的提问来源于stack exchange,提问作者Xiaohe Dong

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 07:30:48