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

Scala中如何向Akka Streams现有Source添加Map中的新值?

解决方案

要实现这个需求,核心思路是先收集原Source中所有元素的ID,再过滤出Map中不存在于原Source的元素,最后将两者合并。以下是具体实现:

import akka.stream.scaladsl.{Sink, Source}
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.Future

case class A(id: String, someValue: String)

def addValuesToExistingSource(mySource: Source[A,_], myMapFuture: Future[Map[String, A]]): Source[A, _] = {
  // 1. 收集原Source中所有元素的ID,存入Set以便快速查找
  val existingIdsFuture: Future[Set[String]] =
    mySource.runWith(Sink.seq).map(_.map(_.id).toSet)

  // 2. 组合两个Future,过滤出Map中不在原Source里的元素
  val newElementsFuture: Future[Iterable[A]] =
    for {
      existingIds <- existingIdsFuture
      map <- myMapFuture
    } yield map.values.filter(a => !existingIds.contains(a.id))

  // 3. 将过滤后的元素转为Source,与原Source合并(用concat保证原元素在前,新增元素在后)
  mySource.concat(Source.fromFuture(newElementsFuture).mapConcat(identity))
}

关键细节说明

  • 收集原Source元素:使用Sink.seq收集原Source的所有元素,仅适用于有限Source。如果原Source是无限流,这种方式会导致内存溢出,此时需要换用带状态的Flow跟踪已出现的ID,逻辑会更复杂。
  • Future组合:通过for推导式组合两个Future,确保在获取到原Source的ID集合和Map数据后,再进行过滤操作。
  • 合并Source:使用concat保证原Source的元素全部输出后,再输出新增的元素;如果不需要严格顺序,也可以用merge并行输出,但元素顺序会被打乱。

内容的提问来源于stack exchange,提问作者Siddharth Shankar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 13:10:50