forked from TencentBlueKing/bk-job
-
Notifications
You must be signed in to change notification settings - Fork 0
/
GseStepListener.java
98 lines (92 loc) · 4.91 KB
/
GseStepListener.java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
/*
* Tencent is pleased to support the open source community by making BK-JOB蓝鲸智云作业平台 available.
*
* Copyright (C) 2021 THL A29 Limited, a Tencent company. All rights reserved.
*
* BK-JOB蓝鲸智云作业平台 is licensed under the MIT License.
*
* License for BK-JOB蓝鲸智云作业平台:
* --------------------------------------------------------------------
* Permission is hereby granted, free of charge, to any person obtaining a copy of this software and associated
* documentation files (the "Software"), to deal in the Software without restriction, including without limitation
* the rights to use, copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the Software, and
* to permit persons to whom the Software is furnished to do so, subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all copies or substantial portions of
* the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO
* THE WARRANTIES OF MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
* AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF
* CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS
* IN THE SOFTWARE.
*/
package com.tencent.bk.job.execute.engine.listener;
import com.tencent.bk.job.execute.common.exception.MessageHandleException;
import com.tencent.bk.job.execute.common.exception.MessageHandlerUnavailableException;
import com.tencent.bk.job.execute.engine.GseTaskManager;
import com.tencent.bk.job.execute.engine.consts.GseStepActionEnum;
import com.tencent.bk.job.execute.engine.message.GseTaskProcessor;
import com.tencent.bk.job.execute.engine.model.StepControlMessage;
import com.tencent.bk.job.execute.monitor.metrics.GseTasksExceptionCounter;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.annotation.StreamListener;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.messaging.handler.annotation.Payload;
import org.springframework.stereotype.Component;
/**
* 执行引擎流程处理-GSE任务
*/
@Component
@EnableBinding({GseTaskProcessor.class})
@Slf4j
public class GseStepListener {
private final GseTaskManager gseTaskManager;
/**
* 任务执行异常Counter
*/
private final GseTasksExceptionCounter gseTasksExceptionCounter;
@Autowired
public GseStepListener(GseTaskManager gseTaskManager,
GseTasksExceptionCounter gseTasksExceptionCounter) {
this.gseTaskManager = gseTaskManager;
this.gseTasksExceptionCounter = gseTasksExceptionCounter;
}
@StreamListener(GseTaskProcessor.INPUT)
public void handleMessage(@Payload StepControlMessage gseStepControlMessage,
@Header("X-B3-TraceId") String traceId, @Header("X-B3-SpanId") String spanId) {
log.info("Receive gse step control message, stepInstanceId={}, action={}, requestId={}, msgSendTime={}, " +
"traceId={}, spanId={}", gseStepControlMessage.getStepInstanceId(),
gseStepControlMessage.getAction(), gseStepControlMessage.getRequestId(), gseStepControlMessage.getTime(),
traceId, spanId);
long stepInstanceId = gseStepControlMessage.getStepInstanceId();
String requestId = gseStepControlMessage.getRequestId();
try {
int action = gseStepControlMessage.getAction();
if (GseStepActionEnum.START.getValue() == action) {
gseTaskManager.startStep(stepInstanceId, gseStepControlMessage.getRequestId());
} else if (GseStepActionEnum.STOP.getValue() == action) {
gseTaskManager.stopStep(stepInstanceId, requestId);
} else if (GseStepActionEnum.RETRY_FAIL.getValue() == action) {
gseTaskManager.retryFail(stepInstanceId, requestId);
} else if (GseStepActionEnum.RETRY_ALL.getValue() == action) {
gseTaskManager.retryAll(stepInstanceId, requestId);
} else {
log.error("Error gse step control action:{}", action);
}
} catch (Throwable e) {
String errorMsg = "Handling gse step control message error,stepInstanceId:" + stepInstanceId;
log.error(errorMsg, e);
handleException(e);
}
}
private void handleException(Throwable e) throws MessageHandleException {
// 服务关闭,消息被拒绝,重新入队列
if (e instanceof MessageHandlerUnavailableException) {
throw (MessageHandlerUnavailableException) e;
}
gseTasksExceptionCounter.increment();
}
}