executorimpl.c 108.8 KB
Newer Older
1 2 3 4 5 6 7 8 9 10 11 12 13 14
/*
 * 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/>.
 */
15

H
Haojun Liao 已提交
16 17
#include "filter.h"
#include "function.h"
18 19
#include "functionMgt.h"
#include "os.h"
H
Haojun Liao 已提交
20
#include "querynodes.h"
21
#include "tfill.h"
dengyihao's avatar
dengyihao 已提交
22
#include "tname.h"
23

H
Haojun Liao 已提交
24
#include "tdatablock.h"
25
#include "tglobal.h"
H
Haojun Liao 已提交
26
#include "tmsg.h"
27
#include "ttime.h"
H
Haojun Liao 已提交
28

29
#include "executorimpl.h"
dengyihao's avatar
dengyihao 已提交
30
#include "index.h"
31
#include "query.h"
32
#include "tcompare.h"
H
Haojun Liao 已提交
33
#include "thash.h"
34
#include "ttypes.h"
dengyihao's avatar
dengyihao 已提交
35
#include "vnode.h"
36

H
Haojun Liao 已提交
37
#define IS_MAIN_SCAN(runtime)          ((runtime)->scanFlag == MAIN_SCAN)
38 39 40 41 42 43
#define SET_REVERSE_SCAN_FLAG(runtime) ((runtime)->scanFlag = REVERSE_SCAN)

#define GET_FORWARD_DIRECTION_FACTOR(ord) (((ord) == TSDB_ORDER_ASC) ? QUERY_ASC_FORWARD_STEP : QUERY_DESC_FORWARD_STEP)

#if 0
static UNUSED_FUNC void *u_malloc (size_t __size) {
wafwerar's avatar
wafwerar 已提交
44
  uint32_t v = taosRand();
45 46 47 48

  if (v % 1000 <= 0) {
    return NULL;
  } else {
wafwerar's avatar
wafwerar 已提交
49
    return taosMemoryMalloc(__size);
50 51 52 53
  }
}

static UNUSED_FUNC void* u_calloc(size_t num, size_t __size) {
wafwerar's avatar
wafwerar 已提交
54
  uint32_t v = taosRand();
55 56 57
  if (v % 1000 <= 0) {
    return NULL;
  } else {
wafwerar's avatar
wafwerar 已提交
58
    return taosMemoryCalloc(num, __size);
59 60 61 62
  }
}

static UNUSED_FUNC void* u_realloc(void* p, size_t __size) {
wafwerar's avatar
wafwerar 已提交
63
  uint32_t v = taosRand();
64 65 66
  if (v % 5 <= 1) {
    return NULL;
  } else {
wafwerar's avatar
wafwerar 已提交
67
    return taosMemoryRealloc(p, __size);
68 69 70 71 72 73 74 75
  }
}

#define calloc  u_calloc
#define malloc  u_malloc
#define realloc u_realloc
#endif

X
Xiaoyu Wang 已提交
76
#define CLEAR_QUERY_STATUS(q, st)   ((q)->status &= (~(st)))
77 78
#define QUERY_IS_INTERVAL_QUERY(_q) ((_q)->interval.interval > 0)

H
Haojun Liao 已提交
79 80 81 82 83 84 85 86 87
typedef struct SAggOperatorInfo {
  SOptrBasicInfo   binfo;
  SAggSupporter    aggSup;
  STableQueryInfo* current;
  uint64_t         groupId;
  SGroupResInfo    groupResInfo;
  SExprSupp        scalarExprSup;
} SAggOperatorInfo;

L
Liu Jicong 已提交
88 89
int32_t getMaximumIdleDurationSec() { return tsShellActivityTimer * 2; }

90
static void setBlockSMAInfo(SqlFunctionCtx* pCtx, SExprInfo* pExpr, SSDataBlock* pBlock);
91

X
Xiaoyu Wang 已提交
92
static void releaseQueryBuf(size_t numOfTables);
93

H
Haojun Liao 已提交
94 95 96 97 98 99 100 101 102 103 104
static void    destroyFillOperatorInfo(void* param);
static void    destroyAggOperatorInfo(void* param);
static void    initCtxOutputBuffer(SqlFunctionCtx* pCtx, int32_t size);
static void    doSetTableGroupOutputBuf(SOperatorInfo* pOperator, int32_t numOfOutput, uint64_t groupId);
static void    doApplyScalarCalculation(SOperatorInfo* pOperator, SSDataBlock* pBlock, int32_t order, int32_t scanFlag);
static int32_t doInitAggInfoSup(SAggSupporter* pAggSup, SqlFunctionCtx* pCtx, int32_t numOfOutput, size_t keyBufSize,
                                const char* pKey);
static void    extractQualifiedTupleByFilterResult(SSDataBlock* pBlock, const SColumnInfoData* p, bool keep,
                                                   int32_t status);
static int32_t doSetInputDataBlock(SExprSupp* pExprSup, SSDataBlock* pBlock, int32_t order, int32_t scanFlag,
                                   bool createDummyCol);
H
Haojun Liao 已提交
105

H
Haojun Liao 已提交
106
void setOperatorCompleted(SOperatorInfo* pOperator) {
107
  pOperator->status = OP_EXEC_DONE;
H
Haojun Liao 已提交
108
  ASSERT(pOperator->pTaskInfo != NULL);
109

110
  pOperator->cost.totalCost = (taosGetTimestampUs() - pOperator->pTaskInfo->cost.start * 1000) / 1000.0;
H
Haojun Liao 已提交
111
  setTaskStatus(pOperator->pTaskInfo, TASK_COMPLETED);
112
}
113

H
Haojun Liao 已提交
114 115 116 117 118 119 120 121 122 123
void setOperatorInfo(SOperatorInfo* pOperator, const char* name, int32_t type, bool blocking, int32_t status,
                     void* pInfo, SExecTaskInfo* pTaskInfo) {
  pOperator->name = (char*)name;
  pOperator->operatorType = type;
  pOperator->blocking = blocking;
  pOperator->status = status;
  pOperator->info = pInfo;
  pOperator->pTaskInfo = pTaskInfo;
}

H
Haojun Liao 已提交
124
int32_t operatorDummyOpenFn(SOperatorInfo* pOperator) {
125
  OPTR_SET_OPENED(pOperator);
126
  pOperator->cost.openCost = 0;
H
Haojun Liao 已提交
127
  return TSDB_CODE_SUCCESS;
H
Haojun Liao 已提交
128 129
}

H
Haojun Liao 已提交
130 131
SOperatorFpSet createOperatorFpSet(__optr_open_fn_t openFn, __optr_fn_t nextFn, __optr_fn_t cleanup,
                                   __optr_close_fn_t closeFn, __optr_explain_fn_t explain) {
132 133 134 135 136 137 138 139 140 141 142
  SOperatorFpSet fpSet = {
      ._openFn = openFn,
      .getNextFn = nextFn,
      .cleanupFn = cleanup,
      .closeFn = closeFn,
      .getExplainFn = explain,
  };

  return fpSet;
}

143 144
static int32_t doCopyToSDataBlock(SExecTaskInfo* pTaskInfo, SSDataBlock* pBlock, SExprSupp* pSup, SDiskbasedBuf* pBuf,
                                  SGroupResInfo* pGroupResInfo);
H
Haojun Liao 已提交
145

146
SResultRow* getNewResultRow(SDiskbasedBuf* pResultBuf, int32_t* currentPageId, int32_t interBufSize) {
L
Liu Jicong 已提交
147
  SFilePage* pData = NULL;
148 149 150

  // in the first scan, new space needed for results
  int32_t pageId = -1;
151
  if (*currentPageId == -1) {
152
    pData = getNewBufPage(pResultBuf, &pageId);
153 154
    pData->num = sizeof(SFilePage);
  } else {
155 156
    pData = getBufPage(pResultBuf, *currentPageId);
    pageId = *currentPageId;
157

wmmhello's avatar
wmmhello 已提交
158
    if (pData->num + interBufSize > getBufPageSize(pResultBuf)) {
159
      // release current page first, and prepare the next one
160
      releaseBufPage(pResultBuf, pData);
161

162
      pData = getNewBufPage(pResultBuf, &pageId);
163 164 165 166 167 168 169 170 171 172
      if (pData != NULL) {
        pData->num = sizeof(SFilePage);
      }
    }
  }

  if (pData == NULL) {
    return NULL;
  }

173 174
  setBufPageDirty(pData, true);

175 176 177 178
  // set the number of rows in current disk page
  SResultRow* pResultRow = (SResultRow*)((char*)pData + pData->num);
  pResultRow->pageId = pageId;
  pResultRow->offset = (int32_t)pData->num;
179
  *currentPageId = pageId;
180

wmmhello's avatar
wmmhello 已提交
181
  pData->num += interBufSize;
182 183 184
  return pResultRow;
}

185 186 187 188 189 190 191
/**
 * the struct of key in hash table
 * +----------+---------------+
 * | group id |   key data    |
 * | 8 bytes  | actual length |
 * +----------+---------------+
 */
192 193 194
SResultRow* doSetResultOutBufByKey(SDiskbasedBuf* pResultBuf, SResultRowInfo* pResultRowInfo, char* pData,
                                   int16_t bytes, bool masterscan, uint64_t groupId, SExecTaskInfo* pTaskInfo,
                                   bool isIntervalQuery, SAggSupporter* pSup) {
195
  SET_RES_WINDOW_KEY(pSup->keyBuf, pData, bytes, groupId);
H
Haojun Liao 已提交
196

dengyihao's avatar
dengyihao 已提交
197
  SResultRowPosition* p1 =
198
      (SResultRowPosition*)tSimpleHashGet(pSup->pResultRowHashTable, pSup->keyBuf, GET_RES_WINDOW_KEY_LEN(bytes));
H
Haojun Liao 已提交
199

200 201
  SResultRow* pResult = NULL;

H
Haojun Liao 已提交
202 203
  // in case of repeat scan/reverse scan, no new time window added.
  if (isIntervalQuery) {
204
    if (masterscan && p1 != NULL) {  // the *p1 may be NULL in case of sliding+offset exists.
205
      pResult = getResultRowByPos(pResultBuf, p1, true);
206
      ASSERT(pResult->pageId == p1->pageId && pResult->offset == p1->offset);
H
Haojun Liao 已提交
207 208
    }
  } else {
dengyihao's avatar
dengyihao 已提交
209 210
    // In case of group by column query, the required SResultRow object must be existInCurrentResusltRowInfo in the
    // pResultRowInfo object.
H
Haojun Liao 已提交
211
    if (p1 != NULL) {
212
      // todo
213
      pResult = getResultRowByPos(pResultBuf, p1, true);
214
      ASSERT(pResult->pageId == p1->pageId && pResult->offset == p1->offset);
H
Haojun Liao 已提交
215 216 217
    }
  }

L
Liu Jicong 已提交
218
  // 1. close current opened time window
219
  if (pResultRowInfo->cur.pageId != -1 && ((pResult == NULL) || (pResult->pageId != pResultRowInfo->cur.pageId))) {
220
    SResultRowPosition pos = pResultRowInfo->cur;
X
Xiaoyu Wang 已提交
221
    SFilePage*         pPage = getBufPage(pResultBuf, pos.pageId);
222 223 224 225 226
    releaseBufPage(pResultBuf, pPage);
  }

  // allocate a new buffer page
  if (pResult == NULL) {
H
Haojun Liao 已提交
227
    ASSERT(pSup->resultRowSize > 0);
228
    pResult = getNewResultRow(pResultBuf, &pSup->currentPageId, pSup->resultRowSize);
229

230 231
    // add a new result set for a new group
    SResultRowPosition pos = {.pageId = pResult->pageId, .offset = pResult->offset};
232
    tSimpleHashPut(pSup->pResultRowHashTable, pSup->keyBuf, GET_RES_WINDOW_KEY_LEN(bytes), &pos,
L
Liu Jicong 已提交
233
                   sizeof(SResultRowPosition));
H
Haojun Liao 已提交
234 235
  }

236 237 238
  // 2. set the new time window to be the new active time window
  pResultRowInfo->cur = (SResultRowPosition){.pageId = pResult->pageId, .offset = pResult->offset};

H
Haojun Liao 已提交
239
  // too many time window in query
240
  if (pTaskInfo->execModel == OPTR_EXEC_MODEL_BATCH &&
241
      tSimpleHashGetSize(pSup->pResultRowHashTable) > MAX_INTERVAL_TIME_WINDOW) {
242
    T_LONG_JMP(pTaskInfo->env, TSDB_CODE_QRY_TOO_MANY_TIMEWINDOW);
H
Haojun Liao 已提交
243 244
  }

H
Haojun Liao 已提交
245
  return pResult;
H
Haojun Liao 已提交
246 247
}

248
// a new buffer page for each table. Needs to opt this design
L
Liu Jicong 已提交
249
static int32_t addNewWindowResultBuf(SResultRow* pWindowRes, SDiskbasedBuf* pResultBuf, int32_t tid, uint32_t size) {
250 251 252 253
  if (pWindowRes->pageId != -1) {
    return 0;
  }

L
Liu Jicong 已提交
254
  SFilePage* pData = NULL;
255 256 257

  // in the first scan, new space needed for results
  int32_t pageId = -1;
258
  SArray* list = getDataBufPagesIdList(pResultBuf);
259 260

  if (taosArrayGetSize(list) == 0) {
261
    pData = getNewBufPage(pResultBuf, &pageId);
262
    pData->num = sizeof(SFilePage);
263 264
  } else {
    SPageInfo* pi = getLastPageInfo(list);
265
    pData = getBufPage(pResultBuf, getPageId(pi));
266
    pageId = getPageId(pi);
267

268
    if (pData->num + size > getBufPageSize(pResultBuf)) {
269
      // release current page first, and prepare the next one
270
      releaseBufPageInfo(pResultBuf, pi);
271

272
      pData = getNewBufPage(pResultBuf, &pageId);
273
      if (pData != NULL) {
274
        pData->num = sizeof(SFilePage);
275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294
      }
    }
  }

  if (pData == NULL) {
    return -1;
  }

  // set the number of rows in current disk page
  if (pWindowRes->pageId == -1) {  // not allocated yet, allocate new buffer
    pWindowRes->pageId = pageId;
    pWindowRes->offset = (int32_t)pData->num;

    pData->num += size;
    assert(pWindowRes->pageId >= 0);
  }

  return 0;
}

295
//  query_range_start, query_range_end, window_duration, window_start, window_end
296
void initExecTimeWindowInfo(SColumnInfoData* pColData, STimeWindow* pQueryWindow) {
297 298 299
  pColData->info.type = TSDB_DATA_TYPE_TIMESTAMP;
  pColData->info.bytes = sizeof(int64_t);

H
Haojun Liao 已提交
300
  colInfoDataEnsureCapacity(pColData, 5, false);
301 302 303 304 305 306 307 308 309
  colDataAppendInt64(pColData, 0, &pQueryWindow->skey);
  colDataAppendInt64(pColData, 1, &pQueryWindow->ekey);

  int64_t interval = 0;
  colDataAppendInt64(pColData, 2, &interval);  // this value may be variable in case of 'n' and 'y'.
  colDataAppendInt64(pColData, 3, &pQueryWindow->skey);
  colDataAppendInt64(pColData, 4, &pQueryWindow->ekey);
}

310 311 312 313 314 315 316
typedef struct {
  bool    hasAgg;
  int32_t numOfRows;
  int32_t startOffset;
} SFunctionCtxStatus;

static void functionCtxSave(SqlFunctionCtx* pCtx, SFunctionCtxStatus* pStatus) {
H
Haojun Liao 已提交
317
  pStatus->hasAgg = pCtx->input.colDataSMAIsSet;
318 319 320 321 322
  pStatus->numOfRows = pCtx->input.numOfRows;
  pStatus->startOffset = pCtx->input.startRowIndex;
}

static void functionCtxRestore(SqlFunctionCtx* pCtx, SFunctionCtxStatus* pStatus) {
H
Haojun Liao 已提交
323
  pCtx->input.colDataSMAIsSet = pStatus->hasAgg;
H
Haojun Liao 已提交
324
  pCtx->input.numOfRows = pStatus->numOfRows;
325 326 327
  pCtx->input.startRowIndex = pStatus->startOffset;
}

H
Haojun Liao 已提交
328 329
void applyAggFunctionOnPartialTuples(SExecTaskInfo* taskInfo, SqlFunctionCtx* pCtx, SColumnInfoData* pTimeWindowData,
                                     int32_t offset, int32_t forwardStep, int32_t numOfTotal, int32_t numOfOutput) {
330
  for (int32_t k = 0; k < numOfOutput; ++k) {
H
Haojun Liao 已提交
331
    // keep it temporarily
332 333
    SFunctionCtxStatus status = {0};
    functionCtxSave(&pCtx[k], &status);
334

335
    pCtx[k].input.startRowIndex = offset;
336
    pCtx[k].input.numOfRows = forwardStep;
337 338 339

    // not a whole block involved in query processing, statistics data can not be used
    // NOTE: the original value of isSet have been changed here
H
Haojun Liao 已提交
340 341
    if (pCtx[k].input.colDataSMAIsSet && forwardStep < numOfTotal) {
      pCtx[k].input.colDataSMAIsSet = false;
342 343
    }

344 345
    if (fmIsWindowPseudoColumnFunc(pCtx[k].functionId)) {
      SResultRowEntryInfo* pEntryInfo = GET_RES_INFO(&pCtx[k]);
346 347

      char* p = GET_ROWCELL_INTERBUF(pEntryInfo);
348

349
      SColumnInfoData idata = {0};
dengyihao's avatar
dengyihao 已提交
350
      idata.info.type = TSDB_DATA_TYPE_BIGINT;
351
      idata.info.bytes = tDataTypes[TSDB_DATA_TYPE_BIGINT].bytes;
dengyihao's avatar
dengyihao 已提交
352
      idata.pData = p;
353 354 355 356

      SScalarParam out = {.columnData = &idata};
      SScalarParam tw = {.numOfRows = 5, .columnData = pTimeWindowData};
      pCtx[k].sfp.process(&tw, 1, &out);
357
      pEntryInfo->numOfRes = 1;
358 359 360 361 362 363 364 365
    } else {
      int32_t code = TSDB_CODE_SUCCESS;
      if (functionNeedToExecute(&pCtx[k]) && pCtx[k].fpSet.process != NULL) {
        code = pCtx[k].fpSet.process(&pCtx[k]);

        if (code != TSDB_CODE_SUCCESS) {
          qError("%s apply functions error, code: %s", GET_TASKID(taskInfo), tstrerror(code));
          taskInfo->code = code;
366
          T_LONG_JMP(taskInfo->env, code);
367
        }
368
      }
369

370
      // restore it
371
      functionCtxRestore(&pCtx[k], &status);
372
    }
373 374 375
  }
}

376 377 378
static void doSetInputDataBlockInfo(SExprSupp* pExprSup, SSDataBlock* pBlock, int32_t order) {
  SqlFunctionCtx* pCtx = pExprSup->pCtx;
  for (int32_t i = 0; i < pExprSup->numOfExprs; ++i) {
379
    pCtx[i].order = order;
380
    pCtx[i].input.numOfRows = pBlock->info.rows;
381
    setBlockSMAInfo(&pCtx[i], &pExprSup->pExprInfo[i], pBlock);
382
    pCtx[i].pSrcBlock = pBlock;
383 384 385
  }
}

386
void setInputDataBlock(SExprSupp* pExprSup, SSDataBlock* pBlock, int32_t order, int32_t scanFlag, bool createDummyCol) {
387
  if (pBlock->pBlockAgg != NULL) {
388
    doSetInputDataBlockInfo(pExprSup, pBlock, order);
389
  } else {
390
    doSetInputDataBlock(pExprSup, pBlock, order, scanFlag, createDummyCol);
H
Haojun Liao 已提交
391
  }
392 393
}

L
Liu Jicong 已提交
394 395
static int32_t doCreateConstantValColumnInfo(SInputColumnInfoData* pInput, SFunctParam* pFuncParam, int32_t paramIndex,
                                             int32_t numOfRows) {
396 397 398 399 400 401 402 403
  SColumnInfoData* pColInfo = NULL;
  if (pInput->pData[paramIndex] == NULL) {
    pColInfo = taosMemoryCalloc(1, sizeof(SColumnInfoData));
    if (pColInfo == NULL) {
      return TSDB_CODE_OUT_OF_MEMORY;
    }

    // Set the correct column info (data type and bytes)
404 405
    pColInfo->info.type = pFuncParam->param.nType;
    pColInfo->info.bytes = pFuncParam->param.nLen;
406 407

    pInput->pData[paramIndex] = pColInfo;
408 409
  } else {
    pColInfo = pInput->pData[paramIndex];
410 411
  }

H
Haojun Liao 已提交
412
  colInfoDataEnsureCapacity(pColInfo, numOfRows, false);
413

414
  int8_t type = pFuncParam->param.nType;
415 416
  if (type == TSDB_DATA_TYPE_BIGINT || type == TSDB_DATA_TYPE_UBIGINT) {
    int64_t v = pFuncParam->param.i;
dengyihao's avatar
dengyihao 已提交
417
    for (int32_t i = 0; i < numOfRows; ++i) {
418 419 420 421
      colDataAppendInt64(pColInfo, i, &v);
    }
  } else if (type == TSDB_DATA_TYPE_DOUBLE) {
    double v = pFuncParam->param.d;
dengyihao's avatar
dengyihao 已提交
422
    for (int32_t i = 0; i < numOfRows; ++i) {
423 424
      colDataAppendDouble(pColInfo, i, &v);
    }
425
  } else if (type == TSDB_DATA_TYPE_VARCHAR) {
L
Liu Jicong 已提交
426
    char* tmp = taosMemoryMalloc(pFuncParam->param.nLen + VARSTR_HEADER_SIZE);
427
    STR_WITH_SIZE_TO_VARSTR(tmp, pFuncParam->param.pz, pFuncParam->param.nLen);
L
Liu Jicong 已提交
428
    for (int32_t i = 0; i < numOfRows; ++i) {
429 430
      colDataAppend(pColInfo, i, tmp, false);
    }
H
Haojun Liao 已提交
431
    taosMemoryFree(tmp);
432 433 434 435 436
  }

  return TSDB_CODE_SUCCESS;
}

