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

Spring Batch如何动态调用Step配置重新初始化Partitioner实例

问题描述

我有一个包含ItemReader、ItemProcessor和NoOpItemWriter的Spring Batch作业,通过Rest控制器传入作业参数启动该作业,具体细节如下:

  • ItemReader从Servlet Context中的列表读取数据
  • ItemProcessor处理数据时执行数据库调用
  • ItemWriter无实际操作
  • 分区器为List Partitioner实现,基于索引范围对列表进行分区

初始状态下Servlet Context仅包含一个占位项(dummy item),因此Bean初始化时List Partitioner会使用该占位项初始化。但实际需求是根据Post请求传入的数据填充列表,让Partitioner使用该列表进行分区。

当前存在的问题:Spring Boot应用启动时,Partitioner已使用占位列表完成初始化,无法识别Rest控制器传入的新列表。


Step与Job配置

@Bean(name = ["controllerStep"])
protected fun controllerStep(
    jobRepository: JobRepository,
    transactionManager: PlatformTransactionManager
): Step {

    servletConfiguration.setServletContext()
    return StepBuilder("controllerStep", jobRepository)
        .partitioner(
            "workerStep",
            ListPartitioner(servletConfiguration.getItemProcessingList()!!.size)
        )
        .step(workerStep(jobRepository, transactionManager))
        .gridSize(6)
        .build()
}


@Bean
fun job(jobRepository: JobRepository?, transactionManager: PlatformTransactionManager?): Job? {
    return JobBuilder("job", jobRepository!!)
        .start(controllerStep)
        .build()
}

Servlet配置类

@Configuration
class ServletConfiguration(private var servletContext: ServletContext) {

    fun setServletContext() {
        servletContext.setAttribute("DatatoBeProcessed", listOf(""))
    }
    @Suppress("UNCHECKED_CAST")
    fun getItemProcessingList(): List<String>? {
        val contextList = servletContext.getAttribute(ATTRIBUTE_NAME_PROCESSING_INPUT_DATA)
        return if (contextList != null) contextList as? List<String> else listOf("")
    }
}

Rest控制器代码

// 设置列表到Servlet Context的方法
dataService.setListsIntoServletContext(jobRequestDto)
@PostMapping("/trigger")
fun startJob(@RequestBody jobDTO: JobDTO): ResponseEntity<String> {

    // 获取待处理列表并保存到Servlet Context
    dataService.setListsIntoServletContext(jobRequestDto)

    // 创建作业参数
    val jobParameters =
        JobParametersBuilder()
            .addString("JobUniquekey", Random.nextInt().toString()) // 待重构
            .addString(
                "status",
                jobDTO.status.toString(),
            )
            .addString("jobType", jobDTO.jobType.toString())
            .toJobParameters()

    // 触发作业
    jobLauncher.run(job, jobParameters)
    return ResponseEntity("JOB trigger:SUCCESS", HttpStatus.ACCEPTED)
}

解决方案

问题核心是ListPartitioner在Spring容器启动时就完成初始化,此时Servlet Context只有占位列表,后续Rest请求传入的新列表无法被已初始化的Partitioner感知。要解决这个问题,需让Partitioner动态获取当前Servlet Context中的列表数据,而非初始化时固定参数。

1. 自定义动态ListPartitioner

创建自定义Partitioner,每次分区操作时从Servlet Context读取最新列表:

@Component
class DynamicListPartitioner(private val servletConfiguration: ServletConfiguration) : Partitioner {

    override fun partition(gridSize: Int): MutableMap<String, ExecutionContext> {
        val itemList = servletConfiguration.getItemProcessingList() ?: emptyList()
        val totalItems = itemList.size
        if (totalItems == 0) {
            return mutableMapOf()
        }

        val partitionSize = Math.ceil(totalItems.toDouble() / gridSize).toInt()
        val partitions = mutableMapOf<String, ExecutionContext>()

        var start = 0
        var partitionNumber = 1
        while (start < totalItems) {
            val end = Math.min(start + partitionSize, totalItems)
            val context = ExecutionContext()
            context.putInt("startIndex", start)
            context.putInt("endIndex", end)
            partitions["partition$partitionNumber"] = context
            start = end
            partitionNumber++
        }

        return partitions
    }
}

2. 修改Step配置,使用自定义Partitioner

不再在Bean初始化时传入固定列表大小,直接注入自定义的DynamicListPartitioner,让它在作业执行时动态计算分区:

@Bean(name = ["controllerStep"])
protected fun controllerStep(
    jobRepository: JobRepository,
    transactionManager: PlatformTransactionManager,
    dynamicListPartitioner: DynamicListPartitioner
): Step {
    return StepBuilder("controllerStep", jobRepository)
        .partitioner("workerStep", dynamicListPartitioner)
        .step(workerStep(jobRepository, transactionManager))
        .gridSize(6)
        .build()
}

3. 调整ServletConfiguration逻辑

移除硬编码的占位项设置,避免覆盖后续传入的真实数据:

@Configuration
class ServletConfiguration(private var servletContext: ServletContext) {

    @Suppress("UNCHECKED_CAST")
    fun getItemProcessingList(): List<String>? {
        val contextList = servletContext.getAttribute(ATTRIBUTE_NAME_PROCESSING_INPUT_DATA)
        return contextList as? List<String> ?: emptyList()
    }

    // 提供设置列表的方法,供dataService调用
    fun setItemProcessingList(list: List<String>) {
        servletContext.setAttribute(ATTRIBUTE_NAME_PROCESSING_INPUT_DATA, list)
    }
}

4. 确保Rest控制器正确设置列表

在dataService.setListsIntoServletContext内部调用servletConfiguration.setItemProcessingList(),保证作业启动前Servlet Context中的列表已更新为传入的真实数据。

修改后,每次作业启动时,DynamicListPartitioner都会读取最新的列表计算分区,彻底解决初始化过早的问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 16:35:14