新增单个下载任务重试接口
This commit is contained in:
parent
ccc05fda19
commit
9475237731
@ -97,4 +97,11 @@ public class GalleryManageController {
|
||||
public String resetUndone(){
|
||||
return galleryManageService.resetUndone();
|
||||
}
|
||||
|
||||
@PostMapping("/retry")
|
||||
public String retryGallery(Integer gid){
|
||||
if(gid == null)
|
||||
return Response._failure("参数不全");
|
||||
return galleryManageService.retryGallery(gid);
|
||||
}
|
||||
}
|
||||
|
||||
@ -25,7 +25,7 @@ public interface GalleryMapper {
|
||||
@Select("select * from gallery where downloader=#{downloader}")
|
||||
Gallery[] selectGalleryByDownloader(int downloader);
|
||||
|
||||
@Select("select * from gallery where status in ('已提交', '下载中', '压缩中')")
|
||||
@Select("select * from gallery where status in ('已提交', '下载中', '等待压缩', '压缩中')")
|
||||
Gallery[] selectUnDoneGalleries();
|
||||
|
||||
@Select("select * from gallery order by createTime")
|
||||
|
||||
@ -444,4 +444,19 @@ public class GalleryManageService {
|
||||
|
||||
return response.toJSONString();
|
||||
}
|
||||
|
||||
public String retryGallery(int gid){
|
||||
Gallery gallery = galleryMapper.selectGalleryByGid(gid);
|
||||
if(gallery == null)
|
||||
return Response._failure("任务不存在");
|
||||
if("下载完成".equals(gallery.getStatus()))
|
||||
return Response._success("下载完成");
|
||||
if(remoteService.isDead())
|
||||
return Response._failure("节点不在线,无法重试");
|
||||
|
||||
RemoteService.RetryResult retryResult = remoteService.retryGallery(gallery);
|
||||
if(retryResult.success())
|
||||
return Response._success(retryResult.message());
|
||||
return Response._failure(retryResult.message());
|
||||
}
|
||||
}
|
||||
|
||||
@ -25,10 +25,14 @@ import java.io.IOException;
|
||||
import java.io.OutputStream;
|
||||
import java.net.*;
|
||||
import java.util.Arrays;
|
||||
import java.util.HashMap;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.CopyOnWriteArrayList;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.TimeoutException;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
@Service
|
||||
@ -49,7 +53,10 @@ public class RemoteService {
|
||||
|
||||
final PushService pushService;
|
||||
|
||||
HashMap<Integer, Promise<AbstractMessage>> promiseHashMap = new HashMap<>();
|
||||
ConcurrentHashMap<Integer, Promise<AbstractMessage>> promiseHashMap = new ConcurrentHashMap<>();
|
||||
|
||||
ConcurrentHashMap<Integer, CopyOnWriteArrayList<CompletableFuture<String>>> retryStatusWaiters =
|
||||
new ConcurrentHashMap<>();
|
||||
|
||||
EventLoop eventLoopGroup = new DefaultEventLoop();
|
||||
|
||||
@ -150,7 +157,7 @@ public class RemoteService {
|
||||
}
|
||||
|
||||
public boolean isDead(){
|
||||
return channelFuture.channel() == null || !channelFuture.channel().isActive();
|
||||
return channelFuture == null || channelFuture.channel() == null || !channelFuture.channel().isActive();
|
||||
}
|
||||
|
||||
public void resetUndone(){
|
||||
@ -174,10 +181,10 @@ public class RemoteService {
|
||||
DownloadPostMessage dpm = new DownloadPostMessage();
|
||||
dpm.messageId = atomicInteger.getAndIncrement();
|
||||
dpm.setGalleryTask(galleryTask);
|
||||
channel.writeAndFlush(dpm);
|
||||
|
||||
DefaultPromise<AbstractMessage> promise = new DefaultPromise<>(eventLoopGroup);
|
||||
promiseHashMap.put(dpm.messageId, promise);
|
||||
channel.writeAndFlush(dpm);
|
||||
try {
|
||||
boolean result = promise.await(10, TimeUnit.SECONDS);
|
||||
if(result){
|
||||
@ -189,9 +196,45 @@ public class RemoteService {
|
||||
log.warn("addGalleryToQueue interrupted", e);
|
||||
Thread.currentThread().interrupt();
|
||||
return -1;
|
||||
}finally {
|
||||
promiseHashMap.remove(dpm.messageId, promise);
|
||||
}
|
||||
}
|
||||
|
||||
public RetryResult retryGallery(Gallery gallery){
|
||||
CompletableFuture<String> statusFuture = new CompletableFuture<>();
|
||||
retryStatusWaiters.computeIfAbsent(gallery.getGid(), ignored -> new CopyOnWriteArrayList<>())
|
||||
.add(statusFuture);
|
||||
try {
|
||||
byte submitResult = addGalleryToQueue(gallery);
|
||||
if(submitResult != 0 && !statusFuture.isDone())
|
||||
return new RetryResult(false, "节点未接受重试请求");
|
||||
return new RetryResult(true, statusFuture.get(10, TimeUnit.SECONDS));
|
||||
}catch (TimeoutException e){
|
||||
return new RetryResult(false, "节点已收到重试请求,但未及时返回任务状态");
|
||||
}catch (ExecutionException e){
|
||||
log.warn("等待重试状态失败, gid={}", gallery.getGid(), e);
|
||||
return new RetryResult(false, "获取任务状态失败");
|
||||
}catch (InterruptedException e){
|
||||
log.warn("等待重试状态被中断, gid={}", gallery.getGid(), e);
|
||||
Thread.currentThread().interrupt();
|
||||
return new RetryResult(false, "获取任务状态被中断");
|
||||
}finally {
|
||||
retryStatusWaiters.computeIfPresent(gallery.getGid(), (gid, waiters) -> {
|
||||
waiters.remove(statusFuture);
|
||||
return waiters.isEmpty() ? null : waiters;
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
private void completeRetryStatusWaiters(int gid, String status){
|
||||
CopyOnWriteArrayList<CompletableFuture<String>> waiters = retryStatusWaiters.remove(gid);
|
||||
if(waiters != null)
|
||||
waiters.forEach(waiter -> waiter.complete(status));
|
||||
}
|
||||
|
||||
public record RetryResult(boolean success, String message) {}
|
||||
|
||||
public byte deleteGallery(Gallery gallery){
|
||||
DeleteGalleryMessage dgm = new DeleteGalleryMessage();
|
||||
dgm.setGalleryName(gallery.getName());
|
||||
@ -270,16 +313,24 @@ public class RemoteService {
|
||||
}
|
||||
else if(galleryTask.getStatus() == GalleryTask.COMPRESSING)
|
||||
gallery.setStatus("压缩中");
|
||||
else if(galleryTask.getProceeding() != 0)
|
||||
else if(galleryTask.getStatus() == GalleryTask.DOWNLOAD_COMPLETE)
|
||||
gallery.setStatus("等待压缩");
|
||||
else if(galleryTask.getStatus() == GalleryTask.DOWNLOADING)
|
||||
gallery.setStatus("下载中");
|
||||
|
||||
log.info(gallery.getName() + "下载进度:" + gallery.getProceeding() + "/" + gallery.getPages());
|
||||
galleryMapper.updateGallery(gallery);
|
||||
completeRetryStatusWaiters(gallery.getGid(), gallery.getStatus());
|
||||
}
|
||||
webSocketService.updateTaskProcessing(galleryTasks);
|
||||
}
|
||||
else if(msg instanceof ResponseMessage rsm)
|
||||
promiseHashMap.get(rsm.messageId).setSuccess(rsm);
|
||||
else if(msg instanceof ResponseMessage rsm) {
|
||||
Promise<AbstractMessage> promise = promiseHashMap.remove(rsm.messageId);
|
||||
if(promise != null)
|
||||
promise.setSuccess(rsm);
|
||||
else
|
||||
log.warn("收到无等待者的响应消息: messageId={}", rsm.messageId);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
Loading…
Reference in New Issue
Block a user