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

Scala+Circe处理大对象转JSON时OOM问题及流处理方案咨询

解决Scala+Circe超大列表JSON编码内存溢出问题

完全可以用迭代器或Akka Reactive Streams解决这个问题,核心思路是避免一次性将整个列表加载到内存并编码,而是采用流式处理,分块或逐个生成JSON片段,最终拼接成完整的JSON数组。

一、用迭代器结合Circe实现轻量流式编码

Circe本身支持基于迭代器的JSON生成,不需要加载整个列表到内存。只需手动处理JSON数组的首尾符号和元素分隔符,配合输出流边生成边写入即可:

import io.circe._
import io.circe.syntax._
import java.io.{BufferedWriter, FileWriter}

// 示例元素类型及对应的Circe编码器
case class Data(id: Int, value: String)
implicit val dataEncoder: Encoder[Data] = Encoder.forProduct2("id", "value")(d => (d.id, d.value))

// 超大迭代器(可从数据库、文件等流式数据源获取)
val largeDataIterator: Iterator[Data] = (1 to 1000000).map(i => Data(i, s"value_$i")).iterator

// 流式写入JSON数组到文件
val writer = new BufferedWriter(new FileWriter("output.json"))
try {
  writer.write("[")
  if (largeDataIterator.hasNext) {
    writer.write(largeDataIterator.next().asJson.noSpaces)
    while (largeDataIterator.hasNext) {
      writer.write(",")
      writer.write(largeDataIterator.next().asJson.noSpaces)
    }
  }
  writer.write("]")
} finally {
  writer.close()
}

这种方式下,内存中始终只保留当前处理的单个元素和少量写入缓冲区,不会出现内存堆积。

二、用Akka Reactive Streams实现带背压的流式处理

如果需要更复杂的流控制(比如背压、错误处理、多阶段数据转换),Akka Streams是更优雅的选择。结合Circe可以实现全链路的流式JSON编码:

import akka.actor.ActorSystem
import akka.stream.scaladsl.{FileIO, Flow, Source}
import akka.util.ByteString
import io.circe._
import io.circe.syntax._
import java.nio.file.Paths

case class Data(id: Int, value: String)
implicit val dataEncoder: Encoder[Data] = Encoder.forProduct2("id", "value")(d => (d.id, d.value))

implicit val system: ActorSystem = ActorSystem("JsonStreaming")

// 从迭代器创建Akka Streams Source
val dataSource: Source[Data, _] = Source.fromIterator(() => (1 to 1000000).map(i => Data(i, s"value_$i")).iterator)

// 构建流式处理逻辑
val jsonStream = dataSource
  // 单个元素编码为JSON字符串
  .map(_.asJson.noSpaces)
  // 为元素添加逗号分隔符(跳过第一个元素)
  .scan(Option.empty[String]) {
    case (None, first) => Some(first)
    case (Some(prev), next) => Some(s"$prev,$next")
  }
  // 过滤初始空值,保留有效JSON片段
  .collect { case Some(str) => str }
  // 拼接JSON数组的首尾符号
  .concat(Source.single("]"))
  .prepend(Source.single("["))
  // 转换为ByteString用于文件写入
  .map(ByteString(_))
  // 写入目标文件
  .runWith(FileIO.toPath(Paths.get("output.json")))

// 流处理完成后关闭ActorSystem
jsonStream.onComplete { _ =>
  system.terminate()
}

Akka Streams会自动处理背压,确保元素生产速度与写入速度匹配,内存占用始终可控。如果你的数据源本身是流式的(比如数据库分页读取、Kafka消费),可以直接替换Source的创建方式,无需先转换为迭代器。

注意事项

  • 绝对避免直接对整个超大列表调用asJson,这会一次性生成完整的JSON字符串,瞬间占用大量内存
  • 若需要格式化JSON,可用asJson.spaces2代替noSpaces,但紧凑格式能减少IO开销,更适合流式场景
  • 优先将结果写入文件或网络流,而非内存中的字符串对象

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 18:53:22