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

H
Haojun Liao 已提交
19
#define STREAM_TASK_INPUT_QUEUEU_CAPACITY 100000
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 (atomic_load_8(&pTask->status.taskStatus) == TASK_STATUS__DROPPING) {
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
  qDebug("vgId:%d s-task:%s receive dispatch req from taskId:%d", pReq->upstreamNodeId, pTask->id.idStr,
L
Liu Jicong 已提交
216
         pReq->upstreamTaskId);
L
Liu Jicong 已提交
217

218
  streamTaskEnqueueBlocks(pTask, pReq, pRsp);
L
Liu Jicong 已提交
219
  tDeleteStreamDispatchReq(pReq);
L
Liu Jicong 已提交
220

L
Liu Jicong 已提交
221
  if (exec) {
L
Liu Jicong 已提交
222 223 224
    if (streamTryExec(pTask) < 0) {
      return -1;
    }
L
Liu Jicong 已提交
225

L
Liu Jicong 已提交
226 227 228
    /*if (pTask->outputType == TASK_OUTPUT__FIXED_DISPATCH || pTask->outputType == TASK_OUTPUT__SHUFFLE_DISPATCH) {*/
    /*streamDispatch(pTask);*/
    /*}*/
L
Liu Jicong 已提交
229
  } else {
L
Liu Jicong 已提交
230
    streamSchedExec(pTask);
231
  }
L
Liu Jicong 已提交
232 233 234 235

  return 0;
}

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

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

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

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

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

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

  streamTaskEnqueueRetrieve(pTask, pReq, pRsp);

276
  ASSERT(pTask->taskLevel != TASK_LEVEL__SINK);
L
Liu Jicong 已提交
277
  streamSchedExec(pTask);
L
Liu Jicong 已提交
278

L
Liu Jicong 已提交
279
  /*streamTryExec(pTask);*/
L
Liu Jicong 已提交
280

L
Liu Jicong 已提交
281
  /*streamDispatch(pTask);*/
L
Liu Jicong 已提交
282 283 284 285

  return 0;
}

286 287 288 289
bool tInputQueueIsFull(const SStreamTask* pTask) {
  return taosQueueItemSize((pTask->inputQueue->queue)) >= STREAM_TASK_INPUT_QUEUEU_CAPACITY;
}

290
int32_t tAppendDataToInputQueue(SStreamTask* pTask, SStreamQueueItem* pItem) {
291 292 293
  int8_t type = pItem->type;

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

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

307 308
    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);
309 310 311 312
      streamDataSubmitDestroy(pSubmitBlock);
      return -1;
    }

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

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

338
  return 0;
339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354
}

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);
  }
355 356 357 358
}

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