stream.c 12.3 KB
Newer Older
L
Liu Jicong 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16
/*
 * 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/>.
 */

#include "streamInc.h"
17
#include "ttimer.h"
L
Liu Jicong 已提交
18

L
liuyao 已提交
19
#define STREAM_TASK_INPUT_QUEUEU_CAPACITY 10240
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
int32_t streamInit() {
  int8_t old;
  while (1) {
    old = atomic_val_compare_exchange_8(&streamEnv.inited, 0, 2);
    if (old != 2) break;
  }

  if (old == 0) {
    streamEnv.timer = taosTmrInit(10000, 100, 10000, "STREAM");
    if (streamEnv.timer == NULL) {
      atomic_store_8(&streamEnv.inited, 0);
      return -1;
    }
    atomic_store_8(&streamEnv.inited, 1);
  }
  return 0;
}

void streamCleanUp() {
  int8_t old;
  while (1) {
    old = atomic_val_compare_exchange_8(&streamEnv.inited, 1, 2);
    if (old != 2) break;
  }

  if (old == 1) {
    taosTmrCleanUp(streamEnv.timer);
    atomic_store_8(&streamEnv.inited, 0);
  }
}

L
Liu Jicong 已提交
52
void streamSchedByTimer(void* param, void* tmrId) {
53 54
  SStreamTask* pTask = (void*)param;

55
  if (streamTaskShouldStop(&pTask->status)) {
L
Liu Jicong 已提交
56
    streamMetaReleaseTask(NULL, pTask);
L
Liu Jicong 已提交
57 58 59
    return;
  }

60
  if (atomic_load_8(&pTask->triggerStatus) == TASK_TRIGGER_STATUS__ACTIVE) {
S
Shengliang Guan 已提交
61
    SStreamTrigger* trigger = taosAllocateQitem(sizeof(SStreamTrigger), DEF_QITEM, 0);
62
    if (trigger == NULL) return;
L
Liu Jicong 已提交
63
    trigger->type = STREAM_INPUT__GET_RES;
64 65 66 67 68 69
    trigger->pBlock = taosMemoryCalloc(1, sizeof(SSDataBlock));
    if (trigger->pBlock == NULL) {
      taosFreeQitem(trigger);
      return;
    }

70
    trigger->pBlock->info.type = STREAM_GET_ALL;
71
    atomic_store_8(&pTask->triggerStatus, TASK_TRIGGER_STATUS__INACTIVE);
72

73
    if (tAppendDataToInputQueue(pTask, (SStreamQueueItem*)trigger) < 0) {
L
Liu Jicong 已提交
74 75 76 77
      taosFreeQitem(trigger);
      taosTmrReset(streamSchedByTimer, (int32_t)pTask->triggerParam, pTask, streamEnv.timer, &pTask->timer);
      return;
    }
78

L
Liu Jicong 已提交
79
    streamSchedExec(pTask);
80 81
  }

L
Liu Jicong 已提交
82
  taosTmrReset(streamSchedByTimer, (int32_t)pTask->triggerParam, pTask, streamEnv.timer, &pTask->timer);
83 84 85 86
}

int32_t streamSetupTrigger(SStreamTask* pTask) {
  if (pTask->triggerParam != 0) {
L
Liu Jicong 已提交
87 88
    int32_t ref = atomic_add_fetch_32(&pTask->refCnt, 1);
    ASSERT(ref == 2);
L
Liu Jicong 已提交
89
    pTask->timer = taosTmrStart(streamSchedByTimer, (int32_t)pTask->triggerParam, pTask, streamEnv.timer);
90
    pTask->triggerStatus = TASK_TRIGGER_STATUS__INACTIVE;
91 92 93 94
  }
  return 0;
}

