schedulerInt.h 10.4 KB
Newer Older
H
refact  
Hongze Cheng 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22
/*
 * Copyright (c) 2019 TAOS Data, Inc. <jhtao@taosdata.com>
 *
 * This program is free software: you can use, redistribute, and/or modify
 * it under the terms of the GNU Affero General Public License, version 3
 * or later ("AGPL"), as published by the Free Software Foundation.
 *
 * This program is distributed in the hope that it will be useful, but WITHOUT
 * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or
 * FITNESS FOR A PARTICULAR PURPOSE.
 *
 * You should have received a copy of the GNU Affero General Public License
 * along with this program. If not, see <http://www.gnu.org/licenses/>.
 */

#ifndef _TD_SCHEDULER_INT_H_
#define _TD_SCHEDULER_INT_H_

#ifdef __cplusplus
extern "C" {
#endif

23 24 25 26
#include "os.h"
#include "tarray.h"
#include "planner.h"
#include "scheduler.h"
27
#include "thash.h"
D
dapan1121 已提交
28
#include "trpc.h"
D
dapan1121 已提交
29
#include "command.h"
30

D
dapan1121 已提交
31 32
#define SCHEDULE_DEFAULT_MAX_JOB_NUM 1000
#define SCHEDULE_DEFAULT_MAX_TASK_NUM 1000
D
dapan1121 已提交
33
#define SCHEDULE_DEFAULT_MAX_NODE_TABLE_NUM 200  // unit is TSDB_TABLE_NUM_UNIT
34

35
#define SCH_MAX_CANDIDATE_EP_NUM TSDB_MAX_REPLICA
D
dapan 已提交
36

D
dapan1121 已提交
37 38 39 40 41
enum {
  SCH_READ = 1,
  SCH_WRITE,
};

D
dapan1121 已提交
42 43 44 45 46 47
typedef enum {
  SCH_RES_TYPE_QUERY,
  SCH_RES_TYPE_FETCH,
} SCH_RES_TYPE;


D
dapan1121 已提交
48 49 50 51 52
typedef struct SSchTrans {
  void *transInst;
  void *transHandle;
} SSchTrans;

D
dapan1121 已提交
53 54
typedef struct SSchHbTrans {
  SRWLatch  lock;
D
dapan1121 已提交
55
  SRpcCtx   rpcCtx;
D
dapan1121 已提交
56 57 58
  SSchTrans trans;
} SSchHbTrans;

D
dapan1121 已提交
59 60
typedef struct SSchApiStat {

wafwerar's avatar
wafwerar 已提交
61 62 63 64
#ifdef WINDOWS
  size_t avoidCompilationErrors;
#endif

D
dapan1121 已提交
65 66 67 68
} SSchApiStat;

typedef struct SSchRuntimeStat {

wafwerar's avatar
wafwerar 已提交
69 70 71 72
#ifdef WINDOWS
  size_t avoidCompilationErrors;
#endif

D
dapan1121 已提交
73 74 75 76
} SSchRuntimeStat;

typedef struct SSchJobStat {

wafwerar's avatar
wafwerar 已提交
77 78 79 80
#ifdef WINDOWS
  size_t avoidCompilationErrors;
#endif

D
dapan1121 已提交
81 82 83 84 85 86 87 88 89
} SSchJobStat;

typedef struct SSchedulerStat {
  SSchApiStat      api;
  SSchRuntimeStat  runtime;
  SSchJobStat      job;
} SSchedulerStat;


90
typedef struct SSchedulerMgmt {
D
dapan1121 已提交
91 92 93
  uint64_t        taskId; // sequential taksId
  uint64_t        sId;    // schedulerId
  SSchedulerCfg   cfg;
D
dapan1121 已提交
94
  SRWLatch        lock;
D
dapan1121 已提交
95
  bool            exit;
D
dapan1121 已提交
96
  int32_t         jobRef;
D
dapan1121 已提交
97
  int32_t         jobNum;
D
dapan1121 已提交
98 99
  SSchedulerStat  stat;
  SHashObj       *hbConnections;
100
} SSchedulerMgmt;
101

D
dapan1121 已提交
102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118
typedef struct SSchCallbackParamHeader {
  bool isHbParam;
} SSchCallbackParamHeader;

typedef struct SSchTaskCallbackParam {
  SSchCallbackParamHeader head;
  uint64_t                queryId;
  int64_t                 refId;
  uint64_t                taskId;
  void                   *transport;
} SSchTaskCallbackParam;

typedef struct SSchHbCallbackParam {
  SSchCallbackParamHeader head;
  SQueryNodeEpId          nodeEpId;
  void                   *transport;
} SSchHbCallbackParam;
D
dapan1121 已提交
119

D
dapan1121 已提交
120 121
typedef struct SSchFlowControl {
  SRWLatch  lock;
D
dapan1121 已提交
122
  bool      sorted;
D
dapan 已提交
123
  int32_t   tableNumSum;
D
dapan1121 已提交
124
  uint32_t  execTaskNum;
D
dapan1121 已提交
125
  SArray   *taskList;      // Element is SSchTask*
D
dapan1121 已提交
126 127
} SSchFlowControl;

D
dapan1121 已提交
128 129 130 131 132
typedef struct SSchNodeInfo {
  SQueryNodeAddr addr;
  void          *handle;
} SSchNodeInfo;

D
dapan 已提交
133
typedef struct SSchLevel {
D
dapan1121 已提交
134 135 136 137 138 139 140 141 142
  int32_t         level;
  int8_t          status;
  SRWLatch        lock;
  int32_t         taskFailed;
  int32_t         taskSucceed;
  int32_t         taskNum;
  int32_t         taskLaunchedNum;
  SHashObj       *flowCtrl;      // key is ep, element is SSchFlowControl
  SArray         *subTasks;      // Element is SQueryTask
D
dapan 已提交
143
} SSchLevel;
D
dapan1121 已提交
144

D
dapan 已提交
145
typedef struct SSchTask {
D
dapan1121 已提交
146
  uint64_t             taskId;         // task id
D
dapan1121 已提交
147
  SRWLatch             lock;           // task lock
D
dapan1121 已提交
148 149 150 151 152
  SSchLevel           *level;          // level
  SSubplan            *plan;           // subplan
  char                *msg;            // operator tree
  int32_t              msgLen;         // msg length
  int8_t               status;         // task status
D
dapan1121 已提交
153
  int32_t              lastMsgType;    // last sent msg type
D
dapan1121 已提交
154
  int32_t              tryTimes;       // task already tried times
D
dapan1121 已提交
155
  SQueryNodeAddr       succeedAddr;    // task executed success node address
156 157
  int8_t               candidateIdx;   // current try condidation index
  SArray              *candidateAddrs; // condidate node addresses, element is SQueryNodeAddr
D
dapan1121 已提交
158
  SArray              *execNodes;      // all tried node for current task, element is SSchNodeInfo
D
dapan1121 已提交
159 160 161 162
  SQueryProfileSummary summary;        // task execution summary
  int32_t              childReady;     // child task ready number
  SArray              *children;       // the datasource tasks,from which to fetch the result, element is SQueryTask*
  SArray              *parents;        // the data destination tasks, get data from current task, element is SQueryTask*
D
dapan1121 已提交
163
  void*                handle;          // task send handle 
D
dapan 已提交
164
} SSchTask;
D
dapan1121 已提交
165

D
dapan 已提交
166
typedef struct SSchJobAttr {
D
dapan1121 已提交
167 168 169 170
  EExplainMode explainMode;
  bool         syncSchedule;
  bool         queryJob;
  bool         needFlowCtrl;
D
dapan 已提交
171
} SSchJobAttr;
D
dapan1121 已提交
172

D
dapan 已提交
173
typedef struct SSchJob {
D
dapan1121 已提交
174
  int64_t          refId;
175
  uint64_t         queryId;
D
dapan1121 已提交
176
  SSchJobAttr      attr;
177
  int32_t          levelNum;
D
dapan1121 已提交
178
  int32_t          taskNum;
D
dapan1121 已提交
179
  void            *transport;
D
dapan1121 已提交
180 181
  SArray          *nodeList;   // qnode/vnode list, SArray<SQueryNodeAddr>
  SArray          *levels;    // starting from 0. SArray<SSchLevel>
X
Xiaoyu Wang 已提交
182
  SNodeList       *subPlans;  // subplan pointer copied from DAG, no need to free it in scheduler
D
dapan1121 已提交
183

184
  int32_t          levelIdx;
D
dapan1121 已提交
185
  SEpSet           dataSrcEps;
D
dapan1121 已提交
186 187 188 189
  SHashObj        *execTasks; // executing tasks, key:taskid, value:SQueryTask*
  SHashObj        *succTasks; // succeed tasks, key:taskid, value:SQueryTask*
  SHashObj        *failTasks; // failed tasks, key:taskid, value:SQueryTask*

D
dapan1121 已提交
190
  SExplainCtx     *explainCtx;
D
dapan1121 已提交
191
  int8_t           status;  
D
dapan1121 已提交
192
  SQueryNodeAddr   resNode;
D
dapan 已提交
193
  tsem_t           rspSem;
D
dapan1121 已提交
194
  int8_t           userFetch;
D
dapan 已提交
195
  int32_t          remoteFetch;
D
dapan1121 已提交
196
  SSchTask        *fetchTask;
D
dapan 已提交
197
  int32_t          errCode;
D
dapan1121 已提交
198
  SArray          *errList;    // SArray<SQueryErrorInfo>
D
dapan 已提交
199
  SRWLatch         resLock;
D
dapan1121 已提交
200
  SCH_RES_TYPE     resType;
D
dapan1121 已提交
201
  void            *resData;         //TODO free it or not
D
dapan1121 已提交
202
  int32_t          resNumOfRows;
203
  const char      *sql;
204
  SQueryProfileSummary summary;
D
dapan 已提交
205
} SSchJob;
D
dapan1121 已提交
206

D
dapan1121 已提交
207 208
extern SSchedulerMgmt schMgmt;

D
dapan 已提交
209
#define SCH_TASK_READY_TO_LUNCH(readyNum, task) ((readyNum) >= taosArrayGetSize((task)->children))
D
dapan1121 已提交
210

D
dapan1121 已提交
211 212 213 214
#define SCH_TASK_ID(_task) ((_task) ? (_task)->taskId : -1)
#define SCH_SET_TASK_LASTMSG_TYPE(_task, _type) do { if(_task) { atomic_store_32(&(_task)->lastMsgType, _type); } } while (0)
#define SCH_GET_TASK_LASTMSG_TYPE(_task) ((_task) ? atomic_load_32(&(_task)->lastMsgType) : -1)

D
dapan1121 已提交
215 216 217
#define SCH_IS_DATA_SRC_QRY_TASK(task) ((task)->plan->subplanType == SUBPLAN_TYPE_SCAN)
#define SCH_IS_DATA_SRC_TASK(task) (((task)->plan->subplanType == SUBPLAN_TYPE_SCAN) || ((task)->plan->subplanType == SUBPLAN_TYPE_MODIFY))
#define SCH_IS_LEAF_TASK(_job, _task) (((_task)->level->level + 1) == (_job)->levelNum)
218

H
Haojun Liao 已提交
219
#define SCH_SET_TASK_STATUS(task, st) atomic_store_8(&(task)->status, st)
D
dapan1121 已提交
220
#define SCH_GET_TASK_STATUS(task) atomic_load_8(&(task)->status)
D
dapan1121 已提交
221 222
#define SCH_GET_TASK_STATUS_STR(task) jobTaskStatusStr(SCH_GET_TASK_STATUS(task))

D
dapan1121 已提交
223 224
#define SCH_GET_TASK_HANDLE(_task) ((_task) ? (_task)->handle : NULL)
#define SCH_SET_TASK_HANDLE(_task, _handle) ((_task)->handle = (_handle))
D
dapan1121 已提交
225

H
Haojun Liao 已提交
226
#define SCH_SET_JOB_STATUS(job, st) atomic_store_8(&(job)->status, st)
D
dapan1121 已提交
227
#define SCH_GET_JOB_STATUS(job) atomic_load_8(&(job)->status)
D
dapan1121 已提交
228
#define SCH_GET_JOB_STATUS_STR(job) jobTaskStatusStr(SCH_GET_JOB_STATUS(job))
D
dapan1121 已提交
229

D
dapan1121 已提交
230 231
#define SCH_SET_JOB_NEED_FLOW_CTRL(_job) (_job)->attr.needFlowCtrl = true
#define SCH_JOB_NEED_FLOW_CTRL(_job) ((_job)->attr.needFlowCtrl)
D
dapan1121 已提交
232
#define SCH_TASK_NEED_FLOW_CTRL(_job, _task) (SCH_IS_DATA_SRC_QRY_TASK(_task) && SCH_JOB_NEED_FLOW_CTRL(_job) && SCH_IS_LEAF_TASK(_job, _task) && SCH_IS_LEVEL_UNFINISHED((_task)->level))
D
dapan1121 已提交
233 234 235 236

#define SCH_SET_JOB_TYPE(_job, type) (_job)->attr.queryJob = ((type) != SUBPLAN_TYPE_MODIFY)
#define SCH_IS_QUERY_JOB(_job) ((_job)->attr.queryJob) 
#define SCH_JOB_NEED_FETCH(_job) SCH_IS_QUERY_JOB(_job)
D
dapan1121 已提交
237 238
#define SCH_IS_WAIT_ALL_JOB(_job) (!SCH_IS_QUERY_JOB(_job))
#define SCH_IS_NEED_DROP_JOB(_job) (SCH_IS_QUERY_JOB(_job))
D
dapan1121 已提交
239
#define SCH_IS_EXPLAIN_JOB(_job) (EXPLAIN_MODE_ANALYZE == (_job)->attr.explainMode)
D
dapan1121 已提交
240

D
dapan1121 已提交
241
#define SCH_IS_LEVEL_UNFINISHED(_level) ((_level)->taskLaunchedNum < (_level)->taskNum)
X
Xiaoyu Wang 已提交
242 243
#define SCH_GET_CUR_EP(_addr) (&(_addr)->epSet.eps[(_addr)->epSet.inUse])
#define SCH_SWITCH_EPSET(_addr) ((_addr)->epSet.inUse = ((_addr)->epSet.inUse + 1) % (_addr)->epSet.numOfEps)
D
dapan1121 已提交
244
#define SCH_TASK_NUM_OF_EPS(_addr) ((_addr)->epSet.numOfEps)
D
dapan1121 已提交
245

H
Haojun Liao 已提交
246 247
#define SCH_JOB_ELOG(param, ...) qError("QID:0x%" PRIx64 " " param, pJob->queryId, __VA_ARGS__)
#define SCH_JOB_DLOG(param, ...) qDebug("QID:0x%" PRIx64 " " param, pJob->queryId, __VA_ARGS__)
S
Shengliang Guan 已提交
248 249

#define SCH_TASK_ELOG(param, ...) \
D
dapan1121 已提交
250
  qError("QID:0x%" PRIx64 ",TID:0x%" PRIx64 " " param, pJob->queryId, SCH_TASK_ID(pTask), __VA_ARGS__)
S
Shengliang Guan 已提交
251
#define SCH_TASK_DLOG(param, ...) \
D
dapan1121 已提交
252
  qDebug("QID:0x%" PRIx64 ",TID:0x%" PRIx64 " " param, pJob->queryId, SCH_TASK_ID(pTask), __VA_ARGS__)
D
dapan1121 已提交
253
#define SCH_TASK_DLOGL(param, ...) \
D
dapan1121 已提交
254
  qDebugL("QID:0x%" PRIx64 ",TID:0x%" PRIx64 " " param, pJob->queryId, SCH_TASK_ID(pTask), __VA_ARGS__)
S
Shengliang Guan 已提交
255
#define SCH_TASK_WLOG(param, ...) \
D
dapan1121 已提交
256
  qWarn("QID:0x%" PRIx64 ",TID:0x%" PRIx64 " " param, pJob->queryId, SCH_TASK_ID(pTask), __VA_ARGS__)
D
dapan1121 已提交
257 258 259 260

#define SCH_ERR_RET(c) do { int32_t _code = c; if (_code != TSDB_CODE_SUCCESS) { terrno = _code; return _code; } } while (0)
#define SCH_RET(c) do { int32_t _code = c; if (_code != TSDB_CODE_SUCCESS) { terrno = _code; } return _code; } while (0)
#define SCH_ERR_JRET(c) do { code = c; if (code != TSDB_CODE_SUCCESS) { terrno = code; goto _return; } } while (0)
261

D
dapan1121 已提交
262 263 264
#define SCH_LOCK(type, _lock) (SCH_READ == (type) ? taosRLockLatch(_lock) : taosWLockLatch(_lock))
#define SCH_UNLOCK(type, _lock) (SCH_READ == (type) ? taosRUnLockLatch(_lock) : taosWUnLockLatch(_lock))

265

D
dapan1121 已提交
266 267
int32_t schLaunchTask(SSchJob *job, SSchTask *task);
int32_t schBuildAndSendMsg(SSchJob *job, SSchTask *task, SQueryNodeAddr *addr, int32_t msgType);
D
dapan1121 已提交
268 269
SSchJob *schAcquireJob(int64_t refId);
int32_t schReleaseJob(int64_t refId);
D
dapan1121 已提交
270 271 272 273 274 275 276
void schFreeFlowCtrl(SSchLevel *pLevel);
int32_t schCheckJobNeedFlowCtrl(SSchJob *pJob, SSchLevel *pLevel);
int32_t schDecTaskFlowQuota(SSchJob *pJob, SSchTask *pTask);
int32_t schCheckIncTaskFlowQuota(SSchJob *pJob, SSchTask *pTask, bool *enough);
int32_t schLaunchTasksInFlowCtrlList(SSchJob *pJob, SSchTask *pTask);
int32_t schLaunchTaskImpl(SSchJob *pJob, SSchTask *pTask);
int32_t schFetchFromRemote(SSchJob *pJob);
D
dapan1121 已提交
277
int32_t schProcessOnTaskFailure(SSchJob *pJob, SSchTask *pTask, int32_t errCode);
D
dapan1121 已提交
278
int32_t schBuildAndSendHbMsg(SQueryNodeEpId *nodeEpId);
D
dapan1121 已提交
279
int32_t schCloneSMsgSendInfo(void *src, void **dst);
D
dapan1121 已提交
280 281
int32_t schValidateAndBuildJob(SQueryPlan *pDag, SSchJob *pJob);
void schFreeJobImpl(void *job);
D
dapan1121 已提交
282

D
dapan 已提交
283

H
refact  
Hongze Cheng 已提交
284 285 286 287
#ifdef __cplusplus
}
#endif

D
dapan1121 已提交
288
#endif /*_TD_SCHEDULER_INT_H_*/