diff --git a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskOperationManager.java b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskOperationManager.java index bba4c42d..9e9ee168 100644 --- a/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskOperationManager.java +++ b/nl-module-task/nl-module-task-server/src/main/java/cn/code/nl/module/task/manage/TransportTaskOperationManager.java @@ -18,6 +18,8 @@ import lombok.extern.slf4j.Slf4j; import org.redisson.api.RLock; import org.redisson.api.RedissonClient; import org.springframework.stereotype.Component; +import org.springframework.transaction.support.TransactionSynchronization; +import org.springframework.transaction.support.TransactionSynchronizationManager; import java.util.EnumMap; import java.util.Map; @@ -116,7 +118,19 @@ public class TransportTaskOperationManager { task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode()); transportTaskMapper.updateById(task); - publishEvent(task, TaskEventTypeEnum.TASK_FINISHED, reqDTO.getPayload()); + // MQ 在事务提交后发送,避免 Consumer 读到未提交的数据 + Map payload = reqDTO.getPayload(); + if (TransactionSynchronizationManager.isSynchronizationActive()) { + TransactionSynchronizationManager.registerSynchronization( + new TransactionSynchronization() { + @Override + public void afterCommit() { + publishEvent(task, TaskEventTypeEnum.TASK_FINISHED, payload); + } + }); + } else { + publishEvent(task, TaskEventTypeEnum.TASK_FINISHED, payload); + } } /** @@ -139,7 +153,19 @@ public class TransportTaskOperationManager { task.setCallbackStatus(CallbackStatusEnum.PENDING.getCode()); transportTaskMapper.updateById(task); - publishEvent(task, TaskEventTypeEnum.TASK_CANCELLED, reqDTO.getPayload()); + // MQ 在事务提交后发送,避免 Consumer 读到未提交的数据 + Map payload = reqDTO.getPayload(); + if (TransactionSynchronizationManager.isSynchronizationActive()) { + TransactionSynchronizationManager.registerSynchronization( + new TransactionSynchronization() { + @Override + public void afterCommit() { + publishEvent(task, TaskEventTypeEnum.TASK_CANCELLED, payload); + } + }); + } else { + publishEvent(task, TaskEventTypeEnum.TASK_CANCELLED, payload); + } } /**