L
Liu Jicong 已提交
95 96
int32_t streamSchedExec(SStreamTask* pTask) {
  int8_t schedStatus =
97
      atomic_val_compare_exchange_8(&pTask->status.schedStatus, TASK_SCHED_STATUS__INACTIVE, TASK_SCHED_STATUS__WAITING);
98

L
Liu Jicong 已提交
99 100 101
  if (schedStatus == TASK_SCHED_STATUS__INACTIVE) {
    SStreamTaskRunReq* pRunReq = rpcMallocCont(sizeof(SStreamTaskRunReq));
    if (pRunReq == NULL) {
102
      terrno = TSDB_CODE_OUT_OF_MEMORY;
103
      atomic_store_8(&pTask->status.schedStatus, TASK_SCHED_STATUS__INACTIVE);
L
Liu Jicong 已提交
104 105
      return -1;
    }
106

L
Liu Jicong 已提交
107
    pRunReq->head.vgId = pTask->nodeId;
108 109
    pRunReq->streamId = pTask->id.streamId;
    pRunReq->taskId = pTask->id.taskId;
110 111

    SRpcMsg msg = { .msgType = TDMT_STREAM_TASK_RUN, .pCont = pRunReq, .contLen = sizeof(SStreamTaskRunReq) };
L
Liu Jicong 已提交
112
    tmsgPutToQueue(pTask->pMsgCb, STREAM_QUEUE, &msg);
113
    qDebug("trigger to run s-task:%s", pTask->id.idStr);
L
Liu Jicong 已提交
114
  }
115

L
Liu Jicong 已提交
116 117 118
  return 0;
}

119
int32_t streamTaskEnqueueBlocks(SStreamTask* pTask, const SStreamDispatchReq* pReq, SRpcMsg* pRsp) {
S
Shengliang Guan 已提交
120
  SStreamDataBlock* pData = taosAllocateQitem(sizeof(SStreamDataBlock), DEF_QITEM, 0);
L
Liu Jicong 已提交
121 122
  int8_t            status;

123
  // enqueue data block
L
Liu Jicong 已提交
124
  if (pData != NULL) {
125
    pData->type = STREAM_INPUT__DATA_BLOCK;
L
Liu Jicong 已提交
126
    pData->srcVgId = pReq->dataSrcVgId;
L
Liu Jicong 已提交
127 128
    // decode
    /*pData->blocks = pReq->data;*/
L
Liu Jicong 已提交
129
    /*pBlock->sourceVer = pReq->sourceVer;*/
L
Liu Jicong 已提交
130
    streamDispatchReqToData(pReq, pData);
131
    if (tAppendDataToInputQueue(pTask, (SStreamQueueItem*)pData) == 0) {
L
Liu Jicong 已提交
132
      status = TASK_INPUT_STATUS__NORMAL;
133 134
    } else {  // input queue is full, upstream is blocked now
      status = TASK_INPUT_STATUS__BLOCKED;
L
Liu Jicong 已提交
135 136 137 138 139 140 141
    }
  } else {
    streamTaskInputFail(pTask);
    status = TASK_INPUT_STATUS__FAILED;
  }

  // rsp by input status
142
  void* buf = rpcMallocCont(sizeof(SMsgHead) + sizeof(SStreamDispatchRsp));
L
Liu Jicong 已提交
143
  ((SMsgHead*)buf)->vgId = htonl(pReq->upstreamNodeId);
144
  SStreamDispatchRsp* pCont = POINTER_SHIFT(buf, sizeof(SMsgHead));
L
Liu Jicong 已提交
145
  pCont->inputStatus = status;
146 147 148 149
  pCont->streamId = htobe64(pReq->streamId);
  pCont->upstreamNodeId = htonl(pReq->upstreamNodeId);
  pCont->upstreamTaskId = htonl(pReq->upstreamTaskId);
  pCont->downstreamNodeId = htonl(pTask->nodeId);
150
  pCont->downstreamTaskId = htonl(pTask->id.taskId);
L
Liu Jicong 已提交
151
  pRsp->pCont = buf;
152

L
Liu Jicong 已提交
153 154
  pRsp->contLen = sizeof(SMsgHead) + sizeof(SStreamDispatchRsp);
  tmsgSendRsp(pRsp);
155

L
Liu Jicong 已提交
156 157 158 159
  return status == TASK_INPUT_STATUS__NORMAL ? 0 : -1;
}

