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

如何用Jackson将大对象序列化为Flux<ByteBuffer>流式上传至AWS S3

流式序列化Java大对象到AWS S3(低内存占用方案)

核心思路

不用先把整个JSON序列化到内存,而是让Jackson直接将JSON数据流式写入自定义输出流,同时把输出流中的字节块转换成Flux<ByteBuffer>,再传给AWS S3的异步上传接口。这样内存里只会保留固定大小的缓冲区,不会加载整个2GB的JSON数据。

具体实现步骤

1. 自定义ByteBuffer输出流

编写一个OutputStream子类,将写入的字节收集到缓冲区,缓冲区满了就通过FluxSink发射对应的ByteBuffer,最后处理剩余字节:

import reactor.core.publisher.FluxSink;
import java.io.IOException;
import java.io.OutputStream;
import java.nio.ByteBuffer;

public class ByteBufferFluxOutputStream extends OutputStream {
    private final FluxSink<ByteBuffer> sink;
    private final ByteBuffer buffer;

    public ByteBufferFluxOutputStream(FluxSink<ByteBuffer> sink, int bufferSize) {
        this.sink = sink;
        this.buffer = ByteBuffer.allocate(bufferSize);
    }

    @Override
    public void write(int b) throws IOException {
        if (!buffer.hasRemaining()) {
            flushBuffer();
        }
        buffer.put((byte) b);
    }

    @Override
    public void write(byte[] b, int off, int len) throws IOException {
        while (len > 0) {
            if (!buffer.hasRemaining()) {
                flushBuffer();
            }
            int writeLen = Math.min(buffer.remaining(), len);
            buffer.put(b, off, writeLen);
            off += writeLen;
            len -= writeLen;
        }
    }

    @Override
    public void flush() throws IOException {
        flushBuffer();
    }

    @Override
    public void close() throws IOException {
        flush();
        sink.complete();
    }

    private void flushBuffer() {
        if (buffer.position() == 0) {
            return;
        }
        buffer.flip();
        sink.next(buffer.asReadOnlyBuffer());
        buffer.clear();
    }
}

2. 生成Flux

用Reactor的Flux.create创建流式序列,将Jackson序列化操作放到单独线程池,避免阻塞Reactor事件循环:

import com.fasterxml.jackson.databind.ObjectMapper;
import reactor.core.publisher.Flux;
import java.nio.ByteBuffer;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class JsonFluxGenerator {
    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
    // 可根据内存情况调整,推荐8KB-64KB
    private static final int BUFFER_SIZE = 8192;
    // 专门的序列化线程池,避免占用IO线程
    private static final ExecutorService SERIALIZATION_POOL = Executors.newSingleThreadExecutor();

    public static Flux<ByteBuffer> serializeToByteBufferFlux(Object largeObject) {
        return Flux.create(sink -> {
            SERIALIZATION_POOL.submit(() -> {
                try (OutputStream outputStream = new ByteBufferFluxOutputStream(sink, BUFFER_SIZE)) {
                    OBJECT_MAPPER.writeValue(outputStream, largeObject);
                } catch (Exception e) {
                    sink.error(e);
                }
            });
        });
    }
}

3. 结合AWS S3异步上传

将生成的Flux<ByteBuffer>传给S3的异步上传接口:

import software.amazon.awssdk.services.s3.S3AsyncClient;
import software.amazon.awssdk.services.s3.model.PutObjectRequest;
import software.amazon.awssdk.core.async.AsyncRequestBody;
import reactor.core.publisher.Mono;

public class S3Uploader {
    private final S3AsyncClient s3AsyncClient;

    public S3Uploader(S3AsyncClient s3AsyncClient) {
        this.s3AsyncClient = s3AsyncClient;
    }

    public Mono<Void> uploadLargeObject(String bucketName, String key, Object largeObject) {
        PutObjectRequest request = PutObjectRequest.builder()
                .bucket(bucketName)
                .key(key)
                .build();

        Flux<ByteBuffer> jsonFlux = JsonFluxGenerator.serializeToByteBufferFlux(largeObject);
        AsyncRequestBody requestBody = AsyncRequestBody.fromPublisher(jsonFlux);

        return Mono.fromFuture(s3AsyncClient.putObject(request, requestBody))
                .then();
    }
}

关键注意事项

  • 缓冲区大小:根据内存情况调整BUFFER_SIZE,太小会频繁发射数据,太大则占用更多内存。
  • 线程池:序列化是CPU密集型操作,必须放到专门线程池,避免影响其他异步任务。
  • 异常处理:序列化出错时,sink.error(e)会将错误传递给订阅者,终止上传流程。
  • 资源关闭:自定义输出流的close方法会触发sink.complete(),通知S3上传完成,Jackson也会自动关闭序列化资源。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 11:52:46