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

如何用Scala FS2流式处理移除文件过滤行并统计过滤行数

Solution Using FS2 Broadcast for Single-Pass Processing

To solve this problem, we need to stream-process a file to remove lines matching a list of values while counting how many lines were removed—all in a single pass over the input file (to avoid redundant I/O operations). Here's a clean, efficient implementation using FS2:

Step-by-Step Explanation

  1. Read and Decode Input: Start by reading the input file as a byte stream, decode it to UTF-8, and split into individual lines.
  2. Broadcast the Stream: Use FS2's broadcast to split the line stream into two identical copies. This ensures we only process the input file once.
  3. Count Discarded Lines: One stream filters lines that match our exclusion list and counts them.
  4. Write Kept Lines: The other stream filters out excluded lines, re-adds newlines, encodes back to bytes, and writes to the output file.
  5. Combine Actions: Run both the count and write operations, then return the final count of discarded lines.

Complete Code

import cats.effect.Sync
import fs2.{Stream, io, text}
import java.nio.file.Paths

// Define our list of lines to remove
val myList = List("chen", "yval")

def converter[F[_]: Sync]: F[Int] = {
  val inputPath = Paths.get("testdata/old.txt")
  val outputPath = Paths.get("testdata/new.txt")

  // Create a stream of lines from the input file
  val lineStream = io.file.readAll[F](inputPath, 4096)
    .through(text.utf8Decode)
    .through(text.lines)

  // Split the stream into two identical copies (single-pass processing)
  val (discardStream, keepStream) = lineStream.broadcast(2).splitAt(2)

  // Count how many lines are in our exclusion list
  val countAction = discardStream
    .filter(myList.contains)
    .compile
    .count
    .map(_.toInt) // Convert Long to Int (adjust to Long for extremely large files)

  // Write lines not in the exclusion list to the output file
  val writeAction = keepStream
    .filterNot(myList.contains)
    .intersperse("\n") // Re-add newlines between lines (preserves original structure)
    .through(text.utf8Encode)
    .through(io.file.writeAll(outputPath))
    .compile
    .drain

  // Run both actions and return the final count
  countAction.product(writeAction).map(_._1)
}

Key Details

  • Single-Pass Efficiency: Using broadcast ensures we only read the input file once, which is critical for performance with large files.
  • Pure Functional Style: All operations align with Cats Effect's Sync type class, avoiding mutable state or unmanaged side effects.
  • Backpressure Handling: FS2 automatically manages backpressure, so we never load the entire file into memory—ideal for processing very large datasets.

If you expect line counts to exceed Int limits, simply change the return type of converter to F[Long] and remove the .map(_.toInt) call.

内容的提问来源于stack exchange,提问作者vkt

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 03:47:11