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

Scala中遍历两个Source并按版本属性筛选最新Record

合并两个Akka Stream Source并筛选每个ID的最新版本Record

问题背景

给定Record类定义:

case class Record(id: String, version: Long)

以及两个包含Record的Akka Stream Source:

val sourceA: Source[Record, _] = <>
val sourceB: Source[Record, _] = <>

两个Source中存在id相同但版本号(version)不同的Record,需要生成一个新的Source,包含每个id对应的最新版本的Record。

你尝试的代码存在逻辑错误:在map操作中嵌套另一个Source的map不符合Akka Stream的操作逻辑,会导致类型不匹配,也无法正确关联相同id的记录。

正确实现方案

实现思路

  1. 合并流:将两个Source合并为一个包含所有Record的流
  2. 按ID分组:把相同id的Record分到同一个子流中
  3. 筛选最新版本:在每个子流中保留版本号最大的Record
  4. 展平流:将所有子流的结果合并回一个主流

代码实现

import akka.stream.scaladsl.Source

case class Record(id: String, version: Long)

def getLatestRecords(sourceA: Source[Record, _], sourceB: Source[Record, _]): Source[Record, _] = {
  // 合并两个Source
  val combined = Source.combine(sourceA, sourceB)(_.merge(_))
  
  combined
    // 按id分组,支持任意数量的不同id
    .groupBy(Int.MaxValue, _.id)
    // 每组内折叠,保留版本号最大的Record
    .fold(Record("", 0L)) { (latest, current) =>
      if (current.version > latest.version) current else latest
    }
    // 合并所有子流结果
    .mergeSubstreams
}

// 使用示例
val sourceA = Source(List(Record("user1", 1), Record("user2", 2)))
val sourceB = Source(List(Record("user1", 3), Record("user2", 1)))

val latestRecordsSource = getLatestRecords(sourceA, sourceB)
// 最终流将输出 Record("user1", 3) 和 Record("user2", 2)

关键操作说明

  • Source.combine(...):将两个Source合并为一个流,这里使用merge策略将元素按到达顺序合并
  • groupBy:按id对元素分组,Int.MaxValue允许处理任意数量的不同id(若已知id数量有限,可设置具体数值优化性能)
  • fold:在每个分组子流中,从初始值开始逐步比较并保留版本号最大的Record
  • mergeSubstreams:将所有分组子流的结果合并回一个单一流

内容的提问来源于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.09 20:55:10