437 438
static int32_t doSetInputDataBlock(SExprSupp* pExprSup, SSDataBlock* pBlock, int32_t order, int32_t scanFlag,
                                   bool createDummyCol) {
439
  int32_t         code = TSDB_CODE_SUCCESS;
440
  SqlFunctionCtx* pCtx = pExprSup->pCtx;
441

442
  for (int32_t i = 0; i < pExprSup->numOfExprs; ++i) {
L
Liu Jicong 已提交
443
    pCtx[i].order = order;
444 445
    pCtx[i].input.numOfRows = pBlock->info.rows;

L
Liu Jicong 已提交
446
    pCtx[i].pSrcBlock = pBlock;
X
Xiaoyu Wang 已提交
447
    pCtx[i].scanFlag = scanFlag;
H
Haojun Liao 已提交
448

449
    SInputColumnInfoData* pInput = &pCtx[i].input;
H
Haojun Liao 已提交
450
    pInput->uid = pBlock->info.id.uid;
H
Haojun Liao 已提交
451
    pInput->colDataSMAIsSet = false;
452

453
    SExprInfo* pOneExpr = &pExprSup->pExprInfo[i];
454
    for (int32_t j = 0; j < pOneExpr->base.numOfParams; ++j) {
dengyihao's avatar
dengyihao 已提交
455
      SFunctParam* pFuncParam = &pOneExpr->base.pParam[j];
G
Ganlin Zhao 已提交
456 457
      if (pFuncParam->type == FUNC_PARAM_TYPE_COLUMN) {
        int32_t slotId = pFuncParam->pCol->slotId;
dengyihao's avatar
dengyihao 已提交
458
        pInput->pData[j] = taosArrayGet(pBlock->pDataBlock, slotId);
459 460 461
        pInput->totalRows = pBlock->info.rows;
        pInput->numOfRows = pBlock->info.rows;
        pInput->startRowIndex = 0;
462

463
        // NOTE: the last parameter is the primary timestamp column
H
Haojun Liao 已提交
464
        // todo: refactor this
465
        if (fmIsImplicitTsFunc(pCtx[i].functionId) && (j == pOneExpr->base.numOfParams - 1)) {
L
Liu Jicong 已提交
466
          pInput->pPTS = pInput->pData[j];  // in case of merge function, this is not always the ts column data.
467
          //          ASSERT(pInput->pPTS->info.type == TSDB_DATA_TYPE_TIMESTAMP);
468
        }
469 470
        ASSERT(pInput->pData[j] != NULL);
      } else if (pFuncParam->type == FUNC_PARAM_TYPE_VALUE) {
471 472 473
        // todo avoid case: top(k, 12), 12 is the value parameter.
        // sum(11), 11 is also the value parameter.
        if (createDummyCol && pOneExpr->base.numOfParams == 1) {
474 475 476 477
          pInput->totalRows = pBlock->info.rows;
          pInput->numOfRows = pBlock->info.rows;
          pInput->startRowIndex = 0;

478
          code = doCreateConstantValColumnInfo(pInput, pFuncParam, j, pBlock->info.rows);
479 480 481
          if (code != TSDB_CODE_SUCCESS) {
            return code;
          }
482
        }
G
Ganlin Zhao 已提交
483 484
      }
    }
H
Haojun Liao 已提交
485
  }
486 487

  return code;
H
Haojun Liao 已提交
488 489
}

490
static int32_t doAggregateImpl(SOperatorInfo* pOperator, SqlFunctionCtx* pCtx) {
491
  for (int32_t k = 0; k < pOperator->exprSupp.numOfExprs; ++k) {
H
Haojun Liao 已提交
492
    if (functionNeedToExecute(&pCtx[k])) {
493
      // todo add a dummy funtion to avoid process check
494 495 496
      if (pCtx[k].fpSet.process == NULL) {
        continue;
      }
H
Haojun Liao 已提交
497

498 499 500 501
      int32_t code = pCtx[k].fpSet.process(&pCtx[k]);
      if (code != TSDB_CODE_SUCCESS) {
        qError("%s aggregate function error happens, code: %s", GET_TASKID(pOperator->pTaskInfo), tstrerror(code));
        return code;
502
      }
503 504
    }
  }
505 506

  return TSDB_CODE_SUCCESS;
507 508
}

5
54liuyao 已提交
509
bool functionNeedToExecute(SqlFunctionCtx* pCtx) {
510
  struct SResultRowEntryInfo* pResInfo = GET_RES_INFO(pCtx);
511

512 513 514 515 516
  // in case of timestamp column, always generated results.
  int32_t functionId = pCtx->functionId;
  if (functionId == -1) {
    return false;
  }
517

518 519
  if (pCtx->scanFlag == REPEAT_SCAN) {
    return fmIsRepeatScanFunc(pCtx->functionId);
520 521
  }

522 523
  if (isRowEntryCompleted(pResInfo)) {
    return false;
524 525
  }

526 527 528
  return true;
}

529 530 531 532 533 534 535
static int32_t doCreateConstantValColumnAggInfo(SInputColumnInfoData* pInput, SFunctParam* pFuncParam, int32_t type,
                                                int32_t paramIndex, int32_t numOfRows) {
  if (pInput->pData[paramIndex] == NULL) {
    pInput->pData[paramIndex] = taosMemoryCalloc(1, sizeof(SColumnInfoData));
    if (pInput->pData[paramIndex] == NULL) {
      return TSDB_CODE_OUT_OF_MEMORY;
    }
536

537 538 539
    // Set the correct column info (data type and bytes)
    pInput->pData[paramIndex]->info.type = type;
    pInput->pData[paramIndex]->info.bytes = tDataTypes[type].bytes;
540
  }
H
Haojun Liao 已提交
541

542 543 544 545 546 547
  SColumnDataAgg* da = NULL;
  if (pInput->pColumnDataAgg[paramIndex] == NULL) {
    da = taosMemoryCalloc(1, sizeof(SColumnDataAgg));
    pInput->pColumnDataAgg[paramIndex] = da;
    if (da == NULL) {
      return TSDB_CODE_OUT_OF_MEMORY;
548 549
    }
  } else {
550
    da = pInput->pColumnDataAgg[paramIndex];
551 552
  }

553
  ASSERT(!IS_VAR_DATA_TYPE(type));
554

555 556
  if (type == TSDB_DATA_TYPE_BIGINT) {
    int64_t v = pFuncParam->param.i;
557
    *da = (SColumnDataAgg){.numOfNull = 0, .min = v, .max = v, .sum = v * numOfRows};
558 559
  } else if (type == TSDB_DATA_TYPE_DOUBLE) {
    double v = pFuncParam->param.d;
560
    *da = (SColumnDataAgg){.numOfNull = 0};
561

562 563 564 565 566 567
    *(double*)&da->min = v;
    *(double*)&da->max = v;
    *(double*)&da->sum = v * numOfRows;
  } else if (type == TSDB_DATA_TYPE_BOOL) {  // todo validate this data type
    bool v = pFuncParam->param.i;

568
    *da = (SColumnDataAgg){.numOfNull = 0};
569 570 571 572 573
    *(bool*)&da->min = 0;
    *(bool*)&da->max = v;
    *(bool*)&da->sum = v * numOfRows;
  } else if (type == TSDB_DATA_TYPE_TIMESTAMP) {
    // do nothing
574
  } else {
575
    ASSERT(0);
576 577
  }

578 579
  return TSDB_CODE_SUCCESS;
}
580

581
void setBlockSMAInfo(SqlFunctionCtx* pCtx, SExprInfo* pExprInfo, SSDataBlock* pBlock) {
582 583 584 585 586 587 588
  int32_t numOfRows = pBlock->info.rows;

  SInputColumnInfoData* pInput = &pCtx->input;
  pInput->numOfRows = numOfRows;
  pInput->totalRows = numOfRows;

  if (pBlock->pBlockAgg != NULL) {
H
Haojun Liao 已提交
589
    pInput->colDataSMAIsSet = true;
590

591 592
    for (int32_t j = 0; j < pExprInfo->base.numOfParams; ++j) {
      SFunctParam* pFuncParam = &pExprInfo->base.pParam[j];
593

594 595
      if (pFuncParam->type == FUNC_PARAM_TYPE_COLUMN) {
        int32_t slotId = pFuncParam->pCol->slotId;
596 597
        pInput->pColumnDataAgg[j] = pBlock->pBlockAgg[slotId];
        if (pInput->pColumnDataAgg[j] == NULL) {
H
Haojun Liao 已提交
598
          pInput->colDataSMAIsSet = false;
599
        }
600 601 602 603

        // Here we set the column info data since the data type for each column data is required, but
        // the data in the corresponding SColumnInfoData will not be used.
        pInput->pData[j] = taosArrayGet(pBlock->pDataBlock, slotId);
604 605
      } else if (pFuncParam->type == FUNC_PARAM_TYPE_VALUE) {
        doCreateConstantValColumnAggInfo(pInput, pFuncParam, pFuncParam->param.nType, j, pBlock->info.rows);
606 607
      }
    }
608
  } else {
H
Haojun Liao 已提交
609
    pInput->colDataSMAIsSet = false;
610 611 612
  }
}

L
Liu Jicong 已提交
613
bool isTaskKilled(SExecTaskInfo* pTaskInfo) {
614 615
  // query has been executed more than tsShellActivityTimer, and the retrieve has not arrived
  // abort current query execution.
L
Liu Jicong 已提交
616 617
  if (pTaskInfo->owner != 0 &&
      ((taosGetTimestampSec() - pTaskInfo->cost.start / 1000) > 10 * getMaximumIdleDurationSec())
618
      /*(!needBuildResAfterQueryComplete(pTaskInfo))*/) {
619
    assert(pTaskInfo->cost.start != 0);
L
Liu Jicong 已提交
620 621 622
    //    qDebug("QInfo:%" PRIu64 " retrieve not arrive beyond %d ms, abort current query execution, start:%" PRId64
    //           ", current:%d", pQInfo->qId, 1, pQInfo->startExecTs, taosGetTimestampSec());
    //    return true;
623 624 625 626 627
  }

  return false;
}

L
Liu Jicong 已提交
628
void setTaskKilled(SExecTaskInfo* pTaskInfo) { pTaskInfo->code = TSDB_CODE_TSC_QUERY_CANCELLED; }
629 630

/////////////////////////////////////////////////////////////////////////////////////////////
631
STimeWindow getAlignQueryTimeWindow(SInterval* pInterval, int32_t precision, int64_t key) {
L
Liu Jicong 已提交
632
  STimeWindow win = {0};
633
  win.skey = taosTimeTruncate(key, pInterval, precision);
634 635

  /*
H
Haojun Liao 已提交
636
   * if the realSkey > INT64_MAX - pInterval->interval, the query duration between
637 638
   * realSkey and realEkey must be less than one interval.Therefore, no need to adjust the query ranges.
   */
639 640 641
  win.ekey = taosTimeAdd(win.skey, pInterval->interval, pInterval->intervalUnit, precision) - 1;
  if (win.ekey < win.skey) {
    win.ekey = INT64_MAX;
642
  }
643 644

  return win;
645 646
}

L
Liu Jicong 已提交
647 648
int32_t loadDataBlockOnDemand(SExecTaskInfo* pTaskInfo, STableScanInfo* pTableScanInfo, SSDataBlock* pBlock,
                              uint32_t* status) {
649
  *status = BLK_DATA_NOT_LOAD;
650

H
Haojun Liao 已提交
651
  pBlock->pDataBlock = NULL;
L
Liu Jicong 已提交
652
  pBlock->pBlockAgg = NULL;
H
Haojun Liao 已提交
653

L
Liu Jicong 已提交
654 655
  //  int64_t groupId = pRuntimeEnv->current->groupIndex;
  //  bool    ascQuery = QUERY_IS_ASC_QUERY(pQueryAttr);
656

H
Haojun Liao 已提交
657
  STaskCostInfo* pCost = &pTaskInfo->cost;
658

659 660
//  pCost->totalBlocks += 1;
//  pCost->totalRows += pBlock->info.rows;
H
Haojun Liao 已提交
661
#if 0
662 663 664
  // Calculate all time windows that are overlapping or contain current data block.
  // If current data block is contained by all possible time window, do not load current data block.
  if (/*pQueryAttr->pFilters || */pQueryAttr->groupbyColumn || pQueryAttr->sw.gap > 0 ||
H
Haojun Liao 已提交
665
      (QUERY_IS_INTERVAL_QUERY(pQueryAttr) && overlapWithTimeWindow(pTaskInfo, &pBlock->info))) {
666
    (*status) = BLK_DATA_DATA_LOAD;
667 668 669
  }

  // check if this data block is required to load
670
  if ((*status) != BLK_DATA_DATA_LOAD) {
671 672 673 674 675 676 677
    bool needFilter = true;

    // the pCtx[i] result is belonged to previous time window since the outputBuf has not been set yet,
    // the filter result may be incorrect. So in case of interval query, we need to set the correct time output buffer
    if (QUERY_IS_INTERVAL_QUERY(pQueryAttr)) {
      SResultRow* pResult = NULL;

H
Haojun Liao 已提交
678
      bool  masterScan = IS_MAIN_SCAN(pRuntimeEnv);
679 680 681 682 683 684
      TSKEY k = ascQuery? pBlock->info.window.skey : pBlock->info.window.ekey;

      STimeWindow win = getActiveTimeWindow(pTableScanInfo->pResultRowInfo, k, pQueryAttr);
      if (pQueryAttr->pointInterpQuery) {
        needFilter = chkWindowOutputBufByKey(pRuntimeEnv, pTableScanInfo->pResultRowInfo, &win, masterScan, &pResult, groupId,
                                    pTableScanInfo->pCtx, pTableScanInfo->numOfOutput,
685
                                    pTableScanInfo->rowEntryInfoOffset);
686
      } else {
H
Haojun Liao 已提交
687
        if (setResultOutputBufByKey(pRuntimeEnv, pTableScanInfo->pResultRowInfo, pBlock->info.id.uid, &win, masterScan, &pResult, groupId,
688
                                    pTableScanInfo->pCtx, pTableScanInfo->numOfOutput,
689
                                    pTableScanInfo->rowEntryInfoOffset) != TSDB_CODE_SUCCESS) {
690
          T_LONG_JMP(pRuntimeEnv->env, TSDB_CODE_QRY_OUT_OF_MEMORY);
691 692 693 694
        }
      }
    } else if (pQueryAttr->stableQuery && (!pQueryAttr->tsCompQuery) && (!pQueryAttr->diffQuery)) { // stable aggregate, not interval aggregate or normal column aggregate
      doSetTableGroupOutputBuf(pRuntimeEnv, pTableScanInfo->pResultRowInfo, pTableScanInfo->pCtx,
695
                               pTableScanInfo->rowEntryInfoOffset, pTableScanInfo->numOfOutput,
696 697 698 699 700 701
                               pRuntimeEnv->current->groupIndex);
    }

    if (needFilter) {
      (*status) = doFilterByBlockTimeWindow(pTableScanInfo, pBlock);
    } else {
702
      (*status) = BLK_DATA_DATA_LOAD;
703 704 705 706
    }
  }

  SDataBlockInfo* pBlockInfo = &pBlock->info;
H
Haojun Liao 已提交
707
//  *status = updateBlockLoadStatus(pRuntimeEnv->pQueryAttr, *status);
708

709
  if ((*status) == BLK_DATA_NOT_LOAD || (*status) == BLK_DATA_FILTEROUT) {
710 711
    //qDebug("QInfo:0x%"PRIx64" data block discard, brange:%" PRId64 "-%" PRId64 ", rows:%d", pQInfo->qId, pBlockInfo->window.skey,
//           pBlockInfo->window.ekey, pBlockInfo->rows);
712
    pCost->skipBlocks += 1;
713
  } else if ((*status) == BLK_DATA_SMA_LOAD) {
714 715
    // this function never returns error?
    pCost->loadBlockStatis += 1;
716
//    tsdbRetrieveDatablockSMA(pTableScanInfo->pTsdbReadHandle, &pBlock->pBlockAgg);
717 718

    if (pBlock->pBlockAgg == NULL) {  // data block statistics does not exist, load data block
719
//      pBlock->pDataBlock = tsdbRetrieveDataBlock(pTableScanInfo->pTsdbReadHandle, NULL);
720 721 722
      pCost->totalCheckedRows += pBlock->info.rows;
    }
  } else {
723
    assert((*status) == BLK_DATA_DATA_LOAD);
724 725 726

    // load the data block statistics to perform further filter
    pCost->loadBlockStatis += 1;
727
//    tsdbRetrieveDatablockSMA(pTableScanInfo->pTsdbReadHandle, &pBlock->pBlockAgg);
728 729 730 731 732 733

    if (pQueryAttr->topBotQuery && pBlock->pBlockAgg != NULL) {
      { // set previous window
        if (QUERY_IS_INTERVAL_QUERY(pQueryAttr)) {
          SResultRow* pResult = NULL;

H
Haojun Liao 已提交
734
          bool  masterScan = IS_MAIN_SCAN(pRuntimeEnv);
735 736 737
          TSKEY k = ascQuery? pBlock->info.window.skey : pBlock->info.window.ekey;

          STimeWindow win = getActiveTimeWindow(pTableScanInfo->pResultRowInfo, k, pQueryAttr);
H
Haojun Liao 已提交
738
          if (setResultOutputBufByKey(pRuntimeEnv, pTableScanInfo->pResultRowInfo, pBlock->info.id.uid, &win, masterScan, &pResult, groupId,
739
                                      pTableScanInfo->pCtx, pTableScanInfo->numOfOutput,
740
                                      pTableScanInfo->rowEntryInfoOffset) != TSDB_CODE_SUCCESS) {
741
            T_LONG_JMP(pRuntimeEnv->env, TSDB_CODE_QRY_OUT_OF_MEMORY);
742 743 744 745 746 747 748 749 750 751
          }
        }
      }
      bool load = false;
      for (int32_t i = 0; i < pQueryAttr->numOfOutput; ++i) {
        int32_t functionId = pTableScanInfo->pCtx[i].functionId;
        if (functionId == FUNCTION_TOP || functionId == FUNCTION_BOTTOM) {
//          load = topbot_datablock_filter(&pTableScanInfo->pCtx[i], (char*)&(pBlock->pBlockAgg[i].min),
//                                         (char*)&(pBlock->pBlockAgg[i].max));
          if (!load) { // current block has been discard due to filter applied
752
            pCost->skipBlocks += 1;
753 754
            //qDebug("QInfo:0x%"PRIx64" data block discard, brange:%" PRId64 "-%" PRId64 ", rows:%d", pQInfo->qId,
//                   pBlockInfo->window.skey, pBlockInfo->window.ekey, pBlockInfo->rows);
755
            (*status) = BLK_DATA_FILTEROUT;
756 757 758 759 760 761 762
            return TSDB_CODE_SUCCESS;
          }
        }
      }
    }

    // current block has been discard due to filter applied
H
Haojun Liao 已提交
763
//    if (!doFilterByBlockSMA(pRuntimeEnv, pBlock->pBlockAgg, pTableScanInfo->pCtx, pBlockInfo->rows)) {
764
//      pCost->skipBlocks += 1;
765 766
//      qDebug("QInfo:0x%"PRIx64" data block discard, brange:%" PRId64 "-%" PRId64 ", rows:%d", pQInfo->qId, pBlockInfo->window.skey,
//             pBlockInfo->window.ekey, pBlockInfo->rows);
767
//      (*status) = BLK_DATA_FILTEROUT;
768 769 770 771 772
//      return TSDB_CODE_SUCCESS;
//    }

    pCost->totalCheckedRows += pBlockInfo->rows;
    pCost->loadBlocks += 1;
773
//    pBlock->pDataBlock = tsdbRetrieveDataBlock(pTableScanInfo->pTsdbReadHandle, NULL);
774 775 776 777 778
//    if (pBlock->pDataBlock == NULL) {
//      return terrno;
//    }

//    if (pQueryAttr->pFilters != NULL) {
779
//      filterSetColFieldData(pQueryAttr->pFilters, taosArrayGetSize(pBlock->pDataBlock), pBlock->pDataBlock);
780
//    }
781

782 783 784 785
//    if (pQueryAttr->pFilters != NULL || pRuntimeEnv->pTsBuf != NULL) {
//      filterColRowsInDataBlock(pRuntimeEnv, pBlock, ascQuery);
//    }
  }
H
Haojun Liao 已提交
786
#endif
787 788 789
  return TSDB_CODE_SUCCESS;
}

L
Liu Jicong 已提交
790
void setTaskStatus(SExecTaskInfo* pTaskInfo, int8_t status) {
791
  if (status == TASK_NOT_COMPLETED) {
H
Haojun Liao 已提交
792
    pTaskInfo->status = status;
793 794
  } else {
    // QUERY_NOT_COMPLETED is not compatible with any other status, so clear its position first
795
    CLEAR_QUERY_STATUS(pTaskInfo, TASK_NOT_COMPLETED);
H
Haojun Liao 已提交
796
    pTaskInfo->status |= status;
797 798 799
  }
}

800
void setResultRowInitCtx(SResultRow* pResult, SqlFunctionCtx* pCtx, int32_t numOfOutput, int32_t* rowEntryInfoOffset) {
5
54liuyao 已提交
801
  bool init = false;
802
  for (int32_t i = 0; i < numOfOutput; ++i) {
803
    pCtx[i].resultInfo = getResultEntryInfo(pResult, i, rowEntryInfoOffset);
5
54liuyao 已提交
804 805 806
    if (init) {
      continue;
    }
807 808 809 810 811

    struct SResultRowEntryInfo* pResInfo = pCtx[i].resultInfo;
    if (isRowEntryCompleted(pResInfo) && isRowEntryInitialized(pResInfo)) {
      continue;
    }
812 813 814 815 816

    if (fmIsWindowPseudoColumnFunc(pCtx[i].functionId)) {
      continue;
    }

817 818 819 820 821 822
    if (!pResInfo->initialized) {
      if (pCtx[i].functionId != -1) {
        pCtx[i].fpSet.init(&pCtx[i], pResInfo);
      } else {
        pResInfo->initialized = true;
      }
5
54liuyao 已提交
823 824
    } else {
      init = true;
825 826 827 828
    }
  }
}

