executorimpl.c 105.9 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
// this is a blocking operator
L
Liu Jicong 已提交
1423
static int32_t doOpenAggregateOptr(SOperatorInfo* pOperator) {
H
Haojun Liao 已提交
1424 1425
  if (OPTR_IS_OPENED(pOperator)) {
    return TSDB_CODE_SUCCESS;
1426 1427
  }

H
Haojun Liao 已提交
1428
  SExecTaskInfo*    pTaskInfo = pOperator->pTaskInfo;
1429
  SAggOperatorInfo* pAggInfo = pOperator->info;
H
Haojun Liao 已提交
1430

1431 1432
  SExprSupp*     pSup = &pOperator->exprSupp;
  SOperatorInfo* downstream = pOperator->pDownstream[0];
1433

1434 1435
  int64_t st = taosGetTimestampUs();

1436 1437 1438
  int32_t order = TSDB_ORDER_ASC;
  int32_t scanFlag = MAIN_SCAN;

H
Haojun Liao 已提交
1439
  while (1) {
1440
    SSDataBlock* pBlock = downstream->fpSet.getNextFn(downstream);
1441 1442 1443 1444
    if (pBlock == NULL) {
      break;
    }

1445 1446
    int32_t code = getTableScanInfo(pOperator, &order, &scanFlag);
    if (code != TSDB_CODE_SUCCESS) {
1447
      T_LONG_JMP(pTaskInfo->env, code);
1448
    }
1449

1450
    // there is an scalar expression that needs to be calculated before apply the group aggregation.
1451 1452 1453
    if (pAggInfo->scalarExprSup.pExprInfo != NULL) {
      SExprSupp* pSup1 = &pAggInfo->scalarExprSup;
      code = projectApplyFunctions(pSup1->pExprInfo, pBlock, pBlock, pSup1->pCtx, pSup1->numOfExprs, NULL);
1454
      if (code != TSDB_CODE_SUCCESS) {
1455
        T_LONG_JMP(pTaskInfo->env, code);
1456
      }
1457 1458
    }

1459
    // the pDataBlock are always the same one, no need to call this again
H
Haojun Liao 已提交
1460
    setExecutionContext(pOperator, pOperator->exprSupp.numOfExprs, pBlock->info.id.groupId);
1461
    setInputDataBlock(pSup, pBlock, order, scanFlag, true);
1462
    code = doAggregateImpl(pOperator, pSup->pCtx);
1463
    if (code != 0) {
1464
      T_LONG_JMP(pTaskInfo->env, code);
1465
    }
1466 1467
  }

1468 1469 1470 1471 1472
  // 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);
  }

1473
  initGroupedResultInfo(&pAggInfo->groupResInfo, pAggInfo->aggSup.pResultRowHashTable, 0);
H
Haojun Liao 已提交
1474
  OPTR_SET_OPENED(pOperator);
1475

1476
  pOperator->cost.openCost = (taosGetTimestampUs() - st) / 1000.0;
1477
  return pTaskInfo->code;
H
Haojun Liao 已提交
1478 1479
}

1480
static SSDataBlock* getAggregateResult(SOperatorInfo* pOperator) {
L
Liu Jicong 已提交
1481
  SAggOperatorInfo* pAggInfo = pOperator->info;
H
Haojun Liao 已提交
1482 1483 1484 1485 1486 1487
  SOptrBasicInfo*   pInfo = &pAggInfo->binfo;

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

L
Liu Jicong 已提交
1488
  SExecTaskInfo* pTaskInfo = pOperator->pTaskInfo;
1489
  pTaskInfo->code = pOperator->fpSet._openFn(pOperator);
H
Haojun Liao 已提交
1490
  if (pTaskInfo->code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
1491
    setOperatorCompleted(pOperator);
H
Haojun Liao 已提交
1492 1493 1494
    return NULL;
  }

H
Haojun Liao 已提交
1495
  blockDataEnsureCapacity(pInfo->pRes, pOperator->resultInfo.capacity);
S
slzhou 已提交
1496 1497
  while (1) {
    doBuildResultDatablock(pOperator, pInfo, &pAggInfo->groupResInfo, pAggInfo->aggSup.pResultBuf);
H
Haojun Liao 已提交
1498
    doFilter(pInfo->pRes, pOperator->exprSupp.pFilterInfo, NULL);
S
slzhou 已提交
1499

1500
    if (!hasRemainResults(&pAggInfo->groupResInfo)) {
H
Haojun Liao 已提交
1501
      setOperatorCompleted(pOperator);
S
slzhou 已提交
1502 1503
      break;
    }
1504

S
slzhou 已提交
1505 1506 1507 1508
    if (pInfo->pRes->info.rows > 0) {
      break;
    }
  }
1509

1510
  size_t rows = blockDataGetNumOfRows(pInfo->pRes);
1511 1512
  pOperator->resultInfo.totalRows += rows;

1513
  return (rows == 0) ? NULL : pInfo->pRes;
1514 1515
}

L
Liu Jicong 已提交
1516 1517
static void doHandleRemainBlockForNewGroupImpl(SOperatorInfo* pOperator, SFillOperatorInfo* pInfo,
                                               SResultInfo* pResultInfo, SExecTaskInfo* pTaskInfo) {
1518
  pInfo->totalInputRows = pInfo->existNewGroupBlock->info.rows;
1519 1520 1521 1522 1523
  SSDataBlock* pResBlock = pInfo->pFinalRes;

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

5
54liuyao 已提交
1525
  int64_t ekey = pInfo->existNewGroupBlock->info.window.ekey;
1526 1527
  taosResetFillInfo(pInfo->pFillInfo, getFillInfoStart(pInfo->pFillInfo));

1528
  blockDataCleanup(pInfo->pRes);
1529 1530 1531 1532
  doApplyScalarCalculation(pOperator, pInfo->existNewGroupBlock, order, scanFlag);

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

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

H
Haojun Liao 已提交
1537
  pInfo->curGroupId = pInfo->existNewGroupBlock->info.id.groupId;
1538 1539 1540
  pInfo->existNewGroupBlock = NULL;
}

L
Liu Jicong 已提交
1541 1542
static void doHandleRemainBlockFromNewGroup(SOperatorInfo* pOperator, SFillOperatorInfo* pInfo,
                                            SResultInfo* pResultInfo, SExecTaskInfo* pTaskInfo) {
1543
  if (taosFillHasMoreResults(pInfo->pFillInfo)) {
H
Haojun Liao 已提交
1544 1545
    int32_t numOfResultRows = pResultInfo->capacity - pInfo->pFinalRes->info.rows;
    taosFillResultDataBlock(pInfo->pFillInfo, pInfo->pFinalRes, numOfResultRows);
H
Haojun Liao 已提交
1546
    pInfo->pRes->info.id.groupId = pInfo->curGroupId;
1547
    return;
1548 1549 1550 1551
  }

  // handle the cached new group data block
  if (pInfo->existNewGroupBlock) {
1552 1553 1554 1555 1556 1557
    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 已提交
1558
  SExprSupp*         pSup = &pOperator->exprSupp;
1559
  setInputDataBlock(pSup, pBlock, order, scanFlag, false);
1560 1561
  projectApplyFunctions(pSup->pExprInfo, pInfo->pRes, pBlock, pSup->pCtx, pSup->numOfExprs, NULL);

1562 1563 1564 1565
  // 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);
1566

1567
  projectApplyFunctions(pNoFillSupp->pExprInfo, pInfo->pRes, pBlock, pNoFillSupp->pCtx, pNoFillSupp->numOfExprs, NULL);
H
Haojun Liao 已提交
1568
  pInfo->pRes->info.id.groupId = pBlock->info.id.groupId;
1569 1570
}

S
slzhou 已提交
1571
static SSDataBlock* doFillImpl(SOperatorInfo* pOperator) {
L
Liu Jicong 已提交
1572 1573
  SFillOperatorInfo* pInfo = pOperator->info;
  SExecTaskInfo*     pTaskInfo = pOperator->pTaskInfo;
1574

H
Haojun Liao 已提交
1575
  SResultInfo* pResultInfo = &pOperator->resultInfo;
H
Haojun Liao 已提交
1576
  SSDataBlock* pResBlock = pInfo->pFinalRes;
1577 1578

  blockDataCleanup(pResBlock);
1579

H
Haojun Liao 已提交
1580 1581
  int32_t order = TSDB_ORDER_ASC;
  int32_t scanFlag = MAIN_SCAN;
1582
  getTableScanInfo(pOperator, &order, &scanFlag);
1583

1584
  doHandleRemainBlockFromNewGroup(pOperator, pInfo, pResultInfo, pTaskInfo);
1585
  if (pResBlock->info.rows > 0) {
H
Haojun Liao 已提交
1586
    pResBlock->info.id.groupId = pInfo->curGroupId;
1587
    return pResBlock;
H
Haojun Liao 已提交
1588
  }
1589

H
Haojun Liao 已提交
1590
  SOperatorInfo* pDownstream = pOperator->pDownstream[0];
L
Liu Jicong 已提交
1591
  while (1) {
1592
    SSDataBlock* pBlock = pDownstream->fpSet.getNextFn(pDownstream);
1593 1594
    if (pBlock == NULL) {
      if (pInfo->totalInputRows == 0) {
H
Haojun Liao 已提交
1595
        setOperatorCompleted(pOperator);
1596 1597
        return NULL;
      }
1598

1599
      taosFillSetStartInfo(pInfo->pFillInfo, 0, pInfo->win.ekey);
1600
    } else {
1601
      blockDataUpdateTsWindow(pBlock, pInfo->primarySrcSlotId);
1602 1603

      blockDataCleanup(pInfo->pRes);
1604 1605
      blockDataEnsureCapacity(pInfo->pRes, pBlock->info.rows);
      blockDataEnsureCapacity(pInfo->pFinalRes, pBlock->info.rows);
1606
      doApplyScalarCalculation(pOperator, pBlock, order, scanFlag);
1607

H
Haojun Liao 已提交
1608 1609
      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 已提交
1610
        pInfo->totalInputRows += pInfo->pRes->info.rows;
1611

H
Haojun Liao 已提交
1612 1613 1614 1615 1616
        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 已提交
1617
        taosFillSetInputDataBlock(pInfo->pFillInfo, pInfo->pRes);
H
Haojun Liao 已提交
1618
      } else if (pInfo->curGroupId != pBlock->info.id.groupId) {  // the new group data block
1619 1620 1621 1622 1623
        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);
1624 1625 1626
      }
    }

1627 1628
    int32_t numOfResultRows = pOperator->resultInfo.capacity - pResBlock->info.rows;
    taosFillResultDataBlock(pInfo->pFillInfo, pResBlock, numOfResultRows);
1629 1630

    // current group has no more result to return
1631
    if (pResBlock->info.rows > 0) {
1632 1633
      // 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
1634
      if (pResBlock->info.rows > pResultInfo->threshold || pBlock == NULL || pInfo->existNewGroupBlock != NULL) {
H
Haojun Liao 已提交
1635
        pResBlock->info.id.groupId = pInfo->curGroupId;
1636
        return pResBlock;
1637 1638
      }

1639
      doHandleRemainBlockFromNewGroup(pOperator, pInfo, pResultInfo, pTaskInfo);
1640
      if (pResBlock->info.rows >= pOperator->resultInfo.threshold || pBlock == NULL) {
H
Haojun Liao 已提交
1641
        pResBlock->info.id.groupId = pInfo->curGroupId;
1642
        return pResBlock;
1643 1644 1645
      }
    } else if (pInfo->existNewGroupBlock) {  // try next group
      assert(pBlock != NULL);
1646 1647 1648 1649

      blockDataCleanup(pResBlock);

      doHandleRemainBlockForNewGroupImpl(pOperator, pInfo, pResultInfo, pTaskInfo);
1650
      if (pResBlock->info.rows > pResultInfo->threshold) {
H
Haojun Liao 已提交
1651
        pResBlock->info.id.groupId = pInfo->curGroupId;
1652
        return pResBlock;
1653 1654 1655 1656 1657 1658 1659
      }
    } else {
      return NULL;
    }
  }
}

