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

如何在Http4s服务中按固定间隔更新变量?

解决方案:基于Cats-Effect和FS2实现定时刷新值

核心思路

用**并发安全的可变引用Ref**存储需要刷新的值(或依赖令牌的case class实例),通过FS2 Stream启动后台定时任务更新Ref中的内容。这样既保留了原对象的纯值使用方式,又能实现定期刷新,完全兼容Http4s的Effect上下文。


第一步:改造基础示例(单个值刷新)

针对你提供的最小示例,修改如下:

package com.example

import scala.concurrent.ExecutionContext
import scala.concurrent.duration._

import cats.effect.{ExitCode, IO, IOApp, Ref, Sync}
import cats.syntax.all._
import cats.Monad
import fs2.Stream
import org.http4s.blaze.server.BlazeServerBuilder
import org.http4s.dsl.Http4sDsl
import org.http4s.{HttpApp, HttpRoutes}

object ExampleService extends IOApp {
  // 生成新值的方法(保持原逻辑)
  def generateNewValue(): Int = scala.util.Random.nextInt(42)

  // 1. 用Ref存储需要刷新的值,初始化时生成第一个值
  private val freshValueRef: IO[Ref[IO, Int]] = Ref.of(generateNewValue())

  // 2. 定义后台定时刷新流:每隔5分钟更新一次Ref的值
  private def refreshStream(ref: Ref[IO, Int], interval: FiniteDuration): Stream[IO, Unit] =
    Stream.awakeEvery[IO](interval)
      .evalMap(_ => IO(generateNewValue()).flatMap(ref.set))
      .onStart(IO(generateNewValue()).flatMap(ref.set)) // 可选:启动时立即刷新一次

  // 3. 路由接收Ref作为参数,读取最新值
  def routes[F[_]: Monad: Sync](ref: Ref[F, Int]): HttpRoutes[F] = {
    val dsl = Http4sDsl[F]
    import dsl._

    HttpRoutes.of[F] {
      case GET -> Root / "ping"        => Ok("pong")
      case GET -> Root / "fresh-value" => ref.get.flatMap(v => Ok(s"$v"))
    }
  }

  override def run(args: List[String]): IO[ExitCode] = {
    val forkJoinPool = ExecutionContext.fromExecutor(new java.util.concurrent.ForkJoinPool(8))

    freshValueRef.flatMap { ref =>
      // 启动后台刷新任务(非阻塞)
      val backgroundRefresh = refreshStream(ref, 5.minutes).compile.drain.start

      BlazeServerBuilder[IO]
        .withExecutionContext(forkJoinPool)
        .bindHttp(1234, "localhost")
        .withHttpApp(routes[IO](ref).orNotFound)
        .resource
        .use(_ => backgroundRefresh.join) // 保持服务运行
        .as(ExitCode.Success)
    }
  }
}

第二步:适配实际场景(带令牌的Case Class)

针对你提到的12小时刷新令牌、更新依赖客户端的场景,只需将Ref存储的对象换成你的Case Class即可,无需修改Case Class本身:

// 原Case Class保持不变,纯值类型
case class ClientConfig(token: String, apiClient: ApiClient)
class ApiClient(token: String) {
  // 客户端的纯方法,无需修改
  def callApi(): String = s"Calling API with token: $token"
}

object TokenRefreshService extends IOApp {
  // 模拟获取新令牌的Effect操作
  private def fetchNewToken(): IO[String] = IO {
    println("Fetching new token...")
    java.util.UUID.randomUUID().toString
  }

  // 生成新的ClientConfig实例
  private def generateNewConfig(): IO[ClientConfig] = 
    fetchNewToken().map(token => ClientConfig(token, new ApiClient(token)))

  // 用Ref存储ClientConfig
  private val configRef: IO[Ref[IO, ClientConfig]] = generateNewConfig().flatMap(Ref.of)

  // 定时刷新流(12小时间隔)
  private def refreshConfigStream(ref: Ref[IO, ClientConfig]): Stream[IO, Unit] =
    Stream.awakeEvery[IO](12.hours)
      .evalMap(_ => generateNewConfig().flatMap(ref.set))
      .onStart(generateNewConfig().flatMap(ref.set))

  // 路由中使用最新的ClientConfig
  def routes[F[_]: Monad: Sync](configRef: Ref[F, ClientConfig]): HttpRoutes[F] = {
    val dsl = Http4sDsl[F]
    import dsl._

    HttpRoutes.of[F] {
      case GET -> Root / "api-status" =>
        configRef.get.map(config => Ok(config.apiClient.callApi()))
    }
  }

  override def run(args: List[String]): IO[ExitCode] = {
    val forkJoinPool = ExecutionContext.fromExecutor(new java.util.concurrent.ForkJoinPool(8))

    configRef.flatMap { ref =>
      val backgroundRefresh = refreshConfigStream(ref).compile.drain.start

      BlazeServerBuilder[IO]
        .withExecutionContext(forkJoinPool)
        .bindHttp(1234, "localhost")
        .withHttpApp(routes[IO](ref).orNotFound)
        .resource
        .use(_ => backgroundRefresh.join)
        .as(ExitCode.Success)
    }
  }
}

关键说明

  1. Ref的作用:Cats-Effect提供的Ref是并发安全的可变引用,保证多线程/Fiber环境下读取和更新的原子性,避免线程安全问题。
  2. 后台流的启动:FS2的Stream.awakeEvery定时触发刷新逻辑,通过compile.drain.start启动为后台Fiber,不阻塞主线程。
  3. 原代码兼容性:你的Case Class和原有纯值逻辑完全不需要修改,只是将实例存储到Ref中,访问时通过ref.get获取最新版本即可。
  4. 避免不安全操作:全程在Effect上下文内处理,不需要调用unsafeRunSync,符合Cats-Effect的纯函数式编程规范。

内容的提问来源于stack exchange,提问作者andres

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.22 00:12:31