如何借助ZIO的ZStream实现请求批量处理并按ID匹配对应响应
ZStream 批量用户信息查询实现方案
核心实现思路
你需要用队列+Promise的模式实现请求聚合和结果分发:队列接收所有单个查询请求,每个请求绑定一个专属Promise;后台ZStream消费队列攒批发起gRPC调用,拿到结果后按ID匹配给对应Promise赋值,调用方通过等待Promise结果拿到返回值。
完整代码实现
import zio._ import zio.stream._ import scala.concurrent.duration._ // 你原有的结构体定义 case class UsersRequest(ids: List[Long]) case class Info(userId: Long, info: String) case class UsersInfoResponse(info: List[Info]) case class User(id: Long, info: String) // 批量查询服务封装 class BatchUserInfoService(service: /* 你的gRPC服务类型 */) { // 队列存储元素:(查询用户ID, 用于返回结果的Promise) private type RequestEntry = (Long, Promise[Throwable, String]) // 队列容量可根据业务QPS调整 private val MAX_QUEUE_SIZE = 1000 // 初始化队列和后台批量处理流 private val (requestQueue, shutdownHook) = Unsafe.unsafe { implicit unsafe => Runtime.default.unsafe.run( for { queue <- Queue.bounded[RequestEntry](MAX_QUEUE_SIZE) // 后台批量处理流 processFiber <- ZStream.fromQueue(queue) // 满足两个条件之一就发送批次:攒够10个ID / 等待满1秒 .groupedWithin(n = 10, d = 1.second) .mapZIO { batch => // 提取当前批次的ID和对应Promise映射 val idPromiseMap = batch.toMap val requestIds = idPromiseMap.keys.toList // 发起批量gRPC请求 service.getUserInfo(UsersRequest(requestIds)) // 把返回结果转成ID到信息的映射表方便匹配 .map(resp => resp.info.map(item => item.userId -> item.info).toMap) .foldZIO( // 请求异常:给批次内所有请求标记失败 err => ZIO.foreachDiscard(idPromiseMap.values)(_.fail(err)), // 请求成功:按ID匹配给每个Promise赋值 resultMap => ZIO.foreachDiscard(idPromiseMap) { case (id, promise) => resultMap.get(id) match { case Some(info) => promise.succeed(info) case None => promise.fail(new RuntimeException(s"用户ID $id 无匹配返回结果")) } } ) } .runDrain .forkDaemon // 作为后台守护进程运行 // 定义服务停止逻辑 stop = processFiber.interrupt *> queue.shutdown } yield (queue, stop) ).getOrThrow() } // 对外暴露的单个查询接口,和你原来的方法签名完全一致 def getUserInfo(id: Long): IO[Throwable, String] = for { promise <- Promise.make[Throwable, String] // 把请求加入队列 _ <- requestQueue.offer((id, promise)) // 等待结果返回 res <- promise.await } yield res // 服务下线时调用释放资源 def shutdown: UIO[Unit] = shutdownHook }
原有业务代码适配
你原来的createUser方法不需要做任何修改,直接替换getUserInfo的实现为上述批量服务的方法即可自动走批量逻辑:
def createUser(id: Long, batchService: BatchUserInfoService): IO[Throwable, User] = { batchService.getUserInfo(id) .map(info => User(id, info)) }
可选优化点
- 可根据业务需要给gRPC调用添加重试逻辑,减少偶发网络错误影响
- 队列满时可根据业务容忍度选择不同队列策略:
dropping丢弃新请求、sliding丢弃最老请求 - 可添加请求超时控制,避免调用方长时间等待无返回
内容的提问来源于stack exchange,提问作者Vlad
相关产品推荐
相关产品推荐