S
slzhou 已提交
1660 1661 1662 1663 1664 1665 1666 1667
static SSDataBlock* doFill(SOperatorInfo* pOperator) {
  SFillOperatorInfo* pInfo = pOperator->info;
  SExecTaskInfo*     pTaskInfo = pOperator->pTaskInfo;

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

S
slzhou 已提交
1668
  SSDataBlock* fillResult = NULL;
S
slzhou 已提交
1669
  while (true) {
S
slzhou 已提交
1670
    fillResult = doFillImpl(pOperator);
S
slzhou 已提交
1671
    if (fillResult == NULL) {
H
Haojun Liao 已提交
1672
      setOperatorCompleted(pOperator);
S
slzhou 已提交
1673 1674 1675
      break;
    }

1676
    doFilter(fillResult, pOperator->exprSupp.pFilterInfo, &pInfo->matchInfo);
S
slzhou 已提交
1677 1678 1679 1680 1681
    if (fillResult->info.rows > 0) {
      break;
    }
  }

S
slzhou 已提交
1682
  if (fillResult != NULL) {
1683
    pOperator->resultInfo.totalRows += fillResult->info.rows;
S
slzhou 已提交
1684
  }
S
slzhou 已提交
1685

S
slzhou 已提交
1686
  return fillResult;
S
slzhou 已提交
1687 1688
}

1689
void destroyExprInfo(SExprInfo* pExpr, int32_t numOfExprs) {
C
Cary Xu 已提交
1690 1691 1692 1693 1694
  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 已提交
1695 1696
      } else if (pExprInfo->base.pParam[j].type == FUNC_PARAM_TYPE_VALUE) {
        taosVariantDestroy(&pExprInfo->base.pParam[j].param);
H
Haojun Liao 已提交
1697
      }
1698
    }
C
Cary Xu 已提交
1699 1700 1701

    taosMemoryFree(pExprInfo->base.pParam);
    taosMemoryFree(pExprInfo->pExpr);
H
Haojun Liao 已提交
1702 1703 1704
  }
}

5
54liuyao 已提交
1705
void destroyOperatorInfo(SOperatorInfo* pOperator) {
1706 1707 1708 1709
  if (pOperator == NULL) {
    return;
  }

1710
  if (pOperator->fpSet.closeFn != NULL) {
1711
    pOperator->fpSet.closeFn(pOperator->info);
1712 1713
  }

H
Haojun Liao 已提交
1714
  if (pOperator->pDownstream != NULL) {
L
Liu Jicong 已提交
1715
    for (int32_t i = 0; i < pOperator->numOfDownstream; ++i) {
H
Haojun Liao 已提交
1716
      destroyOperatorInfo(pOperator->pDownstream[i]);
1717 1718
    }

wafwerar's avatar
wafwerar 已提交
1719
    taosMemoryFreeClear(pOperator->pDownstream);
H
Haojun Liao 已提交
1720
    pOperator->numOfDownstream = 0;
1721 1722
  }

1723
  cleanupExprSupp(&pOperator->exprSupp);
wafwerar's avatar
wafwerar 已提交
1724
  taosMemoryFreeClear(pOperator);
1725 1726
}

1727 1728 1729 1730 1731 1732
int32_t getBufferPgSize(int32_t rowSize, uint32_t* defaultPgsz, uint32_t* defaultBufsz) {
  *defaultPgsz = 4096;
  while (*defaultPgsz < rowSize * 4) {
    *defaultPgsz <<= 1u;
  }

1733
  // The default buffer for each operator in query is 10MB.
1734
  // at least four pages need to be in buffer
1735 1736
  // TODO: make this variable to be configurable.
  *defaultBufsz = 4096 * 2560;
1737 1738 1739 1740 1741 1742 1743
  if ((*defaultBufsz) <= (*defaultPgsz)) {
    (*defaultBufsz) = (*defaultPgsz) * 4;
  }

  return 0;
}

dengyihao's avatar
dengyihao 已提交
1744 1745
int32_t doInitAggInfoSup(SAggSupporter* pAggSup, SqlFunctionCtx* pCtx, int32_t numOfOutput, size_t keyBufSize,
                         const char* pKey) {
L
Liu Jicong 已提交
1746
  int32_t    code = 0;
1747 1748
  _hash_fn_t hashFn = taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY);

1749
  pAggSup->currentPageId = -1;
dengyihao's avatar
dengyihao 已提交
1750 1751
  pAggSup->resultRowSize = getResultRowSize(pCtx, numOfOutput);
  pAggSup->keyBuf = taosMemoryCalloc(1, keyBufSize + POINTER_BYTES + sizeof(int64_t));
1752
  pAggSup->pResultRowHashTable = tSimpleHashInit(10, hashFn);
1753

H
Haojun Liao 已提交
1754
  if (pAggSup->keyBuf == NULL || pAggSup->pResultRowHashTable == NULL) {
1755 1756 1757
    return TSDB_CODE_OUT_OF_MEMORY;
  }

dengyihao's avatar
dengyihao 已提交
1758
  uint32_t defaultPgsz = 0;
1759 1760
  uint32_t defaultBufsz = 0;
  getBufferPgSize(pAggSup->resultRowSize, &defaultPgsz, &defaultBufsz);
H
Haojun Liao 已提交
1761

wafwerar's avatar
wafwerar 已提交
1762
  if (!osTempSpaceAvailable()) {
H
Haojun Liao 已提交
1763 1764 1765
    code = TSDB_CODE_NO_AVAIL_DISK;
    qError("Init stream agg supporter failed since %s, %s", terrstr(code), pKey);
    return code;
wafwerar's avatar
wafwerar 已提交
1766
  }
1767

H
Haojun Liao 已提交
1768
  code = createDiskbasedBuf(&pAggSup->pResultBuf, defaultPgsz, defaultBufsz, pKey, tsTempDir);
H
Haojun Liao 已提交
1769
  if (code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
1770
    qError("Create agg result buf failed since %s, %s", tstrerror(code), pKey);
H
Haojun Liao 已提交
1771 1772 1773
    return code;
  }

H
Haojun Liao 已提交
1774
  return code;
1775 1776
}

1777
void cleanupAggSup(SAggSupporter* pAggSup) {
wafwerar's avatar
wafwerar 已提交
1778
  taosMemoryFreeClear(pAggSup->keyBuf);
1779
  tSimpleHashCleanup(pAggSup->pResultRowHashTable);
H
Haojun Liao 已提交
1780
  destroyDiskbasedBuf(pAggSup->pResultBuf);
1781 1782
}

