如何用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
相关产品推荐
相关产品推荐