H
Haojun Liao 已提交
829 830
void doFilter(SSDataBlock* pBlock, SFilterInfo* pFilterInfo, SColMatchInfo* pColMatchInfo) {
  if (pFilterInfo == NULL || pBlock->info.rows == 0) {
S
shenglian zhou 已提交
831 832
    return;
  }
833

834
  SFilterColumnParam param1 = {.numOfCols = taosArrayGetSize(pBlock->pDataBlock), .pDataBlock = pBlock->pDataBlock};
835
  int32_t            code = filterSetDataFromSlotId(pFilterInfo, &param1);
836

837
  SColumnInfoData* p = NULL;
838
  int32_t          status = 0;
H
Haojun Liao 已提交
839

840
  // todo the keep seems never to be True??
H
Haojun Liao 已提交
841
  bool keep = filterExecute(pFilterInfo, pBlock, &p, NULL, param1.numOfCols, &status);
842
  extractQualifiedTupleByFilterResult(pBlock, p, keep, status);
H
Haojun Liao 已提交
843

844
  if (pColMatchInfo != NULL) {
845
    size_t size = taosArrayGetSize(pColMatchInfo->pList);
H
Haojun Liao 已提交
846
    for (int32_t i = 0; i < size; ++i) {
H
Haojun Liao 已提交
847
      SColMatchItem* pInfo = taosArrayGet(pColMatchInfo->pList, i);
848
      if (pInfo->colId == PRIMARYKEY_TIMESTAMP_COL_ID) {
H
Haojun Liao 已提交
849
        SColumnInfoData* pColData = taosArrayGet(pBlock->pDataBlock, pInfo->dstSlotId);
850
        if (pColData->info.type == TSDB_DATA_TYPE_TIMESTAMP) {
H
Haojun Liao 已提交
851
          blockDataUpdateTsWindow(pBlock, pInfo->dstSlotId);
852 853 854 855 856 857
          break;
        }
      }
    }
  }

858 859
  colDataDestroy(p);
  taosMemoryFree(p);
860 861
}

862
void extractQualifiedTupleByFilterResult(SSDataBlock* pBlock, const SColumnInfoData* p, bool keep, int32_t status) {
863 864 865 866
  if (keep) {
    return;
  }

H
Haojun Liao 已提交
867 868 869
  int32_t totalRows = pBlock->info.rows;

  if (status == FILTER_RESULT_ALL_QUALIFIED) {
870
    // here nothing needs to be done
H
Haojun Liao 已提交
871
  } else if (status == FILTER_RESULT_NONE_QUALIFIED) {
872
    pBlock->info.rows = 0;
H
Haojun Liao 已提交
873
  } else {
874
    SSDataBlock* px = createOneDataBlock(pBlock, true);
875

876 877
    size_t numOfCols = taosArrayGetSize(pBlock->pDataBlock);
    for (int32_t i = 0; i < numOfCols; ++i) {
878 879
      SColumnInfoData* pSrc = taosArrayGet(px->pDataBlock, i);
      SColumnInfoData* pDst = taosArrayGet(pBlock->pDataBlock, i);
880
      // it is a reserved column for scalar function, and no data in this column yet.
881
      if (pDst->pData == NULL || pSrc->pData == NULL) {
882 883 884
        continue;
      }

885 886
      colInfoDataCleanup(pDst, pBlock->info.rows);

887
      int32_t numOfRows = 0;
888
      for (int32_t j = 0; j < totalRows; ++j) {
889
        if (((int8_t*)p->pData)[j] == 0) {
D
dapan1121 已提交
890 891
          continue;
        }
892

D
dapan1121 已提交
893
        if (colDataIsNull_s(pSrc, j)) {
894
          colDataAppendNULL(pDst, numOfRows);
D
dapan1121 已提交
895
        } else {
896
          colDataAppend(pDst, numOfRows, colDataGetData(pSrc, j), false);
D
dapan1121 已提交
897
        }
898
        numOfRows += 1;
H
Haojun Liao 已提交
899
      }
900

901
      // todo this value can be assigned directly
902 903 904 905 906
      if (pBlock->info.rows == totalRows) {
        pBlock->info.rows = numOfRows;
      } else {
        ASSERT(pBlock->info.rows == numOfRows);
      }
907
    }
908

dengyihao's avatar
dengyihao 已提交
909
    blockDataDestroy(px);  // fix memory leak
910 911 912
  }
}

913
void doSetTableGroupOutputBuf(SOperatorInfo* pOperator, int32_t numOfOutput, uint64_t groupId) {
914
  // for simple group by query without interval, all the tables belong to one group result.
915 916 917
  SExecTaskInfo*    pTaskInfo = pOperator->pTaskInfo;
  SAggOperatorInfo* pAggInfo = pOperator->info;

918
  SResultRowInfo* pResultRowInfo = &pAggInfo->binfo.resultRowInfo;
919 920
  SqlFunctionCtx* pCtx = pOperator->exprSupp.pCtx;
  int32_t*        rowEntryInfoOffset = pOperator->exprSupp.rowEntryInfoOffset;
921

922
  SResultRow* pResultRow = doSetResultOutBufByKey(pAggInfo->aggSup.pResultBuf, pResultRowInfo, (char*)&groupId,
L
Liu Jicong 已提交
923
                                                  sizeof(groupId), true, groupId, pTaskInfo, false, &pAggInfo->aggSup);
L
Liu Jicong 已提交
924
  assert(pResultRow != NULL);
925 926 927 928 929 930

  /*
   * not assign result buffer yet, add new result buffer
   * all group belong to one result set, and each group result has different group id so set the id to be one
   */
  if (pResultRow->pageId == -1) {
dengyihao's avatar
dengyihao 已提交
931 932
    int32_t ret =
        addNewWindowResultBuf(pResultRow, pAggInfo->aggSup.pResultBuf, groupId, pAggInfo->binfo.pRes->info.rowSize);
933 934 935 936 937
    if (ret != TSDB_CODE_SUCCESS) {
      return;
    }
  }

938
  setResultRowInitCtx(pResultRow, pCtx, numOfOutput, rowEntryInfoOffset);
939 940
}

941 942 943
static void setExecutionContext(SOperatorInfo* pOperator, int32_t numOfOutput, uint64_t groupId) {
  SAggOperatorInfo* pAggInfo = pOperator->info;
  if (pAggInfo->groupId != UINT64_MAX && pAggInfo->groupId == groupId) {
944 945
    return;
  }
946 947

  doSetTableGroupOutputBuf(pOperator, numOfOutput, groupId);
948 949

  // record the current active group id
H
Haojun Liao 已提交
950
  pAggInfo->groupId = groupId;
951 952
}

dengyihao's avatar
dengyihao 已提交
953 954
static void doUpdateNumOfRows(SqlFunctionCtx* pCtx, SResultRow* pRow, int32_t numOfExprs,
                              const int32_t* rowCellOffset) {
955
  bool returnNotNull = false;
956
  for (int32_t j = 0; j < numOfExprs; ++j) {
957
    struct SResultRowEntryInfo* pResInfo = getResultEntryInfo(pRow, j, rowCellOffset);
958 959 960 961 962 963 964
    if (!isRowEntryInitialized(pResInfo)) {
      continue;
    }

    if (pRow->numOfRows < pResInfo->numOfRes) {
      pRow->numOfRows = pResInfo->numOfRes;
    }
965

966
    if (fmIsNotNullOutputFunc(pCtx[j].functionId)) {
967 968
      returnNotNull = true;
    }
969
  }
S
shenglian zhou 已提交
970 971
  // if all expr skips all blocks, e.g. all null inputs for max function, output one row in final result.
  //  except for first/last, which require not null output, output no rows
972
  if (pRow->numOfRows == 0 && !returnNotNull) {
973
    pRow->numOfRows = 1;
974 975 976
  }
}

977 978
static void doCopyResultToDataBlock(SExprInfo* pExprInfo, int32_t numOfExprs, SResultRow* pRow, SqlFunctionCtx* pCtx,
                                    SSDataBlock* pBlock, const int32_t* rowEntryOffset, SExecTaskInfo* pTaskInfo) {
979 980 981
  for (int32_t j = 0; j < numOfExprs; ++j) {
    int32_t slotId = pExprInfo[j].base.resSchema.slotId;

982
    pCtx[j].resultInfo = getResultEntryInfo(pRow, j, rowEntryOffset);
983
    if (pCtx[j].fpSet.finalize) {
984
      if (strcmp(pCtx[j].pExpr->pExpr->_function.functionName, "_group_key") == 0) {
985 986
        // for groupkey along with functions that output multiple lines(e.g. Histogram)
        // need to match groupkey result for each output row of that function.
987 988 989 990 991
        if (pCtx[j].resultInfo->numOfRes != 0) {
          pCtx[j].resultInfo->numOfRes = pRow->numOfRows;
        }
      }

992 993 994
      int32_t code = pCtx[j].fpSet.finalize(&pCtx[j], pBlock);
      if (TAOS_FAILED(code)) {
        qError("%s build result data block error, code %s", GET_TASKID(pTaskInfo), tstrerror(code));
995
        T_LONG_JMP(pTaskInfo->env, code);
996 997
      }
    } else if (strcmp(pCtx[j].pExpr->pExpr->_function.functionName, "_select_value") == 0) {
998
      // do nothing
999
    } else {
1000 1001
      // expand the result into multiple rows. E.g., _wstart, top(k, 20)
      // the _wstart needs to copy to 20 following rows, since the results of top-k expands to 20 different rows.
1002 1003 1004 1005 1006 1007 1008
      SColumnInfoData* pColInfoData = taosArrayGet(pBlock->pDataBlock, slotId);
      char*            in = GET_ROWCELL_INTERBUF(pCtx[j].resultInfo);
      for (int32_t k = 0; k < pRow->numOfRows; ++k) {
        colDataAppend(pColInfoData, pBlock->info.rows + k, in, pCtx[j].resultInfo->isNullRes);
      }
    }
  }
1009 1010
}

1011 1012 1013
// todo refactor. SResultRow has direct pointer in miainfo
int32_t finalizeResultRows(SDiskbasedBuf* pBuf, SResultRowPosition* resultRowPosition, SExprSupp* pSup,
                           SSDataBlock* pBlock, SExecTaskInfo* pTaskInfo) {
1014 1015 1016 1017 1018 1019 1020 1021 1022 1023 1024 1025 1026 1027 1028 1029 1030 1031 1032 1033 1034 1035 1036 1037 1038 1039
  SFilePage*  page = getBufPage(pBuf, resultRowPosition->pageId);
  SResultRow* pRow = (SResultRow*)((char*)page + resultRowPosition->offset);

  SqlFunctionCtx* pCtx = pSup->pCtx;
  SExprInfo*      pExprInfo = pSup->pExprInfo;
  const int32_t*  rowEntryOffset = pSup->rowEntryInfoOffset;

  doUpdateNumOfRows(pCtx, pRow, pSup->numOfExprs, rowEntryOffset);
  if (pRow->numOfRows == 0) {
    releaseBufPage(pBuf, page);
    return 0;
  }

  int32_t size = pBlock->info.capacity;
  while (pBlock->info.rows + pRow->numOfRows > size) {
    size = size * 1.25;
  }

  int32_t code = blockDataEnsureCapacity(pBlock, size);
  if (TAOS_FAILED(code)) {
    releaseBufPage(pBuf, page);
    qError("%s ensure result data capacity failed, code %s", GET_TASKID(pTaskInfo), tstrerror(code));
    T_LONG_JMP(pTaskInfo->env, code);
  }

  doCopyResultToDataBlock(pExprInfo, pSup->numOfExprs, pRow, pCtx, pBlock, rowEntryOffset, pTaskInfo);
1040 1041

  releaseBufPage(pBuf, page);
1042
  pBlock->info.rows += pRow->numOfRows;
1043 1044 1045
  return 0;
}

1046 1047 1048 1049 1050 1051 1052
int32_t doCopyToSDataBlock(SExecTaskInfo* pTaskInfo, SSDataBlock* pBlock, SExprSupp* pSup, SDiskbasedBuf* pBuf,
                           SGroupResInfo* pGroupResInfo) {
  SExprInfo*      pExprInfo = pSup->pExprInfo;
  int32_t         numOfExprs = pSup->numOfExprs;
  int32_t*        rowEntryOffset = pSup->rowEntryInfoOffset;
  SqlFunctionCtx* pCtx = pSup->pCtx;

1053
  int32_t numOfRows = getNumOfTotalRes(pGroupResInfo);
1054

1055
  for (int32_t i = pGroupResInfo->index; i < numOfRows; i += 1) {
L
Liu Jicong 已提交
1056 1057
    SResKeyPos* pPos = taosArrayGetP(pGroupResInfo->pRows, i);
    SFilePage*  page = getBufPage(pBuf, pPos->pos.pageId);
1058

1059
    SResultRow* pRow = (SResultRow*)((char*)page + pPos->pos.offset);
1060

H
Haojun Liao 已提交
1061
    doUpdateNumOfRows(pCtx, pRow, numOfExprs, rowEntryOffset);
1062 1063

    // no results, continue to check the next one
1064 1065
    if (pRow->numOfRows == 0) {
      pGroupResInfo->index += 1;
1066
      releaseBufPage(pBuf, page);
1067 1068 1069
      continue;
    }

H
Haojun Liao 已提交
1070 1071
    if (pBlock->info.id.groupId == 0) {
      pBlock->info.id.groupId = pPos->groupId;
1072 1073
    } else {
      // current value belongs to different group, it can't be packed into one datablock
H
Haojun Liao 已提交
1074
      if (pBlock->info.id.groupId != pPos->groupId) {
1075
        releaseBufPage(pBuf, page);
1076 1077 1078 1079
        break;
      }
    }

1080
    if (pBlock->info.rows + pRow->numOfRows > pBlock->info.capacity) {
1081
      ASSERT(pBlock->info.rows > 0);
1082
      releaseBufPage(pBuf, page);
1083 1084 1085 1086
      break;
    }

    pGroupResInfo->index += 1;
1087
    doCopyResultToDataBlock(pExprInfo, numOfExprs, pRow, pCtx, pBlock, rowEntryOffset, pTaskInfo);
1088

1089
    releaseBufPage(pBuf, page);
1090
    pBlock->info.rows += pRow->numOfRows;
1091 1092
  }

X
Xiaoyu Wang 已提交
1093
  qDebug("%s result generated, rows:%d, groupId:%" PRIu64, GET_TASKID(pTaskInfo), pBlock->info.rows,
H
Haojun Liao 已提交
1094
         pBlock->info.id.groupId);
1095

1096
  blockDataUpdateTsWindow(pBlock, 0);
1097 1098 1099
  return 0;
}

1100 1101 1102 1103 1104 1105 1106 1107 1108 1109 1110 1111 1112 1113
void doBuildStreamResBlock(SOperatorInfo* pOperator, SOptrBasicInfo* pbInfo, SGroupResInfo* pGroupResInfo,
                           SDiskbasedBuf* pBuf) {
  SExecTaskInfo* pTaskInfo = pOperator->pTaskInfo;
  SSDataBlock*   pBlock = pbInfo->pRes;

  // set output datablock version
  pBlock->info.version = pTaskInfo->version;

  blockDataCleanup(pBlock);
  if (!hasRemainResults(pGroupResInfo)) {
    return;
  }

  // clear the existed group id
H
Haojun Liao 已提交
1114
  pBlock->info.id.groupId = 0;
1115 1116 1117
  ASSERT(!pbInfo->mergeResultBlock);
  doCopyToSDataBlock(pTaskInfo, pBlock, &pOperator->exprSupp, pBuf, pGroupResInfo);

1118
  void* tbname = NULL;
H
Haojun Liao 已提交
1119
  if (streamStateGetParName(pTaskInfo->streamInfo.pState, pBlock->info.id.groupId, &tbname) < 0) {
1120 1121 1122
    pBlock->info.parTbName[0] = 0;
  } else {
    memcpy(pBlock->info.parTbName, tbname, TSDB_TABLE_NAME_LEN);
1123
  }
1124
  tdbFree(tbname);
1125 1126
}

X
Xiaoyu Wang 已提交
1127 1128
void doBuildResultDatablock(SOperatorInfo* pOperator, SOptrBasicInfo* pbInfo, SGroupResInfo* pGroupResInfo,
                            SDiskbasedBuf* pBuf) {
1129
  SExecTaskInfo* pTaskInfo = pOperator->pTaskInfo;
1130
  SSDataBlock*   pBlock = pbInfo->pRes;
1131

1132 1133 1134
  // set output datablock version
  pBlock->info.version = pTaskInfo->version;

1135
  blockDataCleanup(pBlock);
1136
  if (!hasRemainResults(pGroupResInfo)) {
1137 1138 1139
    return;
  }

1140
  // clear the existed group id
H
Haojun Liao 已提交
1141
  pBlock->info.id.groupId = 0;
1142 1143 1144
  if (!pbInfo->mergeResultBlock) {
    doCopyToSDataBlock(pTaskInfo, pBlock, &pOperator->exprSupp, pBuf, pGroupResInfo);
  } else {
dengyihao's avatar
dengyihao 已提交
1145
    while (hasRemainResults(pGroupResInfo)) {
1146 1147 1148
      doCopyToSDataBlock(pTaskInfo, pBlock, &pOperator->exprSupp, pBuf, pGroupResInfo);
      if (pBlock->info.rows >= pOperator->resultInfo.threshold) {
        break;
1149 1150
      }

1151
      // clearing group id to continue to merge data that belong to different groups
H
Haojun Liao 已提交
1152
      pBlock->info.id.groupId = 0;
1153
    }
1154 1155

    // clear the group id info in SSDataBlock, since the client does not need it
H
Haojun Liao 已提交
1156
    pBlock->info.id.groupId = 0;
1157 1158 1159
  }
}

H
Haojun Liao 已提交
1160
void printTaskExecCostInLog(SExecTaskInfo* pTaskInfo) {
L
Liu Jicong 已提交
1161
  STaskCostInfo* pSummary = &pTaskInfo->cost;
1162

1163 1164
  SFileBlockLoadRecorder* pRecorder = pSummary->pRecoder;
  if (pSummary->pRecoder != NULL) {
1165
    qDebug(
H
Haojun Liao 已提交
1166 1167 1168 1169 1170
        "%s :cost summary: elapsed time:%.2f ms, extract tableList:%.2f ms, createGroupIdMap:%.2f ms, total blocks:%d, "
        "load block SMA:%d, load data block:%d, total rows:%" PRId64 ", check rows:%" PRId64,
        GET_TASKID(pTaskInfo), pSummary->elapsedTime / 1000.0, pSummary->extractListTime, pSummary->groupIdMapTime,
        pRecorder->totalBlocks, pRecorder->loadBlockStatis, pRecorder->loadBlocks, pRecorder->totalRows,
        pRecorder->totalCheckedRows);
1171
  }
1172 1173
}

L
Liu Jicong 已提交
1174 1175
// void skipBlocks(STaskRuntimeEnv *pRuntimeEnv) {
//   STaskAttr *pQueryAttr = pRuntimeEnv->pQueryAttr;
1176
//
L
Liu Jicong 已提交
1177 1178 1179
//   if (pQueryAttr->limit.offset <= 0 || pQueryAttr->numOfFilterCols > 0) {
//     return;
//   }
1180
//
L
Liu Jicong 已提交
1181 1182
//   pQueryAttr->pos = 0;
//   int32_t step = GET_FORWARD_DIRECTION_FACTOR(pQueryAttr->order.order);
1183
//
L
Liu Jicong 已提交
1184 1185
//   STableQueryInfo* pTableQueryInfo = pRuntimeEnv->current;
//   TsdbQueryHandleT pTsdbReadHandle = pRuntimeEnv->pTsdbReadHandle;
1186
//
L
Liu Jicong 已提交
1187 1188 1189
//   SDataBlockInfo blockInfo = SDATA_BLOCK_INITIALIZER;
//   while (tsdbNextDataBlock(pTsdbReadHandle)) {
//     if (isTaskKilled(pRuntimeEnv->qinfo)) {
1190
//       T_LONG_JMP(pRuntimeEnv->env, TSDB_CODE_TSC_QUERY_CANCELLED);
L
Liu Jicong 已提交
1191
//     }
1192
//
L
Liu Jicong 已提交
1193
//     tsdbRetrieveDataBlockInfo(pTsdbReadHandle, &blockInfo);
1194
//
L
Liu Jicong 已提交
1195 1196 1197 1198
//     if (pQueryAttr->limit.offset > blockInfo.rows) {
//       pQueryAttr->limit.offset -= blockInfo.rows;
//       pTableQueryInfo->lastKey = (QUERY_IS_ASC_QUERY(pQueryAttr)) ? blockInfo.window.ekey : blockInfo.window.skey;
//       pTableQueryInfo->lastKey += step;
1199
//
L
Liu Jicong 已提交
1200 1201 1202 1203 1204 1205 1206
//       //qDebug("QInfo:0x%"PRIx64" skip rows:%d, offset:%" PRId64, GET_TASKID(pRuntimeEnv), blockInfo.rows,
//              pQuery->limit.offset);
//     } else {  // find the appropriated start position in current block
//       updateOffsetVal(pRuntimeEnv, &blockInfo);
//       break;
//     }
//   }
1207
//
L
Liu Jicong 已提交
1208
//   if (terrno != TSDB_CODE_SUCCESS) {
1209
//     T_LONG_JMP(pRuntimeEnv->env, terrno);
L
Liu Jicong 已提交
1210 1211 1212 1213 1214 1215 1216
//   }
// }

