You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 12:50:22