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

Spring Batch多线程执行时SQL Server事务恢复失败问题求助

Spring Batch多线程批量任务SQL Server事务异常排查与解决

核心问题分析

启用多线程TaskExecutor后出现The server failed to resume the transaction异常,单线程模式正常,说明并发场景下事务或连接管理存在问题,结合SQL Server与Spring Batch特性,从以下方向排查解决:

1. 处理器线程安全与作用域优化

Spring Batch多线程Step中,若处理器实例被共享,可能导致事务上下文混乱。当前PartProcessor无状态变量,但建议给处理器添加@StepScope注解,确保每个线程使用独立实例:

@Bean
@StepScope
public ItemProcessor<ExternalPart, Part> processor(PartRepository partsRepository,
                                                   FactoryRepository factoryRepository,
                                                   PartDetailsRepository partDetailsRepository) {
    return new PartProcessor(partsRepository, factoryRepository, partDetailsRepository);
}

2. SQL Server JDBC驱动与连接池配置

SQL Server JDBC驱动在多线程事务场景下需确保连接管理正确,添加以下配置到application.properties:

# 连接池配置
spring.datasource.hikari.connection-test-query=SELECT 1
spring.datasource.hikari.max-lifetime=1800000
spring.datasource.hikari.leak-detection-threshold=60000
# 禁用Open Session In View避免会话长期持有
spring.jpa.open-in-view=false

同时升级mssql-jdbc到最新稳定版(如12.6.0.jre17),确保与SQL Server版本兼容。

3. 数据查询缓存优化

多个线程频繁查询同一PartDetails会增加DB负载与事务冲突风险,给PartDetailsRepository的查询方法添加缓存:

@Cacheable(value = "partDetails", key = "#partNumber")
PartDetails findByPartNumber(String partNumber);

需提前配置Spring Cache(如使用Caffeine缓存)。

4. TaskExecutor参数调整

当前队列容量过小易导致线程池饱和,调整参数降低线程频繁创建销毁的影响:

@Bean
public TaskExecutor taskExecutor() {
    ThreadPoolTaskExecutor taskExecutor = new ThreadPoolTaskExecutor();
    taskExecutor.setCorePoolSize(4);
    taskExecutor.setMaxPoolSize(6);
    taskExecutor.setQueueCapacity(50);
    taskExecutor.setThreadNamePrefix("parts-");
    taskExecutor.setWaitForTasksToCompleteOnShutdown(true);
    taskExecutor.initialize();
    return taskExecutor;
}

5. Chunk大小调整

当前Chunk大小为500,多线程下单事务数据量过大易引发事务超时或冲突,尝试减小到100-200:

return new StepBuilder("partStep", jobRepository)
        .<ExternalPart, Part>chunk(150, transactionManager)
        .reader(reader)
        .processor(processor)
        .writer(writer)
        .taskExecutor(taskExecutor())
        .build();

问题原代码参考

技术栈与业务场景

基于Spring Boot 3.x、Spring Batch 5.1.1、Spring Data JPA、Hibernate、Microsoft SQL Server、Gradle、JDK 17开发,批量任务从平面文件读取数据,处理后写入数据库(插入/更新Part记录),多个Part可关联同一PartDetails。

实体映射

PartDetails实体

@Getter
@Setter
@ToString
@Audited
@Entity
@Table(name = "part_details")
public class PartDetails {

    @Column(name = "part_number", unique = true)
    @Id
    private String partNumber;
    private String description;
    ...
}

Part实体

@Getter
@Setter
@ToString
@Audited
@Entity
@Table(name = "part")
public class Part {

    private String partNumber;
    private String partStatus;

    @ManyToOne
    @JoinColumn(name = "factory_id", referencedColumnName = "id", nullable = false)
    private Factory factory;

    ...

    @ManyToOne
    @JoinColumn(name = "part_details_id")
    private PartDetails partDetails;

}

实体仓库

PartRepository

@Repository
public interface PartRepository extends JpaRepository<Part, Long> {

    Part findByPartNumberAndFactoryCode(String partNumber, Integer factoryCode);

    ...

}

PartDetailsRepository

@Repository
public interface PartDetailsRepository extends JpaRepository<PartDetails, String> {

    PartDetails findByPartNumber(String partNumber);

}

Spring Batch Item处理器

PartProcessor

@Slf4j
@RequiredArgsConstructor
public class PartProcessor implements ItemProcessor<ExternalPart, Part> {

    private final PartRepository partRepository;
    private final FactoryRepository factoryRepository;
    private final PartDetailsRepository partDetailsRepository;

    @Override
    public Part process(ExternalPart item) {
        Integer factoryCode = Integer.parseInt(item.getFactory());
        String partNumber = item.getArticleNumber();
        Part part = partRepository.findByPartNumberAndFactoryCode(partNumber, factoryCode);

        if (part == null) {
            part = new Part();

            Factory factory = factoryRepository.findByCode(factoryCode);

            if (factory == null) {
                log.error("Unknown factory code {} received. Ignoring this {} article number",
                        item.getFactory(), item.getArticleNumber());

                return null;
            }

            part.setFactory(factory);
            part.setPartNumber(partNumber);
            part.setPartStatus(StringUtils.EMPTY);
        }

        PartDetails partDetails = partDetailsRepository.findByPartNumber(partNumber);

        part.setPartDetails(partDetails);

        return part;
    }
}

批量配置

PartsBatchConfiguration

@Configuration
public class PartsBatchConfiguration {