// static TSKEY doSkipIntervalProcess(STaskRuntimeEnv* pRuntimeEnv, STimeWindow* win, SDataBlockInfo* pBlockInfo,
// STableQueryInfo* pTableQueryInfo) {
//   STaskAttr *pQueryAttr = pRuntimeEnv->pQueryAttr;
//   SResultRowInfo *pWindowResInfo = &pRuntimeEnv->resultRowInfo;
1217
//
L
Liu Jicong 已提交
1218 1219 1220
//   assert(pQueryAttr->limit.offset == 0);
//   STimeWindow tw = *win;
//   getNextTimeWindow(pQueryAttr, &tw);
1221
//
L
Liu Jicong 已提交
1222 1223
//   if ((tw.skey <= pBlockInfo->window.ekey && QUERY_IS_ASC_QUERY(pQueryAttr)) ||
//       (tw.ekey >= pBlockInfo->window.skey && !QUERY_IS_ASC_QUERY(pQueryAttr))) {
1224
//
L
Liu Jicong 已提交
1225 1226 1227 1228
//     // load the data block and check data remaining in current data block
//     // TODO optimize performance
//     SArray *         pDataBlock = tsdbRetrieveDataBlock(pRuntimeEnv->pTsdbReadHandle, NULL);
//     SColumnInfoData *pColInfoData = taosArrayGet(pDataBlock, 0);
1229
//
L
Liu Jicong 已提交
1230 1231 1232 1233
//     tw = *win;
//     int32_t startPos =
//         getNextQualifiedWindow(pQueryAttr, &tw, pBlockInfo, pColInfoData->pData, binarySearchForKey, -1);
//     assert(startPos >= 0);
1234
//
L
Liu Jicong 已提交
1235 1236
//     // set the abort info
//     pQueryAttr->pos = startPos;
1237
//
L
Liu Jicong 已提交
1238 1239 1240 1241
//     // reset the query start timestamp
//     pTableQueryInfo->win.skey = ((TSKEY *)pColInfoData->pData)[startPos];
//     pQueryAttr->window.skey = pTableQueryInfo->win.skey;
//     TSKEY key = pTableQueryInfo->win.skey;
1242
//
L
Liu Jicong 已提交
1243 1244
//     pWindowResInfo->prevSKey = tw.skey;
//     int32_t index = pRuntimeEnv->resultRowInfo.curIndex;
1245
//
L
Liu Jicong 已提交
1246 1247
//     int32_t numOfRes = tableApplyFunctionsOnBlock(pRuntimeEnv, pBlockInfo, NULL, binarySearchForKey, pDataBlock);
//     pRuntimeEnv->resultRowInfo.curIndex = index;  // restore the window index
1248
//
L
Liu Jicong 已提交
1249 1250 1251 1252
//     //qDebug("QInfo:0x%"PRIx64" check data block, brange:%" PRId64 "-%" PRId64 ", numOfRows:%d, numOfRes:%d,
//     lastKey:%" PRId64,
//            GET_TASKID(pRuntimeEnv), pBlockInfo->window.skey, pBlockInfo->window.ekey, pBlockInfo->rows, numOfRes,
//            pQueryAttr->current->lastKey);
1253
//
L
Liu Jicong 已提交
1254 1255 1256 1257 1258
//     return key;
//   } else {  // do nothing
//     pQueryAttr->window.skey      = tw.skey;
//     pWindowResInfo->prevSKey = tw.skey;
//     pTableQueryInfo->lastKey = tw.skey;
1259
//
L
Liu Jicong 已提交
1260 1261
//     return tw.skey;
//   }
1262
//
L
Liu Jicong 已提交
1263 1264 1265 1266 1267 1268 1269 1270 1271 1272
//   return true;
// }

// static bool skipTimeInterval(STaskRuntimeEnv *pRuntimeEnv, TSKEY* start) {
//   STaskAttr *pQueryAttr = pRuntimeEnv->pQueryAttr;
//   if (QUERY_IS_ASC_QUERY(pQueryAttr)) {
//     assert(*start <= pRuntimeEnv->current->lastKey);
//   } else {
//     assert(*start >= pRuntimeEnv->current->lastKey);
//   }
1273
//
L
Liu Jicong 已提交
1274 1275 1276 1277 1278
//   // if queried with value filter, do NOT forward query start position
//   if (pQueryAttr->limit.offset <= 0 || pQueryAttr->numOfFilterCols > 0 || pRuntimeEnv->pTsBuf != NULL ||
//   pRuntimeEnv->pFillInfo != NULL) {
//     return true;
//   }
1279
//
L
Liu Jicong 已提交
1280 1281 1282 1283 1284 1285 1286
//   /*
//    * 1. for interval without interpolation query we forward pQueryAttr->interval.interval at a time for
//    *    pQueryAttr->limit.offset times. Since hole exists, pQueryAttr->interval.interval*pQueryAttr->limit.offset
//    value is
//    *    not valid. otherwise, we only forward pQueryAttr->limit.offset number of points
//    */
//   assert(pRuntimeEnv->resultRowInfo.prevSKey == TSKEY_INITIAL_VAL);
1287
//
L
Liu Jicong 已提交
1288 1289
//   STimeWindow w = TSWINDOW_INITIALIZER;
//   bool ascQuery = QUERY_IS_ASC_QUERY(pQueryAttr);
1290
//
L
Liu Jicong 已提交
1291 1292
//   SResultRowInfo *pWindowResInfo = &pRuntimeEnv->resultRowInfo;
//   STableQueryInfo *pTableQueryInfo = pRuntimeEnv->current;
1293
//
L
Liu Jicong 已提交
1294 1295 1296
//   SDataBlockInfo blockInfo = SDATA_BLOCK_INITIALIZER;
//   while (tsdbNextDataBlock(pRuntimeEnv->pTsdbReadHandle)) {
//     tsdbRetrieveDataBlockInfo(pRuntimeEnv->pTsdbReadHandle, &blockInfo);
1297
//
L
Liu Jicong 已提交
1298 1299 1300 1301 1302 1303 1304 1305 1306
//     if (QUERY_IS_ASC_QUERY(pQueryAttr)) {
//       if (pWindowResInfo->prevSKey == TSKEY_INITIAL_VAL) {
//         getAlignQueryTimeWindow(pQueryAttr, blockInfo.window.skey, blockInfo.window.skey, pQueryAttr->window.ekey,
//         &w); pWindowResInfo->prevSKey = w.skey;
//       }
//     } else {
//       getAlignQueryTimeWindow(pQueryAttr, blockInfo.window.ekey, pQueryAttr->window.ekey, blockInfo.window.ekey, &w);
//       pWindowResInfo->prevSKey = w.skey;
//     }
1307
//
L
Liu Jicong 已提交
1308 1309
//     // the first time window
//     STimeWindow win = getActiveTimeWindow(pWindowResInfo, pWindowResInfo->prevSKey, pQueryAttr);
1310
//
L
Liu Jicong 已提交
1311 1312
//     while (pQueryAttr->limit.offset > 0) {
//       STimeWindow tw = win;
1313
//
L
Liu Jicong 已提交
1314 1315 1316
//       if ((win.ekey <= blockInfo.window.ekey && ascQuery) || (win.ekey >= blockInfo.window.skey && !ascQuery)) {
//         pQueryAttr->limit.offset -= 1;
//         pWindowResInfo->prevSKey = win.skey;
1317
//
L
Liu Jicong 已提交
1318 1319 1320 1321 1322 1323
//         // current time window is aligned with blockInfo.window.ekey
//         // restart it from next data block by set prevSKey to be TSKEY_INITIAL_VAL;
//         if ((win.ekey == blockInfo.window.ekey && ascQuery) || (win.ekey == blockInfo.window.skey && !ascQuery)) {
//           pWindowResInfo->prevSKey = TSKEY_INITIAL_VAL;
//         }
//       }
1324
//
L
Liu Jicong 已提交
1325 1326 1327 1328
//       if (pQueryAttr->limit.offset == 0) {
//         *start = doSkipIntervalProcess(pRuntimeEnv, &win, &blockInfo, pTableQueryInfo);
//         return true;
//       }
1329
//
L
Liu Jicong 已提交
1330 1331
//       // current window does not ended in current data block, try next data block
//       getNextTimeWindow(pQueryAttr, &tw);
1332
//
L
Liu Jicong 已提交
1333 1334 1335 1336 1337 1338 1339 1340 1341
//       /*
//        * If the next time window still starts from current data block,
//        * load the primary timestamp column first, and then find the start position for the next queried time window.
//        * Note that only the primary timestamp column is required.
//        * TODO: Optimize for this cases. All data blocks are not needed to be loaded, only if the first actually
//        required
//        * time window resides in current data block.
//        */
//       if ((tw.skey <= blockInfo.window.ekey && ascQuery) || (tw.ekey >= blockInfo.window.skey && !ascQuery)) {
1342
//
L
Liu Jicong 已提交
1343 1344
//         SArray *pDataBlock = tsdbRetrieveDataBlock(pRuntimeEnv->pTsdbReadHandle, NULL);
//         SColumnInfoData *pColInfoData = taosArrayGet(pDataBlock, 0);
1345
//
L
Liu Jicong 已提交
1346 1347 1348
//         if ((win.ekey > blockInfo.window.ekey && ascQuery) || (win.ekey < blockInfo.window.skey && !ascQuery)) {
//           pQueryAttr->limit.offset -= 1;
//         }
1349
//
L
Liu Jicong 已提交
1350 1351 1352 1353 1354 1355 1356 1357
//         if (pQueryAttr->limit.offset == 0) {
//           *start = doSkipIntervalProcess(pRuntimeEnv, &win, &blockInfo, pTableQueryInfo);
//           return true;
//         } else {
//           tw = win;
//           int32_t startPos =
//               getNextQualifiedWindow(pQueryAttr, &tw, &blockInfo, pColInfoData->pData, binarySearchForKey, -1);
//           assert(startPos >= 0);
1358
//
L
Liu Jicong 已提交
1359 1360 1361 1362 1363 1364 1365 1366 1367 1368 1369
//           // set the abort info
//           pQueryAttr->pos = startPos;
//           pTableQueryInfo->lastKey = ((TSKEY *)pColInfoData->pData)[startPos];
//           pWindowResInfo->prevSKey = tw.skey;
//           win = tw;
//         }
//       } else {
//         break;  // offset is not 0, and next time window begins or ends in the next block.
//       }
//     }
//   }
1370
//
L
Liu Jicong 已提交
1371 1372
//   // check for error
//   if (terrno != TSDB_CODE_SUCCESS) {
1373
//     T_LONG_JMP(pRuntimeEnv->env, terrno);
L
Liu Jicong 已提交
1374
//   }
1375
//
L
Liu Jicong 已提交
1376 1377
//   return true;
// }
1378

1379
int32_t appendDownstream(SOperatorInfo* p, SOperatorInfo** pDownstream, int32_t num) {
H
Haojun Liao 已提交
1380
  if (p->pDownstream == NULL) {
H
Haojun Liao 已提交
1381
    assert(p->numOfDownstream == 0);
1382 1383
  }

wafwerar's avatar
wafwerar 已提交
1384
  p->pDownstream = taosMemoryCalloc(1, num * POINTER_BYTES);
1385 1386 1387 1388 1389 1390 1391
  if (p->pDownstream == NULL) {
    return TSDB_CODE_OUT_OF_MEMORY;
  }

  memcpy(p->pDownstream, pDownstream, num * POINTER_BYTES);
  p->numOfDownstream = num;
  return TSDB_CODE_SUCCESS;
1392 1393
}

X
Xiaoyu Wang 已提交
1394
int32_t getTableScanInfo(SOperatorInfo* pOperator, int32_t* order, int32_t* scanFlag) {
1395
  // todo add more information about exchange operation
1396
  int32_t type = pOperator->operatorType;
X
Xiaoyu Wang 已提交
1397
  if (type == QUERY_NODE_PHYSICAL_PLAN_EXCHANGE || type == QUERY_NODE_PHYSICAL_PLAN_SYSTABLE_SCAN ||
1398
      type == QUERY_NODE_PHYSICAL_PLAN_STREAM_SCAN || type == QUERY_NODE_PHYSICAL_PLAN_TAG_SCAN ||
1399
      type == QUERY_NODE_PHYSICAL_PLAN_BLOCK_DIST_SCAN || type == QUERY_NODE_PHYSICAL_PLAN_LAST_ROW_SCAN) {
1400 1401 1402
    *order = TSDB_ORDER_ASC;
    *scanFlag = MAIN_SCAN;
    return TSDB_CODE_SUCCESS;
1403
  } else if (type == QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN) {
1404
    STableScanInfo* pTableScanInfo = pOperator->info;
H
Haojun Liao 已提交
1405 1406
    *order = pTableScanInfo->base.cond.order;
    *scanFlag = pTableScanInfo->base.scanFlag;
1407
    return TSDB_CODE_SUCCESS;
1408 1409
  } else if (type == QUERY_NODE_PHYSICAL_PLAN_TABLE_MERGE_SCAN) {
    STableMergeScanInfo* pTableScanInfo = pOperator->info;
H
Haojun Liao 已提交
1410 1411
    *order = pTableScanInfo->base.cond.order;
    *scanFlag = pTableScanInfo->base.scanFlag;
1412
    return TSDB_CODE_SUCCESS;
1413
  } else {
H
Haojun Liao 已提交
1414
    if (pOperator->pDownstream == NULL || pOperator->pDownstream[0] == NULL) {
1415
      return TSDB_CODE_INVALID_PARA;
H
Haojun Liao 已提交
1416
    } else {
1417
      return getTableScanInfo(pOperator->pDownstream[0], order, scanFlag);
1418 1419 1420
    }
  }
}
1421

1422 1423 1424 1425 1426
static int32_t createDataBlockForEmptyInput(SOperatorInfo* pOperator, SSDataBlock **ppBlock) {
  if (!tsCountAlwaysReturnValue) {
    return TSDB_CODE_SUCCESS;
  }

1427
  SOperatorInfo* downstream = pOperator->pDownstream[0];
G
Ganlin Zhao 已提交
1428
  if (downstream->operatorType == QUERY_NODE_PHYSICAL_PLAN_PARTITION ||
1429 1430 1431 1432 1433
      (downstream->operatorType == QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN &&
       ((STableScanInfo *)downstream->info)->hasGroupByTag == true)) {
    return TSDB_CODE_SUCCESS;
  }

1434 1435 1436 1437
  SqlFunctionCtx* pCtx = pOperator->exprSupp.pCtx;
  bool hasCountFunc = false;
  for (int32_t i = 0; i < pOperator->exprSupp.numOfExprs; ++i) {
    if ((strcmp(pCtx[i].pExpr->pExpr->_function.functionName, "count") == 0) ||
G
Ganlin Zhao 已提交
1438 1439 1440
        (strcmp(pCtx[i].pExpr->pExpr->_function.functionName, "hyperloglog") == 0) ||
        (strcmp(pCtx[i].pExpr->pExpr->_function.functionName, "_hyperloglog_partial") == 0) ||
        (strcmp(pCtx[i].pExpr->pExpr->_function.functionName, "_hyperloglog_merge") == 0)) {
1441 1442 1443 1444 1445 1446 1447 1448 1449 1450
      hasCountFunc = true;
      break;
    }
  }

  if (!hasCountFunc) {
    return TSDB_CODE_SUCCESS;
  }

  SSDataBlock* pBlock = createDataBlock();
G
Ganlin Zhao 已提交
1451
  pBlock->info.rows = 1;
1452 1453 1454
  pBlock->info.capacity = 0;

  for (int32_t i = 0; i < pOperator->exprSupp.numOfExprs; ++i) {
G
Ganlin Zhao 已提交
1455 1456 1457 1458
    SColumnInfoData colInfo = {0};
    colInfo.hasNull = true;
    colInfo.info.type = TSDB_DATA_TYPE_NULL;
    colInfo.info.bytes = 1;
1459 1460 1461 1462 1463 1464

    SExprInfo* pOneExpr = &pOperator->exprSupp.pExprInfo[i];
    for (int32_t j = 0; j < pOneExpr->base.numOfParams; ++j) {
      SFunctParam* pFuncParam = &pOneExpr->base.pParam[j];
      if (pFuncParam->type == FUNC_PARAM_TYPE_COLUMN) {
        int32_t slotId = pFuncParam->pCol->slotId;
G
Ganlin Zhao 已提交
1465 1466 1467 1468 1469 1470 1471
        int32_t numOfCols = taosArrayGetSize(pBlock->pDataBlock);
        if (slotId >= numOfCols) {
          taosArrayEnsureCap(pBlock->pDataBlock, slotId + 1);
          for (int32_t k = numOfCols; k < slotId + 1; ++k) {
            taosArrayPush(pBlock->pDataBlock, &colInfo);
          }
        }
1472
      } else if (pFuncParam->type == FUNC_PARAM_TYPE_VALUE) {
G
Ganlin Zhao 已提交
1473
        // do nothing
1474 1475 1476 1477
      }
    }
  }

G
Ganlin Zhao 已提交
1478
  blockDataEnsureCapacity(pBlock, pBlock->info.rows);
1479 1480 1481 1482 1483 1484 1485 1486 1487 1488 1489 1490 1491 1492 1493
  *ppBlock = pBlock;

  return TSDB_CODE_SUCCESS;
}

static void destroyDataBlockForEmptyInput(bool blockAllocated, SSDataBlock **ppBlock) {
  if (!blockAllocated) {
    return;
  }

  blockDataDestroy(*ppBlock);
  *ppBlock = NULL;
}


1494
// this is a blocking operator
L
Liu Jicong 已提交
1495
static int32_t doOpenAggregateOptr(SOperatorInfo* pOperator) {
H
Haojun Liao 已提交
1496 1497
  if (OPTR_IS_OPENED(pOperator)) {
    return TSDB_CODE_SUCCESS;
1498 1499
  }

H
Haojun Liao 已提交
1500
  SExecTaskInfo*    pTaskInfo = pOperator->pTaskInfo;
1501
  SAggOperatorInfo* pAggInfo = pOperator->info;
H
Haojun Liao 已提交
1502

1503 1504
  SExprSupp*     pSup = &pOperator->exprSupp;
  SOperatorInfo* downstream = pOperator->pDownstream[0];
1505

1506 1507
  int64_t st = taosGetTimestampUs();

1508 1509 1510
  int32_t order = TSDB_ORDER_ASC;
  int32_t scanFlag = MAIN_SCAN;

1511 1512 1513
  bool    hasValidBlock = false;
  bool    blockAllocated = false;

H
Haojun Liao 已提交
1514
  while (1) {
1515
    SSDataBlock* pBlock = downstream->fpSet.getNextFn(downstream);
1516
    if (pBlock == NULL) {
G
Ganlin Zhao 已提交
1517
      if (!hasValidBlock) {
1518 1519 1520 1521 1522 1523 1524 1525
        createDataBlockForEmptyInput(pOperator, &pBlock);
        if (pBlock == NULL) {
          break;
        }
        blockAllocated = true;
      } else {
        break;
      }
1526
    }
1527
    hasValidBlock = true;
1528

1529 1530
    int32_t code = getTableScanInfo(pOperator, &order, &scanFlag);
    if (code != TSDB_CODE_SUCCESS) {
1531
      destroyDataBlockForEmptyInput(blockAllocated, &pBlock);
1532
      T_LONG_JMP(pTaskInfo->env, code);
1533
    }
1534

1535
    // there is an scalar expression that needs to be calculated before apply the group aggregation.
G
Ganlin Zhao 已提交
1536
    if (pAggInfo->scalarExprSup.pExprInfo != NULL && !blockAllocated) {
1537 1538
      SExprSupp* pSup1 = &pAggInfo->scalarExprSup;
      code = projectApplyFunctions(pSup1->pExprInfo, pBlock, pBlock, pSup1->pCtx, pSup1->numOfExprs, NULL);
1539
      if (code != TSDB_CODE_SUCCESS) {
1540
        destroyDataBlockForEmptyInput(blockAllocated, &pBlock);
1541
        T_LONG_JMP(pTaskInfo->env, code);
1542
      }
1543 1544
    }

1545
    // the pDataBlock are always the same one, no need to call this again
H
Haojun Liao 已提交
1546
    setExecutionContext(pOperator, pOperator->exprSupp.numOfExprs, pBlock->info.id.groupId);
1547
    setInputDataBlock(pSup, pBlock, order, scanFlag, true);
1548
    code = doAggregateImpl(pOperator, pSup->pCtx);
1549
    if (code != 0) {
1550
      destroyDataBlockForEmptyInput(blockAllocated, &pBlock);
1551
      T_LONG_JMP(pTaskInfo->env, code);
1552
    }
1553 1554 1555

    destroyDataBlockForEmptyInput(blockAllocated, &pBlock);

1556 1557
  }

1558 1559 1560 1561 1562
  // the downstream operator may return with error code, so let's check the code before generating results.
  if (pTaskInfo->code != TSDB_CODE_SUCCESS) {
    T_LONG_JMP(pTaskInfo->env, pTaskInfo->code);
  }

1563
  initGroupedResultInfo(&pAggInfo->groupResInfo, pAggInfo->aggSup.pResultRowHashTable, 0);
H
Haojun Liao 已提交
1564
  OPTR_SET_OPENED(pOperator);
1565

1566
  pOperator->cost.openCost = (taosGetTimestampUs() - st) / 1000.0;
1567
  return pTaskInfo->code;
H
Haojun Liao 已提交
1568 1569
}

1570
static SSDataBlock* getAggregateResult(SOperatorInfo* pOperator) {
L
Liu Jicong 已提交
1571
  SAggOperatorInfo* pAggInfo = pOperator->info;
H
Haojun Liao 已提交
1572 1573 1574 1575 1576 1577
  SOptrBasicInfo*   pInfo = &pAggInfo->binfo;

  if (pOperator->status == OP_EXEC_DONE) {
    return NULL;
  }

L
Liu Jicong 已提交
1578
  SExecTaskInfo* pTaskInfo = pOperator->pTaskInfo;
1579
  pTaskInfo->code = pOperator->fpSet._openFn(pOperator);
H
Haojun Liao 已提交
1580
  if (pTaskInfo->code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
1581
    setOperatorCompleted(pOperator);
H
Haojun Liao 已提交
1582 1583 1584
    return NULL;
  }

H
Haojun Liao 已提交
1585
  blockDataEnsureCapacity(pInfo->pRes, pOperator->resultInfo.capacity);
S
slzhou 已提交
1586 1587
  while (1) {
    doBuildResultDatablock(pOperator, pInfo, &pAggInfo->groupResInfo, pAggInfo->aggSup.pResultBuf);
H
Haojun Liao 已提交
1588
    doFilter(pInfo->pRes, pOperator->exprSupp.pFilterInfo, NULL);
S
slzhou 已提交
1589

1590
    if (!hasRemainResults(&pAggInfo->groupResInfo)) {
H
Haojun Liao 已提交
1591
      setOperatorCompleted(pOperator);
S
slzhou 已提交
1592 1593
      break;
    }
1594

S
slzhou 已提交
1595 1596 1597 1598
    if (pInfo->pRes->info.rows > 0) {
      break;
    }
  }
1599

1600
  size_t rows = blockDataGetNumOfRows(pInfo->pRes);
1601 1602
  pOperator->resultInfo.totalRows += rows;

1603
  return (rows == 0) ? NULL : pInfo->pRes;
1604 1605
}

