如何改写ZStream代码实现多CSV文件按时间有序轮询合并?
多文件ZStream轮询读取并按时间有序合并实现方案
核心思路
要实现轮询按指定数量从多文件取数+最终流按时间戳有序,需要把单文件流改造为可按批次拉取的流,再通过调度逻辑轮询拉取各文件的批次,最后对每个轮询周期内的所有数据按时间排序后输出。
步骤1:改造单文件流为可批次拉取的流
基于你已有的fileSource,封装一个能按指定数量拉取单文件数据的流,复用原有数据转换逻辑:
// 假设MsgWithTimestamp是包含unixTimestampMs字段的消息类型,与原有逻辑类型一致 def fileBatchStream(filePath: String, batchSize: Int): Stream[Throwable, List[MsgWithTimestamp]] = fileSource(filePath) .mapZIO(transactionService.mkMsgFromString) .collectSome .grouped(batchSize) // 按指定批次大小分组 .filter(_.nonEmpty) // 过滤文件末尾的空批次
步骤2:实现轮询调度逻辑
针对多文件及对应批次配置,创建循环轮询拉取的流:
// 定义每个文件的轮询配置:路径 + 每次取数数量 case class FilePollConfig(path: String, batchSize: Int) def pollMultiFileStreams(configs: List[FilePollConfig]): Stream[Throwable, MsgWithTimestamp] = { // 为每个配置创建对应批次流 val fileStreams = configs.map { cfg => fileBatchStream(cfg.path, cfg.batchSize).map(chunk => (cfg.path, chunk)) } // 循环轮询所有未读完的文件流,拉取各自批次 ZStream.unfoldLoop(fileStreams) { streams => ZIO.collectAll(streams.map(_.headOption)) // 从每个流拉取一个批次 .map { batches => // 提取所有有效消息 val validMsgs = batches.flatten.flatMap(_._2) // 保留未读完的文件流(已读完的移除) val remainingStreams = streams.zip(batches).flatMap { case (stream, Some(_)) => Some(stream.drop(1)) case (stream, None) => None } // 返回本次轮询的消息和剩余待处理流 (validMsgs, remainingStreams) } .filter(_._1.nonEmpty) // 过滤空轮询周期 }.flatMap(ZStream.fromIterable(_)) }
该逻辑会持续轮询所有未读完的文件,直到全部处理完毕。
步骤3:整合原有处理流程
将轮询得到的流按时间排序后,接入你原有的控速、发送等逻辑:
def multiFileStream(configs: List[FilePollConfig]): Stream[Throwable, Result] = for { // 先对轮询拉取的消息按时间戳全局排序 sortedMsg <- pollMultiFileStreams(configs) .run(Sink.foldLeft( PriorityQueue.empty[MsgWithTimestamp](Ordering.by(_.unixTimestampMs)) )((queue, msg) => queue.enqueue(msg))) .flatMap(q => ZStream.fromIterable(q.dequeueAll)) // 接入原有处理逻辑 res <- ZStream.succeed(sortedMsg) .mapZIO(sleepIfRequired) .mapZIOPar(proxyConf.maxConnections, proxyConf.maxConnections)(emitMessage(_, uriToSendMessages)) } yield res
如果不需要全局严格排序,仅需要轮询批次内排序,可以替换为批次内排序逻辑,减少内存占用:
sortedMsg <- pollMultiFileStreams(configs) .groupAdjacentBy(_.unixTimestampMs) .map { case (_, msgs) => msgs.sortBy(_.unixTimestampMs) } .flatMap(ZStream.fromIterable(_))
关键注意点
- 大文件场景下优先使用流式排序(优先级队列),避免一次性加载所有数据到内存。
- 轮询逻辑会自动移除已读完的文件流,无需额外处理。
- 若需按时间戳聚合业务逻辑,建议先全局排序再分组,保证数据时序正确性。
内容的提问来源于stack exchange,提问作者Arlengur
相关产品推荐
相关产品推荐

