task-29: 合并 scope 状态查询与更新,减少单 chunk 数据库往返

persistResultScope 复用接收前已加载的 scope,不再重复 selectOne ——
每个分片接收的 scope 查询从 2 次降为 1 次;空分片校验前移至任何
scope 查询之前,非法分片零数据库往返。任务级分布式锁保证预读 scope
在锁内不过期。
This commit is contained in:
2026-08-29 19:36:30 +08:00
parent 24ea02f150
commit ec28216984
2 changed files with 615 additions and 8 deletions
@@ -1251,11 +1251,6 @@ public class ShopDataCrawlTaskService {
int chunkTotal = incoming.getChunkTotal();
validateChunkMetadata(chunkIndex, chunkTotal);
String scopeKey = resultChunkScopeKey(shopKey);
String scopeHash = DigestUtil.sha256Hex(scopeKey);
TaskScopeStateEntity scope = findResultScope(taskId, scopeHash);
validateChunkTotal(scope == null ? null : scope.getChunkTotal(), chunkTotal);
List<ShopDataCrawlCountryResultDto> countryResults = copyCountryResults(incoming.getCountryResults());
if (!hasProcessableChunkData(incoming.getCountryResults())) {
throw new BusinessException("店铺数据抓取结果分片内容为空,拒绝接收");
@@ -1263,6 +1258,11 @@ public class ShopDataCrawlTaskService {
String payloadJson = writeJson(countryResults, "序列化店铺数据抓取结果分片失败");
String payloadHash = DigestUtil.sha256Hex(payloadJson);
String scopeKey = resultChunkScopeKey(shopKey);
String scopeHash = DigestUtil.sha256Hex(scopeKey);
TaskScopeStateEntity scope = findResultScope(taskId, scopeHash);
validateChunkTotal(scope == null ? null : scope.getChunkTotal(), chunkTotal);
ensureRustfsPayloadStorageEnabled();
String storedPayload = transientPayloadStorageService.storeChunkPayloadVersioned(
MODULE_TYPE, taskId, scopeHash, chunkIndex, payloadJson);
@@ -1300,7 +1300,7 @@ public class ShopDataCrawlTaskService {
}
try {
persistResultScope(taskId, scopeKey, scopeHash, chunkTotal, receivedChunkCount);
persistResultScope(scope, taskId, scopeKey, scopeHash, chunkTotal, receivedChunkCount);
} catch (RuntimeException ex) {
// 分片行与 scope 计数器是一致性单元:状态写入失败则回滚本次插入的分片行与 payload,
// 客户端重试会走全新插入路径,计数器不会因重试重复计数或永久欠计。
@@ -1380,12 +1380,16 @@ public class ShopDataCrawlTaskService {
return chunkTotal != null && chunkTotal > 0 ? Math.min(chunkTotal, count) : count;
}
private void persistResultScope(Long taskId,
/**
* 复用接收前已加载的 scope(Task 29:合并查询与更新,单 chunk 的 scope 往返 2 次降为 1 次)。
* 任务级分布式锁串行化同一任务的接收,预读 scope 在锁内不会过期。
*/
private void persistResultScope(TaskScopeStateEntity scope,
Long taskId,
String scopeKey,
String scopeHash,
int chunkTotal,
int receivedChunkCount) {
TaskScopeStateEntity scope = findResultScope(taskId, scopeHash);
validateChunkTotal(scope == null ? null : scope.getChunkTotal(), chunkTotal);
LocalDateTime now = LocalDateTime.now();
if (scope == null) {