如何从非Flow函数获取Flow输出并实现持续观测
将非Flow函数包装为可观测Flow(监听文件夹变化)
你需要把返回当前时刻文件夹状态的同步函数,包装成能持续观测文件变化的Flow<T>,核心思路是触发函数重新执行并发射新值,下面提供两种常用实现方案:
方案一:定时轮询(简单易实现)
适合对实时性要求不高的场景,通过固定间隔重复调用原函数并发射结果:
import kotlinx.coroutines.delay import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.flow fun getFlowOfTotalFiles(): Flow<Int> = flow { while (true) { emit(getTotalFiles()) // 发射当前文件数量 delay(1000) // 每秒轮询一次,可根据需求调整间隔 } } fun getFlowOfAllFiles(): Flow<List<File>> = flow { while (true) { emit(getAllFiles()) // 发射当前文件列表 delay(1000) } }
优缺点
- 优点:零额外依赖,代码极简,无需处理文件系统监听逻辑
- 缺点:存在延迟,间隔过短会占用不必要的资源,过长则无法及时感知变化
方案二:文件系统监听(高效实时)
通过监听文件系统的变化事件,仅在文件新增/删除/修改时触发原函数执行,实时性更高且资源消耗低。这里基于Java标准库的WatchService实现:
import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.channels.awaitClose import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.callbackFlow import kotlinx.coroutines.launch import java.nio.file.FileSystems import java.nio.file.Path import java.nio.file.StandardWatchEventKinds.* import java.nio.file.WatchKey // 通用文件夹监听函数:文件变化时发送信号 private fun watchFolder(folderPath: Path): Flow<Unit> = callbackFlow { val watchService = FileSystems.getDefault().newWatchService() // 注册监听:文件创建、删除、修改事件 folderPath.register(watchService, ENTRY_CREATE, ENTRY_DELETE, ENTRY_MODIFY) // 后台循环处理监听事件 val job = launch(Dispatchers.IO) { while (true) { val key: WatchKey = watchService.take() key.pollEvents().forEach { trySend(Unit) } key.reset() // 重置key以继续接收事件 } } // 流取消时清理资源 awaitClose { job.cancel() watchService.close() } } // 包装文件数量函数 fun getFlowOfTotalFiles(folderPath: Path): Flow<Int> = callbackFlow { emit(getTotalFiles()) // 先发射初始值 // 监听文件夹变化,变化时重新获取并发射新值 val job = launch { watchFolder(folderPath).collect { emit(getTotalFiles()) } } awaitClose { job.cancel() } } // 包装文件列表函数 fun getFlowOfAllFiles(folderPath: Path): Flow<List<File>> = callbackFlow { emit(getAllFiles()) // 先发射初始值 val job = launch { watchFolder(folderPath).collect { emit(getAllFiles()) } } awaitClose { job.cancel() } }
补充说明
- 若需要监听子文件夹,需递归遍历所有子目录并注册
WatchService - 使用时需传入目标文件夹的
Path对象,例如Paths.get("/your/target/folder") - 测试示例:
fun main() = runBlocking { val targetFolder = Paths.get("/your/target/folder") getFlowOfTotalFiles(targetFolder).collect { count -> println("当前文件数量:$count") } }
内容的提问来源于stack exchange,提问作者Cyber Avater
相关产品推荐
相关产品推荐