H
Haojun Liao 已提交
1783
int32_t initAggSup(SExprSupp* pSup, SAggSupporter* pAggSup, SExprInfo* pExprInfo, int32_t numOfCols, size_t keyBufSize,
L
Liu Jicong 已提交
1784
                    const char* pkey) {
1785 1786 1787 1788 1789
  int32_t code = initExprSupp(pSup, pExprInfo, numOfCols);
  if (code != TSDB_CODE_SUCCESS) {
    return code;
  }

1790 1791 1792 1793 1794
  code = doInitAggInfoSup(pAggSup, pSup->pCtx, numOfCols, keyBufSize, pkey);
  if (code != TSDB_CODE_SUCCESS) {
    return code;
  }

L
Liu Jicong 已提交
1795
  for (int32_t i = 0; i < numOfCols; ++i) {
1796
    pSup->pCtx[i].saveHandle.pBuf = pAggSup->pResultBuf;
1797 1798
  }

1799
  return TSDB_CODE_SUCCESS;
H
Haojun Liao 已提交
1800 1801
}

L
Liu Jicong 已提交
1802
void initResultSizeInfo(SResultInfo* pResultInfo, int32_t numOfRows) {
wmmhello's avatar
wmmhello 已提交
1803
  ASSERT(numOfRows != 0);
1804 1805
  pResultInfo->capacity = numOfRows;
  pResultInfo->threshold = numOfRows * 0.75;
1806

1807 1808
  if (pResultInfo->threshold == 0) {
    pResultInfo->threshold = numOfRows;
1809 1810 1811
  }
}

1812 1813 1814 1815 1816
void initBasicInfo(SOptrBasicInfo* pInfo, SSDataBlock* pBlock) {
  pInfo->pRes = pBlock;
  initResultRowInfo(&pInfo->resultRowInfo);
}

5
54liuyao 已提交
1817
void* destroySqlFunctionCtx(SqlFunctionCtx* pCtx, int32_t numOfOutput) {
1818 1819 1820 1821 1822 1823 1824 1825 1826 1827
  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);
1828
    taosMemoryFreeClear(pCtx[i].subsidiaries.buf);
1829 1830 1831 1832 1833 1834 1835 1836
    taosMemoryFree(pCtx[i].input.pData);
    taosMemoryFree(pCtx[i].input.pColumnDataAgg);
  }

  taosMemoryFreeClear(pCtx);
  return NULL;
}

1837
int32_t initExprSupp(SExprSupp* pSup, SExprInfo* pExprInfo, int32_t numOfExpr) {
1838 1839 1840 1841
  pSup->pExprInfo = pExprInfo;
  pSup->numOfExprs = numOfExpr;
  if (pSup->pExprInfo != NULL) {
    pSup->pCtx = createSqlFunctionCtx(pExprInfo, numOfExpr, &pSup->rowEntryInfoOffset);
1842 1843 1844
    if (pSup->pCtx == NULL) {
      return TSDB_CODE_OUT_OF_MEMORY;
    }
1845
  }
1846 1847

  return TSDB_CODE_SUCCESS;
1848 1849
}

1850 1851 1852 1853
void cleanupExprSupp(SExprSupp* pSupp) {
  destroySqlFunctionCtx(pSupp->pCtx, pSupp->numOfExprs);
  if (pSupp->pExprInfo != NULL) {
    destroyExprInfo(pSupp->pExprInfo, pSupp->numOfExprs);
C
Cary Xu 已提交
1854
    taosMemoryFreeClear(pSupp->pExprInfo);
1855
  }
H
Haojun Liao 已提交
1856 1857 1858 1859 1860 1861

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

1862 1863 1864
  taosMemoryFree(pSupp->rowEntryInfoOffset);
}

1865 1866
SOperatorInfo* createAggregateOperatorInfo(SOperatorInfo* downstream, SAggPhysiNode* pAggNode,
                                           SExecTaskInfo* pTaskInfo) {
wafwerar's avatar
wafwerar 已提交
1867
  SAggOperatorInfo* pInfo = taosMemoryCalloc(1, sizeof(SAggOperatorInfo));
L
Liu Jicong 已提交
1868
  SOperatorInfo*    pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
H
Haojun Liao 已提交
1869 1870 1871
  if (pInfo == NULL || pOperator == NULL) {
    goto _error;
  }
H
Haojun Liao 已提交
1872

H
Haojun Liao 已提交
1873
  SSDataBlock* pResBlock = createDataBlockFromDescNode(pAggNode->node.pOutputDataBlockDesc);
H
Haojun Liao 已提交
1874 1875 1876
  initBasicInfo(&pInfo->binfo, pResBlock);

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

1879 1880
  int32_t    num = 0;
  SExprInfo* pExprInfo = createExprInfo(pAggNode->pAggFuncs, pAggNode->pGroupKeys, &num);
H
Haojun Liao 已提交
1881
  int32_t    code = initAggSup(&pOperator->exprSupp, &pInfo->aggSup, pExprInfo, num, keyBufSize, pTaskInfo->id.str);
L
Liu Jicong 已提交
1882
  if (code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
1883 1884
    goto _error;
  }
H
Haojun Liao 已提交
1885

H
Haojun Liao 已提交
1886 1887 1888 1889 1890 1891
  int32_t    numOfScalarExpr = 0;
  SExprInfo* pScalarExprInfo = NULL;
  if (pAggNode->pExprs != NULL) {
    pScalarExprInfo = createExprInfo(pAggNode->pExprs, NULL, &numOfScalarExpr);
  }

1892 1893 1894 1895
  code = initExprSupp(&pInfo->scalarExprSup, pScalarExprInfo, numOfScalarExpr);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
1896

H
Haojun Liao 已提交
1897 1898 1899 1900 1901
  code = filterInitFromNode((SNode*)pAggNode->node.pConditions, &pOperator->exprSupp.pFilterInfo, 0);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }

H
Haojun Liao 已提交
1902
  pInfo->binfo.mergeResultBlock = pAggNode->mergeDataBlock;
1903
  pInfo->groupId = UINT64_MAX;
H
Haojun Liao 已提交
1904

1905 1906 1907
  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 已提交
1908

1909 1910
  if (downstream->operatorType == QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN) {
    STableScanInfo* pTableScanInfo = downstream->info;
H
Haojun Liao 已提交
1911 1912
    pTableScanInfo->base.pdInfo.pExprSup = &pOperator->exprSupp;
    pTableScanInfo->base.pdInfo.pAggSup = &pInfo->aggSup;
1913 1914
  }

H
Haojun Liao 已提交
1915 1916 1917 1918
  code = appendDownstream(pOperator, &downstream, 1);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
1919 1920

  return pOperator;
H
Haojun Liao 已提交
1921

1922
_error:
H
Haojun Liao 已提交
1923 1924 1925 1926
  if (pInfo != NULL) {
    destroyAggOperatorInfo(pInfo);
  }

1927 1928 1929
  if (pOperator != NULL) {
    cleanupExprSupp(&pOperator->exprSupp);
  }
H
Haojun Liao 已提交
1930

1931
  taosMemoryFreeClear(pOperator);
H
Haojun Liao 已提交
1932
  pTaskInfo->code = code;
H
Haojun Liao 已提交
1933
  return NULL;
1934 1935
}

1936
void cleanupBasicInfo(SOptrBasicInfo* pInfo) {
1937
  assert(pInfo != NULL);
H
Haojun Liao 已提交
1938
  pInfo->pRes = blockDataDestroy(pInfo->pRes);
1939 1940
}

H
Haojun Liao 已提交
1941 1942 1943 1944 1945 1946 1947
static void freeItem(void* pItem) {
  void** p = pItem;
  if (*p != NULL) {
    taosMemoryFreeClear(*p);
  }
}

1948
void destroyAggOperatorInfo(void* param) {
L
Liu Jicong 已提交
1949
  SAggOperatorInfo* pInfo = (SAggOperatorInfo*)param;
L
Liu Jicong 已提交
1950 1951
  cleanupBasicInfo(&pInfo->binfo);

H
Haojun Liao 已提交
1952
  cleanupAggSup(&pInfo->aggSup);
S
shenglian zhou 已提交
1953
  cleanupExprSupp(&pInfo->scalarExprSup);
H
Haojun Liao 已提交
1954
  cleanupGroupResInfo(&pInfo->groupResInfo);
D
dapan1121 已提交
1955
  taosMemoryFreeClear(param);
1956
}
1957

1958
void destroyFillOperatorInfo(void* param) {
L
Liu Jicong 已提交
1959
  SFillOperatorInfo* pInfo = (SFillOperatorInfo*)param;
1960
  pInfo->pFillInfo = taosDestroyFillInfo(pInfo->pFillInfo);
H
Haojun Liao 已提交
1961
  pInfo->pRes = blockDataDestroy(pInfo->pRes);
H
Haojun Liao 已提交
1962 1963
  pInfo->pFinalRes = blockDataDestroy(pInfo->pFinalRes);

1964
  cleanupExprSupp(&pInfo->noFillExprSupp);
H
Haojun Liao 已提交
1965

wafwerar's avatar
wafwerar 已提交
1966
  taosMemoryFreeClear(pInfo->p);
H
Haojun Liao 已提交
1967
  taosArrayDestroy(pInfo->matchInfo.pList);
D
dapan1121 已提交
1968
  taosMemoryFreeClear(param);
1969 1970
}

H
Haojun Liao 已提交
1971 1972 1973 1974
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 已提交
1975

H
Haojun Liao 已提交
1976 1977 1978
  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 已提交
1979

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

H
Haojun Liao 已提交
1983 1984 1985 1986 1987 1988 1989
  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 已提交