L
Liu Jicong 已提交
1606 1607
static void doHandleRemainBlockForNewGroupImpl(SOperatorInfo* pOperator, SFillOperatorInfo* pInfo,
                                               SResultInfo* pResultInfo, SExecTaskInfo* pTaskInfo) {
1608
  pInfo->totalInputRows = pInfo->existNewGroupBlock->info.rows;
1609 1610 1611 1612 1613
  SSDataBlock* pResBlock = pInfo->pFinalRes;

  int32_t order = TSDB_ORDER_ASC;
  int32_t scanFlag = MAIN_SCAN;
  getTableScanInfo(pOperator, &order, &scanFlag);
H
Haojun Liao 已提交
1614

5
54liuyao 已提交
1615
  int64_t ekey = pInfo->existNewGroupBlock->info.window.ekey;
1616 1617
  taosResetFillInfo(pInfo->pFillInfo, getFillInfoStart(pInfo->pFillInfo));

1618
  blockDataCleanup(pInfo->pRes);
1619 1620 1621 1622
  doApplyScalarCalculation(pOperator, pInfo->existNewGroupBlock, order, scanFlag);

  taosFillSetStartInfo(pInfo->pFillInfo, pInfo->pRes->info.rows, ekey);
  taosFillSetInputDataBlock(pInfo->pFillInfo, pInfo->pRes);
1623

1624 1625
  int32_t numOfResultRows = pResultInfo->capacity - pResBlock->info.rows;
  taosFillResultDataBlock(pInfo->pFillInfo, pResBlock, numOfResultRows);
H
Haojun Liao 已提交
1626

H
Haojun Liao 已提交
1627
  pInfo->curGroupId = pInfo->existNewGroupBlock->info.id.groupId;
1628 1629 1630
  pInfo->existNewGroupBlock = NULL;
}

L
Liu Jicong 已提交
1631 1632
static void doHandleRemainBlockFromNewGroup(SOperatorInfo* pOperator, SFillOperatorInfo* pInfo,
                                            SResultInfo* pResultInfo, SExecTaskInfo* pTaskInfo) {
1633
  if (taosFillHasMoreResults(pInfo->pFillInfo)) {
H
Haojun Liao 已提交
1634 1635
    int32_t numOfResultRows = pResultInfo->capacity - pInfo->pFinalRes->info.rows;
    taosFillResultDataBlock(pInfo->pFillInfo, pInfo->pFinalRes, numOfResultRows);
H
Haojun Liao 已提交
1636
    pInfo->pRes->info.id.groupId = pInfo->curGroupId;
1637
    return;
1638 1639 1640 1641
  }

  // handle the cached new group data block
  if (pInfo->existNewGroupBlock) {
1642 1643 1644 1645 1646 1647
    doHandleRemainBlockForNewGroupImpl(pOperator, pInfo, pResultInfo, pTaskInfo);
  }
}

static void doApplyScalarCalculation(SOperatorInfo* pOperator, SSDataBlock* pBlock, int32_t order, int32_t scanFlag) {
  SFillOperatorInfo* pInfo = pOperator->info;
L
Liu Jicong 已提交
1648
  SExprSupp*         pSup = &pOperator->exprSupp;
1649
  setInputDataBlock(pSup, pBlock, order, scanFlag, false);
1650 1651
  projectApplyFunctions(pSup->pExprInfo, pInfo->pRes, pBlock, pSup->pCtx, pSup->numOfExprs, NULL);

1652 1653 1654 1655
  // reset the row value before applying the no-fill functions to the input data block, which is "pBlock" in this case.
  pInfo->pRes->info.rows = 0;
  SExprSupp* pNoFillSupp = &pInfo->noFillExprSupp;
  setInputDataBlock(pNoFillSupp, pBlock, order, scanFlag, false);
1656

1657
  projectApplyFunctions(pNoFillSupp->pExprInfo, pInfo->pRes, pBlock, pNoFillSupp->pCtx, pNoFillSupp->numOfExprs, NULL);
H
Haojun Liao 已提交
1658
  pInfo->pRes->info.id.groupId = pBlock->info.id.groupId;
1659 1660
}

S
slzhou 已提交
1661
static SSDataBlock* doFillImpl(SOperatorInfo* pOperator) {
L
Liu Jicong 已提交
1662 1663
  SFillOperatorInfo* pInfo = pOperator->info;
  SExecTaskInfo*     pTaskInfo = pOperator->pTaskInfo;
1664

H
Haojun Liao 已提交
1665
  SResultInfo* pResultInfo = &pOperator->resultInfo;
H
Haojun Liao 已提交
1666
  SSDataBlock* pResBlock = pInfo->pFinalRes;
1667 1668

  blockDataCleanup(pResBlock);
1669

H
Haojun Liao 已提交
1670 1671
  int32_t order = TSDB_ORDER_ASC;
  int32_t scanFlag = MAIN_SCAN;
1672
  getTableScanInfo(pOperator, &order, &scanFlag);
1673

1674
  doHandleRemainBlockFromNewGroup(pOperator, pInfo, pResultInfo, pTaskInfo);
1675
  if (pResBlock->info.rows > 0) {
H
Haojun Liao 已提交
1676
    pResBlock->info.id.groupId = pInfo->curGroupId;
1677
    return pResBlock;
H
Haojun Liao 已提交
1678
  }
1679

H
Haojun Liao 已提交
1680
  SOperatorInfo* pDownstream = pOperator->pDownstream[0];
L
Liu Jicong 已提交
1681
  while (1) {
1682
    SSDataBlock* pBlock = pDownstream->fpSet.getNextFn(pDownstream);
1683 1684
    if (pBlock == NULL) {
      if (pInfo->totalInputRows == 0) {
H
Haojun Liao 已提交
1685
        setOperatorCompleted(pOperator);
1686 1687
        return NULL;
      }
1688

1689
      taosFillSetStartInfo(pInfo->pFillInfo, 0, pInfo->win.ekey);
1690
    } else {
1691
      blockDataUpdateTsWindow(pBlock, pInfo->primarySrcSlotId);
1692 1693

      blockDataCleanup(pInfo->pRes);
1694 1695
      blockDataEnsureCapacity(pInfo->pRes, pBlock->info.rows);
      blockDataEnsureCapacity(pInfo->pFinalRes, pBlock->info.rows);
1696
      doApplyScalarCalculation(pOperator, pBlock, order, scanFlag);
1697

H
Haojun Liao 已提交
1698 1699
      if (pInfo->curGroupId == 0 || pInfo->curGroupId == pInfo->pRes->info.id.groupId) {
        pInfo->curGroupId = pInfo->pRes->info.id.groupId;  // the first data block
H
Haojun Liao 已提交
1700
        pInfo->totalInputRows += pInfo->pRes->info.rows;
1701

H
Haojun Liao 已提交
1702 1703 1704 1705 1706
        if (order == pInfo->pFillInfo->order) {
          taosFillSetStartInfo(pInfo->pFillInfo, pInfo->pRes->info.rows, pBlock->info.window.ekey);
        } else {
          taosFillSetStartInfo(pInfo->pFillInfo, pInfo->pRes->info.rows, pBlock->info.window.skey);
        }
H
Haojun Liao 已提交
1707
        taosFillSetInputDataBlock(pInfo->pFillInfo, pInfo->pRes);
H
Haojun Liao 已提交
1708
      } else if (pInfo->curGroupId != pBlock->info.id.groupId) {  // the new group data block
1709 1710 1711 1712 1713
        pInfo->existNewGroupBlock = pBlock;

        // Fill the previous group data block, before handle the data block of new group.
        // Close the fill operation for previous group data block
        taosFillSetStartInfo(pInfo->pFillInfo, 0, pInfo->win.ekey);
1714 1715 1716
      }
    }

1717 1718
    int32_t numOfResultRows = pOperator->resultInfo.capacity - pResBlock->info.rows;
    taosFillResultDataBlock(pInfo->pFillInfo, pResBlock, numOfResultRows);
1719 1720

    // current group has no more result to return
1721
    if (pResBlock->info.rows > 0) {
1722 1723
      // 1. The result in current group not reach the threshold of output result, continue
      // 2. If multiple group results existing in one SSDataBlock is not allowed, return immediately
1724
      if (pResBlock->info.rows > pResultInfo->threshold || pBlock == NULL || pInfo->existNewGroupBlock != NULL) {
H
Haojun Liao 已提交
1725
        pResBlock->info.id.groupId = pInfo->curGroupId;
1726
        return pResBlock;
1727 1728
      }

1729
      doHandleRemainBlockFromNewGroup(pOperator, pInfo, pResultInfo, pTaskInfo);
1730
      if (pResBlock->info.rows >= pOperator->resultInfo.threshold || pBlock == NULL) {
H
Haojun Liao 已提交
1731
        pResBlock->info.id.groupId = pInfo->curGroupId;
1732
        return pResBlock;
1733 1734 1735
      }
    } else if (pInfo->existNewGroupBlock) {  // try next group
      assert(pBlock != NULL);
1736 1737 1738 1739

      blockDataCleanup(pResBlock);

      doHandleRemainBlockForNewGroupImpl(pOperator, pInfo, pResultInfo, pTaskInfo);
1740
      if (pResBlock->info.rows > pResultInfo->threshold) {
H
Haojun Liao 已提交
1741
        pResBlock->info.id.groupId = pInfo->curGroupId;
1742
        return pResBlock;
1743 1744 1745 1746 1747 1748 1749
      }
    } else {
      return NULL;
    }
  }
}

S
slzhou 已提交
1750 1751 1752 1753 1754 1755 1756 1757
static SSDataBlock* doFill(SOperatorInfo* pOperator) {
  SFillOperatorInfo* pInfo = pOperator->info;
  SExecTaskInfo*     pTaskInfo = pOperator->pTaskInfo;

  if (pOperator->status == OP_EXEC_DONE) {
    return NULL;
  }

S
slzhou 已提交
1758
  SSDataBlock* fillResult = NULL;
S
slzhou 已提交
1759
  while (true) {
S
slzhou 已提交
1760
    fillResult = doFillImpl(pOperator);
S
slzhou 已提交
1761
    if (fillResult == NULL) {
H
Haojun Liao 已提交
1762
      setOperatorCompleted(pOperator);
S
slzhou 已提交
1763 1764 1765
      break;
    }

1766
    doFilter(fillResult, pOperator->exprSupp.pFilterInfo, &pInfo->matchInfo);
S
slzhou 已提交
1767 1768 1769 1770 1771
    if (fillResult->info.rows > 0) {
      break;
    }
  }

S
slzhou 已提交
1772
  if (fillResult != NULL) {
1773
    pOperator->resultInfo.totalRows += fillResult->info.rows;
S
slzhou 已提交
1774
  }
S
slzhou 已提交
1775

S
slzhou 已提交
1776
  return fillResult;
S
slzhou 已提交
1777 1778
}

1779
void destroyExprInfo(SExprInfo* pExpr, int32_t numOfExprs) {
C
Cary Xu 已提交
1780 1781 1782 1783 1784
  for (int32_t i = 0; i < numOfExprs; ++i) {
    SExprInfo* pExprInfo = &pExpr[i];
    for (int32_t j = 0; j < pExprInfo->base.numOfParams; ++j) {
      if (pExprInfo->base.pParam[j].type == FUNC_PARAM_TYPE_COLUMN) {
        taosMemoryFreeClear(pExprInfo->base.pParam[j].pCol);
5
54liuyao 已提交
1785 1786
      } else if (pExprInfo->base.pParam[j].type == FUNC_PARAM_TYPE_VALUE) {
        taosVariantDestroy(&pExprInfo->base.pParam[j].param);
H
Haojun Liao 已提交
1787
      }
1788
    }
C
Cary Xu 已提交
1789 1790 1791

    taosMemoryFree(pExprInfo->base.pParam);
    taosMemoryFree(pExprInfo->pExpr);
H
Haojun Liao 已提交
1792 1793 1794
  }
}

5
54liuyao 已提交
1795
void destroyOperatorInfo(SOperatorInfo* pOperator) {
1796 1797 1798 1799
  if (pOperator == NULL) {
    return;
  }

1800
  if (pOperator->fpSet.closeFn != NULL) {
1801
    pOperator->fpSet.closeFn(pOperator->info);
1802 1803
  }

H
Haojun Liao 已提交
1804
  if (pOperator->pDownstream != NULL) {
L
Liu Jicong 已提交
1805
    for (int32_t i = 0; i < pOperator->numOfDownstream; ++i) {
H
Haojun Liao 已提交
1806
      destroyOperatorInfo(pOperator->pDownstream[i]);
1807 1808
    }

wafwerar's avatar
wafwerar 已提交
1809
    taosMemoryFreeClear(pOperator->pDownstream);
H
Haojun Liao 已提交
1810
    pOperator->numOfDownstream = 0;
1811 1812
  }

1813
  cleanupExprSupp(&pOperator->exprSupp);
wafwerar's avatar
wafwerar 已提交
1814
  taosMemoryFreeClear(pOperator);
1815 1816
}

1817 1818 1819 1820 1821 1822
int32_t getBufferPgSize(int32_t rowSize, uint32_t* defaultPgsz, uint32_t* defaultBufsz) {
  *defaultPgsz = 4096;
  while (*defaultPgsz < rowSize * 4) {
    *defaultPgsz <<= 1u;
  }

1823
  // The default buffer for each operator in query is 10MB.
1824
  // at least four pages need to be in buffer
1825 1826
  // TODO: make this variable to be configurable.
  *defaultBufsz = 4096 * 2560;
1827 1828 1829 1830 1831 1832 1833
  if ((*defaultBufsz) <= (*defaultPgsz)) {
    (*defaultBufsz) = (*defaultPgsz) * 4;
  }

  return 0;
}

dengyihao's avatar
dengyihao 已提交
1834 1835
int32_t doInitAggInfoSup(SAggSupporter* pAggSup, SqlFunctionCtx* pCtx, int32_t numOfOutput, size_t keyBufSize,
                         const char* pKey) {
L
Liu Jicong 已提交
1836
  int32_t    code = 0;
1837 1838
  _hash_fn_t hashFn = taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY);

1839
  pAggSup->currentPageId = -1;
dengyihao's avatar
dengyihao 已提交
1840 1841
  pAggSup->resultRowSize = getResultRowSize(pCtx, numOfOutput);
  pAggSup->keyBuf = taosMemoryCalloc(1, keyBufSize + POINTER_BYTES + sizeof(int64_t));
1842
  pAggSup->pResultRowHashTable = tSimpleHashInit(10, hashFn);
1843

H
Haojun Liao 已提交
1844
  if (pAggSup->keyBuf == NULL || pAggSup->pResultRowHashTable == NULL) {
1845 1846 1847
    return TSDB_CODE_OUT_OF_MEMORY;
  }

dengyihao's avatar
dengyihao 已提交
1848
  uint32_t defaultPgsz = 0;
1849 1850
  uint32_t defaultBufsz = 0;
  getBufferPgSize(pAggSup->resultRowSize, &defaultPgsz, &defaultBufsz);
H
Haojun Liao 已提交
1851

wafwerar's avatar
wafwerar 已提交
1852
  if (!osTempSpaceAvailable()) {
H
Haojun Liao 已提交
1853 1854 1855
    code = TSDB_CODE_NO_AVAIL_DISK;
    qError("Init stream agg supporter failed since %s, %s", terrstr(code), pKey);
    return code;
wafwerar's avatar
wafwerar 已提交
1856
  }
1857

H
Haojun Liao 已提交
1858
  code = createDiskbasedBuf(&pAggSup->pResultBuf, defaultPgsz, defaultBufsz, pKey, tsTempDir);
H
Haojun Liao 已提交
1859
  if (code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
1860
    qError("Create agg result buf failed since %s, %s", tstrerror(code), pKey);
H
Haojun Liao 已提交
1861 1862 1863
    return code;
  }

H
Haojun Liao 已提交
1864
  return code;
1865 1866
}

1867
void cleanupAggSup(SAggSupporter* pAggSup) {
wafwerar's avatar
wafwerar 已提交
1868
  taosMemoryFreeClear(pAggSup->keyBuf);
1869
  tSimpleHashCleanup(pAggSup->pResultRowHashTable);
H
Haojun Liao 已提交
1870
  destroyDiskbasedBuf(pAggSup->pResultBuf);
1871 1872
}

H
Haojun Liao 已提交
1873
int32_t initAggSup(SExprSupp* pSup, SAggSupporter* pAggSup, SExprInfo* pExprInfo, int32_t numOfCols, size_t keyBufSize,
L
Liu Jicong 已提交
1874
                    const char* pkey) {
1875 1876 1877 1878 1879
  int32_t code = initExprSupp(pSup, pExprInfo, numOfCols);
  if (code != TSDB_CODE_SUCCESS) {
    return code;
  }

1880 1881 1882 1883 1884
  code = doInitAggInfoSup(pAggSup, pSup->pCtx, numOfCols, keyBufSize, pkey);
  if (code != TSDB_CODE_SUCCESS) {
    return code;
  }

L
Liu Jicong 已提交
1885
  for (int32_t i = 0; i < numOfCols; ++i) {
1886
    pSup->pCtx[i].saveHandle.pBuf = pAggSup->pResultBuf;
1887 1888
  }

1889
  return TSDB_CODE_SUCCESS;
H
Haojun Liao 已提交
1890 1891
}

L
Liu Jicong 已提交
1892
void initResultSizeInfo(SResultInfo* pResultInfo, int32_t numOfRows) {
wmmhello's avatar
wmmhello 已提交
1893
  ASSERT(numOfRows != 0);
1894 1895
  pResultInfo->capacity = numOfRows;
  pResultInfo->threshold = numOfRows * 0.75;
1896

1897 1898
  if (pResultInfo->threshold == 0) {
    pResultInfo->threshold = numOfRows;
1899 1900 1901
  }
}

1902 1903 1904 1905 1906
void initBasicInfo(SOptrBasicInfo* pInfo, SSDataBlock* pBlock) {
  pInfo->pRes = pBlock;
  initResultRowInfo(&pInfo->resultRowInfo);
}

5
54liuyao 已提交
1907
void* destroySqlFunctionCtx(SqlFunctionCtx* pCtx, int32_t numOfOutput) {
1908 1909 1910 1911 1912 1913 1914 1915 1916 1917
  if (pCtx == NULL) {
    return NULL;
  }

  for (int32_t i = 0; i < numOfOutput; ++i) {
    for (int32_t j = 0; j < pCtx[i].numOfParams; ++j) {
      taosVariantDestroy(&pCtx[i].param[j].param);
    }

    taosMemoryFreeClear(pCtx[i].subsidiaries.pCtx);
1918
    taosMemoryFreeClear(pCtx[i].subsidiaries.buf);
1919 1920 1921 1922 1923 1924 1925 1926
    taosMemoryFree(pCtx[i].input.pData);
    taosMemoryFree(pCtx[i].input.pColumnDataAgg);
  }

  taosMemoryFreeClear(pCtx);
  return NULL;
}

1927
int32_t initExprSupp(SExprSupp* pSup, SExprInfo* pExprInfo, int32_t numOfExpr) {
1928 1929 1930 1931
  pSup->pExprInfo = pExprInfo;
  pSup->numOfExprs = numOfExpr;
  if (pSup->pExprInfo != NULL) {
    pSup->pCtx = createSqlFunctionCtx(pExprInfo, numOfExpr, &pSup->rowEntryInfoOffset);
1932 1933 1934
    if (pSup->pCtx == NULL) {
      return TSDB_CODE_OUT_OF_MEMORY;
    }
1935
  }
1936 1937

  return TSDB_CODE_SUCCESS;
1938 1939
}

1940 1941 1942 1943
void cleanupExprSupp(SExprSupp* pSupp) {
  destroySqlFunctionCtx(pSupp->pCtx, pSupp->numOfExprs);
  if (pSupp->pExprInfo != NULL) {
    destroyExprInfo(pSupp->pExprInfo, pSupp->numOfExprs);
C
Cary Xu 已提交
1944
    taosMemoryFreeClear(pSupp->pExprInfo);
1945
  }
H
Haojun Liao 已提交
1946 1947 1948 1949 1950 1951

  if (pSupp->pFilterInfo != NULL) {
    filterFreeInfo(pSupp->pFilterInfo);
    pSupp->pFilterInfo = NULL;
  }

1952 1953 1954
  taosMemoryFree(pSupp->rowEntryInfoOffset);
}

1955 1956
SOperatorInfo* createAggregateOperatorInfo(SOperatorInfo* downstream, SAggPhysiNode* pAggNode,
                                           SExecTaskInfo* pTaskInfo) {
wafwerar's avatar
wafwerar 已提交
1957
  SAggOperatorInfo* pInfo = taosMemoryCalloc(1, sizeof(SAggOperatorInfo));
L
Liu Jicong 已提交
1958
  SOperatorInfo*    pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
H
Haojun Liao 已提交
1959 1960 1961
  if (pInfo == NULL || pOperator == NULL) {
    goto _error;
  }
H
Haojun Liao 已提交
1962

H
Haojun Liao 已提交
1963
  SSDataBlock* pResBlock = createDataBlockFromDescNode(pAggNode->node.pOutputDataBlockDesc);
H
Haojun Liao 已提交
1964 1965 1966
  initBasicInfo(&pInfo->binfo, pResBlock);

  size_t keyBufSize = sizeof(int64_t) + sizeof(int64_t) + POINTER_BYTES;
1967
  initResultSizeInfo(&pOperator->resultInfo, 4096);
H
Haojun Liao 已提交
1968

1969 1970
  int32_t    num = 0;
  SExprInfo* pExprInfo = createExprInfo(pAggNode->pAggFuncs, pAggNode->pGroupKeys, &num);
H
Haojun Liao 已提交
1971
  int32_t    code = initAggSup(&pOperator->exprSupp, &pInfo->aggSup, pExprInfo, num, keyBufSize, pTaskInfo->id.str);
L
Liu Jicong 已提交
1972
  if (code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
1973 1974
    goto _error;
  }
H
Haojun Liao 已提交
1975

H
Haojun Liao 已提交
1976 1977 1978 1979 1980 1981
  int32_t    numOfScalarExpr = 0;
  SExprInfo* pScalarExprInfo = NULL;
  if (pAggNode->pExprs != NULL) {
    pScalarExprInfo = createExprInfo(pAggNode->pExprs, NULL, &numOfScalarExpr);
  }

1982 1983 1984 1985
  code = initExprSupp(&pInfo->scalarExprSup, pScalarExprInfo, numOfScalarExpr);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
1986

H
Haojun Liao 已提交
1987 1988 1989 1990 1991
  code = filterInitFromNode((SNode*)pAggNode->node.pConditions, &pOperator->exprSupp.pFilterInfo, 0);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }

H
Haojun Liao 已提交
1992
  pInfo->binfo.mergeResultBlock = pAggNode->mergeDataBlock;
1993
  pInfo->groupId = UINT64_MAX;
H
Haojun Liao 已提交
1994

1995 1996 1997
  setOperatorInfo(pOperator, "TableAggregate", QUERY_NODE_PHYSICAL_PLAN_HASH_AGG, true, OP_NOT_OPENED, pInfo,
                  pTaskInfo);
  pOperator->fpSet = createOperatorFpSet(doOpenAggregateOptr, getAggregateResult, NULL, destroyAggOperatorInfo, NULL);
H
Haojun Liao 已提交
1998

1999 2000
  if (downstream->operatorType == QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN) {
    STableScanInfo* pTableScanInfo = downstream->info;
H
Haojun Liao 已提交
2001 2002
    pTableScanInfo->base.pdInfo.pExprSup = &pOperator->exprSupp;
    pTableScanInfo->base.pdInfo.pAggSup = &pInfo->aggSup;
2003 2004
  }

H
Haojun Liao 已提交
2005 2006 2007 2008
  code = appendDownstream(pOperator, &downstream, 1);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
2009 2010

  return pOperator;
H
Haojun Liao 已提交
2011

2012
_error:
H
Haojun Liao 已提交
2013 2014 2015 2016
  if (pInfo != NULL) {
    destroyAggOperatorInfo(pInfo);
  }

2017 2018 2019
  if (pOperator != NULL) {
    cleanupExprSupp(&pOperator->exprSupp);
  }
H
Haojun Liao 已提交
2020

2021
  taosMemoryFreeClear(pOperator);
H
Haojun Liao 已提交
2022
  pTaskInfo->code = code;
H
Haojun Liao 已提交
2023
  return NULL;
2024 2025
}

2026
void cleanupBasicInfo(SOptrBasicInfo* pInfo) {
2027
  assert(pInfo != NULL);
H
Haojun Liao 已提交
2028
  pInfo->pRes = blockDataDestroy(pInfo->pRes);
2029 2030
}

H
Haojun Liao 已提交
2031 2032 2033 2034 2035 2036 2037
static void freeItem(void* pItem) {
  void** p = pItem;
  if (*p != NULL) {
    taosMemoryFreeClear(*p);
  }
}

2038
void destroyAggOperatorInfo(void* param) {
L
Liu Jicong 已提交
2039
  SAggOperatorInfo* pInfo = (SAggOperatorInfo*)param;
L
Liu Jicong 已提交
2040 2041
  cleanupBasicInfo(&pInfo->binfo);

H
Haojun Liao 已提交
2042
  cleanupAggSup(&pInfo->aggSup);
S
shenglian zhou 已提交
2043
  cleanupExprSupp(&pInfo->scalarExprSup);
H
Haojun Liao 已提交
2044
  cleanupGroupResInfo(&pInfo->groupResInfo);
D
dapan1121 已提交
2045
  taosMemoryFreeClear(param);
2046
}
2047

2048
void destroyFillOperatorInfo(void* param) {
L
Liu Jicong 已提交
2049
  SFillOperatorInfo* pInfo = (SFillOperatorInfo*)param;
2050
  pInfo->pFillInfo = taosDestroyFillInfo(pInfo->pFillInfo);
H
Haojun Liao 已提交
2051
  pInfo->pRes = blockDataDestroy(pInfo->pRes);
H
Haojun Liao 已提交
2052 2053
  pInfo->pFinalRes = blockDataDestroy(pInfo->pFinalRes);

2054
  cleanupExprSupp(&pInfo->noFillExprSupp);
H
Haojun Liao 已提交
2055

wafwerar's avatar
wafwerar 已提交
2056
  taosMemoryFreeClear(pInfo->p);
H
Haojun Liao 已提交
2057
  taosArrayDestroy(pInfo->matchInfo.pList);
D
dapan1121 已提交
2058
  taosMemoryFreeClear(param);
2059 2060
}

H
Haojun Liao 已提交
2061 2062 2063 2064
static int32_t initFillInfo(SFillOperatorInfo* pInfo, SExprInfo* pExpr, int32_t numOfCols, SExprInfo* pNotFillExpr,
                            int32_t numOfNotFillCols, SNodeListNode* pValNode, STimeWindow win, int32_t capacity,
                            const char* id, SInterval* pInterval, int32_t fillType, int32_t order) {
  SFillColInfo* pColInfo = createFillColInfo(pExpr, numOfCols, pNotFillExpr, numOfNotFillCols, pValNode);
H
Haojun Liao 已提交
2065

H
Haojun Liao 已提交
2066 2067 2068
  int64_t     startKey = (order == TSDB_ORDER_ASC) ? win.skey : win.ekey;
  STimeWindow w = getAlignQueryTimeWindow(pInterval, pInterval->precision, startKey);
  w = getFirstQualifiedTimeWindow(startKey, &w, pInterval, order);
H
Haojun Liao 已提交
2069

L
Liu Jicong 已提交
2070 2071
  pInfo->pFillInfo = taosCreateFillInfo(w.skey, numOfCols, numOfNotFillCols, capacity, pInterval, fillType, pColInfo,
                                        pInfo->primaryTsCol, order, id);
H
Haojun Liao 已提交
2072

H
Haojun Liao 已提交
2073 2074 2075 2076 2077 2078 2079
  if (order == TSDB_ORDER_ASC) {
    pInfo->win.skey = win.skey;
    pInfo->win.ekey = win.ekey;
  } else {
    pInfo->win.skey = win.ekey;
    pInfo->win.ekey = win.skey;
  }
L
Liu Jicong 已提交
2080
  pInfo->p = taosMemoryCalloc(numOfCols, POINTER_BYTES);
2081

H
Haojun Liao 已提交
2082
  if (pInfo->pFillInfo == NULL || pInfo->p == NULL) {
H
Haojun Liao 已提交
2083 2084
    taosMemoryFree(pInfo->pFillInfo);
    taosMemoryFree(pInfo->p);
H
Haojun Liao 已提交
2085 2086 2087 2088 2089 2090
    return TSDB_CODE_OUT_OF_MEMORY;
  } else {
    return TSDB_CODE_SUCCESS;
  }
}

2091
static bool isWstartColumnExist(SFillOperatorInfo* pInfo) {
2092
  if (pInfo->noFillExprSupp.numOfExprs == 0) {
2093 2094
    return false;
  }
2095 2096 2097

  for (int32_t i = 0; i < pInfo->noFillExprSupp.numOfExprs; ++i) {
    SExprInfo* exprInfo = pInfo->noFillExprSupp.pExprInfo + i;
2098
    if (exprInfo->pExpr->nodeType == QUERY_NODE_COLUMN && exprInfo->base.numOfParams == 1 &&
2099 2100 2101 2102 2103 2104 2105
        exprInfo->base.pParam[0].pCol->colType == COLUMN_TYPE_WINDOW_START) {
      return true;
    }
  }
  return false;
}

2106 2107
static int32_t createPrimaryTsExprIfNeeded(SFillOperatorInfo* pInfo, SFillPhysiNode* pPhyFillNode, SExprSupp* pExprSupp,
                                           const char* idStr) {
2108
  bool wstartExist = isWstartColumnExist(pInfo);
2109

2110 2111
  if (wstartExist == false) {
    if (pPhyFillNode->pWStartTs->type != QUERY_NODE_TARGET) {
2112
      qError("pWStartTs of fill physical node is not a target node, %s", idStr);
2113 2114 2115
      return TSDB_CODE_QRY_SYS_ERROR;
    }

2116 2117
    SExprInfo* pExpr = taosMemoryRealloc(pExprSupp->pExprInfo, (pExprSupp->numOfExprs + 1) * sizeof(SExprInfo));
    if (pExpr == NULL) {
2118 2119 2120
      return TSDB_CODE_OUT_OF_MEMORY;
    }

2121 2122 2123
    createExprFromTargetNode(&pExpr[pExprSupp->numOfExprs], (STargetNode*)pPhyFillNode->pWStartTs);
    pExprSupp->numOfExprs += 1;
    pExprSupp->pExprInfo = pExpr;
2124
  }
2125

2126 2127 2128
  return TSDB_CODE_SUCCESS;
}

L
Liu Jicong 已提交
2129 2130
SOperatorInfo* createFillOperatorInfo(SOperatorInfo* downstream, SFillPhysiNode* pPhyFillNode,
                                      SExecTaskInfo* pTaskInfo) {
2131 2132 2133 2134 2135 2136
  SFillOperatorInfo* pInfo = taosMemoryCalloc(1, sizeof(SFillOperatorInfo));
  SOperatorInfo*     pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
  if (pInfo == NULL || pOperator == NULL) {
    goto _error;
  }

H
Haojun Liao 已提交
2137
  pInfo->pRes = createDataBlockFromDescNode(pPhyFillNode->node.pOutputDataBlockDesc);
H
Haojun Liao 已提交
2138 2139 2140
  SExprInfo* pExprInfo = createExprInfo(pPhyFillNode->pFillExprs, NULL, &pInfo->numOfExpr);
  pOperator->exprSupp.pExprInfo = pExprInfo;

2141 2142 2143 2144 2145 2146 2147 2148
  SExprSupp* pNoFillSupp = &pInfo->noFillExprSupp;
  pNoFillSupp->pExprInfo = createExprInfo(pPhyFillNode->pNotFillExprs, NULL, &pNoFillSupp->numOfExprs);
  int32_t code = createPrimaryTsExprIfNeeded(pInfo, pPhyFillNode, pNoFillSupp, pTaskInfo->id.str);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }

  code = initExprSupp(pNoFillSupp, pNoFillSupp->pExprInfo, pNoFillSupp->numOfExprs);
2149 2150 2151
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
H
Haojun Liao 已提交
2152

L
Liu Jicong 已提交
2153
  SInterval* pInterval =
2154
      QUERY_NODE_PHYSICAL_PLAN_MERGE_ALIGNED_INTERVAL == downstream->operatorType
2155 2156
          ? &((SMergeAlignedIntervalAggOperatorInfo*)downstream->info)->intervalAggOperatorInfo->interval
          : &((SIntervalAggOperatorInfo*)downstream->info)->interval;
2157

2158
  int32_t order = (pPhyFillNode->inputTsOrder == ORDER_ASC) ? TSDB_ORDER_ASC : TSDB_ORDER_DESC;
2159
  int32_t type = convertFillType(pPhyFillNode->mode);
2160

H
Haojun Liao 已提交
2161
  SResultInfo* pResultInfo = &pOperator->resultInfo;
H
Haojun Liao 已提交
2162

2163
  initResultSizeInfo(&pOperator->resultInfo, 4096);
H
Haojun Liao 已提交
2164 2165 2166 2167 2168
  blockDataEnsureCapacity(pInfo->pRes, pOperator->resultInfo.capacity);
  code = initExprSupp(&pOperator->exprSupp, pExprInfo, pInfo->numOfExpr);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
H
Haojun Liao 已提交
2169

H
Haojun Liao 已提交
2170 2171
  pInfo->primaryTsCol = ((STargetNode*)pPhyFillNode->pWStartTs)->slotId;
  pInfo->primarySrcSlotId = ((SColumnNode*)((STargetNode*)pPhyFillNode->pWStartTs)->pExpr)->slotId;
2172

2173
  int32_t numOfOutputCols = 0;
H
Haojun Liao 已提交
2174 2175
  code = extractColMatchInfo(pPhyFillNode->pFillExprs, pPhyFillNode->node.pOutputDataBlockDesc, &numOfOutputCols,
                             COL_MATCH_FROM_SLOT_ID, &pInfo->matchInfo);
2176

2177
  code = initFillInfo(pInfo, pExprInfo, pInfo->numOfExpr, pNoFillSupp->pExprInfo, pNoFillSupp->numOfExprs,
2178 2179
                      (SNodeListNode*)pPhyFillNode->pValues, pPhyFillNode->timeRange, pResultInfo->capacity,
                      pTaskInfo->id.str, pInterval, type, order);
2180 2181 2182
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
2183

H
Haojun Liao 已提交
2184
  pInfo->pFinalRes = createOneDataBlock(pInfo->pRes, false);
H
Haojun Liao 已提交
2185 2186
  blockDataEnsureCapacity(pInfo->pFinalRes, pOperator->resultInfo.capacity);

H
Haojun Liao 已提交
2187 2188 2189 2190 2191
  code = filterInitFromNode((SNode*)pPhyFillNode->node.pConditions, &pOperator->exprSupp.pFilterInfo, 0);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }

H
Haojun Liao 已提交
2192
  setOperatorInfo(pOperator, "FillOperator", QUERY_NODE_PHYSICAL_PLAN_FILL, false, OP_NOT_OPENED, pInfo, pTaskInfo);
2193
  pOperator->exprSupp.numOfExprs = pInfo->numOfExpr;
H
Haojun Liao 已提交
2194
  pOperator->fpSet = createOperatorFpSet(operatorDummyOpenFn, doFill, NULL, destroyFillOperatorInfo, NULL);
2195

2196
  code = appendDownstream(pOperator, &downstream, 1);
2197
  return pOperator;
H
Haojun Liao 已提交
2198

2199
_error:
H
Haojun Liao 已提交
2200 2201 2202 2203
  if (pInfo != NULL) {
    destroyFillOperatorInfo(pInfo);
  }

2204
  pTaskInfo->code = code;
wafwerar's avatar
wafwerar 已提交
2205
  taosMemoryFreeClear(pOperator);
H
Haojun Liao 已提交
2206
  return NULL;
2207 2208
}

D
dapan1121 已提交
2209
static SExecTaskInfo* createExecTaskInfo(uint64_t queryId, uint64_t taskId, EOPTR_EXEC_MODEL model, char* dbFName) {
wafwerar's avatar
wafwerar 已提交
2210
  SExecTaskInfo* pTaskInfo = taosMemoryCalloc(1, sizeof(SExecTaskInfo));
H
Haojun Liao 已提交
2211 2212 2213 2214 2215
  if (pTaskInfo == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return NULL;
  }

2216
  setTaskStatus(pTaskInfo, TASK_NOT_COMPLETED);
H
Haojun Liao 已提交
2217

2218
  pTaskInfo->schemaInfo.dbname = strdup(dbFName);
2219
  pTaskInfo->cost.created = taosGetTimestampMs();
H
Haojun Liao 已提交
2220
  pTaskInfo->id.queryId = queryId;
dengyihao's avatar
dengyihao 已提交
2221
  pTaskInfo->execModel = model;
H
Haojun Liao 已提交
2222
  pTaskInfo->pTableInfoList = tableListCreate();
D
dapan1121 已提交
2223
  pTaskInfo->stopInfo.pStopInfo = taosArrayInit(4, sizeof(SExchangeOpStopInfo));
H
Haojun Liao 已提交
2224
  pTaskInfo->pResultBlockList = taosArrayInit(128, POINTER_BYTES);
H
Haojun Liao 已提交
2225

wafwerar's avatar
wafwerar 已提交
2226
  char* p = taosMemoryCalloc(1, 128);
L
Liu Jicong 已提交
2227
  snprintf(p, 128, "TID:0x%" PRIx64 " QID:0x%" PRIx64, taskId, queryId);
H
Haojun Liao 已提交
2228
  pTaskInfo->id.str = p;
H
Haojun Liao 已提交
2229

2230 2231
  return pTaskInfo;
}
H
Haojun Liao 已提交
2232

H
Haojun Liao 已提交
2233 2234
SSchemaWrapper* extractQueriedColumnSchema(SScanPhysiNode* pScanNode);

2235
int32_t extractTableSchemaInfo(SReadHandle* pHandle, SScanPhysiNode* pScanNode, SExecTaskInfo* pTaskInfo) {
2236 2237
  SMetaReader mr = {0};
  metaReaderInit(&mr, pHandle->meta, 0);
H
Haojun Liao 已提交
2238
  int32_t code = metaGetTableEntryByUidCache(&mr, pScanNode->uid);
2239
  if (code != TSDB_CODE_SUCCESS) {
L
Liu Jicong 已提交
2240 2241
    qError("failed to get the table meta, uid:0x%" PRIx64 ", suid:0x%" PRIx64 ", %s", pScanNode->uid, pScanNode->suid,
           GET_TASKID(pTaskInfo));
H
Haojun Liao 已提交
2242

D
dapan1121 已提交
2243
    metaReaderClear(&mr);
2244
    return terrno;
D
dapan1121 已提交
2245
  }
2246

2247 2248
  SSchemaInfo* pSchemaInfo = &pTaskInfo->schemaInfo;
  pSchemaInfo->tablename = strdup(mr.me.name);
2249 2250

  if (mr.me.type == TSDB_SUPER_TABLE) {
2251 2252
    pSchemaInfo->sw = tCloneSSchemaWrapper(&mr.me.stbEntry.schemaRow);
    pSchemaInfo->tversion = mr.me.stbEntry.schemaTag.version;
2253
  } else if (mr.me.type == TSDB_CHILD_TABLE) {
2254 2255
    tDecoderClear(&mr.coder);

2256
    tb_uid_t suid = mr.me.ctbEntry.suid;
H
Haojun Liao 已提交
2257
    metaGetTableEntryByUidCache(&mr, suid);
2258 2259
    pSchemaInfo->sw = tCloneSSchemaWrapper(&mr.me.stbEntry.schemaRow);
    pSchemaInfo->tversion = mr.me.stbEntry.schemaTag.version;
2260
  } else {
2261
    pSchemaInfo->sw = tCloneSSchemaWrapper(&mr.me.ntbEntry.schemaRow);
2262
  }
2263 2264

  metaReaderClear(&mr);
2265

H
Haojun Liao 已提交
2266 2267 2268 2269 2270
  pSchemaInfo->qsw = extractQueriedColumnSchema(pScanNode);
  return TSDB_CODE_SUCCESS;
}

SSchemaWrapper* extractQueriedColumnSchema(SScanPhysiNode* pScanNode) {
2271 2272 2273
  int32_t numOfCols = LIST_LENGTH(pScanNode->pScanCols);
  int32_t numOfTags = LIST_LENGTH(pScanNode->pScanPseudoCols);

2274
  SSchemaWrapper* pqSw = taosMemoryCalloc(1, sizeof(SSchemaWrapper));
2275
  pqSw->pSchema = taosMemoryCalloc(numOfCols + numOfTags, sizeof(SSchema));
2276

L
Liu Jicong 已提交
2277
  for (int32_t i = 0; i < numOfCols; ++i) {
H
Haojun Liao 已提交
2278
    STargetNode* pNode = (STargetNode*)nodesListGetNode(pScanNode->pScanCols, i);
2279 2280
    SColumnNode* pColNode = (SColumnNode*)pNode->pExpr;

H
Haojun Liao 已提交
2281 2282 2283
    SSchema* pSchema = &pqSw->pSchema[pqSw->nCols++];
    pSchema->colId = pColNode->colId;
    pSchema->type = pColNode->node.resType.type;
H
Haojun Liao 已提交
2284 2285
    pSchema->bytes = pColNode->node.resType.bytes;
    tstrncpy(pSchema->name, pColNode->colName, tListLen(pSchema->name));
2286 2287
  }

2288
  // this the tags and pseudo function columns, we only keep the tag columns
2289
  for (int32_t i = 0; i < numOfTags; ++i) {
2290 2291 2292 2293 2294 2295 2296 2297 2298
    STargetNode* pNode = (STargetNode*)nodesListGetNode(pScanNode->pScanPseudoCols, i);

    int32_t type = nodeType(pNode->pExpr);
    if (type == QUERY_NODE_COLUMN) {
      SColumnNode* pColNode = (SColumnNode*)pNode->pExpr;

      SSchema* pSchema = &pqSw->pSchema[pqSw->nCols++];
      pSchema->colId = pColNode->colId;
      pSchema->type = pColNode->node.resType.type;
H
Haojun Liao 已提交
2299
      pSchema->bytes = pColNode->node.resType.bytes;
H
Haojun Liao 已提交
2300
      tstrncpy(pSchema->name, pColNode->colName, tListLen(pSchema->name));
2301 2302 2303
    }
  }

H
Haojun Liao 已提交
2304
  return pqSw;
2305 2306
}

2307 2308
static void cleanupTableSchemaInfo(SSchemaInfo* pSchemaInfo) {
  taosMemoryFreeClear(pSchemaInfo->dbname);
2309
  taosMemoryFreeClear(pSchemaInfo->tablename);
2310 2311
  tDeleteSSchemaWrapper(pSchemaInfo->sw);
  tDeleteSSchemaWrapper(pSchemaInfo->qsw);
2312 2313
}

2314
static void cleanupStreamInfo(SStreamTaskInfo* pStreamInfo) { tDeleteSSchemaWrapper(pStreamInfo->schema); }
2315

2316
bool groupbyTbname(SNodeList* pGroupList) {
2317
  bool bytbname = false;
2318
  if (LIST_LENGTH(pGroupList) == 1) {
2319 2320 2321 2322 2323 2324 2325 2326 2327 2328
    SNode* p = nodesListGetNode(pGroupList, 0);
    if (p->type == QUERY_NODE_FUNCTION) {
      // partition by tbname/group by tbname
      bytbname = (strcmp(((struct SFunctionNode*)p)->functionName, "tbname") == 0);
    }
  }

  return bytbname;
}

