Spring Batch处理器修改状态引发Chunk漏处理数据的解决方案咨询
ItemReader按状态从表中选取数据,Processor将数据状态从READY修改为PROCESSED。已处理的数据不再被下一个Chunk的查询包含,导致部分数据完全未被处理。
第一个Chunk正确选取并处理了ID为1-10的数据,将其状态设置为PROCESSED,执行的SQL为:
SELECT t.* FROM the_table t WHERE t.status = READY OFFSET 0 FETCH NEXT 10 ROWS ONLY;
下一个Chunk尝试通过以下SQL选取ID为11-20的数据:
SELECT t.* FROM the_table t WHERE t.status = READY OFFSET 10 FETCH NEXT 10 ROWS ONLY;
但由于1-10的数据状态已变更,不再被查询包含,实际选取的是ID为21-30的数据,ID为11-20的数据从未被处理。
将Chunk size设置为足够大以覆盖所有READY数据可得到预期结果,但这不适用于存在大量数据的生产环境。这种情况会导致部分数据未被处理,具体取决于设置的并行Chunk数量。
核心疑问
- 是否可以预选取待处理的ID,基于前一步的结果进行Chunk处理?
- 如何在整个Step中使用仅选取一次的固定数据列表?
注:上述SQL为简化版,实际ItemReader使用的复杂SQL会按特定列(ssn)分组选取最新可处理数据,且该分组下无SUCCESS状态的更新数据,SQL如下:
SELECT bi.*, ri.* FROM ri_item ri JOIN item bi ON ri.item_id = bi.id WHERE (ri.ssn, bi.created_timestamp) IN (SELECT ri2.ssn, MAX(bi2.created_timestamp) FROM ri_item ri2 JOIN item bi2 ON ri2.item_id = bi2.id WHERE ri2.status IN ('READY', 'FAILED_RECOVERABLE') AND bi2.created_timestamp > COALESCE((SELECT MAX(bi3.created_timestamp) FROM ri_item ri3 JOIN item bi3 ON ri3.item_id = bi3.id WHERE ri3.status IN ('SUCCESS') AND ri3.ssn = ri2.ssn), TO_TIMESTAMP('1970-01-01', 'YYYY-MM-DD')) GROUP BY ri2.ssn )
尝试过相关建议但结果相同,希望获取Process Indicator模式的简单示例,并确认该模式是否适用于当前场景。
1. 预选取待处理ID并分Chunk处理
这是解决该问题的常用方案,核心思路是先一次性锁定所有待处理的数据ID,再基于这些ID分批次读取处理,避免后续查询因数据状态变更导致的偏移错误。
实现步骤:
- 第一步:预查询所有符合条件的ID:在Step启动时,执行一次查询获取所有需要处理的数据ID,将这些ID存入一个临时存储(比如内存列表、临时表)。
示例SQL(对应业务场景):SELECT ri.id, bi.id FROM ri_item ri JOIN item bi ON ri.item_id = bi.id WHERE (ri.ssn, bi.created_timestamp) IN (SELECT ri2.ssn, MAX(bi2.created_timestamp) FROM ri_item ri2 JOIN item bi2 ON ri2.item_id = bi2.id WHERE ri2.status IN ('READY', 'FAILED_RECOVERABLE') AND bi2.created_timestamp > COALESCE((SELECT MAX(bi3.created_timestamp) FROM ri_item ri3 JOIN item bi3 ON ri3.item_id = bi3.id WHERE ri3.status IN ('SUCCESS') AND ri3.ssn = ri2.ssn), TO_TIMESTAMP('1970-01-01', 'YYYY-MM-DD')) GROUP BY ri2.ssn ) - 第二步:分Chunk读取ID对应的完整数据:ItemReader不再直接查询
READY状态的数据,而是从预存的ID列表中每次取Chunk Size数量的ID,再通过ID查询完整数据进行处理。 - 第三步:处理后更新状态:Processor依旧将数据状态从
READY改为PROCESSED,但此时后续Chunk的读取基于预存ID,不会受状态变更影响。
2. 使用Process Indicator模式
该模式适用于需要避免重复处理、同时解决分页偏移问题的场景,核心是在查询时先将数据标记为「正在处理」状态,而非直接修改为最终状态,确保后续查询不会重复选取,同时分页逻辑不受状态变更影响。
简单示例(基于业务场景):
第一步:原子标记+查询数据
假设表的status字段新增PROCESSING状态,通过原子操作将符合条件的READY/FAILED_RECOVERABLE数据标记为PROCESSING,同时返回这些数据:
WITH candidate_data AS ( SELECT ri.id, bi.id FROM ri_item ri JOIN item bi ON ri.item_id = bi.id WHERE (ri.ssn, bi.created_timestamp) IN (SELECT ri2.ssn, MAX(bi2.created_timestamp) FROM ri_item ri2 JOIN item bi2 ON ri2.item_id = bi2.id WHERE ri2.status IN ('READY', 'FAILED_RECOVERABLE') AND bi2.created_timestamp > COALESCE((SELECT MAX(bi3.created_timestamp) FROM ri_item ri3 JOIN item bi3 ON ri3.item_id = bi3.id WHERE ri3.status IN ('SUCCESS') AND ri3.ssn = ri2.ssn), TO_TIMESTAMP('1970-01-01', 'YYYY-MM-DD')) GROUP BY ri2.ssn ) LIMIT 10 -- Chunk Size ) UPDATE ri_item ri SET status = 'PROCESSING' FROM candidate_data cd WHERE ri.id = cd.id RETURNING ri.*, bi.*; -- 返回处理中的完整数据
第二步:处理后更新最终状态
Processor处理完数据后,将PROCESSING状态更新为PROCESSED(或SUCCESS):
UPDATE ri_item SET status = 'PROCESSED' WHERE id = ?;
第三步:异常回滚处理
如果Chunk处理失败,需将PROCESSING状态回滚为FAILED_RECOVERABLE,确保数据可被后续重试处理。
模式适配性说明
该模式完全适用于当前场景:
- 原子性的「标记+查询」操作确保每个数据只会被一个Chunk选取,避免分页偏移问题;
- 区分「处理中」和「已处理」状态,避免未完成处理的数据被排除在后续查询之外;
- 支持大数量数据的分批处理,无需设置超大Chunk Size。
3. 整个Step使用固定数据列表
如果数据量在内存可承受范围内,可在Step启动时一次性查询所有符合条件的完整数据,存入内存列表,然后ItemReader从该列表中按Chunk Size读取数据。这种方式简单直接,但仅适用于数据量不大的场景。
实现思路:
- 在Step的
beforeStep方法中执行完整查询,将结果存入线程安全的列表; - 自定义ItemReader,每次从列表中读取Chunk Size数量的元素;
- Processor处理数据后更新状态,由于读取基于内存列表,不会受数据库状态变更影响。
内容的提问来源于stack exchange,提问作者Tomáš Mika

