基于Spring Batch实现参数化查询与SOAP调用的数据迁移咨询
解决方案实现指南
一、Repository带参数查询实现
Repository只负责数据访问逻辑,别把SOAP调用塞进来,违反单一职责。用Spring Data JPA定义带参数查询很简单:
示例代码
// 旧库客户Repository public interface OldCustomerRepository extends JpaRepository<OldCustomer, Long> { } // 旧库预订Repository,根据客户ID查询所有预订 public interface OldBookingRepository extends JpaRepository<OldBooking, Long> { List<OldBooking> findByCustomerId(Long customerId); } // 新库对应Repository public interface NewCustomerRepository extends JpaRepository<NewCustomer, Long> { } public interface NewBookingRepository extends JpaRepository<NewBooking, Long> { } public interface NewBookingInfoRepository extends JpaRepository<NewBookingInfo, Long> { }
二、SOAP调用单独封装成Service
把SOAP请求、响应映射逻辑抽成独立的Service,避免和数据层耦合:
示例代码
@Service public class SoapBookingClient { private final WebServiceTemplate webServiceTemplate; // 注入Spring WS的WebServiceTemplate,提前配置好SOAP端点地址 public SoapBookingClient(WebServiceTemplate webServiceTemplate) { this.webServiceTemplate = webServiceTemplate; } public NewBookingInfo getAndMapBookingDetails(String bookingType, String bookingRef) { // 构造SOAP请求体 BookingRequest soapRequest = new BookingRequest(); soapRequest.setBookingType(bookingType); soapRequest.setBookingReference(bookingRef); // 调用SOAP接口 Object rawResponse = webServiceTemplate.marshalSendAndReceive(soapRequest); // 根据预订类型做响应格式转换 if ("SPECIAL_TYPE".equals(bookingType)) { return convertSpecialTypeResponse((SpecialTypeBookingResponse) rawResponse); } else { return convertNormalTypeResponse((NormalTypeBookingResponse) rawResponse); } } // 特殊类型响应转换 private NewBookingInfo convertSpecialTypeResponse(SpecialTypeBookingResponse response) { NewBookingInfo info = new NewBookingInfo(); info.setBookingId(response.getSpecialBookingId()); info.setCustomerContact(response.getSpecialContactInfo()); info.setBookingDate(response.getSpecialCreateDate()); // 其他字段映射 return info; } // 普通类型响应转换 private NewBookingInfo convertNormalTypeResponse(NormalTypeBookingResponse response) { NewBookingInfo info = new NewBookingInfo(); info.setBookingId(response.getBookingId()); info.setCustomerContact(response.getCustomerPhone()); info.setBookingDate(response.getCreateTime()); // 其他字段映射 return info; } }
三、Spring Batch作业整合
分三个步骤实现需求,用Chunk模式处理高延迟场景:
1. 客户与预订数据迁移到新Schema
@Bean public Step migrateDataStep(OldCustomerRepository oldCustomerRepo, OldBookingRepository oldBookingRepo, NewCustomerRepository newCustomerRepo, NewBookingRepository newBookingRepo) { return stepBuilderFactory.get("migrateDataStep") .<OldCustomer, NewCustomer>chunk(20) // 批量迁移,减少数据库交互 .reader(new RepositoryItemReaderBuilder<OldCustomer>() .repository(oldCustomerRepo) .methodName("findAll") .build()) .processor(oldCustomer -> { // 转换旧客户实体到新实体 NewCustomer newCustomer = new NewCustomer(); newCustomer.setId(oldCustomer.getId()); newCustomer.setName(oldCustomer.getFullName()); newCustomer.setEmail(oldCustomer.getEmailAddress()); // 批量迁移该客户的所有预订 List<OldBooking> oldBookings = oldBookingRepo.findByCustomerId(oldCustomer.getId()); List<NewBooking> newBookings = oldBookings.stream() .map(oldBooking -> { NewBooking newBooking = new NewBooking(); newBooking.setId(oldBooking.getId()); newBooking.setCustomerId(oldBooking.getCustomerId()); newBooking.setBookingType(oldBooking.getType()); newBooking.setBookingRef(oldBooking.getReference()); return newBooking; }) .collect(Collectors.toList()); newBookingRepo.saveAll(newBookings); return newCustomer; }) .writer(new RepositoryItemWriterBuilder<NewCustomer>() .repository(newCustomerRepo) .methodName("save") .build()) .build(); }
2. 调用SOAP接口并持久化结果
针对超过5条延迟高的问题,把Chunk size设为5,控制并发调用数量:
@Bean public Step processSoapBookingStep(NewBookingRepository newBookingRepo, SoapBookingClient soapClient, NewBookingInfoRepository newBookingInfoRepo) { return stepBuilderFactory.get("processSoapBookingStep") .<NewBooking, NewBookingInfo>chunk(5) .reader(new RepositoryItemReaderBuilder<NewBooking>() .repository(newBookingRepo) .methodName("findAll") .build()) .processor(newBooking -> { // 调用SOAP客户端获取并转换数据 return soapClient.getAndMapBookingDetails(newBooking.getBookingType(), newBooking.getBookingRef()); }) .writer(new RepositoryItemWriterBuilder<NewBookingInfo>() .repository(newBookingInfoRepo) .methodName("save") .build()) // 配置重试策略,处理SOAP调用超时 .faultTolerant() .retry(SoapFaultClientException.class) .retryLimit(3) .build(); }
3. 整合作业
@Bean public Job bookingMigrationJob(Step migrateDataStep, Step processSoapBookingStep) { return jobBuilderFactory.get("bookingMigrationJob") .start(migrateDataStep) .next(processSoapBookingStep) .build(); }
关键注意事项
- Repository职责单一:只做数据库CRUD和查询,业务逻辑(比如SOAP调用、格式转换)放到Service或Processor里。
- 延迟优化:通过Chunk size控制SOAP调用并发数,避免同时发起大量请求导致超时;也可以配置异步Processor进一步优化。
- 异常处理:给SOAP调用加重试、超时配置,避免单个请求失败导致整个作业中断。
内容的提问来源于stack exchange,提问作者Tyali Nexi
相关产品推荐
相关产品推荐