1990
  pInfo->p = taosMemoryCalloc(numOfCols, POINTER_BYTES);
1991

H
Haojun Liao 已提交
1992
  if (pInfo->pFillInfo == NULL || pInfo->p == NULL) {
H
Haojun Liao 已提交
1993 1994
    taosMemoryFree(pInfo->pFillInfo);
    taosMemoryFree(pInfo->p);
H
Haojun Liao 已提交
1995 1996 1997 1998 1999 2000
    return TSDB_CODE_OUT_OF_MEMORY;
  } else {
    return TSDB_CODE_SUCCESS;
  }
}

2001
static bool isWstartColumnExist(SFillOperatorInfo* pInfo) {
2002
  if (pInfo->noFillExprSupp.numOfExprs == 0) {
2003 2004
    return false;
  }
2005 2006 2007

  for (int32_t i = 0; i < pInfo->noFillExprSupp.numOfExprs; ++i) {
    SExprInfo* exprInfo = pInfo->noFillExprSupp.pExprInfo + i;
2008
    if (exprInfo->pExpr->nodeType == QUERY_NODE_COLUMN && exprInfo->base.numOfParams == 1 &&
2009 2010 2011 2012 2013 2014 2015
        exprInfo->base.pParam[0].pCol->colType == COLUMN_TYPE_WINDOW_START) {
      return true;
    }
  }
  return false;
}

2016 2017
static int32_t createPrimaryTsExprIfNeeded(SFillOperatorInfo* pInfo, SFillPhysiNode* pPhyFillNode, SExprSupp* pExprSupp,
                                           const char* idStr) {
2018
  bool wstartExist = isWstartColumnExist(pInfo);
2019

2020 2021
  if (wstartExist == false) {
    if (pPhyFillNode->pWStartTs->type != QUERY_NODE_TARGET) {
2022
      qError("pWStartTs of fill physical node is not a target node, %s", idStr);
2023 2024 2025
      return TSDB_CODE_QRY_SYS_ERROR;
    }

2026 2027
    SExprInfo* pExpr = taosMemoryRealloc(pExprSupp->pExprInfo, (pExprSupp->numOfExprs + 1) * sizeof(SExprInfo));
    if (pExpr == NULL) {
2028 2029 2030
      return TSDB_CODE_OUT_OF_MEMORY;
    }

2031 2032 2033
    createExprFromTargetNode(&pExpr[pExprSupp->numOfExprs], (STargetNode*)pPhyFillNode->pWStartTs);
    pExprSupp->numOfExprs += 1;
    pExprSupp->pExprInfo = pExpr;
2034
  }
2035

2036 2037 2038
  return TSDB_CODE_SUCCESS;
}

L
Liu Jicong 已提交
2039 2040
SOperatorInfo* createFillOperatorInfo(SOperatorInfo* downstream, SFillPhysiNode* pPhyFillNode,
                                      SExecTaskInfo* pTaskInfo) {
2041 2042 2043 2044 2045 2046
  SFillOperatorInfo* pInfo = taosMemoryCalloc(1, sizeof(SFillOperatorInfo));
  SOperatorInfo*     pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
  if (pInfo == NULL || pOperator == NULL) {
    goto _error;
  }

H
Haojun Liao 已提交
2047
  pInfo->pRes = createDataBlockFromDescNode(pPhyFillNode->node.pOutputDataBlockDesc);
H
Haojun Liao 已提交
2048 2049 2050
  SExprInfo* pExprInfo = createExprInfo(pPhyFillNode->pFillExprs, NULL, &pInfo->numOfExpr);
  pOperator->exprSupp.pExprInfo = pExprInfo;

2051 2052 2053 2054 2055 2056 2057 2058
  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);
2059 2060 2061
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
H
Haojun Liao 已提交
2062

L
Liu Jicong 已提交
2063
  SInterval* pInterval =
2064
      QUERY_NODE_PHYSICAL_PLAN_MERGE_ALIGNED_INTERVAL == downstream->operatorType
2065 2066
          ? &((SMergeAlignedIntervalAggOperatorInfo*)downstream->info)->intervalAggOperatorInfo->interval
          : &((SIntervalAggOperatorInfo*)downstream->info)->interval;
2067

2068
  int32_t order = (pPhyFillNode->inputTsOrder == ORDER_ASC) ? TSDB_ORDER_ASC : TSDB_ORDER_DESC;
2069
  int32_t type = convertFillType(pPhyFillNode->mode);
2070

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

2073
  initResultSizeInfo(&pOperator->resultInfo, 4096);
H
Haojun Liao 已提交
2074 2075 2076 2077 2078
  blockDataEnsureCapacity(pInfo->pRes, pOperator->resultInfo.capacity);
  code = initExprSupp(&pOperator->exprSupp, pExprInfo, pInfo->numOfExpr);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
H
Haojun Liao 已提交
2079

H
Haojun Liao 已提交
2080 2081
  pInfo->primaryTsCol = ((STargetNode*)pPhyFillNode->pWStartTs)->slotId;
  pInfo->primarySrcSlotId = ((SColumnNode*)((STargetNode*)pPhyFillNode->pWStartTs)->pExpr)->slotId;
2082

2083
  int32_t numOfOutputCols = 0;
H
Haojun Liao 已提交
2084 2085
  code = extractColMatchInfo(pPhyFillNode->pFillExprs, pPhyFillNode->node.pOutputDataBlockDesc, &numOfOutputCols,
                             COL_MATCH_FROM_SLOT_ID, &pInfo->matchInfo);
2086

2087
  code = initFillInfo(pInfo, pExprInfo, pInfo->numOfExpr, pNoFillSupp->pExprInfo, pNoFillSupp->numOfExprs,
2088 2089
                      (SNodeListNode*)pPhyFillNode->pValues, pPhyFillNode->timeRange, pResultInfo->capacity,
                      pTaskInfo->id.str, pInterval, type, order);
2090 2091 2092
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
2093

H
Haojun Liao 已提交
2094
  pInfo->pFinalRes = createOneDataBlock(pInfo->pRes, false);
H
Haojun Liao 已提交
2095 2096
  blockDataEnsureCapacity(pInfo->pFinalRes, pOperator->resultInfo.capacity);

H
Haojun Liao 已提交
2097 2098 2099 2100 2101
  code = filterInitFromNode((SNode*)pPhyFillNode->node.pConditions, &pOperator->exprSupp.pFilterInfo, 0);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }

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

2106
  code = appendDownstream(pOperator, &downstream, 1);
2107
  return pOperator;
H
Haojun Liao 已提交
2108

2109
_error:
H
Haojun Liao 已提交
2110 2111 2112 2113
  if (pInfo != NULL) {
    destroyFillOperatorInfo(pInfo);
  }

2114
  pTaskInfo->code = code;
wafwerar's avatar
wafwerar 已提交
2115
  taosMemoryFreeClear(pOperator);
H
Haojun Liao 已提交
2116
  return NULL;
2117 2118
}

D
dapan1121 已提交
2119
static SExecTaskInfo* createExecTaskInfo(uint64_t queryId, uint64_t taskId, EOPTR_EXEC_MODEL model, char* dbFName) {
wafwerar's avatar
wafwerar 已提交
2120
  SExecTaskInfo* pTaskInfo = taosMemoryCalloc(1, sizeof(SExecTaskInfo));
H
Haojun Liao 已提交
2121 2122 2123 2124 2125
  if (pTaskInfo == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return NULL;
  }

2126
  setTaskStatus(pTaskInfo, TASK_NOT_COMPLETED);
H
Haojun Liao 已提交
2127

2128
  pTaskInfo->schemaInfo.dbname = strdup(dbFName);
2129
  pTaskInfo->cost.created = taosGetTimestampMs();
H
Haojun Liao 已提交
2130
  pTaskInfo->id.queryId = queryId;
dengyihao's avatar
dengyihao 已提交
2131
  pTaskInfo->execModel = model;
H
Haojun Liao 已提交
2132
  pTaskInfo->pTableInfoList = tableListCreate();
D
dapan1121 已提交
2133
  pTaskInfo->stopInfo.pStopInfo = taosArrayInit(4, sizeof(SExchangeOpStopInfo));
H
Haojun Liao 已提交
2134
  pTaskInfo->pResultBlockList = taosArrayInit(128, POINTER_BYTES);
H
Haojun Liao 已提交
2135

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

2140 2141
  return pTaskInfo;
}
H
Haojun Liao 已提交
2142

H
Haojun Liao 已提交
2143 2144
SSchemaWrapper* extractQueriedColumnSchema(SScanPhysiNode* pScanNode);

