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

如何从非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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 06:11:17