Kotlin中如何将Flow<Byte>转换为Flow<String>(Kotlin/JS环境)
Kotlin/JS 中将 Flow 转换为按行分割的 Flow
要实现将UTF-8字节流(Flow<Byte>)转换为按行分割的字符串流,核心是流式处理字节,正确解码UTF-8并按换行符分割,同时避免一次性加载所有数据(符合冷数据源的特性)。以下是具体实现:
实现代码
import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.transform import org.w3c.dom.TextDecoder fun Flow<Byte>.toUtf8LineFlow(): Flow<String> = transform { byte -> // 每个收集器独立维护解码和行缓存状态 val decoder = TextDecoder("utf-8", js("{ stream: true }")) val currentLine = StringBuilder() val tempBytes = mutableListOf<Byte>() // 处理单个字节,解码后分割行 fun processByte(byte: Byte) { tempBytes.add(byte) val input = Uint8Array(tempBytes.toByteArray()) tempBytes.clear() val decoded = decoder.decode(input, js("{ stream: true }")) decoded.forEach { char -> if (char == '\n') { emit(currentLine.toString()) currentLine.clear() } else { currentLine.append(char) } } } processByte(byte) // 流结束时处理剩余未完成的内容 awaitClose { val finalDecoded = decoder.decode(Uint8Array(0), js("{ stream: false }")) finalDecoded.forEach { char -> if (char == '\n') { emit(currentLine.toString()) currentLine.clear() } else { currentLine.append(char) } } if (currentLine.isNotEmpty()) { emit(currentLine.toString()) } } }
代码说明
- UTF-8 流式解码:使用
TextDecoder的stream模式,确保多字节UTF-8字符(如中文、俄文)被正确解码,即使字节分开发射。 - 行缓存与分割:用
StringBuilder累积当前行内容,遇到\n时立即发射该行并清空缓存。 - 冷流兼容:每个流收集器都会初始化独立的解码器和缓存,避免不同收集操作之间的状态干扰。
- 收尾处理:流结束时触发解码器的最终解码,处理剩余未完成的字节,确保最后一行(无结尾换行的情况)被正确发射。
测试使用
将示例中的TODO()替换为source.toUtf8LineFlow()即可通过测试:
val transformed: Flow<String> = source.toUtf8LineFlow()
内容的提问来源于stack exchange,提问作者Андрей Щеглов
相关产品推荐
相关产品推荐

