Spring Batch同作业中将数据库导出至多个文件的实现方案
按属性分组导出JSON文件的自动化实现方案
完全可以实现这个需求,基于你正在使用的Spring Batch组件,通过自定义分组ItemWriter就能完成全量数据的自动分组导出,以下是具体实现方案:
核心思路
用RepositoryItemReader分页读取全量数据,通过自定义的ItemWriter根据指定属性(attribute1)的值动态创建对应文件的JsonFileItemWriter,将同一属性值的记录写入对应JSON文件。
实现步骤
1. 自定义分组ItemWriter
这个Writer会维护一个属性值到JsonFileItemWriter的映射,每次处理数据时自动匹配或创建对应的文件写入器:
import org.springframework.batch.item.ItemWriter; import org.springframework.batch.item.json.JsonFileItemWriter; import org.springframework.batch.item.json.builder.JsonFileItemWriterBuilder; import org.springframework.core.io.FileSystemResource; import com.fasterxml.jackson.databind.ObjectMapper; import org.springframework.batch.item.json.JacksonJsonObjectMarshaller; import java.io.IOException; import java.util.HashMap; import java.util.List; import java.util.Map; public class GroupedJsonItemWriter<T> implements ItemWriter<T> { private final String groupAttribute; private final Map<Object, JsonFileItemWriter<T>> writerMap = new HashMap<>(); private final Class<T> itemClass; private final ObjectMapper objectMapper; public GroupedJsonItemWriter(String groupAttribute, Class<T> itemClass, ObjectMapper objectMapper) { this.groupAttribute = groupAttribute; this.itemClass = itemClass; this.objectMapper = objectMapper; } @Override public void write(List<? extends T> items) throws Exception { for (T item : items) { // 获取分组属性值(这里假设实体有对应的getter方法) Object attributeValue = getItemAttributeValue(item); // 不存在则创建对应文件的Writer JsonFileItemWriter<T> writer = writerMap.computeIfAbsent(attributeValue, this::createJsonWriter); writer.write(List.of(item)); } } private Object getItemAttributeValue(T item) throws ReflectiveOperationException { // 通过反射调用getter方法获取属性值,也可以用BeanUtils简化 String getterName = "get" + groupAttribute.substring(0, 1).toUpperCase() + groupAttribute.substring(1); return item.getClass().getMethod(getterName).invoke(item); } private JsonFileItemWriter<T> createJsonWriter(Object attributeValue) { String fileName = attributeValue + ".json"; return new JsonFileItemWriterBuilder<T>() .resource(new FileSystemResource(fileName)) .jsonObjectMarshaller(new JacksonJsonObjectMarshaller<>(objectMapper)) .append(true) // 追加模式,避免覆盖已有内容 .build(); } // Step结束时关闭所有Writer,释放资源 public void closeAllWriters() throws IOException { for (JsonFileItemWriter<T> writer : writerMap.values()) { writer.close(); } } }
2. 配置Spring Batch作业
配置分页读取器、自定义Writer和作业流程,确保18万条数据分批处理,避免内存溢出:
import org.springframework.batch.core.Job; import org.springframework.batch.core.Step; import org.springframework.batch.core.configuration.annotation.EnableBatchProcessing; import org.springframework.batch.core.configuration.annotation.JobBuilderFactory; import org.springframework.batch.core.configuration.annotation.StepBuilderFactory; import org.springframework.batch.core.listener.StepExecutionListenerSupport; import org.springframework.batch.item.data.RepositoryItemReader; import org.springframework.batch.item.data.builder.RepositoryItemReaderBuilder; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.data.domain.Sort; import java.util.Collections; @Configuration @EnableBatchProcessing public class BatchExportConfig { private final JobBuilderFactory jobBuilderFactory; private final StepBuilderFactory stepBuilderFactory; private final YourEntityRepository entityRepository; // 替换为你的Repository private final ObjectMapper objectMapper; public BatchExportConfig(JobBuilderFactory jobBuilderFactory, StepBuilderFactory stepBuilderFactory, YourEntityRepository entityRepository, ObjectMapper objectMapper) { this.jobBuilderFactory = jobBuilderFactory; this.stepBuilderFactory = stepBuilderFactory; this.entityRepository = entityRepository; this.objectMapper = objectMapper; } @Bean public RepositoryItemReader<YourEntity> entityReader() { return new RepositoryItemReaderBuilder<YourEntity>() .repository(entityRepository) .methodName("findAll") .pageSize(1000) // 分页读取,控制内存占用 .sort(Collections.singletonMap("id", Sort.Direction.ASC)) // 必须指定排序,保证分页稳定性 .build(); } @Bean public GroupedJsonItemWriter<YourEntity> groupedJsonWriter() { return new GroupedJsonItemWriter<>("attribute1", YourEntity.class, objectMapper); // 替换为你的分组属性名 } @Bean public Step exportStep() { return stepBuilderFactory.get("groupedExportStep") .<YourEntity, YourEntity>chunk(1000) // Chunk大小与分页大小匹配 .reader(entityReader()) .writer(groupedJsonWriter()) .listener(new StepExecutionListenerSupport() { @Override public void afterStep(org.springframework.batch.core.StepExecution stepExecution) { try { groupedJsonWriter().closeAllWriters(); } catch (IOException e) { throw new RuntimeException("Failed to close JSON writers", e); } } }) .build(); } @Bean public Job exportJob() { return jobBuilderFactory.get("groupedJsonExportJob") .start(exportStep()) .build(); } }
关键注意事项
- 内存控制:18万条数据必须分页读取,
pageSize建议设置为1000-2000,避免一次性加载大量数据到内存。 - 资源释放:必须在Step结束时调用
closeAllWriters()关闭所有文件写入器,防止文件句柄泄漏和数据未完全写入。 - 属性获取优化:如果觉得反射繁琐,可以用
BeanUtils.getProperty(item, groupAttribute)简化属性值的获取。 - 多分组场景:如果
attribute1的取值非常多(比如上万种),需要检查操作系统的文件句柄限制,必要时调整系统参数。 - 异常处理:可以在自定义Writer中加入异常捕获逻辑,记录失败的记录和对应的属性值,方便后续排查。
内容的提问来源于stack exchange,提问作者cvetan
相关产品推荐
相关产品推荐

