Spring?Batch?數(shù)據(jù)處理的實(shí)現(xiàn)
一、Spring Batch 核心概念
Spring Batch 是 Spring 生態(tài)系統(tǒng)中用于批處理的框架,它提供了強(qiáng)大的批處理功能,支持大規(guī)模數(shù)據(jù)處理。
1.1 核心概念
- Job:批處理作業(yè),是批處理的頂層概念
- Step:作業(yè)的步驟,一個(gè)作業(yè)由一個(gè)或多個(gè)步驟組成
- ItemReader:讀取數(shù)據(jù)的組件
- ItemProcessor:處理數(shù)據(jù)的組件
- ItemWriter:寫入數(shù)據(jù)的組件
- JobRepository:存儲(chǔ)作業(yè)執(zhí)行狀態(tài)的倉庫
- JobLauncher:啟動(dòng)作業(yè)的組件
- JobExecution:作業(yè)執(zhí)行實(shí)例
- StepExecution:步驟執(zhí)行實(shí)例
1.2 Spring Batch 的優(yōu)勢
- 可擴(kuò)展性:支持大規(guī)模數(shù)據(jù)處理
- 可靠性:支持事務(wù)管理和重啟機(jī)制
- 可監(jiān)控性:提供詳細(xì)的執(zhí)行狀態(tài)和日志
- 靈活性:支持多種數(shù)據(jù)源和處理方式
- 集成性:與 Spring 生態(tài)系統(tǒng)無縫集成
二、Spring Batch 配置
2.1 基本配置
@Configuration
@EnableBatchProcessing
public class BatchConfig {
@Autowired
private JobBuilderFactory jobBuilderFactory;
@Autowired
private StepBuilderFactory stepBuilderFactory;
@Bean
public ItemReader<User> userItemReader() {
// 從數(shù)據(jù)庫讀取數(shù)據(jù)
return new JdbcCursorItemReaderBuilder<User>()
.name("userItemReader")
.dataSource(dataSource)
.sql("SELECT id, name, email FROM users WHERE status = 'ACTIVE'")
.rowMapper(new BeanPropertyRowMapper<>(User.class))
.build();
}
@Bean
public ItemProcessor<User, UserDTO> userItemProcessor() {
return user -> {
UserDTO dto = new UserDTO();
dto.setId(user.getId());
dto.setName(user.getName().toUpperCase());
dto.setEmail(user.getEmail().toLowerCase());
return dto;
};
}
@Bean
public ItemWriter<UserDTO> userItemWriter() {
// 寫入到文件
return items -> {
for (UserDTO item : items) {
System.out.println("Processing user: " + item.getName());
// 寫入到文件或其他目標(biāo)
}
};
}
@Bean
public Step processUserStep() {
return stepBuilderFactory.get("processUserStep")
.<User, UserDTO>chunk(10)
.reader(userItemReader())
.processor(userItemProcessor())
.writer(userItemWriter())
.build();
}
@Bean
public Job processUserJob() {
return jobBuilderFactory.get("processUserJob")
.incrementer(new RunIdIncrementer())
.flow(processUserStep())
.end()
.build();
}
}2.2 數(shù)據(jù)源配置
@Configuration
public class DataSourceConfig {
@Bean
public DataSource dataSource() {
HikariConfig config = new HikariConfig();
config.setJdbcUrl("jdbc:mysql://localhost:3306/batch_db");
config.setUsername("root");
config.setPassword("password");
config.setMaximumPoolSize(10);
return new HikariDataSource(config);
}
@Bean
public JdbcTemplate jdbcTemplate(DataSource dataSource) {
return new JdbcTemplate(dataSource);
}
}2.3 作業(yè)倉庫配置
@Configuration
public class JobRepositoryConfig {
@Bean
public JobRepository jobRepository(DataSource dataSource, PlatformTransactionManager transactionManager) throws Exception {
JobRepositoryFactoryBean factory = new JobRepositoryFactoryBean();
factory.setDataSource(dataSource);
factory.setTransactionManager(transactionManager);
factory.setIsolationLevelForCreate("ISOLATION_SERIALIZABLE");
factory.setTablePrefix("BATCH_");
factory.setMaxVarCharLength(1000);
return factory.getObject();
}
@Bean
public PlatformTransactionManager transactionManager(DataSource dataSource) {
return new DataSourceTransactionManager(dataSource);
}
@Bean
public JobLauncher jobLauncher(JobRepository jobRepository) throws Exception {
SimpleJobLauncher launcher = new SimpleJobLauncher();
launcher.setJobRepository(jobRepository);
launcher.setTaskExecutor(new SimpleAsyncTaskExecutor());
return launcher;
}
}三、ItemReader 實(shí)現(xiàn)
3.1 數(shù)據(jù)庫讀取
@Bean
public ItemReader<Customer> customerItemReader(DataSource dataSource) {
return new JdbcPagingItemReaderBuilder<Customer>()
.name("customerItemReader")
.dataSource(dataSource)
.selectClause("SELECT id, first_name, last_name, email, phone")
.fromClause("FROM customers")
.whereClause("WHERE last_updated > :lastUpdated")
.parameterValues(Collections.singletonMap("lastUpdated", LocalDateTime.now().minusDays(1)))
.sortKeys(Collections.singletonMap("id", Order.ASCENDING))
.rowMapper(new BeanPropertyRowMapper<>(Customer.class))
.pageSize(100)
.build();
}3.2 文件讀取
@Bean
public ItemReader<Product> productItemReader() {
return new FlatFileItemReaderBuilder<Product>()
.name("productItemReader")
.resource(new ClassPathResource("products.csv"))
.delimited()
.names("id", "name", "price", "quantity")
.fieldSetMapper(fieldSet -> {
Product product = new Product();
product.setId(fieldSet.readLong("id"));
product.setName(fieldSet.readString("name"));
product.setPrice(fieldSet.readBigDecimal("price"));
product.setQuantity(fieldSet.readInt("quantity"));
return product;
})
.build();
}3.3 自定義讀取器
public class CustomItemReader implements ItemReader<String> {
private final List<String> items;
private int index = 0;
public CustomItemReader(List<String> items) {
this.items = items;
}
@Override
public String read() {
if (index < items.size()) {
return items.get(index++);
}
return null;
}
}
@Bean
public ItemReader<String> customItemReader() {
List<String> items = Arrays.asList("item1", "item2", "item3", "item4", "item5");
return new CustomItemReader(items);
}四、ItemProcessor 實(shí)現(xiàn)
4.1 基本處理器
public class ProductProcessor implements ItemProcessor<Product, ProductDTO> {
@Override
public ProductDTO process(Product item) {
ProductDTO dto = new ProductDTO();
dto.setId(item.getId());
dto.setName(item.getName());
dto.setPrice(item.getPrice());
dto.setQuantity(item.getQuantity());
dto.setTotalValue(item.getPrice().multiply(BigDecimal.valueOf(item.getQuantity())));
return dto;
}
}
@Bean
public ItemProcessor<Product, ProductDTO> productProcessor() {
return new ProductProcessor();
}4.2 條件處理
public class OrderProcessor implements ItemProcessor<Order, Order> {
@Override
public Order process(Order item) {
if (item.getStatus().equals(OrderStatus.PENDING)) {
item.setStatus(OrderStatus.PROCESSED);
item.setProcessedAt(LocalDateTime.now());
return item;
}
return null; // 跳過非待處理訂單
}
}
@Bean
public ItemProcessor<Order, Order> orderProcessor() {
return new OrderProcessor();
}4.3 復(fù)合處理器
public class CompositeItemProcessor<T, R> implements ItemProcessor<T, R> {
private final List<ItemProcessor> processors;
public CompositeItemProcessor(List<ItemProcessor> processors) {
this.processors = processors;
}
@Override
public R process(T item) {
Object result = item;
for (ItemProcessor processor : processors) {
result = processor.process(result);
if (result == null) {
return null;
}
}
return (R) result;
}
}
@Bean
public ItemProcessor<Customer, CustomerDTO> customerProcessor() {
List<ItemProcessor> processors = new ArrayList<>();
processors.add(new ValidationProcessor());
processors.add(new TransformationProcessor());
processors.add(new EnrichmentProcessor());
return new CompositeItemProcessor<>(processors);
}五、ItemWriter 實(shí)現(xiàn)
5.1 數(shù)據(jù)庫寫入
@Bean
public ItemWriter<CustomerDTO> customerItemWriter(DataSource dataSource) {
return new JdbcBatchItemWriterBuilder<CustomerDTO>()
.dataSource(dataSource)
.sql("INSERT INTO customer_processed (id, first_name, last_name, email, processed_at) VALUES (:id, :firstName, :lastName, :email, :processedAt)")
.itemSqlParameterSourceProvider(new BeanPropertyItemSqlParameterSourceProvider<>())
.build();
}
5.2 文件寫入
@Bean
public ItemWriter<ProductDTO> productItemWriter() {
return new FlatFileItemWriterBuilder<ProductDTO>()
.name("productItemWriter")
.resource(new FileSystemResource("output/products-processed.csv"))
.delimited()
.names("id", "name", "price", "quantity", "totalValue")
.headerCallback(writer -> writer.write("ID,Name,Price,Quantity,Total Value"))
.build();
}
5.3 自定義寫入器
public class CustomItemWriter implements ItemWriter<User> {
private final Logger logger = LoggerFactory.getLogger(CustomItemWriter.class);
@Override
public void write(List<? extends User> items) {
for (User item : items) {
logger.info("Writing user: {}", item.getName());
// 寫入到外部系統(tǒng)或其他目標(biāo)
}
}
}
@Bean
public ItemWriter<User> customItemWriter() {
return new CustomItemWriter();
}六、作業(yè)執(zhí)行與監(jiān)控
6.1 作業(yè)啟動(dòng)
@Service
public class BatchService {
@Autowired
private JobLauncher jobLauncher;
@Autowired
private Job processUserJob;
public void runProcessUserJob() throws Exception {
JobParameters jobParameters = new JobParametersBuilder()
.addString("jobName", "processUserJob")
.addLong("time", System.currentTimeMillis())
.toJobParameters();
JobExecution execution = jobLauncher.run(processUserJob, jobParameters);
System.out.println("Job execution status: " + execution.getStatus());
}
}6.2 作業(yè)監(jiān)控
@RestController
@RequestMapping("/api/batch")
public class BatchController {
@Autowired
private JobExplorer jobExplorer;
@GetMapping("/jobs")
public List<JobInstance> getJobs() {
return jobExplorer.getJobInstances("processUserJob", 0, 10);
}
@GetMapping("/executions/{jobInstanceId}")
public List<JobExecution> getExecutions(@PathVariable Long jobInstanceId) {
JobInstance jobInstance = jobExplorer.getJobInstance(jobInstanceId);
return jobExplorer.getJobExecutions(jobInstance);
}
@GetMapping("/steps/{jobExecutionId}")
public List<StepExecution> getSteps(@PathVariable Long jobExecutionId) {
JobExecution jobExecution = jobExplorer.getJobExecution(jobExecutionId);
return jobExecution.getStepExecutions();
}
}6.3 作業(yè)調(diào)度
@Configuration
@EnableScheduling
public class BatchScheduler {
@Autowired
private JobLauncher jobLauncher;
@Autowired
private Job processUserJob;
@Scheduled(cron = "0 0 0 * * ?") // 每天凌晨執(zhí)行
public void runDailyJob() throws Exception {
JobParameters jobParameters = new JobParametersBuilder()
.addString("jobName", "processUserJob")
.addLong("time", System.currentTimeMillis())
.toJobParameters();
jobLauncher.run(processUserJob, jobParameters);
}
}七、Spring Batch 最佳實(shí)踐
7.1 性能優(yōu)化
- 合理設(shè)置 chunk 大小:根據(jù)數(shù)據(jù)量和系統(tǒng)資源設(shè)置合適的 chunk 大小
- 使用并行處理:對于大規(guī)模數(shù)據(jù)處理,使用并行步驟
- 優(yōu)化數(shù)據(jù)庫操作:使用批量操作,減少數(shù)據(jù)庫連接次數(shù)
- 使用異步處理:對于IO密集型操作,使用異步處理
7.2 錯(cuò)誤處理
- 跳過策略:設(shè)置合理的跳過策略,處理錯(cuò)誤數(shù)據(jù)
- 重試機(jī)制:對于臨時(shí)錯(cuò)誤,使用重試機(jī)制
- 錯(cuò)誤日志:詳細(xì)記錄錯(cuò)誤信息,便于排查
- 死信隊(duì)列:將無法處理的數(shù)據(jù)放入死信隊(duì)列
@Bean
public Step processOrderStep() {
return stepBuilderFactory.get("processOrderStep")
.<Order, Order>chunk(10)
.reader(orderItemReader())
.processor(orderItemProcessor())
.writer(orderItemWriter())
.faultTolerant()
.skipLimit(10)
.skip(OrderProcessingException.class)
.retryLimit(3)
.retry(ConnectionException.class)
.build();
}7.3 事務(wù)管理
- 合理設(shè)置事務(wù)邊界:根據(jù)業(yè)務(wù)需求設(shè)置合適的事務(wù)邊界
- 使用局部事務(wù):對于不需要全局事務(wù)的步驟,使用局部事務(wù)
- 事務(wù)隔離級(jí)別:根據(jù)業(yè)務(wù)需求設(shè)置合適的事務(wù)隔離級(jí)別
7.4 監(jiān)控與告警
- 作業(yè)執(zhí)行監(jiān)控:監(jiān)控作業(yè)執(zhí)行狀態(tài)和性能
- 錯(cuò)誤告警:對作業(yè)執(zhí)行錯(cuò)誤進(jìn)行告警
- 性能指標(biāo):收集作業(yè)執(zhí)行的性能指標(biāo)
八、生產(chǎn)環(huán)境案例分析
8.1 案例一:電商平臺(tái)數(shù)據(jù)同步
某電商平臺(tái)使用 Spring Batch 實(shí)現(xiàn)了從線下系統(tǒng)到線上系統(tǒng)的數(shù)據(jù)同步。主要功能包括:
- 從線下數(shù)據(jù)庫讀取商品信息
- 處理和轉(zhuǎn)換數(shù)據(jù)格式
- 寫入到線上數(shù)據(jù)庫
- 生成同步報(bào)告
通過 Spring Batch,該平臺(tái)實(shí)現(xiàn)了每天同步超過 100 萬條商品數(shù)據(jù),同步時(shí)間從原來的 4 小時(shí)減少到 30 分鐘,數(shù)據(jù)準(zhǔn)確率達(dá)到 99.99%。
8.2 案例二:金融系統(tǒng)批處理
某銀行使用 Spring Batch 實(shí)現(xiàn)了每日 批處理作業(yè),包括:
- 賬戶余額計(jì)算
- 交易對賬
- 報(bào)表生成
- 風(fēng)險(xiǎn)評(píng)估
通過 Spring Batch,該銀行實(shí)現(xiàn)了每天處理超過 1000 萬筆交易,批處理時(shí)間從原來的 6 小時(shí)減少到 1.5 小時(shí),系統(tǒng)穩(wěn)定性顯著提高。
九、常見誤區(qū)與解決方案
9.1 內(nèi)存溢出
問題:處理大量數(shù)據(jù)時(shí)出現(xiàn)內(nèi)存溢出
解決方案:合理設(shè)置 chunk 大小,使用分頁讀取,避免一次性加載所有數(shù)據(jù)
9.2 事務(wù)管理不當(dāng)
問題:事務(wù)范圍過大,導(dǎo)致鎖定時(shí)間過長
解決方案:合理設(shè)置事務(wù)邊界,使用局部事務(wù)
9.3 錯(cuò)誤處理不完善
問題:錯(cuò)誤處理機(jī)制不完善,導(dǎo)致作業(yè)頻繁失敗
解決方案:設(shè)置合理的跳過策略和重試機(jī)制
9.4 監(jiān)控不足
問題:缺乏對作業(yè)執(zhí)行狀態(tài)的監(jiān)控
解決方案:建立完善的監(jiān)控體系,及時(shí)發(fā)現(xiàn)和解決問題
十、總結(jié)與展望
Spring Batch 是一個(gè)強(qiáng)大的批處理框架,它為企業(yè)級(jí)應(yīng)用提供了可靠、高效的數(shù)據(jù)處理能力。通過合理配置和使用 Spring Batch,可以顯著提高數(shù)據(jù)處理效率,減少人工干預(yù),提高系統(tǒng)可靠性。
在云原生時(shí)代,Spring Batch 也在不斷演進(jìn)。未來,我們將看到 Spring Batch 與云原生技術(shù)的深度融合,如與 Kubernetes 的集成,以及對 Serverless 架構(gòu)的支持,為批處理作業(yè)提供更加靈活、高效的運(yùn)行環(huán)境。
記住,批處理作業(yè)的設(shè)計(jì)應(yīng)該根據(jù)業(yè)務(wù)需求和數(shù)據(jù)特點(diǎn)進(jìn)行合理規(guī)劃。這其實(shí)可以更優(yōu)雅一點(diǎn)。
到此這篇關(guān)于Spring Batch 數(shù)據(jù)處理的實(shí)現(xiàn)的文章就介紹到這了,更多相關(guān)Spring Batch 數(shù)據(jù)處理內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
java算法入門之有效的括號(hào)刪除有序數(shù)組中的重復(fù)項(xiàng)實(shí)現(xiàn)strStr
大家好,我是哪吒,一個(gè)熱愛編碼的Java工程師,本著"欲速則不達(dá),欲達(dá)則欲速"的學(xué)習(xí)態(tài)度,在程序猿這條不歸路上不斷成長,所謂成長,不過是用時(shí)間慢慢擦亮你的眼睛,少時(shí)看重的,年長后卻視若鴻毛,少時(shí)看輕的,年長后卻視若泰山,成長之路,亦是漸漸放下執(zhí)念,內(nèi)心歸于平靜的旅程2021-08-08
SpringBoot讀寫xml上傳到AWS存儲(chǔ)服務(wù)S3的示例
這篇文章主要介紹了SpringBoot讀寫xml上傳到S3的示例,幫助大家更好的理解和使用springboot框架,感興趣的朋友可以了解下2020-10-10
maven多個(gè)plugin相同phase的執(zhí)行順序
這篇文章主要介紹了maven多個(gè)plugin相同phase的執(zhí)行順序,本文給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2020-12-12
SpringCloud-Alibaba-Sentinel服務(wù)降級(jí),熱點(diǎn)限流,服務(wù)熔斷
這篇文章主要介紹了SpringCloud-Alibaba-Sentinel服務(wù)降級(jí),熱點(diǎn)限流,服務(wù)熔斷,本文給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2020-12-12
springboot集成swagger3與knife4j的詳細(xì)代碼
這篇文章主要介紹了springboot集成swagger3與knife4j,本文通過實(shí)例代碼給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2022-08-08
Java鏈表中添加元素的原理與實(shí)現(xiàn)方法詳解
這篇文章主要介紹了Java鏈表中添加元素的原理與實(shí)現(xiàn)方法,結(jié)合實(shí)例形式詳細(xì)分析了Java實(shí)現(xiàn)鏈表中添加元素的相關(guān)原理、操作技巧與注意事項(xiàng),需要的朋友可以參考下2020-03-03
SpringBoot封裝JDBC的實(shí)現(xiàn)步驟
本文主要介紹了SpringBoot封裝JDBC的實(shí)現(xiàn)步驟,文中通過示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2021-12-12
詳解Spring 中如何控制2個(gè)bean中的初始化順序
本篇文章主要介紹了Spring 中如何控制2個(gè)bean中的初始化順序,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2017-10-10
springboot集成mybatisplus實(shí)例詳解
這篇文章主要介紹了springboot集成mybatisplus實(shí)例詳解,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2019-09-09