2145
int32_t extractTableSchemaInfo(SReadHandle* pHandle, SScanPhysiNode* pScanNode, SExecTaskInfo* pTaskInfo) {
2146 2147
  SMetaReader mr = {0};
  metaReaderInit(&mr, pHandle->meta, 0);
H
Haojun Liao 已提交
2148
  int32_t code = metaGetTableEntryByUidCache(&mr, pScanNode->uid);
2149
  if (code != TSDB_CODE_SUCCESS) {
L
Liu Jicong 已提交
2150 2151
    qError("failed to get the table meta, uid:0x%" PRIx64 ", suid:0x%" PRIx64 ", %s", pScanNode->uid, pScanNode->suid,
           GET_TASKID(pTaskInfo));
H
Haojun Liao 已提交
2152

D
dapan1121 已提交
2153
    metaReaderClear(&mr);
2154
    return terrno;
D
dapan1121 已提交
2155
  }
2156

2157 2158
  SSchemaInfo* pSchemaInfo = &pTaskInfo->schemaInfo;
  pSchemaInfo->tablename = strdup(mr.me.name);
2159 2160

  if (mr.me.type == TSDB_SUPER_TABLE) {
2161 2162
    pSchemaInfo->sw = tCloneSSchemaWrapper(&mr.me.stbEntry.schemaRow);
    pSchemaInfo->tversion = mr.me.stbEntry.schemaTag.version;
2163
  } else if (mr.me.type == TSDB_CHILD_TABLE) {
2164 2165
    tDecoderClear(&mr.coder);

2166
    tb_uid_t suid = mr.me.ctbEntry.suid;
H
Haojun Liao 已提交
2167
    metaGetTableEntryByUidCache(&mr, suid);
2168 2169
    pSchemaInfo->sw = tCloneSSchemaWrapper(&mr.me.stbEntry.schemaRow);
    pSchemaInfo->tversion = mr.me.stbEntry.schemaTag.version;
2170
  } else {
2171
    pSchemaInfo->sw = tCloneSSchemaWrapper(&mr.me.ntbEntry.schemaRow);
2172
  }
2173 2174

  metaReaderClear(&mr);
2175

H
Haojun Liao 已提交
2176 2177 2178 2179 2180
  pSchemaInfo->qsw = extractQueriedColumnSchema(pScanNode);
  return TSDB_CODE_SUCCESS;
}

SSchemaWrapper* extractQueriedColumnSchema(SScanPhysiNode* pScanNode) {
2181 2182 2183
  int32_t numOfCols = LIST_LENGTH(pScanNode->pScanCols);
  int32_t numOfTags = LIST_LENGTH(pScanNode->pScanPseudoCols);

2184
  SSchemaWrapper* pqSw = taosMemoryCalloc(1, sizeof(SSchemaWrapper));
2185
  pqSw->pSchema = taosMemoryCalloc(numOfCols + numOfTags, sizeof(SSchema));
2186

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

H
Haojun Liao 已提交
2191 2192 2193
    SSchema* pSchema = &pqSw->pSchema[pqSw->nCols++];
    pSchema->colId = pColNode->colId;
    pSchema->type = pColNode->node.resType.type;
H
Haojun Liao 已提交
2194 2195
    pSchema->bytes = pColNode->node.resType.bytes;
    tstrncpy(pSchema->name, pColNode->colName, tListLen(pSchema->name));
2196 2197
  }

2198
  // this the tags and pseudo function columns, we only keep the tag columns
2199
  for (int32_t i = 0; i < numOfTags; ++i) {
2200 2201 2202 2203 2204 2205 2206 2207 2208
    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 已提交
2209
      pSchema->bytes = pColNode->node.resType.bytes;
H
Haojun Liao 已提交
2210
      tstrncpy(pSchema->name, pColNode->colName, tListLen(pSchema->name));
2211 2212 2213
    }
  }

H
Haojun Liao 已提交
2214
  return pqSw;
2215 2216
}

2217 2218
static void cleanupTableSchemaInfo(SSchemaInfo* pSchemaInfo) {
  taosMemoryFreeClear(pSchemaInfo->dbname);
2219
  taosMemoryFreeClear(pSchemaInfo->tablename);
2220 2221
  tDeleteSSchemaWrapper(pSchemaInfo->sw);
  tDeleteSSchemaWrapper(pSchemaInfo->qsw);
2222 2223
}

2224
static void cleanupStreamInfo(SStreamTaskInfo* pStreamInfo) { tDeleteSSchemaWrapper(pStreamInfo->schema); }
2225

