refactor(openapi): 重构回调数据结构和方法命名

master
xuxueli 3 months ago
parent 3ff307180d
commit 8920422ba9

@ -16,7 +16,7 @@ import com.xxl.job.core.context.XxlJobContext;
import com.xxl.job.core.openapi.executor.ExecutorBiz; import com.xxl.job.core.openapi.executor.ExecutorBiz;
import com.xxl.job.core.openapi.executor.dto.KillRequest; import com.xxl.job.core.openapi.executor.dto.KillRequest;
import com.xxl.job.core.openapi.executor.dto.LogRequest; import com.xxl.job.core.openapi.executor.dto.LogRequest;
import com.xxl.job.core.openapi.executor.dto.LogResult; import com.xxl.job.core.openapi.executor.dto.LogData;
import com.xxl.tool.core.CollectionTool; import com.xxl.tool.core.CollectionTool;
import com.xxl.tool.core.DateTool; import com.xxl.tool.core.DateTool;
import com.xxl.tool.core.StringTool; import com.xxl.tool.core.StringTool;
@ -284,9 +284,9 @@ public class JobLogController {
@RequestMapping("/logDetailCat") @RequestMapping("/logDetailCat")
@ResponseBody @ResponseBody
public Response<LogResult> logDetailCat(HttpServletRequest request, public Response<LogData> logDetailCat(HttpServletRequest request,
@RequestParam("logId") long logId, @RequestParam("logId") long logId,
@RequestParam("fromLineNum") int fromLineNum){ @RequestParam("fromLineNum") int fromLineNum){
try { try {
// valid // valid
XxlJobLog jobLog = xxlJobLogMapper.load(logId); XxlJobLog jobLog = xxlJobLogMapper.load(logId);
@ -299,7 +299,7 @@ public class JobLogController {
// log cat // log cat
ExecutorBiz executorBiz = XxlJobAdminBootstrap.getExecutorBiz(jobLog.getExecutorAddress()); ExecutorBiz executorBiz = XxlJobAdminBootstrap.getExecutorBiz(jobLog.getExecutorAddress());
Response<LogResult> logResult = executorBiz.log(new LogRequest(jobLog.getTriggerTime().getTime(), logId, fromLineNum)); Response<LogData> logResult = executorBiz.log(new LogRequest(logId, jobLog.getTriggerTime().getTime(), fromLineNum));
// is end // is end
if (logResult.getData()!=null && logResult.getData().getFromLineNum() > logResult.getData().getToLineNum()) { if (logResult.getData()!=null && logResult.getData().getFromLineNum() > logResult.getData().getToLineNum()) {

@ -14,8 +14,6 @@ import jakarta.servlet.http.HttpServletRequest;
import org.springframework.stereotype.Controller; import org.springframework.stereotype.Controller;
import org.springframework.web.bind.annotation.*; import org.springframework.web.bind.annotation.*;
import java.util.List;
/** /**
* Created by xuxueli on 17/5/10. * Created by xuxueli on 17/5/10.
*/ */
@ -57,8 +55,8 @@ public class OpenApiController {
try { try {
switch (uri) { switch (uri) {
case "callback": { case "callback": {
List<CallbackRequest> callbackParamList = GsonTool.fromJson(requestBody, List.class, CallbackRequest.class); CallbackRequest callbackParam = GsonTool.fromJson(requestBody, CallbackRequest.class);
return adminBiz.callback(callbackParamList); return adminBiz.callback(callbackParam);
} }
case "registry": { case "registry": {
RegistryRequest registryParam = GsonTool.fromJson(requestBody, RegistryRequest.class); RegistryRequest registryParam = GsonTool.fromJson(requestBody, RegistryRequest.class);

@ -4,7 +4,7 @@ import com.xxl.job.admin.business.model.XxlJobLog;
import com.xxl.job.admin.business.scheduler.config.XxlJobAdminBootstrap; import com.xxl.job.admin.business.scheduler.config.XxlJobAdminBootstrap;
import com.xxl.job.admin.framework.util.I18nUtil; import com.xxl.job.admin.framework.util.I18nUtil;
import com.xxl.job.core.context.XxlJobContext; import com.xxl.job.core.context.XxlJobContext;
import com.xxl.job.core.openapi.admin.dto.CallbackRequest; import com.xxl.job.core.openapi.admin.dto.CallbackData;
import com.xxl.tool.concurrent.CyclicThread; import com.xxl.tool.concurrent.CyclicThread;
import com.xxl.tool.core.DateTool; import com.xxl.tool.core.DateTool;
import com.xxl.tool.response.Response; import com.xxl.tool.response.Response;
@ -104,12 +104,12 @@ public class JobCompleteHelper {
* @param callbackParamList callback param * @param callbackParamList callback param
* @return callback result * @return callback result
*/ */
public Response<String> callback(List<CallbackRequest> callbackParamList) { public Response<String> callback(List<CallbackData> callbackParamList) {
callbackThreadPool.execute(new Runnable() { callbackThreadPool.execute(new Runnable() {
@Override @Override
public void run() { public void run() {
for (CallbackRequest callbackRequest: callbackParamList) { for (CallbackData callbackRequest: callbackParamList) {
Response<String> callbackResult = doCallback(callbackRequest); Response<String> callbackResult = doCallback(callbackRequest);
logger.debug(">>>>>>>>> JobApiController.callback {}, callbackRequest={}, callbackResult={}", logger.debug(">>>>>>>>> JobApiController.callback {}, callbackRequest={}, callbackResult={}",
(callbackResult.isSuccess()?"success":"fail"), callbackRequest, callbackResult); (callbackResult.isSuccess()?"success":"fail"), callbackRequest, callbackResult);
@ -120,7 +120,7 @@ public class JobCompleteHelper {
return Response.ofSuccess(); return Response.ofSuccess();
} }
private Response<String> doCallback(CallbackRequest handleCallbackParam) { private Response<String> doCallback(CallbackData handleCallbackParam) {
// valid log item // valid log item
XxlJobLog log = XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().load(handleCallbackParam.getLogId()); XxlJobLog log = XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().load(handleCallbackParam.getLogId());
if (log == null) { if (log == null) {

@ -261,7 +261,7 @@ public class JobTrigger {
ExecutorBiz executorBiz = XxlJobAdminBootstrap.getExecutorBiz(address); ExecutorBiz executorBiz = XxlJobAdminBootstrap.getExecutorBiz(address);
// invoke // invoke
Response<String> runResult = executorBiz.run(triggerParam); Response<String> runResult = executorBiz.trigger(triggerParam);
// build result // build result
StringBuffer runResultSB = new StringBuffer(I18nUtil.getString("jobconf_trigger_run") + ":"); StringBuffer runResultSB = new StringBuffer(I18nUtil.getString("jobconf_trigger_run") + ":");

@ -7,8 +7,6 @@ import com.xxl.job.core.openapi.admin.dto.RegistryRequest;
import com.xxl.tool.response.Response; import com.xxl.tool.response.Response;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import java.util.List;
/** /**
* @author xuxueli 2017-07-27 21:54:20 * @author xuxueli 2017-07-27 21:54:20
*/ */
@ -16,8 +14,8 @@ import java.util.List;
public class AdminBizImpl implements AdminBiz { public class AdminBizImpl implements AdminBiz {
@Override @Override
public Response<String> callback(List<CallbackRequest> callbackRequestList) { public Response<String> callback(CallbackRequest callbackRequest) {
return XxlJobAdminBootstrap.getInstance().getJobCompleteHelper().callback(callbackRequestList); return XxlJobAdminBootstrap.getInstance().getJobCompleteHelper().callback(callbackRequest.getCallbackList());
} }
@Override @Override

@ -2,6 +2,7 @@ package com.xxl.job.openapi;
import com.xxl.job.core.constant.RegistTypeEnum; import com.xxl.job.core.constant.RegistTypeEnum;
import com.xxl.job.core.openapi.admin.AdminBiz; import com.xxl.job.core.openapi.admin.AdminBiz;
import com.xxl.job.core.openapi.admin.dto.CallbackData;
import com.xxl.job.core.openapi.admin.dto.CallbackRequest; import com.xxl.job.core.openapi.admin.dto.CallbackRequest;
import com.xxl.job.core.openapi.admin.dto.RegistryRequest; import com.xxl.job.core.openapi.admin.dto.RegistryRequest;
import com.xxl.job.core.context.XxlJobContext; import com.xxl.job.core.context.XxlJobContext;
@ -42,13 +43,13 @@ public class AdminBizTest {
public void callback() throws Exception { public void callback() throws Exception {
AdminBiz adminBiz = buildClient(); AdminBiz adminBiz = buildClient();
CallbackRequest param = new CallbackRequest(); CallbackData param = new CallbackData();
param.setLogId(1); param.setLogId(1);
param.setHandleCode(XxlJobContext.HANDLE_CODE_SUCCESS); param.setHandleCode(XxlJobContext.HANDLE_CODE_SUCCESS);
List<CallbackRequest> callbackParamList = Arrays.asList(param); CallbackRequest callbackParam = new CallbackRequest(List.of(param));
Response<String> returnT = adminBiz.callback(callbackParamList); Response<String> returnT = adminBiz.callback(callbackParam);
assertTrue(returnT.isSuccess()); assertTrue(returnT.isSuccess());
} }

@ -61,7 +61,7 @@ public class ExecutorBizTest {
} }
@Test @Test
public void run(){ public void trigger(){
ExecutorBiz executorBiz = buildClient(); ExecutorBiz executorBiz = buildClient();
// trigger data // trigger data
@ -77,7 +77,7 @@ public class ExecutorBizTest {
triggerParam.setLogDateTime(System.currentTimeMillis()); triggerParam.setLogDateTime(System.currentTimeMillis());
// Act // Act
final Response<String> retval = executorBiz.run(triggerParam); final Response<String> retval = executorBiz.trigger(triggerParam);
// Assert result // Assert result
Assertions.assertNotNull(retval); Assertions.assertNotNull(retval);
@ -104,12 +104,12 @@ public class ExecutorBizTest {
public void log(){ public void log(){
ExecutorBiz executorBiz = buildClient(); ExecutorBiz executorBiz = buildClient();
final long logDateTim = 0L;
final long logId = 0; final long logId = 0;
final long logDateTim = 0L;
final int fromLineNum = 0; final int fromLineNum = 0;
// Act // Act
final Response<LogResult> retval = executorBiz.log(new LogRequest(logDateTim, logId, fromLineNum)); final Response<LogData> retval = executorBiz.log(new LogRequest(logId, logDateTim, fromLineNum));
// Assert result // Assert result
Assertions.assertNotNull(retval); Assertions.assertNotNull(retval);

@ -1,6 +1,6 @@
package com.xxl.job.core.log; package com.xxl.job.core.log;
import com.xxl.job.core.openapi.executor.dto.LogResult; import com.xxl.job.core.openapi.executor.dto.LogData;
import com.xxl.tool.core.DateTool; import com.xxl.tool.core.DateTool;
import com.xxl.tool.core.StringTool; import com.xxl.tool.core.StringTool;
import com.xxl.tool.io.FileTool; import com.xxl.tool.io.FileTool;
@ -123,14 +123,14 @@ public class XxlJobFileAppender {
* @param fromLineNum from line num * @param fromLineNum from line num
* @return log content * @return log content
*/ */
public static LogResult readLog(String logFileName, final int fromLineNum){ public static LogData readLog(String logFileName, final int fromLineNum){
// valid // valid
if (StringTool.isBlank(logFileName)) { if (StringTool.isBlank(logFileName)) {
return new LogResult(fromLineNum, 0, "readLog fail, logFile not found", true); return new LogData(fromLineNum, 0, "readLog fail, logFile not found", true);
} }
if (!FileTool.exists(logFileName)) { if (!FileTool.exists(logFileName)) {
return new LogResult(fromLineNum, 0, "readLog fail, logFile not exists", true); return new LogData(fromLineNum, 0, "readLog fail, logFile not exists", true);
} }
// read data // read data
@ -168,7 +168,7 @@ public class XxlJobFileAppender {
} }
// result // result
return new LogResult(fromLineNum, toLineNum.get(), logContentBuilder.toString(), false); return new LogData(fromLineNum, toLineNum.get(), logContentBuilder.toString(), false);
} }
} }

@ -4,8 +4,6 @@ import com.xxl.job.core.openapi.admin.dto.CallbackRequest;
import com.xxl.job.core.openapi.admin.dto.RegistryRequest; import com.xxl.job.core.openapi.admin.dto.RegistryRequest;
import com.xxl.tool.response.Response; import com.xxl.tool.response.Response;
import java.util.List;
/** /**
* @author xuxueli 2017-07-27 21:52:49 * @author xuxueli 2017-07-27 21:52:49
*/ */
@ -17,10 +15,10 @@ public interface AdminBiz {
/** /**
* callback * callback
* *
* @param callbackRequestList callback request list * @param callbackRequest callback request
* @return response * @return response
*/ */
public Response<String> callback(List<CallbackRequest> callbackRequestList); public Response<String> callback(CallbackRequest callbackRequest);
// ---------------------- registry ---------------------- // ---------------------- registry ----------------------

@ -0,0 +1,67 @@
package com.xxl.job.core.openapi.admin.dto;
import java.io.Serializable;
/**
* Created by xuxueli on 17/3/2.
*/
public class CallbackData implements Serializable {
private static final long serialVersionUID = 42L;
private long logId;
private long logDateTime;
private int handleCode;
private String handleMsg;
public CallbackData(){}
public CallbackData(long logId, long logDateTime, int handleCode, String handleMsg) {
this.logId = logId;
this.logDateTime = logDateTime;
this.handleCode = handleCode;
this.handleMsg = handleMsg;
}
public long getLogId() {
return logId;
}
public void setLogId(long logId) {
this.logId = logId;
}
public long getLogDateTime() {
return logDateTime;
}
public void setLogDateTime(long logDateTime) {
this.logDateTime = logDateTime;
}
public int getHandleCode() {
return handleCode;
}
public void setHandleCode(int handleCode) {
this.handleCode = handleCode;
}
public String getHandleMsg() {
return handleMsg;
}
public void setHandleMsg(String handleMsg) {
this.handleMsg = handleMsg;
}
@Override
public String toString() {
return "CallbackRequest{" +
"logId=" + logId +
", logDateTime=" + logDateTime +
", handleCode=" + handleCode +
", handleMsg='" + handleMsg + '\'' +
'}';
}
}

@ -1,66 +1,32 @@
package com.xxl.job.core.openapi.admin.dto; package com.xxl.job.core.openapi.admin.dto;
import java.io.Serializable; import java.io.Serializable;
import java.util.List;
/**
* Created by xuxueli on 17/3/2.
*/
public class CallbackRequest implements Serializable { public class CallbackRequest implements Serializable {
private static final long serialVersionUID = 42L; private static final long serialVersionUID = 42L;
private long logId; private List<CallbackData> callbackList;
private long logDateTime;
private int handleCode; public CallbackRequest() {
private String handleMsg;
public CallbackRequest(){}
public CallbackRequest(long logId, long logDateTime, int handleCode, String handleMsg) {
this.logId = logId;
this.logDateTime = logDateTime;
this.handleCode = handleCode;
this.handleMsg = handleMsg;
}
public long getLogId() {
return logId;
}
public void setLogId(long logId) {
this.logId = logId;
}
public long getLogDateTime() {
return logDateTime;
}
public void setLogDateTime(long logDateTime) {
this.logDateTime = logDateTime;
}
public int getHandleCode() {
return handleCode;
} }
public void setHandleCode(int handleCode) { public CallbackRequest(List<CallbackData> callbackList) {
this.handleCode = handleCode; this.callbackList = callbackList;
} }
public String getHandleMsg() { public List<CallbackData> getCallbackList() {
return handleMsg; return callbackList;
} }
public void setHandleMsg(String handleMsg) { public void setCallbackList(List<CallbackData> callbackList) {
this.handleMsg = handleMsg; this.callbackList = callbackList;
} }
@Override @Override
public String toString() { public String toString() {
return "CallbackRequest{" + return "CallbackRequest{" +
"logId=" + logId + "callbackList=" + callbackList +
", logDateTime=" + logDateTime +
", handleCode=" + handleCode +
", handleMsg='" + handleMsg + '\'' +
'}'; '}';
} }

@ -29,7 +29,7 @@ public interface ExecutorBiz {
* @param triggerRequest triggerRequest * @param triggerRequest triggerRequest
* @return response * @return response
*/ */
public Response<String> run(TriggerRequest triggerRequest); public Response<String> trigger(TriggerRequest triggerRequest);
/** /**
* kill * kill
@ -45,6 +45,6 @@ public interface ExecutorBiz {
* @param logRequest logRequest * @param logRequest logRequest
* @return response * @return response
*/ */
public Response<LogResult> log(LogRequest logRequest); public Response<LogData> log(LogRequest logRequest);
} }

@ -8,15 +8,14 @@ import java.io.Serializable;
public class IdleBeatRequest implements Serializable { public class IdleBeatRequest implements Serializable {
private static final long serialVersionUID = 42L; private static final long serialVersionUID = 42L;
private int jobId;
public IdleBeatRequest() { public IdleBeatRequest() {
} }
public IdleBeatRequest(int jobId) { public IdleBeatRequest(int jobId) {
this.jobId = jobId; this.jobId = jobId;
} }
private int jobId;
public int getJobId() { public int getJobId() {
return jobId; return jobId;
} }

@ -8,15 +8,14 @@ import java.io.Serializable;
public class KillRequest implements Serializable { public class KillRequest implements Serializable {
private static final long serialVersionUID = 42L; private static final long serialVersionUID = 42L;
private int jobId;
public KillRequest() { public KillRequest() {
} }
public KillRequest(int jobId) { public KillRequest(int jobId) {
this.jobId = jobId; this.jobId = jobId;
} }
private int jobId;
public int getJobId() { public int getJobId() {
return jobId; return jobId;
} }

@ -5,7 +5,7 @@ import java.io.Serializable;
/** /**
* Created by xuxueli on 17/3/23. * Created by xuxueli on 17/3/23.
*/ */
public class LogResult implements Serializable { public class LogData implements Serializable {
private static final long serialVersionUID = 42L; private static final long serialVersionUID = 42L;
private int fromLineNum; private int fromLineNum;
@ -13,9 +13,9 @@ public class LogResult implements Serializable {
private String logContent; private String logContent;
private boolean isEnd; private boolean isEnd;
public LogResult() { public LogData() {
} }
public LogResult(int fromLineNum, int toLineNum, String logContent, boolean isEnd) { public LogData(int fromLineNum, int toLineNum, String logContent, boolean isEnd) {
this.fromLineNum = fromLineNum; this.fromLineNum = fromLineNum;
this.toLineNum = toLineNum; this.toLineNum = toLineNum;
this.logContent = logContent; this.logContent = logContent;

@ -8,18 +8,18 @@ import java.io.Serializable;
public class LogRequest implements Serializable { public class LogRequest implements Serializable {
private static final long serialVersionUID = 42L; private static final long serialVersionUID = 42L;
private long logId;
private long logDateTime;
private int fromLineNum;
public LogRequest() { public LogRequest() {
} }
public LogRequest(long logDateTime, long logId, int fromLineNum) { public LogRequest(long logId, long logDateTime, int fromLineNum) {
this.logDateTime = logDateTime;
this.logId = logId; this.logId = logId;
this.logDateTime = logDateTime;
this.fromLineNum = fromLineNum; this.fromLineNum = fromLineNum;
} }
private long logDateTime;
private long logId;
private int fromLineNum;
public long getLogDateTime() { public long getLogDateTime() {
return logDateTime; return logDateTime;
} }

@ -46,7 +46,7 @@ public class ExecutorBizImpl implements ExecutorBiz {
} }
@Override @Override
public Response<String> run(TriggerRequest triggerRequest) { public Response<String> trigger(TriggerRequest triggerRequest) {
// load job info:jobHandler + jobThread + glueTypeEnum // load job info:jobHandler + jobThread + glueTypeEnum
JobThread jobThread = XxlJobExecutor.getInstance().loadJobThread(triggerRequest.getJobId()); JobThread jobThread = XxlJobExecutor.getInstance().loadJobThread(triggerRequest.getJobId());
@ -172,11 +172,11 @@ public class ExecutorBizImpl implements ExecutorBiz {
} }
@Override @Override
public Response<LogResult> log(LogRequest logRequest) { public Response<LogData> log(LogRequest logRequest) {
// log filename: logPath/yyyy-MM-dd/9999.log // log filename: logPath/yyyy-MM-dd/9999.log
String logFileName = XxlJobFileAppender.makeLogFileName(new Date(logRequest.getLogDateTime()), logRequest.getLogId()); String logFileName = XxlJobFileAppender.makeLogFileName(new Date(logRequest.getLogDateTime()), logRequest.getLogId());
LogResult logResult = XxlJobFileAppender.readLog(logFileName, logRequest.getFromLineNum()); LogData logResult = XxlJobFileAppender.readLog(logFileName, logRequest.getFromLineNum());
return Response.ofSuccess(logResult); return Response.ofSuccess(logResult);
} }

@ -203,7 +203,7 @@ public class EmbedServer {
return executorBiz.idleBeat(idleBeatParam); return executorBiz.idleBeat(idleBeatParam);
case "/run": case "/run":
TriggerRequest triggerParam = GsonTool.fromJson(requestData, TriggerRequest.class); TriggerRequest triggerParam = GsonTool.fromJson(requestData, TriggerRequest.class);
return executorBiz.run(triggerParam); return executorBiz.trigger(triggerParam);
case "/kill": case "/kill":
KillRequest killParam = GsonTool.fromJson(requestData, KillRequest.class); KillRequest killParam = GsonTool.fromJson(requestData, KillRequest.class);
return executorBiz.kill(killParam); return executorBiz.kill(killParam);

@ -1,6 +1,6 @@
package com.xxl.job.core.thread; package com.xxl.job.core.thread;
import com.xxl.job.core.openapi.admin.dto.CallbackRequest; import com.xxl.job.core.openapi.admin.dto.CallbackData;
import com.xxl.job.core.openapi.executor.dto.TriggerRequest; import com.xxl.job.core.openapi.executor.dto.TriggerRequest;
import com.xxl.job.core.context.XxlJobContext; import com.xxl.job.core.context.XxlJobContext;
import com.xxl.job.core.context.XxlJobHelper; import com.xxl.job.core.context.XxlJobHelper;
@ -202,7 +202,7 @@ public class JobThread extends Thread{
// callback handler info // callback handler info
if (!toStop) { if (!toStop) {
// common // common
XxlJobExecutor.getInstance().getTriggerCallbackThreadHelper().pushCallBack(new CallbackRequest( XxlJobExecutor.getInstance().getTriggerCallbackThreadHelper().pushCallBack(new CallbackData(
triggerParam.getLogId(), triggerParam.getLogId(),
triggerParam.getLogDateTime(), triggerParam.getLogDateTime(),
XxlJobContext.getXxlJobContext().getHandleCode(), XxlJobContext.getXxlJobContext().getHandleCode(),
@ -210,7 +210,7 @@ public class JobThread extends Thread{
); );
} else { } else {
// is killed // is killed
XxlJobExecutor.getInstance().getTriggerCallbackThreadHelper().pushCallBack(new CallbackRequest( XxlJobExecutor.getInstance().getTriggerCallbackThreadHelper().pushCallBack(new CallbackData(
triggerParam.getLogId(), triggerParam.getLogId(),
triggerParam.getLogDateTime(), triggerParam.getLogDateTime(),
XxlJobContext.HANDLE_CODE_FAIL, XxlJobContext.HANDLE_CODE_FAIL,
@ -226,7 +226,7 @@ public class JobThread extends Thread{
TriggerRequest triggerParam = triggerQueue.poll(); TriggerRequest triggerParam = triggerQueue.poll();
if (triggerParam!=null) { if (triggerParam!=null) {
// is killed // is killed
XxlJobExecutor.getInstance().getTriggerCallbackThreadHelper().pushCallBack(new CallbackRequest( XxlJobExecutor.getInstance().getTriggerCallbackThreadHelper().pushCallBack(new CallbackData(
triggerParam.getLogId(), triggerParam.getLogId(),
triggerParam.getLogDateTime(), triggerParam.getLogDateTime(),
XxlJobContext.HANDLE_CODE_FAIL, XxlJobContext.HANDLE_CODE_FAIL,

@ -6,6 +6,7 @@ import com.xxl.job.core.context.XxlJobHelper;
import com.xxl.job.core.executor.XxlJobExecutor; import com.xxl.job.core.executor.XxlJobExecutor;
import com.xxl.job.core.log.XxlJobFileAppender; import com.xxl.job.core.log.XxlJobFileAppender;
import com.xxl.job.core.openapi.admin.AdminBiz; import com.xxl.job.core.openapi.admin.AdminBiz;
import com.xxl.job.core.openapi.admin.dto.CallbackData;
import com.xxl.job.core.openapi.admin.dto.CallbackRequest; import com.xxl.job.core.openapi.admin.dto.CallbackRequest;
import com.xxl.tool.concurrent.CyclicThread; import com.xxl.tool.concurrent.CyclicThread;
import com.xxl.tool.concurrent.MessageQueue; import com.xxl.tool.concurrent.MessageQueue;
@ -38,7 +39,7 @@ public class TriggerCallbackThreadHelper {
/** /**
* callback message-queue * callback message-queue
*/ */
private volatile MessageQueue<CallbackRequest> callbackMessageQueue; private volatile MessageQueue<CallbackData> callbackMessageQueue;
/** /**
* retry callback-file thread * retry callback-file thread
@ -61,7 +62,7 @@ public class TriggerCallbackThreadHelper {
/** /**
* 1、callback message-queue * 1、callback message-queue
*/ */
callbackMessageQueue = new MessageQueue<CallbackRequest>( callbackMessageQueue = new MessageQueue<CallbackData>(
"TriggerCallbackThreadHelper#callbackMessageQueue", "TriggerCallbackThreadHelper#callbackMessageQueue",
messages -> { messages -> {
@ -104,7 +105,7 @@ public class TriggerCallbackThreadHelper {
} }
// parse callback param // parse callback param
List<CallbackRequest> callbackParamList = GsonTool.fromJsonList(callbackData, CallbackRequest.class); List<CallbackData> callbackParamList = GsonTool.fromJsonList(callbackData, CallbackData.class);
FileTool.delete(callbackLogFile); FileTool.delete(callbackLogFile);
// retry callback // retry callback
@ -138,7 +139,7 @@ public class TriggerCallbackThreadHelper {
/** /**
* submit callback message * submit callback message
*/ */
public void pushCallBack(CallbackRequest callback){ public void pushCallBack(CallbackData callback){
if (!callbackMessageQueue.produce(callback)) { if (!callbackMessageQueue.produce(callback)) {
doCallback(new ArrayList<>(Collections.singletonList(callback)), XxlJobExecutor.getInstance()); doCallback(new ArrayList<>(Collections.singletonList(callback)), XxlJobExecutor.getInstance());
} }
@ -151,38 +152,38 @@ public class TriggerCallbackThreadHelper {
/** /**
* do callback, will retry if error * do callback, will retry if error
* *
* @param callbackParamList callback param list * @param callbackDataList callback data list
*/ */
private void doCallback(List<CallbackRequest> callbackParamList, final XxlJobExecutor xxlJobExecutor){ private void doCallback(List<CallbackData> callbackDataList, final XxlJobExecutor xxlJobExecutor){
boolean callbackRet = false; boolean callbackRet = false;
// callback request, will retry + append-log if fail // callback request, will retry + append-log if fail
for (AdminBiz adminBiz: xxlJobExecutor.getAdminBizList()) { for (AdminBiz adminBiz: xxlJobExecutor.getAdminBizList()) {
try { try {
Response<String> callbackResult = adminBiz.callback(callbackParamList); Response<String> callbackResult = adminBiz.callback(new CallbackRequest(callbackDataList));
if (callbackResult!=null && callbackResult.isSuccess()) { if (callbackResult!=null && callbackResult.isSuccess()) {
appendCallbackResult(callbackParamList, "<br>----------- xxl-job job callback finish."); appendCallbackResult(callbackDataList, "<br>----------- xxl-job job callback finish.");
callbackRet = true; callbackRet = true;
break; break;
} else { } else {
appendCallbackResult(callbackParamList, "<br>----------- xxl-job job callback fail, callbackResult:" + callbackResult); appendCallbackResult(callbackDataList, "<br>----------- xxl-job job callback fail, callbackResult:" + callbackResult);
} }
} catch (Throwable e) { } catch (Throwable e) {
appendCallbackResult(callbackParamList, "<br>----------- xxl-job job callback error, errorMsg:" + e.getMessage()); appendCallbackResult(callbackDataList, "<br>----------- xxl-job job callback error, errorMsg:" + e.getMessage());
} }
} }
// write callback-file, will retry later // write callback-file, will retry later
if (!callbackRet) { if (!callbackRet) {
writeCallbackLog(callbackParamList); writeCallbackLog(callbackDataList);
} }
} }
/** /**
* append callback result, to each joblog * append callback result, to each joblog
*/ */
private void appendCallbackResult(List<CallbackRequest> callbackParamList, String logContent){ private void appendCallbackResult(List<CallbackData> callbackParamList, String logContent){
for (CallbackRequest callbackParam: callbackParamList) { for (CallbackData callbackParam: callbackParamList) {
String logFileName = XxlJobFileAppender.makeLogFileName(new Date(callbackParam.getLogDateTime()), callbackParam.getLogId()); String logFileName = XxlJobFileAppender.makeLogFileName(new Date(callbackParam.getLogDateTime()), callbackParam.getLogId());
XxlJobContext.setXxlJobContext(new XxlJobContext( XxlJobContext.setXxlJobContext(new XxlJobContext(
-1, -1,
@ -213,7 +214,7 @@ public class TriggerCallbackThreadHelper {
* *
* @param callbackParamList callback param list * @param callbackParamList callback param list
*/ */
private void writeCallbackLog(List<CallbackRequest> callbackParamList) { private void writeCallbackLog(List<CallbackData> callbackParamList) {
// valid // valid
if (CollectionTool.isEmpty(callbackParamList)) { if (CollectionTool.isEmpty(callbackParamList)) {
return; return;

Loading…
Cancel
Save