使用Akka Streams+Slick从PostgreSQL导大量数据到S3内存不足求助
解决Akka Streams + Slick处理大数据时的内存不足问题
你的问题根源很明确:当前代码是一次性把所有200万条记录加载到内存中,MyRepository.getAllRecordss()返回的Future[Seq[MyEntityImpl]]会让Slick把整个查询结果拉取到JVM堆里,当数据量超过10万条时,内存自然就不够用了。
要解决这个问题,我们需要改成流式处理——让Slick分批从数据库拉取数据,Akka Streams边处理边写入S3,全程不把所有数据存在内存里。
第一步:修改Repository,返回流式Source
Slick支持流式查询,我们可以把getAllRecordss()从返回Future[Seq]改成返回Akka Streams的Source,这样就能分批拉取数据:
import akka.stream.scaladsl.Source import slick.jdbc.PostgresProfile.api._ class MyRepository(db: Database) { // 假设你的表查询是MyEntityTable.query def getAllRecords(): Source[MyEntityImpl, _] = { // 用withStatementParameters设置fetchSize,控制每次从数据库拉取的行数 val streamingQuery = MyEntityTable.query.result .withStatementParameters(fetchSize = 1000) // 可根据内存情况调整,比如5000 // 将Slick的流式结果转换成Akka Streams Source Source.fromPublisher(db.stream(streamingQuery)) } }
这里的fetchSize很关键,它告诉PostgreSQL每次返回多少条数据,避免一次性拉取全量。
第二步:重构Akka Streams处理流程
现在不需要先把所有数据加载到内存,直接用上面的Source来构建流:
val results: Future[MultipartUploadResult] = MyRepository.getAllRecords() .map(myEntity => myEntity.toPSV + "\n") // 转换为PSV格式 .map(ByteString(_)) // 转为ByteString供S3 Sink使用 .runWith(s3Sink)
这样修改后,整个流程是完全响应式的:
- Slick每次从数据库拉取1000条记录
- Akka Streams逐条转换为PSV字符串,再转为ByteString
- 数据被S3 Sink写入到S3,处理完的记录会被GC回收
- 基于Akka Streams的背压机制,下游S3写入速度慢的话,上游会自动暂停拉取数据,不会导致内存积压
额外优化建议
- 优化PSV字符串拼接:如果你的实体字段很多,用
StringBuilder代替直接的字符串插值会更高效,减少临时对象创建:case class MyPartOne(field1: String, field2: String) { def toPSV: String = new StringBuilder() .append(field1) .append('|') .append(field2) .toString() } - 调整S3 Sink的分块大小:如果你用的是Alpakka S3的
MultipartUploadSink,可以配置partSize(默认5MB),比如设置为10MB,减少S3的分块上传请求次数,提升效率。 - 监控内存使用:可以用JVM监控工具(比如VisualVM)观察处理过程中的内存占用,根据实际情况调整
fetchSize的值。
内容的提问来源于stack exchange,提问作者Anthony
相关产品推荐
相关产品推荐