    @Bean
    @StepScope
    public SynchronizedItemStreamReader<ExternalPart> reader() {
        ...
    }

    @Bean
    public TaskExecutor taskExecutor() {
        ThreadPoolTaskExecutor taskExecutor = new ThreadPoolTaskExecutor();

        taskExecutor.setCorePoolSize(6);
        taskExecutor.setMaxPoolSize(8);
        taskExecutor.setQueueCapacity(10);
        taskExecutor.setThreadNamePrefix("parts-");

        return taskExecutor;
    }

    @Bean
    public RepositoryItemWriter<Part> writer(PartRepository partsRepository) {
        RepositoryItemWriter<Part> writer = new RepositoryItemWriter<>();
        writer.setRepository(partsRepository);
        return writer;
    }

    @Bean
    public Job job(Step step, JobRepository jobRepository) {
        return new JobBuilder("importParts", jobRepository)
                .incrementer(new RunIdIncrementer())
                .flow(step)
                .end()
                .build();
    }

    @Bean
    public Step step(SynchronizedItemStreamReader<ExternalPart> reader, JobRepository jobRepository,
                          PlatformTransactionManager transactionManager,
                          ItemProcessor<ExternalPart, Part> processor,
                          RepositoryItemWriter<Part> writer) {
        return new StepBuilder("partStep", jobRepository)
                .<ExternalPart, Part>chunk(500, transactionManager)
                .reader(reader)
                .processor(processor)
                .writer(writer)
                .taskExecutor(taskExecutor())
                .build();
    }

    @Bean
    public ItemProcessor<ExternalPart, Part> processor(PartRepository partsRepository,
                                                           FactoryRepository factoryRepository,
                                                           PartDetailsRepository partDetailsRepository) {
        return new PartProcessor(partsRepository, factoryRepository, partDetailsRepository);
    }

}

Gradle构建文件

plugins {
    id 'java'
    id 'org.springframework.boot' version '3.2.3'
}

apply plugin: 'io.spring.dependency-management'

configurations {
    compileOnly {
        extendsFrom annotationProcessor
    }
}

group 'my.parts'
version = '1.0'

java {
    sourceCompatibility = JavaVersion.VERSION_17
    targetCompatibility = JavaVersion.VERSION_17
}

repositories {
    mavenCentral()
}

dependencyManagement {
    imports {
        mavenBom "org.springframework.cloud:spring-cloud-dependencies:2023.0.0"
    }
}

// other variables
def lombokVersion = '1.18.30'

dependencies {

    // spring
    implementation 'org.springframework.boot:spring-boot-starter-data-jpa'
    implementation 'org.springframework.boot:spring-boot-starter-web'
    implementation 'org.springframework.boot:spring-boot-starter-validation'
    implementation 'org.springframework.boot:spring-boot-starter-batch'

    implementation 'org.springframework.cloud:spring-cloud-starter-openfeign'

    implementation 'org.springdoc:springdoc-openapi-starter-webmvc-ui:2.3.0'

    // db
    runtimeOnly 'com.microsoft.sqlserver:mssql-jdbc'
    implementation 'org.hibernate:hibernate-envers:6.4.3.Final'

    // other libs
    implementation 'org.apache.poi:poi:5.2.5'
    implementation 'org.apache.poi:poi-ooxml:5.2.5'
    implementation 'org.apache.commons:commons-lang3:3.14.0'
    implementation 'org.apache.commons:commons-collections4:4.4'
    implementation 'io.github.openfeign:feign-okhttp:13.1'
    annotationProcessor "org.projectlombok:lombok:$lombokVersion"

    // tests
    testImplementation 'org.springframework.boot:spring-boot-starter-test'
    testImplementation 'org.springframework.batch:spring-batch-test'
    testAnnotationProcessor "org.projectlombok:lombok:$lombokVersion"
    testCompileOnly "org.projectlombok:lombok:$lombokVersion"
}

test {
    useJUnitPlatform()
}

异常信息

启用多线程TaskExecutor时出现以下错误:

2024-03-07T15:43:18.296+02:00  WARN 22487 --- [        parts-4] o.h.engine.jdbc.spi.SqlExceptionHelper   : SQL Error: 3971, SQLState: S0001
2024-03-07T15:43:18.296+02:00 ERROR 22487 --- [        parts-4] o.h.engine.jdbc.spi.SqlExceptionHelper   : The server failed to resume the transaction. Desc:4300000001.
2024-03-07T15:43:18.299+02:00  INFO 22487 --- [        parts-4] o.h.e.internal.DefaultLoadEventListener  : HHH000327: Error performing load command

org.hibernate.exception.SQLGrammarException: JDBC exception executing SQL [select pd1_0.part_number,pd1_0.created_at,pd1_0.created_by,pd1_0.desc_english,pd1_0.desc_swedish,pd1_0.part_color,pd1_0.part_sort,pd1_0.part_type,pd1_0.updated_at,pd1_0.updated_by,pd1_0.weight from part_details pd1_0 where pd1_0.part_number=?] [The server failed to resume the transaction. Desc:4300000001.] [n/a]
    at org.hibernate.exception.internal.SQLStateConversionDelegate.convert(SQLStateConversionDelegate.java:91) ~[hibernate-core-6.4.4.Final.jar:6.4.4.Final]
    ...(省略部分栈信息)
Caused by: com.microsoft.sqlserver.jdbc.SQLServerException: The server failed to resume the transaction. Desc:4300000001.
    ...(省略部分栈信息)

内容的提问来源于stack exchange,提问作者ydrozhdzhal

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 10:07:02