2329 2330
SOperatorInfo* createOperatorTree(SPhysiNode* pPhyNode, SExecTaskInfo* pTaskInfo, SReadHandle* pHandle, SNode* pTagCond,
                                  SNode* pTagIndexCond, const char* pUser) {
2331
  int32_t         type = nodeType(pPhyNode);
2332
  STableListInfo* pTableListInfo = pTaskInfo->pTableInfoList;
2333
  const char*     idstr = GET_TASKID(pTaskInfo);
2334

X
Xiaoyu Wang 已提交
2335
  if (pPhyNode->pChildren == NULL || LIST_LENGTH(pPhyNode->pChildren) == 0) {
2336
    SOperatorInfo* pOperator = NULL;
H
Haojun Liao 已提交
2337
    if (QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN == type) {
dengyihao's avatar
dengyihao 已提交
2338
      STableScanPhysiNode* pTableScanNode = (STableScanPhysiNode*)pPhyNode;
H
Haojun Liao 已提交
2339

2340 2341 2342 2343 2344 2345
      // NOTE: this is an patch to fix the physical plan
      // TODO remove it later
      if (pTableScanNode->scan.node.pLimit != NULL) {
        pTableScanNode->groupSort = true;
      }

L
Liu Jicong 已提交
2346 2347
      int32_t code =
          createScanTableListInfo(&pTableScanNode->scan, pTableScanNode->pGroupTags, pTableScanNode->groupSort, pHandle,
H
Haojun Liao 已提交
2348
                                  pTableListInfo, pTagCond, pTagIndexCond, pTaskInfo);
2349
      if (code) {
wmmhello's avatar
wmmhello 已提交
2350
        pTaskInfo->code = code;
2351
        qError("failed to createScanTableListInfo, code:%s, %s", tstrerror(code), idstr);
D
dapan1121 已提交
2352 2353
        return NULL;
      }
wmmhello's avatar
wmmhello 已提交
2354

2355
      code = extractTableSchemaInfo(pHandle, &pTableScanNode->scan, pTaskInfo);
S
slzhou 已提交
2356
      if (code) {
2357
        pTaskInfo->code = terrno;
wmmhello's avatar
wmmhello 已提交
2358 2359 2360
        return NULL;
      }

2361
      pOperator = createTableScanOperatorInfo(pTableScanNode, pHandle, pTaskInfo);
D
dapan1121 已提交
2362 2363 2364 2365 2366
      if (NULL == pOperator) {
        pTaskInfo->code = terrno;
        return NULL;
      }

2367
      STableScanInfo* pScanInfo = pOperator->info;
H
Haojun Liao 已提交
2368
      pTaskInfo->cost.pRecoder = &pScanInfo->base.readRecorder;
S
slzhou 已提交
2369 2370
    } else if (QUERY_NODE_PHYSICAL_PLAN_TABLE_MERGE_SCAN == type) {
      STableMergeScanPhysiNode* pTableScanNode = (STableMergeScanPhysiNode*)pPhyNode;
H
Haojun Liao 已提交
2371 2372 2373

      int32_t code = createScanTableListInfo(&pTableScanNode->scan, pTableScanNode->pGroupTags, true, pHandle,
                                             pTableListInfo, pTagCond, pTagIndexCond, pTaskInfo);
L
Liu Jicong 已提交
2374
      if (code) {
wmmhello's avatar
wmmhello 已提交
2375
        pTaskInfo->code = code;
H
Haojun Liao 已提交
2376
        qError("failed to createScanTableListInfo, code: %s", tstrerror(code));
wmmhello's avatar
wmmhello 已提交
2377 2378
        return NULL;
      }
2379

2380
      code = extractTableSchemaInfo(pHandle, &pTableScanNode->scan, pTaskInfo);
wmmhello's avatar
wmmhello 已提交
2381 2382 2383 2384
      if (code) {
        pTaskInfo->code = terrno;
        return NULL;
      }
wmmhello's avatar
wmmhello 已提交
2385

H
Haojun Liao 已提交
2386
      pOperator = createTableMergeScanOperatorInfo(pTableScanNode, pHandle, pTaskInfo);
D
dapan1121 已提交
2387 2388 2389 2390
      if (NULL == pOperator) {
        pTaskInfo->code = terrno;
        return NULL;
      }
wmmhello's avatar
wmmhello 已提交
2391

2392
      STableScanInfo* pScanInfo = pOperator->info;
H
Haojun Liao 已提交
2393
      pTaskInfo->cost.pRecoder = &pScanInfo->base.readRecorder;
H
Haojun Liao 已提交
2394
    } else if (QUERY_NODE_PHYSICAL_PLAN_EXCHANGE == type) {
2395 2396
      pOperator = createExchangeOperatorInfo(pHandle ? pHandle->pMsgCb->clientRpc : NULL, (SExchangePhysiNode*)pPhyNode,
                                             pTaskInfo);
H
Haojun Liao 已提交
2397
    } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_SCAN == type) {
5
54liuyao 已提交
2398
      STableScanPhysiNode* pTableScanNode = (STableScanPhysiNode*)pPhyNode;
5
54liuyao 已提交
2399
      if (pHandle->vnode) {
L
Liu Jicong 已提交
2400 2401
        int32_t code =
            createScanTableListInfo(&pTableScanNode->scan, pTableScanNode->pGroupTags, pTableScanNode->groupSort,
H
Haojun Liao 已提交
2402
                                    pHandle, pTableListInfo, pTagCond, pTagIndexCond, pTaskInfo);
L
Liu Jicong 已提交
2403
        if (code) {
wmmhello's avatar
wmmhello 已提交
2404
          pTaskInfo->code = code;
H
Haojun Liao 已提交
2405
          qError("failed to createScanTableListInfo, code: %s", tstrerror(code));
wmmhello's avatar
wmmhello 已提交
2406 2407
          return NULL;
        }
L
Liu Jicong 已提交
2408 2409

#ifndef NDEBUG
H
Haojun Liao 已提交
2410
        int32_t sz = tableListGetSize(pTableListInfo);
H
Haojun Liao 已提交
2411 2412
        qDebug("create stream task, total:%d", sz);

L
Liu Jicong 已提交
2413
        for (int32_t i = 0; i < sz; i++) {
H
Haojun Liao 已提交
2414
          STableKeyInfo* pKeyInfo = tableListGetInfo(pTableListInfo, i);
2415
          qDebug("add table uid:%" PRIu64 ", gid:%" PRIu64, pKeyInfo->uid, pKeyInfo->groupId);
L
Liu Jicong 已提交
2416 2417
        }
#endif
2418
      }
2419

H
Haojun Liao 已提交
2420
      pTaskInfo->schemaInfo.qsw = extractQueriedColumnSchema(&pTableScanNode->scan);
2421
      pOperator = createStreamScanOperatorInfo(pHandle, pTableScanNode, pTagCond, pTaskInfo);
H
Haojun Liao 已提交
2422
    } else if (QUERY_NODE_PHYSICAL_PLAN_SYSTABLE_SCAN == type) {
L
Liu Jicong 已提交
2423
      SSystemTableScanPhysiNode* pSysScanPhyNode = (SSystemTableScanPhysiNode*)pPhyNode;
2424
      pOperator = createSysTableScanOperatorInfo(pHandle, pSysScanPhyNode, pUser, pTaskInfo);
2425
    } else if (QUERY_NODE_PHYSICAL_PLAN_TAG_SCAN == type) {
X
Xiaoyu Wang 已提交
2426
      STagScanPhysiNode* pScanPhyNode = (STagScanPhysiNode*)pPhyNode;
2427 2428

      int32_t code = createScanTableListInfo(pScanPhyNode, NULL, false, pHandle, pTableListInfo, pTagCond,
H
Haojun Liao 已提交
2429
                                             pTagIndexCond, pTaskInfo);
2430
      if (code != TSDB_CODE_SUCCESS) {
2431
        pTaskInfo->code = code;
H
Haojun Liao 已提交
2432
        qError("failed to getTableList, code: %s", tstrerror(code));
2433 2434 2435
        return NULL;
      }

H
Haojun Liao 已提交
2436
      pOperator = createTagScanOperatorInfo(pHandle, pScanPhyNode, pTaskInfo);
2437
    } else if (QUERY_NODE_PHYSICAL_PLAN_BLOCK_DIST_SCAN == type) {
2438
      SBlockDistScanPhysiNode* pBlockNode = (SBlockDistScanPhysiNode*)pPhyNode;
2439 2440

      if (pBlockNode->tableType == TSDB_SUPER_TABLE) {
H
Haojun Liao 已提交
2441 2442
        SArray* pList = taosArrayInit(4, sizeof(STableKeyInfo));
        int32_t code = vnodeGetAllTableList(pHandle->vnode, pBlockNode->uid, pList);
2443 2444 2445 2446
        if (code != TSDB_CODE_SUCCESS) {
          pTaskInfo->code = terrno;
          return NULL;
        }
H
Haojun Liao 已提交
2447

H
Haojun Liao 已提交
2448 2449
        size_t num = taosArrayGetSize(pList);
        for (int32_t i = 0; i < num; ++i) {
H
Haojun Liao 已提交
2450
          STableKeyInfo* p = taosArrayGet(pList, i);
H
Haojun Liao 已提交
2451
          tableListAddTableInfo(pTableListInfo, p->uid, 0);
H
Haojun Liao 已提交
2452
        }
H
Haojun Liao 已提交
2453

H
Haojun Liao 已提交
2454
        taosArrayDestroy(pList);
2455
      } else {  // Create group with only one table
H
Haojun Liao 已提交
2456
        tableListAddTableInfo(pTableListInfo, pBlockNode->uid, 0);
2457 2458
      }

2459
      pOperator = createDataBlockInfoScanOperator(pHandle, pBlockNode, pTaskInfo);
H
Haojun Liao 已提交
2460 2461 2462
    } else if (QUERY_NODE_PHYSICAL_PLAN_LAST_ROW_SCAN == type) {
      SLastRowScanPhysiNode* pScanNode = (SLastRowScanPhysiNode*)pPhyNode;

L
Liu Jicong 已提交
2463
      int32_t code = createScanTableListInfo(&pScanNode->scan, pScanNode->pGroupTags, true, pHandle, pTableListInfo,
H
Haojun Liao 已提交
2464
                                             pTagCond, pTagIndexCond, pTaskInfo);
2465 2466 2467 2468
      if (code != TSDB_CODE_SUCCESS) {
        pTaskInfo->code = code;
        return NULL;
      }
2469

2470
      code = extractTableSchemaInfo(pHandle, &pScanNode->scan, pTaskInfo);
2471 2472 2473
      if (code != TSDB_CODE_SUCCESS) {
        pTaskInfo->code = code;
        return NULL;
H
Haojun Liao 已提交
2474 2475
      }

2476
      pOperator = createCacherowsScanOperator(pScanNode, pHandle, pTaskInfo);
2477
    } else if (QUERY_NODE_PHYSICAL_PLAN_PROJECT == type) {
2478
      pOperator = createProjectOperatorInfo(NULL, (SProjectPhysiNode*)pPhyNode, pTaskInfo);
H
Haojun Liao 已提交
2479 2480
    } else {
      ASSERT(0);
H
Haojun Liao 已提交
2481
    }
2482 2483 2484 2485 2486

    if (pOperator != NULL) {
      pOperator->resultDataBlockId = pPhyNode->pOutputDataBlockDesc->dataBlockId;
    }

2487
    return pOperator;
H
Haojun Liao 已提交
2488 2489
  }

2490
  size_t          size = LIST_LENGTH(pPhyNode->pChildren);
2491
  SOperatorInfo** ops = taosMemoryCalloc(size, POINTER_BYTES);
2492 2493 2494 2495
  if (ops == NULL) {
    return NULL;
  }

dengyihao's avatar
dengyihao 已提交
2496
  for (int32_t i = 0; i < size; ++i) {
2497
    SPhysiNode* pChildNode = (SPhysiNode*)nodesListGetNode(pPhyNode->pChildren, i);
2498
    ops[i] = createOperatorTree(pChildNode, pTaskInfo, pHandle, pTagCond, pTagIndexCond, pUser);
2499
    if (ops[i] == NULL) {
H
Haojun Liao 已提交
2500
      taosMemoryFree(ops);
2501 2502
      return NULL;
    }
2503
  }
H
Haojun Liao 已提交
2504

2505
  SOperatorInfo* pOptr = NULL;
H
Haojun Liao 已提交
2506
  if (QUERY_NODE_PHYSICAL_PLAN_PROJECT == type) {
2507
    pOptr = createProjectOperatorInfo(ops[0], (SProjectPhysiNode*)pPhyNode, pTaskInfo);
2508
  } else if (QUERY_NODE_PHYSICAL_PLAN_HASH_AGG == type) {
H
Haojun Liao 已提交
2509 2510
    SAggPhysiNode* pAggNode = (SAggPhysiNode*)pPhyNode;
    if (pAggNode->pGroupKeys != NULL) {
H
Haojun Liao 已提交
2511
      pOptr = createGroupOperatorInfo(ops[0], pAggNode, pTaskInfo);
H
Haojun Liao 已提交
2512
    } else {
H
Haojun Liao 已提交
2513
      pOptr = createAggregateOperatorInfo(ops[0], pAggNode, pTaskInfo);
H
Haojun Liao 已提交
2514
    }
2515
  } else if (QUERY_NODE_PHYSICAL_PLAN_HASH_INTERVAL == type) {
H
Haojun Liao 已提交
2516
    SIntervalPhysiNode* pIntervalPhyNode = (SIntervalPhysiNode*)pPhyNode;
H
Haojun Liao 已提交
2517

H
Haojun Liao 已提交
2518 2519
    bool isStream = (QUERY_NODE_PHYSICAL_PLAN_STREAM_INTERVAL == type);
    pOptr = createIntervalOperatorInfo(ops[0], pIntervalPhyNode, pTaskInfo, isStream);
2520
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_INTERVAL == type) {
2521
    pOptr = createStreamIntervalOperatorInfo(ops[0], pPhyNode, pTaskInfo);
2522 2523
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE_ALIGNED_INTERVAL == type) {
    SMergeAlignedIntervalPhysiNode* pIntervalPhyNode = (SMergeAlignedIntervalPhysiNode*)pPhyNode;
2524
    pOptr = createMergeAlignedIntervalOperatorInfo(ops[0], pIntervalPhyNode, pTaskInfo);
S
shenglian zhou 已提交
2525
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE_INTERVAL == type) {
X
Xiaoyu Wang 已提交
2526
    SMergeIntervalPhysiNode* pIntervalPhyNode = (SMergeIntervalPhysiNode*)pPhyNode;
2527
    pOptr = createMergeIntervalOperatorInfo(ops[0], pIntervalPhyNode, pTaskInfo);
5
54liuyao 已提交
2528
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_SEMI_INTERVAL == type) {
2529
    int32_t children = 0;
5
54liuyao 已提交
2530 2531
    pOptr = createStreamFinalIntervalOperatorInfo(ops[0], pPhyNode, pTaskInfo, children);
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_FINAL_INTERVAL == type) {
5
54liuyao 已提交
2532
    int32_t children = pHandle->numOfVgroups;
5
54liuyao 已提交
2533
    pOptr = createStreamFinalIntervalOperatorInfo(ops[0], pPhyNode, pTaskInfo, children);
H
Haojun Liao 已提交
2534
  } else if (QUERY_NODE_PHYSICAL_PLAN_SORT == type) {
2535
    pOptr = createSortOperatorInfo(ops[0], (SSortPhysiNode*)pPhyNode, pTaskInfo);
S
shenglian zhou 已提交
2536 2537
  } else if (QUERY_NODE_PHYSICAL_PLAN_GROUP_SORT == type) {
    pOptr = createGroupSortOperatorInfo(ops[0], (SGroupSortPhysiNode*)pPhyNode, pTaskInfo);
X
Xiaoyu Wang 已提交
2538
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE == type) {
2539
    SMergePhysiNode* pMergePhyNode = (SMergePhysiNode*)pPhyNode;
2540
    pOptr = createMultiwayMergeOperatorInfo(ops, size, pMergePhyNode, pTaskInfo);
2541
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE_SESSION == type) {
H
Haojun Liao 已提交
2542
    SSessionWinodwPhysiNode* pSessionNode = (SSessionWinodwPhysiNode*)pPhyNode;
H
Haojun Liao 已提交
2543
    pOptr = createSessionAggOperatorInfo(ops[0], pSessionNode, pTaskInfo);
2544
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_SESSION == type) {
2545 2546 2547 2548 2549
    pOptr = createStreamSessionAggOperatorInfo(ops[0], pPhyNode, pTaskInfo);
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_SEMI_SESSION == type) {
    int32_t children = 0;
    pOptr = createStreamFinalSessionAggOperatorInfo(ops[0], pPhyNode, pTaskInfo, children);
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_FINAL_SESSION == type) {
2550
    int32_t children = pHandle->numOfVgroups;
2551
    pOptr = createStreamFinalSessionAggOperatorInfo(ops[0], pPhyNode, pTaskInfo, children);
H
Haojun Liao 已提交
2552
  } else if (QUERY_NODE_PHYSICAL_PLAN_PARTITION == type) {
2553
    pOptr = createPartitionOperatorInfo(ops[0], (SPartitionPhysiNode*)pPhyNode, pTaskInfo);
2554
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_PARTITION == type) {
2555
    pOptr = createStreamPartitionOperatorInfo(ops[0], (SStreamPartitionPhysiNode*)pPhyNode, pTaskInfo);
2556
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE_STATE == type) {
dengyihao's avatar
dengyihao 已提交
2557
    SStateWinodwPhysiNode* pStateNode = (SStateWinodwPhysiNode*)pPhyNode;
2558
    pOptr = createStatewindowOperatorInfo(ops[0], pStateNode, pTaskInfo);
2559
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_STATE == type) {
5
54liuyao 已提交
2560
    pOptr = createStreamStateAggOperatorInfo(ops[0], pPhyNode, pTaskInfo);
2561
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE_JOIN == type) {
2562
    pOptr = createMergeJoinOperatorInfo(ops, size, (SSortMergeJoinPhysiNode*)pPhyNode, pTaskInfo);
2563
  } else if (QUERY_NODE_PHYSICAL_PLAN_FILL == type) {
H
Haojun Liao 已提交
2564
    pOptr = createFillOperatorInfo(ops[0], (SFillPhysiNode*)pPhyNode, pTaskInfo);
5
54liuyao 已提交
2565 2566
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_FILL == type) {
    pOptr = createStreamFillOperatorInfo(ops[0], (SStreamFillPhysiNode*)pPhyNode, pTaskInfo);
H
Haojun Liao 已提交
2567 2568
  } else if (QUERY_NODE_PHYSICAL_PLAN_INDEF_ROWS_FUNC == type) {
    pOptr = createIndefinitOutputOperatorInfo(ops[0], pPhyNode, pTaskInfo);
2569 2570
  } else if (QUERY_NODE_PHYSICAL_PLAN_INTERP_FUNC == type) {
    pOptr = createTimeSliceOperatorInfo(ops[0], pPhyNode, pTaskInfo);
H
Haojun Liao 已提交
2571 2572
  } else {
    ASSERT(0);
H
Haojun Liao 已提交
2573
  }
2574

2575
  taosMemoryFree(ops);
2576 2577 2578 2579
  if (pOptr) {
    pOptr->resultDataBlockId = pPhyNode->pOutputDataBlockDesc->dataBlockId;
  }

2580
  return pOptr;
2581
}
H
Haojun Liao 已提交
2582

L
Liu Jicong 已提交
2583 2584 2585 2586 2587 2588 2589 2590 2591 2592 2593 2594 2595
static int32_t extractTbscanInStreamOpTree(SOperatorInfo* pOperator, STableScanInfo** ppInfo) {
  if (pOperator->operatorType != QUERY_NODE_PHYSICAL_PLAN_STREAM_SCAN) {
    if (pOperator->numOfDownstream == 0) {
      qError("failed to find stream scan operator");
      return TSDB_CODE_QRY_APP_ERROR;
    }

    if (pOperator->numOfDownstream > 1) {
      qError("join not supported for stream block scan");
      return TSDB_CODE_QRY_APP_ERROR;
    }
    return extractTbscanInStreamOpTree(pOperator->pDownstream[0], ppInfo);
  } else {
2596 2597 2598
    SStreamScanInfo* pInfo = pOperator->info;
    ASSERT(pInfo->pTableScanOp->operatorType == QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN);
    *ppInfo = pInfo->pTableScanOp->info;
L
Liu Jicong 已提交
2599 2600 2601 2602
    return 0;
  }
}

2603 2604 2605 2606 2607 2608 2609 2610 2611 2612 2613 2614 2615 2616 2617 2618 2619 2620 2621 2622 2623 2624
int32_t extractTableScanNode(SPhysiNode* pNode, STableScanPhysiNode** ppNode) {
  if (pNode->pChildren == NULL || LIST_LENGTH(pNode->pChildren) == 0) {
    if (QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN == pNode->type) {
      *ppNode = (STableScanPhysiNode*)pNode;
      return 0;
    } else {
      ASSERT(0);
      terrno = TSDB_CODE_QRY_APP_ERROR;
      return -1;
    }
  } else {
    if (LIST_LENGTH(pNode->pChildren) != 1) {
      ASSERT(0);
      terrno = TSDB_CODE_QRY_APP_ERROR;
      return -1;
    }
    SPhysiNode* pChildNode = (SPhysiNode*)nodesListGetNode(pNode->pChildren, 0);
    return extractTableScanNode(pChildNode, ppNode);
  }
  return -1;
}

2625
#if 0
L
Liu Jicong 已提交
2626 2627 2628 2629 2630
int32_t rebuildReader(SOperatorInfo* pOperator, SSubplan* plan, SReadHandle* pHandle, int64_t uid, int64_t ts) {
  STableScanInfo* pTableScanInfo = NULL;
  if (extractTbscanInStreamOpTree(pOperator, &pTableScanInfo) < 0) {
    return -1;
  }
2631

L
Liu Jicong 已提交
2632 2633 2634 2635
  STableScanPhysiNode* pNode = NULL;
  if (extractTableScanNode(plan->pNode, &pNode) < 0) {
    ASSERT(0);
  }
2636

H
Haojun Liao 已提交
2637
  tsdbReaderClose(pTableScanInfo->dataReader);
2638

L
Liu Jicong 已提交
2639
  STableListInfo info = {0};
H
Haojun Liao 已提交
2640
  pTableScanInfo->dataReader = doCreateDataReader(pNode, pHandle, &info, NULL);
L
Liu Jicong 已提交
2641 2642 2643 2644
  if (pTableScanInfo->dataReader == NULL) {
    ASSERT(0);
    qError("failed to create data reader");
    return TSDB_CODE_QRY_APP_ERROR;
2645
  }
L
Liu Jicong 已提交
2646
  // TODO: set uid and ts to data reader
2647 2648
  return 0;
}
2649
#endif
2650