int32_t streamTaskEnqueueRetrieve(SStreamTask* pTask, SStreamRetrieveReq* pReq, SRpcMsg* pRsp) {
S
Shengliang Guan 已提交
160
  SStreamDataBlock* pData = taosAllocateQitem(sizeof(SStreamDataBlock), DEF_QITEM, 0);
L
Liu Jicong 已提交
161 162 163 164
  int8_t            status = TASK_INPUT_STATUS__NORMAL;

  // enqueue
  if (pData != NULL) {
165
    qDebug("task %d(child %d) recv retrieve req from task %d, reqId %" PRId64, pTask->id.taskId, pTask->selfChildId,
L
Liu Jicong 已提交
166 167
           pReq->srcTaskId, pReq->reqId);

168
    pData->type = STREAM_INPUT__DATA_RETRIEVE;
L
Liu Jicong 已提交
169 170 171 172 173
    pData->srcVgId = 0;
    // decode
    /*pData->blocks = pReq->data;*/
    /*pBlock->sourceVer = pReq->sourceVer;*/
    streamRetrieveReqToData(pReq, pData);
174
    if (tAppendDataToInputQueue(pTask, (SStreamQueueItem*)pData) == 0) {
L
Liu Jicong 已提交
175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190
      status = TASK_INPUT_STATUS__NORMAL;
    } else {
      status = TASK_INPUT_STATUS__FAILED;
    }
  } else {
    /*streamTaskInputFail(pTask);*/
    /*status = TASK_INPUT_STATUS__FAILED;*/
  }

  // rsp by input status
  void* buf = rpcMallocCont(sizeof(SMsgHead) + sizeof(SStreamRetrieveRsp));
  ((SMsgHead*)buf)->vgId = htonl(pReq->srcNodeId);
  SStreamRetrieveRsp* pCont = POINTER_SHIFT(buf, sizeof(SMsgHead));
  pCont->streamId = pReq->streamId;
  pCont->rspToTaskId = pReq->srcTaskId;
  pCont->rspFromTaskId = pReq->dstTaskId;
191
  pRsp->pCont = buf;
192
  pRsp->contLen = sizeof(SMsgHead) + sizeof(SStreamRetrieveRsp);
L
Liu Jicong 已提交
193 194 195 196
  tmsgSendRsp(pRsp);
  return status == TASK_INPUT_STATUS__NORMAL ? 0 : -1;
}

L
Liu Jicong 已提交
197
int32_t streamTaskOutput(SStreamTask* pTask, SStreamDataBlock* pBlock) {
dengyihao's avatar
dengyihao 已提交
198
  int32_t code = 0;
L
Liu Jicong 已提交
199 200 201 202 203 204 205 206 207 208
  if (pTask->outputType == TASK_OUTPUT__TABLE) {
    pTask->tbSink.tbSinkFunc(pTask, pTask->tbSink.vnode, 0, pBlock->blocks);
    taosArrayDestroyEx(pBlock->blocks, (FDelete)blockDataFreeRes);
    taosFreeQitem(pBlock);
  } else if (pTask->outputType == TASK_OUTPUT__SMA) {
    pTask->smaSink.smaSink(pTask->smaSink.vnode, pTask->smaSink.smaId, pBlock->blocks);
    taosArrayDestroyEx(pBlock->blocks, (FDelete)blockDataFreeRes);
    taosFreeQitem(pBlock);
  } else {
    ASSERT(pTask->outputType == TASK_OUTPUT__FIXED_DISPATCH || pTask->outputType == TASK_OUTPUT__SHUFFLE_DISPATCH);
dengyihao's avatar
dengyihao 已提交
209 210 211 212
    code = taosWriteQitem(pTask->outputQueue->queue, pBlock);
    if (code != 0) {
      return code;
    }
L
Liu Jicong 已提交
213 214 215 216 217
    streamDispatch(pTask);
  }
  return 0;
}

L
Liu Jicong 已提交
218
int32_t streamProcessDispatchReq(SStreamTask* pTask, SStreamDispatchReq* pReq, SRpcMsg* pRsp, bool exec) {
219
  qDebug("vgId:%d s-task:%s receive dispatch req from taskId:%d", pReq->upstreamNodeId, pTask->id.idStr,
L
Liu Jicong 已提交
220
         pReq->upstreamTaskId);
L
Liu Jicong 已提交
221

222
  streamTaskEnqueueBlocks(pTask, pReq, pRsp);
L
Liu Jicong 已提交
223
  tDeleteStreamDispatchReq(pReq);
L
Liu Jicong 已提交
224

L
Liu Jicong 已提交
225
  if (exec) {
L
Liu Jicong 已提交
226 227 228
    if (streamTryExec(pTask) < 0) {
      return -1;
    }
L
Liu Jicong 已提交
229

L
Liu Jicong 已提交
230 231 232
    /*if (pTask->outputType == TASK_OUTPUT__FIXED_DISPATCH || pTask->outputType == TASK_OUTPUT__SHUFFLE_DISPATCH) {*/
    /*streamDispatch(pTask);*/
    /*}*/
L
Liu Jicong 已提交
233
  } else {
L
Liu Jicong 已提交
234
    streamSchedExec(pTask);
235
  }
L
Liu Jicong 已提交
236 237 238 239

  return 0;
}