2226
bool groupbyTbname(SNodeList* pGroupList) {
2227
  bool bytbname = false;
2228
  if (LIST_LENGTH(pGroupList) == 1) {
2229 2230 2231 2232 2233 2234 2235 2236 2237 2238
    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;
}

2239 2240
SOperatorInfo* createOperatorTree(SPhysiNode* pPhyNode, SExecTaskInfo* pTaskInfo, SReadHandle* pHandle, SNode* pTagCond,
                                  SNode* pTagIndexCond, const char* pUser) {
2241
  int32_t         type = nodeType(pPhyNode);
2242
  STableListInfo* pTableListInfo = pTaskInfo->pTableInfoList;
2243
  const char*     idstr = GET_TASKID(pTaskInfo);
2244

X
Xiaoyu Wang 已提交
2245
  if (pPhyNode->pChildren == NULL || LIST_LENGTH(pPhyNode->pChildren) == 0) {
2246
    SOperatorInfo* pOperator = NULL;
H
Haojun Liao 已提交
2247
    if (QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN == type) {
dengyihao's avatar
dengyihao 已提交
2248
      STableScanPhysiNode* pTableScanNode = (STableScanPhysiNode*)pPhyNode;
H
Haojun Liao 已提交
2249

2250 2251 2252 2253 2254 2255
      // 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 已提交
2256 2257
      int32_t code =
          createScanTableListInfo(&pTableScanNode->scan, pTableScanNode->pGroupTags, pTableScanNode->groupSort, pHandle,
H
Haojun Liao 已提交
2258
                                  pTableListInfo, pTagCond, pTagIndexCond, pTaskInfo);
2259
      if (code) {
wmmhello's avatar
wmmhello 已提交
2260
        pTaskInfo->code = code;
2261
        qError("failed to createScanTableListInfo, code:%s, %s", tstrerror(code), idstr);
D
dapan1121 已提交
2262 2263
        return NULL;
      }
wmmhello's avatar
wmmhello 已提交
2264

2265
      code = extractTableSchemaInfo(pHandle, &pTableScanNode->scan, pTaskInfo);
S
slzhou 已提交
2266
      if (code) {
2267
        pTaskInfo->code = terrno;
wmmhello's avatar
wmmhello 已提交
2268 2269 2270
        return NULL;
      }

2271
      pOperator = createTableScanOperatorInfo(pTableScanNode, pHandle, pTaskInfo);
D
dapan1121 已提交
2272 2273 2274 2275 2276
      if (NULL == pOperator) {
        pTaskInfo->code = terrno;
        return NULL;
      }

2277
      STableScanInfo* pScanInfo = pOperator->info;
H
Haojun Liao 已提交
2278
      pTaskInfo->cost.pRecoder = &pScanInfo->base.readRecorder;
S
slzhou 已提交
2279 2280
    } else if (QUERY_NODE_PHYSICAL_PLAN_TABLE_MERGE_SCAN == type) {
      STableMergeScanPhysiNode* pTableScanNode = (STableMergeScanPhysiNode*)pPhyNode;
H
Haojun Liao 已提交
2281 2282 2283

      int32_t code = createScanTableListInfo(&pTableScanNode->scan, pTableScanNode->pGroupTags, true, pHandle,
                                             pTableListInfo, pTagCond, pTagIndexCond, pTaskInfo);
L
Liu Jicong 已提交
2284
      if (code) {
wmmhello's avatar
wmmhello 已提交
2285
        pTaskInfo->code = code;
H
Haojun Liao 已提交
2286
        qError("failed to createScanTableListInfo, code: %s", tstrerror(code));
wmmhello's avatar
wmmhello 已提交
2287 2288
        return NULL;
      }
2289

2290
      code = extractTableSchemaInfo(pHandle, &pTableScanNode->scan, pTaskInfo);
wmmhello's avatar
wmmhello 已提交
2291 2292 2293 2294
      if (code) {
        pTaskInfo->code = terrno;
        return NULL;
      }
wmmhello's avatar
wmmhello 已提交
2295

H
Haojun Liao 已提交
2296
      pOperator = createTableMergeScanOperatorInfo(pTableScanNode, pHandle, pTaskInfo);
D
dapan1121 已提交
2297 2298 2299 2300
      if (NULL == pOperator) {
        pTaskInfo->code = terrno;
        return NULL;
      }
wmmhello's avatar
wmmhello 已提交
2301

2302
      STableScanInfo* pScanInfo = pOperator->info;
H
Haojun Liao 已提交
2303
      pTaskInfo->cost.pRecoder = &pScanInfo->base.readRecorder;
H
Haojun Liao 已提交
2304
    } else if (QUERY_NODE_PHYSICAL_PLAN_EXCHANGE == type) {
2305 2306
      pOperator = createExchangeOperatorInfo(pHandle ? pHandle->pMsgCb->clientRpc : NULL, (SExchangePhysiNode*)pPhyNode,
                                             pTaskInfo);
H
Haojun Liao 已提交
2307
    } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_SCAN == type) {
5
54liuyao 已提交
2308
      STableScanPhysiNode* pTableScanNode = (STableScanPhysiNode*)pPhyNode;
5
54liuyao 已提交
2309
      if (pHandle->vnode) {
L
Liu Jicong 已提交
2310 2311
        int32_t code =
            createScanTableListInfo(&pTableScanNode->scan, pTableScanNode->pGroupTags, pTableScanNode->groupSort,
H
Haojun Liao 已提交
2312
                                    pHandle, pTableListInfo, pTagCond, pTagIndexCond, pTaskInfo);
L
Liu Jicong 已提交
2313
        if (code) {
wmmhello's avatar
wmmhello 已提交
2314
          pTaskInfo->code = code;
H
Haojun Liao 已提交
2315
          qError("failed to createScanTableListInfo, code: %s", tstrerror(code));
wmmhello's avatar
wmmhello 已提交
2316 2317
          return NULL;
        }
L
Liu Jicong 已提交
2318 2319

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

L
Liu Jicong 已提交
2323
        for (int32_t i = 0; i < sz; i++) {
H
Haojun Liao 已提交
2324
          STableKeyInfo* pKeyInfo = tableListGetInfo(pTableListInfo, i);
2325
          qDebug("add table uid:%" PRIu64 ", gid:%" PRIu64, pKeyInfo->uid, pKeyInfo->groupId);
L
Liu Jicong 已提交
2326 2327
        }
#endif
2328
      }
2329

H
Haojun Liao 已提交
2330
      pTaskInfo->schemaInfo.qsw = extractQueriedColumnSchema(&pTableScanNode->scan);
2331
      pOperator = createStreamScanOperatorInfo(pHandle, pTableScanNode, pTagCond, pTaskInfo);
H
Haojun Liao 已提交
2332
    } else if (QUERY_NODE_PHYSICAL_PLAN_SYSTABLE_SCAN == type) {
L
Liu Jicong 已提交
2333
      SSystemTableScanPhysiNode* pSysScanPhyNode = (SSystemTableScanPhysiNode*)pPhyNode;
2334
      pOperator = createSysTableScanOperatorInfo(pHandle, pSysScanPhyNode, pUser, pTaskInfo);
2335
    } else if (QUERY_NODE_PHYSICAL_PLAN_TAG_SCAN == type) {
X
Xiaoyu Wang 已提交
2336
      STagScanPhysiNode* pScanPhyNode = (STagScanPhysiNode*)pPhyNode;
2337 2338

      int32_t code = createScanTableListInfo(pScanPhyNode, NULL, false, pHandle, pTableListInfo, pTagCond,
H
Haojun Liao 已提交
2339
                                             pTagIndexCond, pTaskInfo);
2340
      if (code != TSDB_CODE_SUCCESS) {
2341
        pTaskInfo->code = code;
H
Haojun Liao 已提交
2342
        qError("failed to getTableList, code: %s", tstrerror(code));
2343 2344 2345
        return NULL;
      }

H
Haojun Liao 已提交
2346
      pOperator = createTagScanOperatorInfo(pHandle, pScanPhyNode, pTaskInfo);
2347
    } else if (QUERY_NODE_PHYSICAL_PLAN_BLOCK_DIST_SCAN == type) {
2348
      SBlockDistScanPhysiNode* pBlockNode = (SBlockDistScanPhysiNode*)pPhyNode;
2349 2350

      if (pBlockNode->tableType == TSDB_SUPER_TABLE) {
H
Haojun Liao 已提交
2351 2352
        SArray* pList = taosArrayInit(4, sizeof(STableKeyInfo));
        int32_t code = vnodeGetAllTableList(pHandle->vnode, pBlockNode->uid, pList);
2353 2354 2355 2356
        if (code != TSDB_CODE_SUCCESS) {
          pTaskInfo->code = terrno;
          return NULL;
        }
H
Haojun Liao 已提交
2357

2358
        for (int32_t i = 0; i < tableListGetSize(pTableListInfo); ++i) {
H
Haojun Liao 已提交
2359
          STableKeyInfo* p = taosArrayGet(pList, i);
H
Haojun Liao 已提交
2360
          tableListAddTableInfo(pTableListInfo, p->uid, 0);
H
Haojun Liao 已提交
2361 2362
        }
        taosArrayDestroy(pList);
2363
      } else {  // Create group with only one table
H
Haojun Liao 已提交
2364
        tableListAddTableInfo(pTableListInfo, pBlockNode->uid, 0);
2365 2366
      }

2367
      pOperator = createDataBlockInfoScanOperator(pHandle, pBlockNode, pTaskInfo);
H
Haojun Liao 已提交
2368 2369 2370
    } else if (QUERY_NODE_PHYSICAL_PLAN_LAST_ROW_SCAN == type) {
      SLastRowScanPhysiNode* pScanNode = (SLastRowScanPhysiNode*)pPhyNode;

L
Liu Jicong 已提交
2371
      int32_t code = createScanTableListInfo(&pScanNode->scan, pScanNode->pGroupTags, true, pHandle, pTableListInfo,
H
Haojun Liao 已提交
2372
                                             pTagCond, pTagIndexCond, pTaskInfo);
2373 2374 2375 2376
      if (code != TSDB_CODE_SUCCESS) {
        pTaskInfo->code = code;
        return NULL;
      }
2377

2378
      code = extractTableSchemaInfo(pHandle, &pScanNode->scan, pTaskInfo);
2379 2380 2381
      if (code != TSDB_CODE_SUCCESS) {
        pTaskInfo->code = code;
        return NULL;
H
Haojun Liao 已提交
2382 2383
      }

2384
      pOperator = createCacherowsScanOperator(pScanNode, pHandle, pTaskInfo);
2385
    } else if (QUERY_NODE_PHYSICAL_PLAN_PROJECT == type) {
2386
      pOperator = createProjectOperatorInfo(NULL, (SProjectPhysiNode*)pPhyNode, pTaskInfo);
H
Haojun Liao 已提交
2387 2388
    } else {
      ASSERT(0);
H
Haojun Liao 已提交
2389
    }
2390 2391 2392 2393 2394

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

2395
    return pOperator;
H
Haojun Liao 已提交
2396 2397
  }

2398
  size_t          size = LIST_LENGTH(pPhyNode->pChildren);
2399
  SOperatorInfo** ops = taosMemoryCalloc(size, POINTER_BYTES);
2400 2401 2402 2403
  if (ops == NULL) {
    return NULL;
  }

dengyihao's avatar
dengyihao 已提交
2404
  for (int32_t i = 0; i < size; ++i) {
2405
    SPhysiNode* pChildNode = (SPhysiNode*)nodesListGetNode(pPhyNode->pChildren, i);
2406
    ops[i] = createOperatorTree(pChildNode, pTaskInfo, pHandle, pTagCond, pTagIndexCond, pUser);
2407
    if (ops[i] == NULL) {
H
Haojun Liao 已提交
2408
      taosMemoryFree(ops);
2409 2410
      return NULL;
    }
2411
  }
H
Haojun Liao 已提交
2412

2413
  SOperatorInfo* pOptr = NULL;
H
Haojun Liao 已提交
2414
  if (QUERY_NODE_PHYSICAL_PLAN_PROJECT == type) {
2415
    pOptr = createProjectOperatorInfo(ops[0], (SProjectPhysiNode*)pPhyNode, pTaskInfo);
2416
  } else if (QUERY_NODE_PHYSICAL_PLAN_HASH_AGG == type) {
H
Haojun Liao 已提交
2417 2418
    SAggPhysiNode* pAggNode = (SAggPhysiNode*)pPhyNode;
    if (pAggNode->pGroupKeys != NULL) {
H
Haojun Liao 已提交
2419
      pOptr = createGroupOperatorInfo(ops[0], pAggNode, pTaskInfo);
H
Haojun Liao 已提交
2420
    } else {
H
Haojun Liao 已提交
2421
      pOptr = createAggregateOperatorInfo(ops[0], pAggNode, pTaskInfo);
H
Haojun Liao 已提交
2422
    }
2423
  } else if (QUERY_NODE_PHYSICAL_PLAN_HASH_INTERVAL == type) {
H
Haojun Liao 已提交
2424
    SIntervalPhysiNode* pIntervalPhyNode = (SIntervalPhysiNode*)pPhyNode;
H
Haojun Liao 已提交
2425

H
Haojun Liao 已提交
2426 2427
    bool isStream = (QUERY_NODE_PHYSICAL_PLAN_STREAM_INTERVAL == type);
    pOptr = createIntervalOperatorInfo(ops[0], pIntervalPhyNode, pTaskInfo, isStream);
2428
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_INTERVAL == type) {
2429
    pOptr = createStreamIntervalOperatorInfo(ops[0], pPhyNode, pTaskInfo);
2430 2431
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE_ALIGNED_INTERVAL == type) {
    SMergeAlignedIntervalPhysiNode* pIntervalPhyNode = (SMergeAlignedIntervalPhysiNode*)pPhyNode;
2432
    pOptr = createMergeAlignedIntervalOperatorInfo(ops[0], pIntervalPhyNode, pTaskInfo);
S
shenglian zhou 已提交
2433
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE_INTERVAL == type) {
X
Xiaoyu Wang 已提交
2434
    SMergeIntervalPhysiNode* pIntervalPhyNode = (SMergeIntervalPhysiNode*)pPhyNode;
2435
    pOptr = createMergeIntervalOperatorInfo(ops[0], pIntervalPhyNode, pTaskInfo);
5
54liuyao 已提交
2436
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_SEMI_INTERVAL == type) {
2437
    int32_t children = 0;
5
54liuyao 已提交
2438 2439
    pOptr = createStreamFinalIntervalOperatorInfo(ops[0], pPhyNode, pTaskInfo, children);
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_FINAL_INTERVAL == type) {
5
54liuyao 已提交
2440
    int32_t children = pHandle->numOfVgroups;
5
54liuyao 已提交
2441
    pOptr = createStreamFinalIntervalOperatorInfo(ops[0], pPhyNode, pTaskInfo, children);
H
Haojun Liao 已提交
2442
  } else if (QUERY_NODE_PHYSICAL_PLAN_SORT == type) {
2443
    pOptr = createSortOperatorInfo(ops[0], (SSortPhysiNode*)pPhyNode, pTaskInfo);
S
shenglian zhou 已提交
2444 2445
  } else if (QUERY_NODE_PHYSICAL_PLAN_GROUP_SORT == type) {
    pOptr = createGroupSortOperatorInfo(ops[0], (SGroupSortPhysiNode*)pPhyNode, pTaskInfo);
X
Xiaoyu Wang 已提交
2446
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE == type) {
2447
    SMergePhysiNode* pMergePhyNode = (SMergePhysiNode*)pPhyNode;
2448
    pOptr = createMultiwayMergeOperatorInfo(ops, size, pMergePhyNode, pTaskInfo);
2449
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE_SESSION == type) {
H
Haojun Liao 已提交
2450
    SSessionWinodwPhysiNode* pSessionNode = (SSessionWinodwPhysiNode*)pPhyNode;
H
Haojun Liao 已提交
2451
    pOptr = createSessionAggOperatorInfo(ops[0], pSessionNode, pTaskInfo);
2452
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_SESSION == type) {
2453 2454 2455 2456 2457
    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) {
2458
    int32_t children = pHandle->numOfVgroups;
2459
    pOptr = createStreamFinalSessionAggOperatorInfo(ops[0], pPhyNode, pTaskInfo, children);
H
Haojun Liao 已提交
2460
  } else if (QUERY_NODE_PHYSICAL_PLAN_PARTITION == type) {
2461
    pOptr = createPartitionOperatorInfo(ops[0], (SPartitionPhysiNode*)pPhyNode, pTaskInfo);
2462
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_PARTITION == type) {
2463
    pOptr = createStreamPartitionOperatorInfo(ops[0], (SStreamPartitionPhysiNode*)pPhyNode, pTaskInfo);
2464
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE_STATE == type) {
dengyihao's avatar
dengyihao 已提交
2465
    SStateWinodwPhysiNode* pStateNode = (SStateWinodwPhysiNode*)pPhyNode;
2466
    pOptr = createStatewindowOperatorInfo(ops[0], pStateNode, pTaskInfo);
2467
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_STATE == type) {
5
54liuyao 已提交
2468
    pOptr = createStreamStateAggOperatorInfo(ops[0], pPhyNode, pTaskInfo);
2469
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE_JOIN == type) {
2470
    pOptr = createMergeJoinOperatorInfo(ops, size, (SSortMergeJoinPhysiNode*)pPhyNode, pTaskInfo);
2471
  } else if (QUERY_NODE_PHYSICAL_PLAN_FILL == type) {
H
Haojun Liao 已提交
2472
    pOptr = createFillOperatorInfo(ops[0], (SFillPhysiNode*)pPhyNode, pTaskInfo);
5
54liuyao 已提交
2473 2474
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_FILL == type) {
    pOptr = createStreamFillOperatorInfo(ops[0], (SStreamFillPhysiNode*)pPhyNode, pTaskInfo);
H
Haojun Liao 已提交
2475 2476
  } else if (QUERY_NODE_PHYSICAL_PLAN_INDEF_ROWS_FUNC == type) {
    pOptr = createIndefinitOutputOperatorInfo(ops[0], pPhyNode, pTaskInfo);
2477 2478
  } else if (QUERY_NODE_PHYSICAL_PLAN_INTERP_FUNC == type) {
    pOptr = createTimeSliceOperatorInfo(ops[0], pPhyNode, pTaskInfo);
H
Haojun Liao 已提交
2479 2480
  } else {
    ASSERT(0);
H
Haojun Liao 已提交
2481
  }
2482

2483
  taosMemoryFree(ops);
2484 2485 2486 2487
  if (pOptr) {
    pOptr->resultDataBlockId = pPhyNode->pOutputDataBlockDesc->dataBlockId;
  }

2488
  return pOptr;
2489
}
H
Haojun Liao 已提交
2490

