feat:取消完成通过配置决定是否使用mq或者feign
This commit is contained in:
@@ -5,17 +5,22 @@ import cn.code.nl.framework.execute.core.AbstractTask;
|
||||
import cn.code.nl.framework.execute.core.TaskFactory;
|
||||
import cn.code.nl.framework.execute.core.dto.TaskExecuteDTO;
|
||||
import cn.code.nl.module.task.api.TransportTaskApi;
|
||||
import cn.code.nl.module.task.dto.TaskCallbackResultReqDTO;
|
||||
import cn.code.nl.module.task.dto.TaskInfoDTO;
|
||||
import cn.code.nl.module.task.enums.CallbackStatusEnum;
|
||||
import cn.code.nl.module.task.enums.TransportTaskStatusEnum;
|
||||
import cn.code.nl.module.task.message.TaskEventMessage;
|
||||
import jakarta.annotation.Resource;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
|
||||
import org.apache.rocketmq.spring.core.RocketMQListener;
|
||||
import org.redisson.api.RBucket;
|
||||
import org.redisson.api.RLock;
|
||||
import org.redisson.api.RedissonClient;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import java.time.Duration;
|
||||
|
||||
import static cn.code.nl.module.task.message.LockKeyConstants.TASK_STATUS_CHANGE_LOCK_KEY;
|
||||
|
||||
/**
|
||||
@@ -32,6 +37,12 @@ import static cn.code.nl.module.task.message.LockKeyConstants.TASK_STATUS_CHANGE
|
||||
)
|
||||
public class WmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMessage> {
|
||||
|
||||
/** 任务业务处理成功标识前缀 */
|
||||
private static final String TASK_BUSINESS_HANDLED_KEY = "task:business:handled:";
|
||||
|
||||
/** 任务业务处理成功标识有效期 */
|
||||
private static final Duration TASK_BUSINESS_HANDLED_TIMEOUT = Duration.ofDays(7);
|
||||
|
||||
@Resource
|
||||
private TaskFactory taskFactory;
|
||||
|
||||
@@ -56,7 +67,19 @@ public class WmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMe
|
||||
if (!isTaskCallbackPending(message)) {
|
||||
return;
|
||||
}
|
||||
executeTaskHandler(message, handleCode, eventType);
|
||||
RBucket<Boolean> handledBucket = redissonClient.getBucket(
|
||||
TASK_BUSINESS_HANDLED_KEY + message.getEventId());
|
||||
if (!Boolean.TRUE.equals(handledBucket.get())) {
|
||||
try {
|
||||
executeTaskHandler(message, handleCode, eventType);
|
||||
handledBucket.set(true, TASK_BUSINESS_HANDLED_TIMEOUT);
|
||||
} catch (RuntimeException exception) {
|
||||
callbackFailedResult(message, exception);
|
||||
throw exception;
|
||||
}
|
||||
}
|
||||
callbackTaskResult(message, CallbackStatusEnum.SUCCESS.getCode(), "任务业务处理成功");
|
||||
handledBucket.delete();
|
||||
} finally {
|
||||
if (lock.isHeldByCurrentThread()) {
|
||||
lock.unlock();
|
||||
@@ -99,12 +122,13 @@ public class WmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMe
|
||||
private void executeTaskHandler(TaskEventMessage message, String handleCode, String eventType) {
|
||||
AbstractTask task = taskFactory.getTask(handleCode);
|
||||
if (task == null) {
|
||||
log.warn("未找到对应任务处理器, handleCode={}", handleCode);
|
||||
return;
|
||||
throw new IllegalStateException("未找到对应任务处理器:" + handleCode);
|
||||
}
|
||||
|
||||
TaskExecuteDTO dto = new TaskExecuteDTO();
|
||||
dto.setTaskId(message.getTaskId());
|
||||
dto.setEventId(message.getEventId());
|
||||
dto.setEventType(message.getEventType());
|
||||
dto.setPayload(message.getPayload());
|
||||
|
||||
if (TaskEventMessage.EVENT_TYPE_FINISHED.equals(eventType)) {
|
||||
@@ -115,4 +139,31 @@ public class WmsTaskStatusChangeConsumer implements RocketMQListener<TaskEventMe
|
||||
log.warn("未知事件类型, eventType={}, handleCode={}", eventType, handleCode);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 回调 Task 服务记录业务处理结果
|
||||
*/
|
||||
private void callbackTaskResult(TaskEventMessage message, String result, String resultMessage) {
|
||||
TaskCallbackResultReqDTO reqDTO = new TaskCallbackResultReqDTO();
|
||||
reqDTO.setTaskId(message.getTaskId());
|
||||
reqDTO.setEventId(message.getEventId());
|
||||
reqDTO.setResult(result);
|
||||
reqDTO.setMessage(resultMessage);
|
||||
CommonResult<Boolean> callbackResult = transportTaskApi.receiveCallbackResult(reqDTO);
|
||||
if (!callbackResult.isSuccess()) {
|
||||
log.error("任务业务结果回调失败, reqDTO={}, callbackResult={}", reqDTO, callbackResult);
|
||||
callbackResult.checkError();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 尝试回调 Task 服务记录业务处理失败
|
||||
*/
|
||||
private void callbackFailedResult(TaskEventMessage message, RuntimeException exception) {
|
||||
try {
|
||||
callbackTaskResult(message, CallbackStatusEnum.FAILED.getCode(), exception.getMessage());
|
||||
} catch (RuntimeException callbackException) {
|
||||
log.error("任务业务失败结果回调异常, message={}", message, callbackException);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user