构建可恢复可中断的FS2定时数据流的技术疑问及代码咨询
FS2流实现问题解答
需求与现有实现
需求
- 每3秒执行一次数据获取函数,并将结果写入文件
- 发生错误时,自动重新认证并重启流
- 支持通过独立流触发中断
现有实现代码
private def CREATE_STREAM_RECURSIVE( session: AuthSession, myService: MyService[IO], sessionStore: SessionStore[IO], ): Stream[IO, Unit] = { Stream .repeatEval(service.fetchData(session)) .flatMap { data => val jsonStream = Stream.eval(IO.delay(data.asJson.noSpaces)) persistToFile(jsonStream) <-- another fs2 stream that writes to a file } .metered(3.seconds) .handleErrorWith { _ => Stream.force { myService.login() .flatMap(sessionStore.setSession) .map(CREATE_STREAM_RECURSIVE(_, myService, sessionStore)) //recursive } } .interruptWhen(isItOkToInterruptStream) }
疑问解答
1. 错误持续发生时,递归是否会造成栈溢出?
不会造成栈溢出。FS2的handleErrorWith是在流的运行时循环中处理错误,递归调用CREATE_STREAM_RECURSIVE并不会直接在当前调用栈上执行,而是返回一个新的Stream实例,由FS2 runtime负责调度执行。这种递归属于尾递归风格的流构造,不会累积调用栈帧,即便错误反复出现,也不会触发栈溢出。
2. 中断发生时,流是否可能未完成文件写入就停止?
是的,有可能。默认情况下,interruptWhen会立即终止流的所有当前运行任务,包括正在执行的文件写入操作。如果persistToFile没有处理中断的逻辑,正在进行的写入可能被强制中断,导致文件内容不完整或损坏。
3. 如何确保退出前完成文件写入?
可以通过以下方式解决:
- 给写入操作添加中断安全处理:在
persistToFile流中,用IO.interruptible标记写入操作,确保中断时能完成当前写入;或者通过onFinalize添加收尾逻辑,在流终止前完成文件的flush或关闭操作。 - 调整中断作用范围:将
interruptWhen的作用范围限定在数据获取循环部分,确保当前文件写入流执行完毕后再响应中断:Stream .repeatEval(service.fetchData(session)) .interruptWhen(isItOkToInterruptStream) // 仅中断数据获取循环 .flatMap { data => val jsonStream = Stream.eval(IO.delay(data.asJson.noSpaces)) persistToFile(jsonStream) // 写入操作不受中断直接影响,会执行完成 } .metered(3.seconds) .handleErrorWith { _ => Stream.force { myService.login() .flatMap(sessionStore.setSession) .map(CREATE_STREAM_RECURSIVE(_, myService, sessionStore)) } } - 使用
gracefulShutdown:FS2的gracefulShutdown操作符可以在收到中断信号后,允许当前正在处理的元素完成处理,再终止流。可将其应用到整个流或写入流部分,确保写入操作完成后再退出。
内容的提问来源于stack exchange,提问作者tharindu_DG
相关产品推荐
相关产品推荐

