Scala中使用Stream处理大文件的技术疑问
有一个存储用户数据的大型CSV文件,需要实现接口接收用户信息并返回该用户是否存在于文件中的布尔值。为避免内存溢出,尝试采用流式读取方式,使用懒加载的getLines方法,编写的代码如下:
def readUsers(filePath: Path): IO[Seq[User]] = { Stream .eval(IO(scala.io.Source.fromFile(filePath.toFile))) .flatMap(source => Stream.fromIterator[IO](source.getLines(), 64)) .map(line => line.split(",") .map(cols => User(cols(0), cols(1), cols(2), cols(3))) .compile .toList } def isUserExists(user: User, users: IO[Seq[User]]): IO[Boolean] = users.map(_.contains(user))
应用启动时执行val users = readUsers(filePath)赋值,每次接口调用时将其传入isUserExists函数。有以下疑问:
toList是否会将整个文件内容加载到内存中?- 若是,该如何避免?是否需要移除
toList并在isUserExists中逐行遍历? - 当文件内容变更时,由于
users是IO值,接口返回结果是否会同步更新?
1. toList是否会将整个文件内容加载到内存中?
是的。compile.toList会把Stream中的所有元素全部收集到List里,这意味着CSV文件的所有用户数据都会被一次性加载到内存中,完全违背了流式读取避免内存溢出的初衷——不管Stream的构建逻辑有多懒,一旦调用compile.toList,就会强制拉取所有元素到内存存储。
2. 如何避免内存溢出?
需要移除toList,并把用户存在性检查逻辑融入流式处理流程,不需要加载全部数据到内存。可以重构代码,让流在遍历过程中匹配目标用户,找到后立即终止遍历,只保留当前处理的行数据:
def isUserExists(user: User, filePath: Path): IO[Boolean] = { Stream .eval(IO(scala.io.Source.fromFile(filePath.toFile))) .flatMap(source => Stream.fromIterator[IO](source.getLines(), 64)) .map { line => val cols = line.split(",") User(cols(0), cols(1), cols(2), cols(3)) } .exists(_ == user) // 找到匹配用户后立即停止遍历 .compile .lastOrError // 取出判断结果,流为空时默认返回false }
每次接口调用时直接传入文件路径和目标用户,全程仅在内存中保留当前处理的行,不会加载整个文件。另外注意:原代码中map(line => line.split(",")...存在语法错误,重构代码已修正该问题。
3. 文件内容变更时,接口结果是否同步更新?
不会。你在应用启动时执行val users = readUsers(filePath),这个IO值只会被执行一次(首次触发时),之后每次isUserExists调用都是复用之前加载到内存的List数据,文件后续的变更不会被感知到。
如果需要同步文件变更,不能在启动时提前加载数据,而是要每次接口调用都重新读取文件(即上面重构后的isUserExists逻辑),这样每次检查都会基于最新文件内容。但这种方式会增加每次请求的文件IO开销,需要根据业务场景权衡性能与数据一致性:如果文件更新频繁且要求强一致性,这种方式合适;如果文件很少更新,也可以考虑定时重新加载数据,或者用文件变更监听机制(如Java的WatchService)触发数据刷新。
内容的提问来源于stack exchange,提问作者fullstack

