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的记录。
正确实现方案
实现思路
- 合并流:将两个Source合并为一个包含所有Record的流
- 按ID分组:把相同id的Record分到同一个子流中
- 筛选最新版本:在每个子流中保留版本号最大的Record
- 展平流:将所有子流的结果合并回一个主流
代码实现
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:在每个分组子流中,从初始值开始逐步比较并保留版本号最大的RecordmergeSubstreams:将所有分组子流的结果合并回一个单一流
内容的提问来源于stack exchange,提问作者Siddharth Shankar
相关产品推荐
相关产品推荐