L
Liu Jicong 已提交
2491 2492 2493 2494 2495 2496 2497 2498 2499 2500 2501 2502 2503
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 {
2504 2505 2506
    SStreamScanInfo* pInfo = pOperator->info;
    ASSERT(pInfo->pTableScanOp->operatorType == QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN);
    *ppInfo = pInfo->pTableScanOp->info;
L
Liu Jicong 已提交
2507 2508 2509 2510
    return 0;
  }
}

2511 2512 2513 2514 2515 2516 2517 2518 2519 2520 2521 2522 2523 2524 2525 2526 2527 2528 2529 2530 2531 2532
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;
}

2533
#if 0
L
Liu Jicong 已提交
2534 2535 2536 2537 2538
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;
  }
2539

L
Liu Jicong 已提交
2540 2541 2542 2543
  STableScanPhysiNode* pNode = NULL;
  if (extractTableScanNode(plan->pNode, &pNode) < 0) {
    ASSERT(0);
  }
2544

H
Haojun Liao 已提交
2545
  tsdbReaderClose(pTableScanInfo->dataReader);
2546

L
Liu Jicong 已提交
2547
  STableListInfo info = {0};
H
Haojun Liao 已提交
2548
  pTableScanInfo->dataReader = doCreateDataReader(pNode, pHandle, &info, NULL);
L
Liu Jicong 已提交
2549 2550 2551 2552
  if (pTableScanInfo->dataReader == NULL) {
    ASSERT(0);
    qError("failed to create data reader");
    return TSDB_CODE_QRY_APP_ERROR;
2553
  }
L
Liu Jicong 已提交
2554
  // TODO: set uid and ts to data reader
2555 2556
  return 0;
}
2557
#endif
2558

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

D
dapan1121 已提交
2562
  switch (pNode->type) {
D
dapan1121 已提交
2563 2564 2565 2566 2567 2568
    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 已提交
2569

D
dapan1121 已提交
2570 2571 2572
      *pParam = pInserterParam;
      break;
    }
D
dapan1121 已提交
2573
    case QUERY_NODE_PHYSICAL_PLAN_DELETE: {
2574
      SDeleterParam* pDeleterParam = taosMemoryCalloc(1, sizeof(SDeleterParam));
D
dapan1121 已提交
2575 2576 2577
      if (NULL == pDeleterParam) {
        return TSDB_CODE_OUT_OF_MEMORY;
      }
H
Haojun Liao 已提交
2578 2579 2580 2581
      int32_t tbNum = tableListGetSize(pTask->pTableInfoList);
      pDeleterParam->suid = tableListGetSuid(pTask->pTableInfoList);

      // TODO extract uid list
D
dapan1121 已提交
2582 2583 2584 2585 2586
      pDeleterParam->pUidList = taosArrayInit(tbNum, sizeof(uint64_t));
      if (NULL == pDeleterParam->pUidList) {
        taosMemoryFree(pDeleterParam);
        return TSDB_CODE_OUT_OF_MEMORY;
      }
H
Haojun Liao 已提交
2587

D
dapan1121 已提交
2588
      for (int32_t i = 0; i < tbNum; ++i) {
H
Haojun Liao 已提交
2589
        STableKeyInfo* pTable = tableListGetInfo(pTask->pTableInfoList, i);
D
dapan1121 已提交
2590 2591 2592 2593 2594 2595 2596 2597 2598 2599 2600 2601 2602
        taosArrayPush(pDeleterParam->pUidList, &pTable->uid);
      }

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

  return TSDB_CODE_SUCCESS;
}

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

D
dapan1121 已提交
2607
  *pTaskInfo = createExecTaskInfo(queryId, taskId, model, pPlan->dbFName);
H
Haojun Liao 已提交
2608 2609 2610
  if (*pTaskInfo == NULL) {
    goto _complete;
  }
H
Haojun Liao 已提交
2611

2612
  if (pHandle) {
L
Liu Jicong 已提交
2613
    /*(*pTaskInfo)->streamInfo.fillHistoryVer1 = pHandle->fillHistoryVer1;*/
2614 2615 2616
    if (pHandle->pStateBackend) {
      (*pTaskInfo)->streamInfo.pState = pHandle->pStateBackend;
    }
H
Haojun Liao 已提交
2617 2618
  }

2619
  (*pTaskInfo)->sql = sql;
D
dapan1121 已提交
2620
  sql = NULL;
H
Haojun Liao 已提交
2621

2622
  (*pTaskInfo)->pSubplan = pPlan;
2623 2624
  (*pTaskInfo)->pRoot =
      createOperatorTree(pPlan->pNode, *pTaskInfo, pHandle, pPlan->pTagCond, pPlan->pTagIndexCond, pPlan->user);
L
Liu Jicong 已提交
2625

D
dapan1121 已提交
2626
  if (NULL == (*pTaskInfo)->pRoot) {
H
Haojun Liao 已提交
2627
    terrno = (*pTaskInfo)->code;
D
dapan1121 已提交
2628
    goto _complete;
2629 2630
  }

H
Haojun Liao 已提交
2631
  return TSDB_CODE_SUCCESS;
H
Haojun Liao 已提交
2632

H
Haojun Liao 已提交
2633
_complete:
D
dapan1121 已提交
2634
  taosMemoryFree(sql);
H
Haojun Liao 已提交
2635
  doDestroyTask(*pTaskInfo);
H
Haojun Liao 已提交
2636
  return terrno;
H
Haojun Liao 已提交
2637 2638
}

H
Haojun Liao 已提交
2639 2640 2641 2642 2643
static void freeBlock(void* pParam) {
  SSDataBlock* pBlock = *(SSDataBlock**)pParam;
  blockDataDestroy(pBlock);
}

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

H
Haojun Liao 已提交
2647
  pTaskInfo->pTableInfoList = tableListDestroy(pTaskInfo->pTableInfoList);
