如何配置Spring Batch CSV Writer输出到S3并实现内容追加
Spring Batch写入S3追加CSV数据解决方案
问题根源
S3是对象存储服务,不支持本地文件系统的随机追加写入能力:
- 常规PutObject上传每次都会直接覆盖目标路径的已有对象,这是你第二个方案每次覆盖数据的核心原因
- Spring提供的S3 WritableResource底层封装的就是PutObject接口,所以每次调用
getOutputStream()都会发起一次新的覆盖上传请求 - 你最后全局持有输出流的写法不符合S3的上传逻辑:S3单次上传的输出流有超时限制,长时间不关闭会被服务端主动断开,所以会抛出流已关闭的异常,且只有流正常关闭后S3才会生成最终对象,这也是你看不到文件生成的原因
推荐解决方案
方案一:本地临时文件缓存,作业完成后统一上传(最适配你的需求)
你的需求明确提到待批处理作业完成后将文件存储到S3桶中,该方案完全匹配场景,性能最高,实现最简单:
- 写入阶段直接往本地临时文件追加写入,和你本地开发的逻辑完全一致
- 作业执行完成后,一次性将完整的CSV文件上传到S3
代码实现
- 批处理Writer逻辑:
// 用系统临时目录存储中间文件,避免权限问题 String tempCsvPath = System.getProperty("java.io.tmpdir") + "/subscription_output.csv" @Override void write(List<? extends Subscription> items) throws Exception { // 开启追加模式写入本地临时文件 try (FileWriter fileWriter = new FileWriter(tempCsvPath, true); CSVPrinter csvPrinter = new CSVPrinter(fileWriter, CSVFormat.DEFAULT.withDelimiter('|' as char))) { items.each { subscription -> csvPrinter.printRecord(subscription.id, subscription.subscription) } } }
- 新增作业监听器,作业完成后上传S3:
@Component class S3UploadJobListener implements JobExecutionListener { @Autowired S3Client s3Client private final String tempCsvPath = System.getProperty("java.io.tmpdir") + "/subscription_output.csv" private final String s3Bucket = "你的S3桶名" private final String s3FileKey = "output.csv" @Override void beforeJob(JobExecution jobExecution) { // 作业启动前删除旧的临时文件,避免残留脏数据 Files.deleteIfExists(Paths.get(tempCsvPath)) } @Override void afterJob(JobExecution jobExecution) { // 仅作业执行成功时上传文件 if (jobExecution.getStatus() == BatchStatus.COMPLETED) { PutObjectRequest uploadRequest = PutObjectRequest.builder() .bucket(s3Bucket) .key(s3FileKey) .build() s3Client.putObject(uploadRequest, Paths.get(tempCsvPath)) // 上传完成后清理临时文件 Files.deleteIfExists(Paths.get(tempCsvPath)) } } }
- 将监听器绑定到你的Spring Batch作业配置中即可。
方案二:S3分片上传实现边处理边写入
如果你需要作业运行过程中就把数据同步到S3,可以使用S3的分片上传能力,每批次数据作为一个分片上传,作业完成后合并分片:
// 作业启动时初始化分片上传,仅执行一次 String uploadId List<CompletedPart> partList = [] private final String s3Bucket = "你的S3桶名" private final String s3FileKey = "output.csv" @PostConstruct void initUpload() { CreateMultipartUploadRequest initReq = CreateMultipartUploadRequest.builder() .bucket(s3Bucket) .key(s3FileKey) .build() uploadId = s3Client.createMultipartUpload(initReq).uploadId() } @Override void write(List<? extends Subscription> items) throws Exception { // 把当前批次数据转为字节数组 ByteArrayOutputStream baos = new ByteArrayOutputStream() try (CSVPrinter csvPrinter = new CSVPrinter(new OutputStreamWriter(baos), CSVFormat.DEFAULT.withDelimiter('|' as char))) { items.each { subscription -> csvPrinter.printRecord(subscription.id, subscription.subscription) } } byte[] partData = baos.toByteArray() // 上传当前分片 UploadPartRequest partReq = UploadPartRequest.builder() .bucket(s3Bucket) .key(s3FileKey) .uploadId(uploadId) .partNumber(partList.size() + 1) .contentLength(partData.length as long) .build() String etag = s3Client.uploadPart(partReq, RequestBody.fromBytes(partData)).eTag() partList.add(CompletedPart.builder().partNumber(partList.size() + 1).eTag(etag).build()) } // 作业完成后合并分片,生成最终文件 void completeUpload() { CompletedMultipartUpload completedUpload = CompletedMultipartUpload.builder() .parts(partList) .build() CompleteMultipartUploadRequest completeReq = CompleteMultipartUploadRequest.builder() .bucket(s3Bucket) .key(s3FileKey) .uploadId(uploadId) .multipartUpload(completedUpload) .build() s3Client.completeMultipartUpload(completeReq) }
内容的提问来源于stack exchange,提问作者Moon Rise
相关产品推荐
相关产品推荐