D
dapan1121 已提交
2651
int32_t createDataSinkParam(SDataSinkNode* pNode, void** pParam, qTaskInfo_t* pTaskInfo, SReadHandle* readHandle) {
D
dapan1121 已提交
2652
  SExecTaskInfo* pTask = *(SExecTaskInfo**)pTaskInfo;
2653

D
dapan1121 已提交
2654
  switch (pNode->type) {
D
dapan1121 已提交
2655 2656 2657 2658 2659 2660
    case QUERY_NODE_PHYSICAL_PLAN_QUERY_INSERT: {
      SInserterParam* pInserterParam = taosMemoryCalloc(1, sizeof(SInserterParam));
      if (NULL == pInserterParam) {
        return TSDB_CODE_OUT_OF_MEMORY;
      }
      pInserterParam->readHandle = readHandle;
L
Liu Jicong 已提交
2661

D
dapan1121 已提交
2662 2663 2664
      *pParam = pInserterParam;
      break;
    }
D
dapan1121 已提交
2665
    case QUERY_NODE_PHYSICAL_PLAN_DELETE: {
2666
      SDeleterParam* pDeleterParam = taosMemoryCalloc(1, sizeof(SDeleterParam));
D
dapan1121 已提交
2667 2668 2669
      if (NULL == pDeleterParam) {
        return TSDB_CODE_OUT_OF_MEMORY;
      }
H
Haojun Liao 已提交
2670 2671 2672 2673
      int32_t tbNum = tableListGetSize(pTask->pTableInfoList);
      pDeleterParam->suid = tableListGetSuid(pTask->pTableInfoList);

      // TODO extract uid list
D
dapan1121 已提交
2674 2675 2676 2677 2678
      pDeleterParam->pUidList = taosArrayInit(tbNum, sizeof(uint64_t));
      if (NULL == pDeleterParam->pUidList) {
        taosMemoryFree(pDeleterParam);
        return TSDB_CODE_OUT_OF_MEMORY;
      }
H
Haojun Liao 已提交
2679

D
dapan1121 已提交
2680
      for (int32_t i = 0; i < tbNum; ++i) {
H
Haojun Liao 已提交
2681
        STableKeyInfo* pTable = tableListGetInfo(pTask->pTableInfoList, i);
D
dapan1121 已提交
2682 2683 2684 2685 2686 2687 2688 2689 2690 2691 2692 2693 2694
        taosArrayPush(pDeleterParam->pUidList, &pTable->uid);
      }

      *pParam = pDeleterParam;
      break;
    }
    default:
      break;
  }

  return TSDB_CODE_SUCCESS;
}

dengyihao's avatar
dengyihao 已提交
2695
int32_t createExecTaskInfoImpl(SSubplan* pPlan, SExecTaskInfo** pTaskInfo, SReadHandle* pHandle, uint64_t taskId,
D
dapan1121 已提交
2696
                               char* sql, EOPTR_EXEC_MODEL model) {
H
Haojun Liao 已提交
2697 2698
  uint64_t queryId = pPlan->id.queryId;

D
dapan1121 已提交
2699
  *pTaskInfo = createExecTaskInfo(queryId, taskId, model, pPlan->dbFName);
H
Haojun Liao 已提交
2700 2701 2702
  if (*pTaskInfo == NULL) {
    goto _complete;
  }
H
Haojun Liao 已提交
2703

2704
  if (pHandle) {
L
Liu Jicong 已提交
2705
    /*(*pTaskInfo)->streamInfo.fillHistoryVer1 = pHandle->fillHistoryVer1;*/
2706 2707 2708
    if (pHandle->pStateBackend) {
      (*pTaskInfo)->streamInfo.pState = pHandle->pStateBackend;
    }
H
Haojun Liao 已提交
2709 2710
  }

2711
  (*pTaskInfo)->sql = sql;
D
dapan1121 已提交
2712
  sql = NULL;
H
Haojun Liao 已提交
2713

2714
  (*pTaskInfo)->pSubplan = pPlan;
2715 2716
  (*pTaskInfo)->pRoot =
      createOperatorTree(pPlan->pNode, *pTaskInfo, pHandle, pPlan->pTagCond, pPlan->pTagIndexCond, pPlan->user);
L
Liu Jicong 已提交
2717

D
dapan1121 已提交
2718
  if (NULL == (*pTaskInfo)->pRoot) {
H
Haojun Liao 已提交
2719
    terrno = (*pTaskInfo)->code;
D
dapan1121 已提交
2720
    goto _complete;
2721 2722
  }

H
Haojun Liao 已提交
2723
  return TSDB_CODE_SUCCESS;
H
Haojun Liao 已提交
2724

H
Haojun Liao 已提交
2725
_complete:
D
dapan1121 已提交
2726
  taosMemoryFree(sql);
H
Haojun Liao 已提交
2727
  doDestroyTask(*pTaskInfo);
H
Haojun Liao 已提交
2728
  return terrno;
H
Haojun Liao 已提交
2729 2730
}

H
Haojun Liao 已提交
2731 2732 2733 2734 2735
static void freeBlock(void* pParam) {
  SSDataBlock* pBlock = *(SSDataBlock**)pParam;
  blockDataDestroy(pBlock);
}

L
Liu Jicong 已提交
2736
void doDestroyTask(SExecTaskInfo* pTaskInfo) {
H
Haojun Liao 已提交
2737 2738
  qDebug("%s execTask is freed", GET_TASKID(pTaskInfo));

H
Haojun Liao 已提交
2739
  pTaskInfo->pTableInfoList = tableListDestroy(pTaskInfo->pTableInfoList);
H
Haojun Liao 已提交
2740
  destroyOperatorInfo(pTaskInfo->pRoot);
2741
  cleanupTableSchemaInfo(&pTaskInfo->schemaInfo);
2742
  cleanupStreamInfo(&pTaskInfo->streamInfo);
2743

D
dapan1121 已提交
2744
  if (!pTaskInfo->localFetch.localExec) {
D
dapan1121 已提交
2745 2746
    nodesDestroyNode((SNode*)pTaskInfo->pSubplan);
  }
2747

H
Haojun Liao 已提交
2748
  taosArrayDestroyEx(pTaskInfo->pResultBlockList, freeBlock);
D
dapan1121 已提交
2749
  taosArrayDestroy(pTaskInfo->stopInfo.pStopInfo);
wafwerar's avatar
wafwerar 已提交
2750 2751 2752
  taosMemoryFreeClear(pTaskInfo->sql);
  taosMemoryFreeClear(pTaskInfo->id.str);
  taosMemoryFreeClear(pTaskInfo);
2753 2754 2755 2756
}

static int64_t getQuerySupportBufSize(size_t numOfTables) {
  size_t s1 = sizeof(STableQueryInfo);
L
Liu Jicong 已提交
2757 2758
  //  size_t s3 = sizeof(STableCheckInfo);  buffer consumption in tsdb
  return (int64_t)(s1 * 1.5 * numOfTables);
2759 2760 2761 2762 2763 2764 2765
}

int32_t checkForQueryBuf(size_t numOfTables) {
  int64_t t = getQuerySupportBufSize(numOfTables);
  if (tsQueryBufferSizeBytes < 0) {
    return TSDB_CODE_SUCCESS;
  } else if (tsQueryBufferSizeBytes > 0) {
L
Liu Jicong 已提交
2766
    while (1) {
2767 2768 2769 2770 2771 2772 2773 2774 2775 2776 2777 2778 2779 2780 2781 2782 2783 2784 2785 2786 2787 2788 2789 2790 2791 2792
      int64_t s = tsQueryBufferSizeBytes;
      int64_t remain = s - t;
      if (remain >= 0) {
        if (atomic_val_compare_exchange_64(&tsQueryBufferSizeBytes, s, remain) == s) {
          return TSDB_CODE_SUCCESS;
        }
      } else {
        return TSDB_CODE_QRY_NOT_ENOUGH_BUFFER;
      }
    }
  }

  // disable query processing if the value of tsQueryBufferSize is zero.
  return TSDB_CODE_QRY_NOT_ENOUGH_BUFFER;
}

void releaseQueryBuf(size_t numOfTables) {
  if (tsQueryBufferSizeBytes < 0) {
    return;
  }

  int64_t t = getQuerySupportBufSize(numOfTables);

  // restore value is not enough buffer available
  atomic_add_fetch_64(&tsQueryBufferSizeBytes, t);
}
D
dapan1121 已提交
2793

H
Haojun Liao 已提交
2794
int32_t getOperatorExplainExecInfo(SOperatorInfo* operatorInfo, SArray* pExecInfoList) {
2795
  SExplainExecInfo  execInfo = {0};
H
Haojun Liao 已提交
2796
  SExplainExecInfo* pExplainInfo = taosArrayPush(pExecInfoList, &execInfo);
2797

H
Haojun Liao 已提交
2798 2799 2800 2801 2802
  pExplainInfo->numOfRows = operatorInfo->resultInfo.totalRows;
  pExplainInfo->startupCost = operatorInfo->cost.openCost;
  pExplainInfo->totalCost = operatorInfo->cost.totalCost;
  pExplainInfo->verboseLen = 0;
  pExplainInfo->verboseInfo = NULL;
D
dapan1121 已提交
2803

2804
  if (operatorInfo->fpSet.getExplainFn) {
2805 2806
    int32_t code =
        operatorInfo->fpSet.getExplainFn(operatorInfo, &pExplainInfo->verboseInfo, &pExplainInfo->verboseLen);
D
dapan1121 已提交
2807
    if (code) {
2808
      qError("%s operator getExplainFn failed, code:%s", GET_TASKID(operatorInfo->pTaskInfo), tstrerror(code));
D
dapan1121 已提交
2809 2810 2811
      return code;
    }
  }
dengyihao's avatar
dengyihao 已提交
2812

D
dapan1121 已提交
2813
  int32_t code = 0;
D
dapan1121 已提交
2814
  for (int32_t i = 0; i < operatorInfo->numOfDownstream; ++i) {
H
Haojun Liao 已提交
2815 2816
    code = getOperatorExplainExecInfo(operatorInfo->pDownstream[i], pExecInfoList);
    if (code != TSDB_CODE_SUCCESS) {
2817
      //      taosMemoryFreeClear(*pRes);
D
dapan1121 已提交
2818 2819 2820 2821 2822
      return TSDB_CODE_QRY_OUT_OF_MEMORY;
    }
  }

  return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
2823
}
5
54liuyao 已提交
2824

2825 2826
int32_t setOutputBuf(SStreamState* pState, STimeWindow* win, SResultRow** pResult, int64_t tableGroupId,
                     SqlFunctionCtx* pCtx, int32_t numOfOutput, int32_t* rowEntryInfoOffset, SAggSupporter* pAggSup) {
2827 2828 2829 2830 2831 2832
  SWinKey key = {
      .ts = win->skey,
      .groupId = tableGroupId,
  };
  char*   value = NULL;
  int32_t size = pAggSup->resultRowSize;
5
54liuyao 已提交
2833

2834
  if (streamStateAddIfNotExist(pState, &key, (void**)&value, &size) < 0) {
2835 2836 2837 2838 2839 2840 2841 2842 2843 2844
    return TSDB_CODE_QRY_OUT_OF_MEMORY;
  }
  *pResult = (SResultRow*)value;
  ASSERT(*pResult);
  // set time window for current result
  (*pResult)->win = (*win);
  setResultRowInitCtx(*pResult, pCtx, numOfOutput, rowEntryInfoOffset);
  return TSDB_CODE_SUCCESS;
}

2845 2846
int32_t releaseOutputBuf(SStreamState* pState, SWinKey* pKey, SResultRow* pResult) {
  streamStateReleaseBuf(pState, pKey, pResult);
2847 2848 2849
  return TSDB_CODE_SUCCESS;
}

2850 2851
int32_t saveOutputBuf(SStreamState* pState, SWinKey* pKey, SResultRow* pResult, int32_t resSize) {
  streamStatePut(pState, pKey, pResult, resSize);
2852 2853 2854
  return TSDB_CODE_SUCCESS;
}

2855
int32_t buildDataBlockFromGroupRes(SOperatorInfo* pOperator, SStreamState* pState, SSDataBlock* pBlock, SExprSupp* pSup,
2856
                                   SGroupResInfo* pGroupResInfo) {
2857
  SExecTaskInfo*  pTaskInfo = pOperator->pTaskInfo;
2858 2859 2860 2861 2862 2863 2864 2865 2866 2867 2868 2869
  SExprInfo*      pExprInfo = pSup->pExprInfo;
  int32_t         numOfExprs = pSup->numOfExprs;
  int32_t*        rowEntryOffset = pSup->rowEntryInfoOffset;
  SqlFunctionCtx* pCtx = pSup->pCtx;

  int32_t numOfRows = getNumOfTotalRes(pGroupResInfo);

  for (int32_t i = pGroupResInfo->index; i < numOfRows; i += 1) {
    SResKeyPos* pPos = taosArrayGetP(pGroupResInfo->pRows, i);
    int32_t     size = 0;
    void*       pVal = NULL;
    SWinKey     key = {
2870 2871
            .ts = *(TSKEY*)pPos->key,
            .groupId = pPos->groupId,
2872
    };
2873
    int32_t code = streamStateGet(pState, &key, &pVal, &size);
2874 2875 2876 2877 2878 2879
    ASSERT(code == 0);
    SResultRow* pRow = (SResultRow*)pVal;
    doUpdateNumOfRows(pCtx, pRow, numOfExprs, rowEntryOffset);
    // no results, continue to check the next one
    if (pRow->numOfRows == 0) {
      pGroupResInfo->index += 1;
2880
      releaseOutputBuf(pState, &key, pRow);
2881 2882 2883
      continue;
    }

H
Haojun Liao 已提交
2884 2885
    if (pBlock->info.id.groupId == 0) {
      pBlock->info.id.groupId = pPos->groupId;
2886
      void* tbname = NULL;
H
Haojun Liao 已提交
2887
      if (streamStateGetParName(pTaskInfo->streamInfo.pState, pBlock->info.id.groupId, &tbname) < 0) {
L
Liu Jicong 已提交
2888
        pBlock->info.parTbName[0] = 0;
2889 2890
      } else {
        memcpy(pBlock->info.parTbName, tbname, TSDB_TABLE_NAME_LEN);
2891
      }
2892
      tdbFree(tbname);
2893 2894
    } else {
      // current value belongs to different group, it can't be packed into one datablock
H
Haojun Liao 已提交
2895
      if (pBlock->info.id.groupId != pPos->groupId) {
2896
        releaseOutputBuf(pState, &key, pRow);
2897 2898 2899 2900 2901 2902
        break;
      }
    }

    if (pBlock->info.rows + pRow->numOfRows > pBlock->info.capacity) {
      ASSERT(pBlock->info.rows > 0);
2903
      releaseOutputBuf(pState, &key, pRow);
2904 2905 2906 2907 2908 2909 2910 2911 2912 2913
      break;
    }

    pGroupResInfo->index += 1;

    for (int32_t j = 0; j < numOfExprs; ++j) {
      int32_t slotId = pExprInfo[j].base.resSchema.slotId;

      pCtx[j].resultInfo = getResultEntryInfo(pRow, j, rowEntryOffset);
      if (pCtx[j].fpSet.finalize) {
2914 2915 2916 2917
        int32_t code1 = pCtx[j].fpSet.finalize(&pCtx[j], pBlock);
        if (TAOS_FAILED(code1)) {
          qError("%s build result data block error, code %s", GET_TASKID(pTaskInfo), tstrerror(code1));
          T_LONG_JMP(pTaskInfo->env, code1);
2918 2919 2920 2921 2922 2923 2924 2925 2926 2927 2928 2929 2930
        }
      } else if (strcmp(pCtx[j].pExpr->pExpr->_function.functionName, "_select_value") == 0) {
        // do nothing, todo refactor
      } else {
        // expand the result into multiple rows. E.g., _wstart, top(k, 20)
        // the _wstart needs to copy to 20 following rows, since the results of top-k expands to 20 different rows.
        SColumnInfoData* pColInfoData = taosArrayGet(pBlock->pDataBlock, slotId);
        char*            in = GET_ROWCELL_INTERBUF(pCtx[j].resultInfo);
        for (int32_t k = 0; k < pRow->numOfRows; ++k) {
          colDataAppend(pColInfoData, pBlock->info.rows + k, in, pCtx[j].resultInfo->isNullRes);
        }
      }
    }
5
54liuyao 已提交
2931

2932
    pBlock->info.rows += pRow->numOfRows;
2933
    releaseOutputBuf(pState, &key, pRow);
2934 2935 2936 2937
  }
  blockDataUpdateTsWindow(pBlock, 0);
  return TSDB_CODE_SUCCESS;
}
5
54liuyao 已提交
2938 2939 2940 2941 2942 2943 2944

int32_t saveSessionDiscBuf(SStreamState* pState, SSessionKey* key, void* buf, int32_t size) {
  streamStateSessionPut(pState, key, (const void*)buf, size);
  releaseOutputBuf(pState, NULL, (SResultRow*)buf);
  return TSDB_CODE_SUCCESS;
}

2945
int32_t buildSessionResultDataBlock(SOperatorInfo* pOperator, SStreamState* pState, SSDataBlock* pBlock,
5
54liuyao 已提交
2946
                                    SExprSupp* pSup, SGroupResInfo* pGroupResInfo) {
2947
  SExecTaskInfo*  pTaskInfo = pOperator->pTaskInfo;
5
54liuyao 已提交
2948 2949 2950 2951 2952 2953 2954 2955 2956 2957 2958 2959 2960
  SExprInfo*      pExprInfo = pSup->pExprInfo;
  int32_t         numOfExprs = pSup->numOfExprs;
  int32_t*        rowEntryOffset = pSup->rowEntryInfoOffset;
  SqlFunctionCtx* pCtx = pSup->pCtx;

  int32_t numOfRows = getNumOfTotalRes(pGroupResInfo);

  for (int32_t i = pGroupResInfo->index; i < numOfRows; i += 1) {
    SSessionKey* pKey = taosArrayGet(pGroupResInfo->pRows, i);
    int32_t      size = 0;
    void*        pVal = NULL;
    int32_t      code = streamStateSessionGet(pState, pKey, &pVal, &size);
    ASSERT(code == 0);
2961 2962
    if (code == -1) {
      // coverity scan
5
54liuyao 已提交
2963
      pGroupResInfo->index += 1;
2964 2965
      continue;
    }
5
54liuyao 已提交
2966 2967 2968 2969 2970 2971 2972 2973 2974
    SResultRow* pRow = (SResultRow*)pVal;
    doUpdateNumOfRows(pCtx, pRow, numOfExprs, rowEntryOffset);
    // no results, continue to check the next one
    if (pRow->numOfRows == 0) {
      pGroupResInfo->index += 1;
      releaseOutputBuf(pState, NULL, pRow);
      continue;
    }

H
Haojun Liao 已提交
2975 2976
    if (pBlock->info.id.groupId == 0) {
      pBlock->info.id.groupId = pKey->groupId;
2977

2978
      void* tbname = NULL;
H
Haojun Liao 已提交
2979
      if (streamStateGetParName(pTaskInfo->streamInfo.pState, pBlock->info.id.groupId, &tbname) < 0) {
2980
        pBlock->info.parTbName[0] = 0;
2981
      } else {
2982
        memcpy(pBlock->info.parTbName, tbname, TSDB_TABLE_NAME_LEN);
2983
      }
2984
      tdbFree(tbname);
5
54liuyao 已提交
2985 2986
    } else {
      // current value belongs to different group, it can't be packed into one datablock
H
Haojun Liao 已提交
2987
      if (pBlock->info.id.groupId != pKey->groupId) {
5
54liuyao 已提交
2988 2989 2990 2991 2992 2993 2994 2995 2996 2997 2998 2999 3000 3001 3002 3003 3004 3005 3006 3007 3008 3009 3010 3011 3012 3013 3014 3015 3016 3017 3018 3019 3020 3021 3022 3023 3024 3025 3026 3027 3028 3029
        releaseOutputBuf(pState, NULL, pRow);
        break;
      }
    }

    if (pBlock->info.rows + pRow->numOfRows > pBlock->info.capacity) {
      ASSERT(pBlock->info.rows > 0);
      releaseOutputBuf(pState, NULL, pRow);
      break;
    }

    pGroupResInfo->index += 1;

    for (int32_t j = 0; j < numOfExprs; ++j) {
      int32_t slotId = pExprInfo[j].base.resSchema.slotId;

      pCtx[j].resultInfo = getResultEntryInfo(pRow, j, rowEntryOffset);
      if (pCtx[j].fpSet.finalize) {
        int32_t code1 = pCtx[j].fpSet.finalize(&pCtx[j], pBlock);
        if (TAOS_FAILED(code1)) {
          qError("%s build result data block error, code %s", GET_TASKID(pTaskInfo), tstrerror(code1));
          T_LONG_JMP(pTaskInfo->env, code1);
        }
      } else if (strcmp(pCtx[j].pExpr->pExpr->_function.functionName, "_select_value") == 0) {
        // do nothing, todo refactor
      } else {
        // expand the result into multiple rows. E.g., _wstart, top(k, 20)
        // the _wstart needs to copy to 20 following rows, since the results of top-k expands to 20 different rows.
        SColumnInfoData* pColInfoData = taosArrayGet(pBlock->pDataBlock, slotId);
        char*            in = GET_ROWCELL_INTERBUF(pCtx[j].resultInfo);
        for (int32_t k = 0; k < pRow->numOfRows; ++k) {
          colDataAppend(pColInfoData, pBlock->info.rows + k, in, pCtx[j].resultInfo->isNullRes);
        }
      }
    }

    pBlock->info.rows += pRow->numOfRows;
    // saveSessionDiscBuf(pState, pKey, pVal, size);
    releaseOutputBuf(pState, NULL, pRow);
  }
  blockDataUpdateTsWindow(pBlock, 0);
  return TSDB_CODE_SUCCESS;
3030
}