H
Haojun Liao 已提交
2648
  destroyOperatorInfo(pTaskInfo->pRoot);
2649
  cleanupTableSchemaInfo(&pTaskInfo->schemaInfo);
2650
  cleanupStreamInfo(&pTaskInfo->streamInfo);
2651

D
dapan1121 已提交
2652
  if (!pTaskInfo->localFetch.localExec) {
D
dapan1121 已提交
2653 2654
    nodesDestroyNode((SNode*)pTaskInfo->pSubplan);
  }
2655

H
Haojun Liao 已提交
2656
  taosArrayDestroyEx(pTaskInfo->pResultBlockList, freeBlock);
D
dapan1121 已提交
2657
  taosArrayDestroy(pTaskInfo->stopInfo.pStopInfo);
wafwerar's avatar
wafwerar 已提交
2658 2659 2660
  taosMemoryFreeClear(pTaskInfo->sql);
  taosMemoryFreeClear(pTaskInfo->id.str);
  taosMemoryFreeClear(pTaskInfo);
2661 2662 2663 2664
}

static int64_t getQuerySupportBufSize(size_t numOfTables) {
  size_t s1 = sizeof(STableQueryInfo);
L
Liu Jicong 已提交
2665 2666
  //  size_t s3 = sizeof(STableCheckInfo);  buffer consumption in tsdb
  return (int64_t)(s1 * 1.5 * numOfTables);
2667 2668 2669 2670 2671 2672 2673
}

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 已提交
2674
    while (1) {
2675 2676 2677 2678 2679 2680 2681 2682 2683 2684 2685 2686 2687 2688 2689 2690 2691 2692 2693 2694 2695 2696 2697 2698 2699 2700
      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 已提交
2701

H
Haojun Liao 已提交
2702
int32_t getOperatorExplainExecInfo(SOperatorInfo* operatorInfo, SArray* pExecInfoList) {
2703
  SExplainExecInfo  execInfo = {0};
H
Haojun Liao 已提交
2704
  SExplainExecInfo* pExplainInfo = taosArrayPush(pExecInfoList, &execInfo);
2705

H
Haojun Liao 已提交
2706 2707 2708 2709 2710
  pExplainInfo->numOfRows = operatorInfo->resultInfo.totalRows;
  pExplainInfo->startupCost = operatorInfo->cost.openCost;
  pExplainInfo->totalCost = operatorInfo->cost.totalCost;
  pExplainInfo->verboseLen = 0;
  pExplainInfo->verboseInfo = NULL;
D
dapan1121 已提交
2711

2712
  if (operatorInfo->fpSet.getExplainFn) {
2713 2714
    int32_t code =
        operatorInfo->fpSet.getExplainFn(operatorInfo, &pExplainInfo->verboseInfo, &pExplainInfo->verboseLen);
D
dapan1121 已提交
2715
    if (code) {
2716
      qError("%s operator getExplainFn failed, code:%s", GET_TASKID(operatorInfo->pTaskInfo), tstrerror(code));
D
dapan1121 已提交
2717 2718 2719
      return code;
    }
  }
dengyihao's avatar
dengyihao 已提交
2720

D
dapan1121 已提交
2721
  int32_t code = 0;
D
dapan1121 已提交
2722
  for (int32_t i = 0; i < operatorInfo->numOfDownstream; ++i) {
H
Haojun Liao 已提交
2723 2724
    code = getOperatorExplainExecInfo(operatorInfo->pDownstream[i], pExecInfoList);
    if (code != TSDB_CODE_SUCCESS) {
2725
      //      taosMemoryFreeClear(*pRes);
D
dapan1121 已提交
2726 2727 2728 2729 2730
      return TSDB_CODE_QRY_OUT_OF_MEMORY;
    }
  }

  return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
2731
}
5
54liuyao 已提交
2732

2733 2734
int32_t setOutputBuf(SStreamState* pState, STimeWindow* win, SResultRow** pResult, int64_t tableGroupId,
                     SqlFunctionCtx* pCtx, int32_t numOfOutput, int32_t* rowEntryInfoOffset, SAggSupporter* pAggSup) {
2735 2736 2737 2738 2739 2740
  SWinKey key = {
      .ts = win->skey,
      .groupId = tableGroupId,
  };
  char*   value = NULL;
  int32_t size = pAggSup->resultRowSize;
5
54liuyao 已提交
2741

2742
  if (streamStateAddIfNotExist(pState, &key, (void**)&value, &size) < 0) {
2743 2744 2745 2746 2747 2748 2749 2750 2751 2752
    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;
}

2753 2754
int32_t releaseOutputBuf(SStreamState* pState, SWinKey* pKey, SResultRow* pResult) {
  streamStateReleaseBuf(pState, pKey, pResult);
2755 2756 2757
  return TSDB_CODE_SUCCESS;
}

2758 2759
int32_t saveOutputBuf(SStreamState* pState, SWinKey* pKey, SResultRow* pResult, int32_t resSize) {
  streamStatePut(pState, pKey, pResult, resSize);
2760 2761 2762
  return TSDB_CODE_SUCCESS;
}

2763
int32_t buildDataBlockFromGroupRes(SOperatorInfo* pOperator, SStreamState* pState, SSDataBlock* pBlock, SExprSupp* pSup,
2764
                                   SGroupResInfo* pGroupResInfo) {
2765
  SExecTaskInfo*  pTaskInfo = pOperator->pTaskInfo;
2766 2767 2768 2769 2770 2771 2772 2773 2774 2775 2776 2777
  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 = {
2778 2779
            .ts = *(TSKEY*)pPos->key,
            .groupId = pPos->groupId,
2780
    };
2781
    int32_t code = streamStateGet(pState, &key, &pVal, &size);
2782 2783 2784 2785 2786 2787
    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;
2788
      releaseOutputBuf(pState, &key, pRow);
2789 2790 2791
      continue;
    }

H
Haojun Liao 已提交
2792 2793
    if (pBlock->info.id.groupId == 0) {
      pBlock->info.id.groupId = pPos->groupId;
2794
      void* tbname = NULL;
H
Haojun Liao 已提交
2795
      if (streamStateGetParName(pTaskInfo->streamInfo.pState, pBlock->info.id.groupId, &tbname) < 0) {
L
Liu Jicong 已提交
2796
        pBlock->info.parTbName[0] = 0;
2797 2798
      } else {
        memcpy(pBlock->info.parTbName, tbname, TSDB_TABLE_NAME_LEN);
2799
      }
2800
      tdbFree(tbname);
2801 2802
    } else {
      // current value belongs to different group, it can't be packed into one datablock
H
Haojun Liao 已提交
2803
      if (pBlock->info.id.groupId != pPos->groupId) {
2804
        releaseOutputBuf(pState, &key, pRow);
2805 2806 2807 2808 2809 2810
        break;
      }
    }

    if (pBlock->info.rows + pRow->numOfRows > pBlock->info.capacity) {
      ASSERT(pBlock->info.rows > 0);
2811
      releaseOutputBuf(pState, &key, pRow);
2812 2813 2814 2815 2816 2817 2818 2819 2820 2821
      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) {
2822 2823 2824 2825
        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);
2826 2827 2828 2829 2830 2831 2832 2833 2834 2835 2836 2837 2838
        }
      } 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 已提交
2839

2840
    pBlock->info.rows += pRow->numOfRows;
2841
    releaseOutputBuf(pState, &key, pRow);
2842 2843 2844 2845
  }
  blockDataUpdateTsWindow(pBlock, 0);
  return TSDB_CODE_SUCCESS;
}
5
54liuyao 已提交
2846 2847 2848 2849 2850 2851 2852

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;
}

2853
int32_t buildSessionResultDataBlock(SOperatorInfo* pOperator, SStreamState* pState, SSDataBlock* pBlock,
5
54liuyao 已提交
2854
                                    SExprSupp* pSup, SGroupResInfo* pGroupResInfo) {
2855
  SExecTaskInfo*  pTaskInfo = pOperator->pTaskInfo;
5
54liuyao 已提交
2856 2857 2858 2859 2860 2861 2862 2863 2864 2865 2866 2867 2868
  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);
2869 2870
    if (code == -1) {
      // coverity scan
5
54liuyao 已提交
2871
      pGroupResInfo->index += 1;
2872 2873
      continue;
    }
5
54liuyao 已提交
2874 2875 2876 2877 2878 2879 2880 2881 2882
    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 已提交
2883 2884
    if (pBlock->info.id.groupId == 0) {
      pBlock->info.id.groupId = pKey->groupId;
2885

2886
      void* tbname = NULL;
H
Haojun Liao 已提交
2887
      if (streamStateGetParName(pTaskInfo->streamInfo.pState, pBlock->info.id.groupId, &tbname) < 0) {
2888
        pBlock->info.parTbName[0] = 0;
2889
      } else {
2890
        memcpy(pBlock->info.parTbName, tbname, TSDB_TABLE_NAME_LEN);
2891
      }
2892
      tdbFree(tbname);
5
54liuyao 已提交
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 != pKey->groupId) {
5
54liuyao 已提交
2896 2897 2898 2899 2900 2901 2902 2903 2904 2905 2906 2907 2908 2909 2910 2911 2912 2913 2914 2915 2916 2917 2918 2919 2920 2921 2922 2923 2924 2925 2926 2927 2928 2929 2930 2931 2932 2933 2934 2935 2936 2937
        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;
2938
}