240
int32_t streamProcessDispatchRsp(SStreamTask* pTask, SStreamDispatchRsp* pRsp, int32_t code) {
241
  ASSERT(pRsp->inputStatus == TASK_OUTPUT_STATUS__NORMAL || pRsp->inputStatus == TASK_OUTPUT_STATUS__BLOCKED);
242
  qDebug("s-task:%s receive dispatch rsp, code: %x", pTask->id.idStr, code);
L
Liu Jicong 已提交
243

244
  if (pTask->outputType == TASK_OUTPUT__SHUFFLE_DISPATCH) {
L
Liu Jicong 已提交
245
    int32_t leftRsp = atomic_sub_fetch_32(&pTask->shuffleDispatcher.waitingRspCnt, 1);
246
    qDebug("task %d is shuffle, left waiting rsp %d", pTask->id.taskId, leftRsp);
247 248 249
    if (leftRsp > 0) {
      return 0;
    }
L
Liu Jicong 已提交
250 251
  }

252 253
  int8_t old = atomic_exchange_8(&pTask->outputStatus, pRsp->inputStatus);
  ASSERT(old == TASK_OUTPUT_STATUS__WAIT);
L
Liu Jicong 已提交
254 255
  if (pRsp->inputStatus == TASK_INPUT_STATUS__BLOCKED) {
    // TODO: init recover timer
L
Liu Jicong 已提交
256
    ASSERT(0);
257
    return 0;
L
Liu Jicong 已提交
258 259
  }
  // continue dispatch
L
Liu Jicong 已提交
260
  streamDispatch(pTask);
L
Liu Jicong 已提交
261 262 263
  return 0;
}

L
Liu Jicong 已提交
264
int32_t streamProcessRunReq(SStreamTask* pTask) {
L
Liu Jicong 已提交
265 266 267
  if (streamTryExec(pTask) < 0) {
    return -1;
  }
L
Liu Jicong 已提交
268

L
Liu Jicong 已提交
269 270 271
  /*if (pTask->outputType == TASK_OUTPUT__FIXED_DISPATCH || pTask->outputType == TASK_OUTPUT__SHUFFLE_DISPATCH) {*/
  /*streamDispatch(pTask);*/
  /*}*/
L
Liu Jicong 已提交
272 273 274
  return 0;
}

L
Liu Jicong 已提交
275
int32_t streamProcessRetrieveReq(SStreamTask* pTask, SStreamRetrieveReq* pReq, SRpcMsg* pRsp) {
276
  qDebug("task %d receive retrieve req from node %d task %d", pTask->id.taskId, pReq->srcNodeId, pReq->srcTaskId);
L
Liu Jicong 已提交
277 278 279

  streamTaskEnqueueRetrieve(pTask, pReq, pRsp);

280
  ASSERT(pTask->taskLevel != TASK_LEVEL__SINK);
L
Liu Jicong 已提交
281
  streamSchedExec(pTask);
L
Liu Jicong 已提交
282

L
Liu Jicong 已提交
283
  /*streamTryExec(pTask);*/
L
Liu Jicong 已提交
284

L
Liu Jicong 已提交
285
  /*streamDispatch(pTask);*/
L
Liu Jicong 已提交
286 287 288 289

  return 0;
}

290 291 292 293
bool tInputQueueIsFull(const SStreamTask* pTask) {
  return taosQueueItemSize((pTask->inputQueue->queue)) >= STREAM_TASK_INPUT_QUEUEU_CAPACITY;
}

