stream.c 12.1 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

19
#define STREAM_TASK_INPUT_QUEUEU_CAPACITY 3000
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 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213
int32_t streamTaskOutput(SStreamTask* pTask, SStreamDataBlock* pBlock) {
  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);
    taosWriteQitem(pTask->outputQueue->queue, pBlock);
    streamDispatch(pTask);
  }
  return 0;
}

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

218
  // todo add the input queue buffer limitation
219
  streamTaskEnqueueBlocks(pTask, pReq, pRsp);
L
Liu Jicong 已提交
220
  tDeleteStreamDispatchReq(pReq);
L
Liu Jicong 已提交
221

L
Liu Jicong 已提交
222
  if (exec) {
L
Liu Jicong 已提交
223 224 225
    if (streamTryExec(pTask) < 0) {
      return -1;
    }
L
Liu Jicong 已提交
226
  } else {
L
Liu Jicong 已提交
227
    streamSchedExec(pTask);
228
  }
L
Liu Jicong 已提交
229 230 231 232

  return 0;
}

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

237
  if (pTask->outputType == TASK_OUTPUT__SHUFFLE_DISPATCH) {
L
Liu Jicong 已提交
238
    int32_t leftRsp = atomic_sub_fetch_32(&pTask->shuffleDispatcher.waitingRspCnt, 1);
239
    qDebug("task %d is shuffle, left waiting rsp %d", pTask->id.taskId, leftRsp);
240 241 242
    if (leftRsp > 0) {
      return 0;
    }
L
Liu Jicong 已提交
243 244
  }

245 246
  int8_t old = atomic_exchange_8(&pTask->outputStatus, pRsp->inputStatus);
  ASSERT(old == TASK_OUTPUT_STATUS__WAIT);
L
Liu Jicong 已提交
247 248
  if (pRsp->inputStatus == TASK_INPUT_STATUS__BLOCKED) {
    // TODO: init recover timer
L
Liu Jicong 已提交
249
    ASSERT(0);
250
    return 0;
L
Liu Jicong 已提交
251 252
  }
  // continue dispatch
L
Liu Jicong 已提交
253
  streamDispatch(pTask);
L
Liu Jicong 已提交
254 255 256
  return 0;
}

L
Liu Jicong 已提交
257
int32_t streamProcessRunReq(SStreamTask* pTask) {
L
Liu Jicong 已提交
258 259 260
  if (streamTryExec(pTask) < 0) {
    return -1;
  }
L
Liu Jicong 已提交
261

L
Liu Jicong 已提交
262 263 264
  /*if (pTask->outputType == TASK_OUTPUT__FIXED_DISPATCH || pTask->outputType == TASK_OUTPUT__SHUFFLE_DISPATCH) {*/
  /*streamDispatch(pTask);*/
  /*}*/
L
Liu Jicong 已提交
265 266 267
  return 0;
}

L
Liu Jicong 已提交
268
int32_t streamProcessRetrieveReq(SStreamTask* pTask, SStreamRetrieveReq* pReq, SRpcMsg* pRsp) {
269
  qDebug("task %d receive retrieve req from node %d task %d", pTask->id.taskId, pReq->srcNodeId, pReq->srcTaskId);
L
Liu Jicong 已提交
270 271 272

  streamTaskEnqueueRetrieve(pTask, pReq, pRsp);

273
  ASSERT(pTask->taskLevel != TASK_LEVEL__SINK);
L
Liu Jicong 已提交
274
  streamSchedExec(pTask);
L
Liu Jicong 已提交
275

L
Liu Jicong 已提交
276
  /*streamTryExec(pTask);*/
L
Liu Jicong 已提交
277

L
Liu Jicong 已提交
278
  /*streamDispatch(pTask);*/
L
Liu Jicong 已提交
279 280 281 282

  return 0;
}

283 284 285 286
bool tInputQueueIsFull(const SStreamTask* pTask) {
  return taosQueueItemSize((pTask->inputQueue->queue)) >= STREAM_TASK_INPUT_QUEUEU_CAPACITY;
}

287
int32_t tAppendDataToInputQueue(SStreamTask* pTask, SStreamQueueItem* pItem) {
288 289 290
  int8_t type = pItem->type;

  if (type == STREAM_INPUT__DATA_SUBMIT) {
291 292
    SStreamDataSubmit2* pSubmitBlock = streamSubmitBlockClone((SStreamDataSubmit2*)pItem);
    if (pSubmitBlock == NULL) {
293
      qDebug("task %d %p submit enqueue failed since out of memory", pTask->id.taskId, pTask);
294 295 296 297 298
      terrno = TSDB_CODE_OUT_OF_MEMORY;
      atomic_store_8(&pTask->inputStatus, TASK_INPUT_STATUS__FAILED);
      return -1;
    }

H
Haojun Liao 已提交
299
    int32_t total = taosQueueItemSize(pTask->inputQueue->queue) + 1;
300 301
    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,
302
           pSubmitBlock->submit.ver, total);
H
Haojun Liao 已提交
303

304 305
    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);
306 307 308 309
      streamDataSubmitDestroy(pSubmitBlock);
      return -1;
    }

H
Haojun Liao 已提交
310
    taosWriteQitem(pTask->inputQueue->queue, pSubmitBlock);
311 312
  } else if (type == STREAM_INPUT__DATA_BLOCK || type == STREAM_INPUT__DATA_RETRIEVE ||
             type == STREAM_INPUT__REF_DATA_BLOCK) {
313
    int32_t total = taosQueueItemSize(pTask->inputQueue->queue) + 1;
314 315
    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);
316 317 318 319
      return -1;
    }

    qDebug("s-task:%s data block enqueue, total in queue:%d", pTask->id.idStr, total);
320 321 322 323 324 325 326 327 328 329 330 331 332 333
    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
334

335
  return 0;
336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351
}

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);
  }
352 353 354 355
}

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