基于相同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 实现方案
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
相关产品推荐
相关产品推荐