294
int32_t tAppendDataToInputQueue(SStreamTask* pTask, SStreamQueueItem* pItem) {
295 296 297
  int8_t type = pItem->type;

  if (type == STREAM_INPUT__DATA_SUBMIT) {
298 299
    SStreamDataSubmit2* pSubmitBlock = streamSubmitBlockClone((SStreamDataSubmit2*)pItem);
    if (pSubmitBlock == NULL) {
300
      qDebug("task %d %p submit enqueue failed since out of memory", pTask->id.taskId, pTask);
301 302 303 304 305
      terrno = TSDB_CODE_OUT_OF_MEMORY;
      atomic_store_8(&pTask->inputStatus, TASK_INPUT_STATUS__FAILED);
      return -1;
    }

H
Haojun Liao 已提交
306
    int32_t total = taosQueueItemSize(pTask->inputQueue->queue) + 1;
307 308
    qDebug("s-task:%s submit enqueue %p %p msgLen:%d ver:%" PRId64 ", total in queue:%d", pTask->id.idStr,
           pItem, pSubmitBlock->submit.msgStr, pSubmitBlock->submit.msgLen,
309
           pSubmitBlock->submit.ver, total);
H
Haojun Liao 已提交
310

311 312
    if ((pTask->taskLevel == TASK_LEVEL__SOURCE) && total > STREAM_TASK_INPUT_QUEUEU_CAPACITY) {
      qError("s-task:%s input queue is full, capacity:%d, abort", pTask->id.idStr, STREAM_TASK_INPUT_QUEUEU_CAPACITY);
313 314 315 316
      streamDataSubmitDestroy(pSubmitBlock);
      return -1;
    }

H
Haojun Liao 已提交
317
    taosWriteQitem(pTask->inputQueue->queue, pSubmitBlock);
318 319
  } else if (type == STREAM_INPUT__DATA_BLOCK || type == STREAM_INPUT__DATA_RETRIEVE ||
             type == STREAM_INPUT__REF_DATA_BLOCK) {
320
    int32_t total = taosQueueItemSize(pTask->inputQueue->queue) + 1;
321 322
    if ((pTask->taskLevel == TASK_LEVEL__SOURCE) && total > STREAM_TASK_INPUT_QUEUEU_CAPACITY) {
      qError("s-task:%s input queue is full, capacity:%d, abort", pTask->id.idStr, STREAM_TASK_INPUT_QUEUEU_CAPACITY);
323 324 325 326
      return -1;
    }

    qDebug("s-task:%s data block enqueue, total in queue:%d", pTask->id.idStr, total);
327 328 329 330 331 332 333 334 335 336 337 338 339 340
    taosWriteQitem(pTask->inputQueue->queue, pItem);
  } else if (type == STREAM_INPUT__CHECKPOINT) {
    taosWriteQitem(pTask->inputQueue->queue, pItem);
  } else if (type == STREAM_INPUT__GET_RES) {
    taosWriteQitem(pTask->inputQueue->queue, pItem);
  }

  if (type != STREAM_INPUT__GET_RES && type != STREAM_INPUT__CHECKPOINT && pTask->triggerParam != 0) {
    atomic_val_compare_exchange_8(&pTask->triggerStatus, TASK_TRIGGER_STATUS__INACTIVE, TASK_TRIGGER_STATUS__ACTIVE);
  }

#if 0
  atomic_store_8(&pTask->inputStatus, TASK_INPUT_STATUS__NORMAL);
#endif
341

342
  return 0;
343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358
}

void* streamQueueNextItem(SStreamQueue* queue) {
  int8_t dequeueFlag = atomic_exchange_8(&queue->status, STREAM_QUEUE__PROCESSING);
  if (dequeueFlag == STREAM_QUEUE__FAILED) {
    ASSERT(queue->qItem != NULL);
    return streamQueueCurItem(queue);
  } else {
    queue->qItem = NULL;
    taosGetQitem(queue->qall, &queue->qItem);
    if (queue->qItem == NULL) {
      taosReadAllQitems(queue->queue, queue->qall);
      taosGetQitem(queue->qall, &queue->qItem);
    }
    return streamQueueCurItem(queue);
  }
359 360 361 362
}

void streamTaskInputFail(SStreamTask* pTask) {
  atomic_store_8(&pTask->inputStatus, TASK_INPUT_STATUS__FAILED);
363
}