task-68: payload 引用删除改为批量引用检查与异步物理删除

- 新增 TransientPayloadDeleteOrchestrator:提交去重入队(maxPendingDeletes 上限),
  flush 时一次性 IN 批量反查 biz_task_chunk / biz_task_scope_state,
  未引用对象交由后台线程池异步物理删除(剥 rustfs: 前缀)
- 幂等:重复提交只入队一次;保守:查询失败整批保留可重试,不误删
- 空值/非指针忽略,本地/OSS 指针不在批量删除范围(各有归属与清理路径)
This commit is contained in:
2026-08-30 19:13:43 +08:00
parent 45b23b559c
commit 2de58b7afc
2 changed files with 426 additions and 0 deletions
@@ -0,0 +1,205 @@
package com.nanri.aiimage.modules.task.service;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.nanri.aiimage.modules.file.service.object.RustfsObjectStorageService;
import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper;
import com.nanri.aiimage.modules.task.mapper.TaskScopeStateMapper;
import com.nanri.aiimage.modules.task.model.entity.TaskChunkEntity;
import com.nanri.aiimage.modules.task.model.entity.TaskScopeStateEntity;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ExecutorService;
import java.util.function.Supplier;
/**
* P2-9 扩展(task-68):payload 引用删除改为批量引用检查与异步物理删除。
* <p>调用方把待删 payload(指针或 JSON 编码指针)批量提交,本组件去重入队;
* {@link #flushPendingDeletes()} 对整批 pending 一次性 IN 反查
* biz_task_chunk / biz_task_scope_state(比逐条两次查询少一个数量级的 DB 往返),
* 确认不再被引用的对象交由后台线程池异步物理删除,不阻塞业务线程。
* <p>幂等:重复提交相同对象只入队一次;保守:引用检查失败时本批保留,
* 调用方(如周期清理任务)可再次 flush 重试。
*/
@Service
@Slf4j
public class TransientPayloadDeleteOrchestrator {
private static final String RUSTFS_POINTER_PREFIX = "rustfs:";
private final TransientPayloadStorageService transientPayloadStorageService;
private final RustfsObjectStorageService rustfsObjectStorageService;
private final TaskChunkMapper taskChunkMapper;
private final TaskScopeStateMapper taskScopeStateMapper;
private final ObjectMapper objectMapper;
private final ExecutorService asyncDeleteExecutor;
private final Set<String> pendingPointers = new LinkedHashSet<>();
@Value("${aiimage.transient-storage.max-pending-deletes:1000}")
private long maxPendingDeletes = 1000;
public TransientPayloadDeleteOrchestrator(TransientPayloadStorageService transientPayloadStorageService,
RustfsObjectStorageService rustfsObjectStorageService,
TaskChunkMapper taskChunkMapper,
TaskScopeStateMapper taskScopeStateMapper,
ObjectMapper objectMapper,
ExecutorService asyncDeleteExecutor) {
this.transientPayloadStorageService = transientPayloadStorageService;
this.rustfsObjectStorageService = rustfsObjectStorageService;
this.taskChunkMapper = taskChunkMapper;
this.taskScopeStateMapper = taskScopeStateMapper;
this.objectMapper = objectMapper;
this.asyncDeleteExecutor = asyncDeleteExecutor;
}
/**
* 批量提交待删 payload 值(指针或 JSON 编码指针)。空值与非指针值忽略;
* 已在 pending 中的对象幂等跳过;pending 达到上限后拒绝新提交并返回实际入队数。
*/
public int submitDeletes(List<String> values) {
if (values == null || values.isEmpty()) {
return 0;
}
synchronized (pendingPointers) {
int accepted = 0;
for (String value : values) {
String pointer = transientPayloadStorageService.extractPointer(value);
if (pointer == null) {
continue;
}
if (pendingPointers.size() >= maxPendingDeletes) {
log.warn("[transient-payload] delete queue full, drop {} pointer={}", "submit", pointer);
break;
}
if (pendingPointers.add(pointer)) {
accepted++;
}
}
return accepted;
}
}
/**
* 对 pending 中的对象批量做引用检查,未引用的异步物理删除,返回本批删除数。
* 引用检查异常时保守跳过整批(不删),调用方可再次 flush 重试。
*/
public int flushPendingDeletes() {
List<String> batch;
synchronized (pendingPointers) {
if (pendingPointers.isEmpty()) {
return 0;
}
batch = new ArrayList<>(pendingPointers);
pendingPointers.clear();
}
try {
Set<String> stillReferenced = batchReferencedPointers(batch);
List<String> toDelete = new ArrayList<>(batch);
toDelete.removeAll(stillReferenced);
if (!toDelete.isEmpty()) {
asyncDeleteExecutor.submit(() -> deleteObjects(toDelete));
}
if (!stillReferenced.isEmpty()) {
log.info("[transient-payload] skip delete, still referenced count={}", stillReferenced.size());
}
return toDelete.size();
} catch (Exception ex) {
log.warn("[transient-payload] batch reference check failed, keep pending count={} err={}",
batch.size(), ex.getMessage());
synchronized (pendingPointers) {
pendingPointers.addAll(batch);
}
return 0;
}
}
public int pendingCount() {
synchronized (pendingPointers) {
return pendingPointers.size();
}
}
private void deleteObjects(List<String> pointers) {
for (String pointer : pointers) {
if (pointer == null || !pointer.startsWith(RUSTFS_POINTER_PREFIX)) {
// 本地/OSS 指针有实例归属与其它清理路径,不在此批量删除范围。
continue;
}
String objectKey = pointer.substring(RUSTFS_POINTER_PREFIX.length());
try {
rustfsObjectStorageService.deleteObject(objectKey);
} catch (Exception ex) {
log.warn("[transient-payload] async delete failed objectKey={} err={}",
objectKey, ex.getMessage());
}
}
}
/** 批量 IN 反查两张引用表,返回仍被引用的指针集合(查询异常抛给调用方)。 */
private Set<String> batchReferencedPointers(List<String> pointers) {
List<String> jsonEncoded = pointers.stream()
.map(this::jsonEncodePointer)
.filter(java.util.Objects::nonNull)
.toList();
List<String> candidates = new ArrayList<>(pointers);
candidates.addAll(jsonEncoded);
Set<String> referenced = new LinkedHashSet<>();
List<TaskChunkEntity> chunks = taskChunkMapper.selectList(new LambdaQueryWrapper<TaskChunkEntity>()
.in(TaskChunkEntity::getPayloadJson, candidates));
for (TaskChunkEntity chunk : chunks) {
referenced.addAll(matchingPointers(chunk.getPayloadJson(), pointers));
}
List<TaskScopeStateEntity> scopeStates = taskScopeStateMapper.selectList(new LambdaQueryWrapper<TaskScopeStateEntity>()
.and(w -> w.in(TaskScopeStateEntity::getParsedPayloadJson, candidates)
.or()
.in(TaskScopeStateEntity::getStateJson, candidates)));
for (TaskScopeStateEntity scopeState : scopeStates) {
referenced.addAll(matchingPointers(scopeState.getParsedPayloadJson(), pointers));
referenced.addAll(matchingPointers(scopeState.getStateJson(), pointers));
}
return referenced;
}
private Set<String> matchingPointers(String dbValue, List<String> pointers) {
Set<String> matches = new LinkedHashSet<>();
if (dbValue == null || dbValue.isBlank()) {
return matches;
}
String candidate = dbValue.trim();
for (int i = 0; i < 4; i++) {
String pointer = transientPayloadStorageService.extractPointer(candidate);
if (pointer != null && pointers.contains(pointer)) {
matches.add(pointer);
}
String decoded;
try {
decoded = objectMapper.readValue(candidate, String.class);
} catch (Exception ex) {
break;
}
if (decoded == null || decoded.equals(candidate)) {
break;
}
candidate = decoded.trim();
}
return matches;
}
private String jsonEncodePointer(String pointer) {
try {
return objectMapper.writeValueAsString(pointer);
} catch (Exception ex) {
return null;
}
}
}