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
相关产品推荐
相关产品推荐

