executorimpl.c 108.5 KB
Newer Older
1 2 3 4 5 6 7 8 9 10 11 12 13 14
/*
 * Copyright (c) 2019 TAOS Data, Inc. <jhtao@taosdata.com>
 *
 * This program is free software: you can use, redistribute, and/or modify
 * it under the terms of the GNU Affero General Public License, version 3
 * or later ("AGPL"), as published by the Free Software Foundation.
 *
 * This program is distributed in the hope that it will be useful, but WITHOUT
 * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or
 * FITNESS FOR A PARTICULAR PURPOSE.
 *
 * You should have received a copy of the GNU Affero General Public License
 * along with this program. If not, see <http://www.gnu.org/licenses/>.
 */
15

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

  return fpSet;
}

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

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

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

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

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

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

173 174
  setBufPageDirty(pData, true);

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

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

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

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

200 201
  SResultRow* pResult = NULL;

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

  return 0;
}

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

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

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

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

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

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

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

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

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

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

      char* p = GET_ROWCELL_INTERBUF(pEntryInfo);
348

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

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

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

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

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

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

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

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

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

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

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

  return TSDB_CODE_SUCCESS;
}

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

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

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

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

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

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

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

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

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

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

  return TSDB_CODE_SUCCESS;
507 508
}

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

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

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

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

526 527 528
  return true;
}

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

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

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

553
  ASSERT(!IS_VAR_DATA_TYPE(type));
554

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

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

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

578 579
  return TSDB_CODE_SUCCESS;
}
580

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

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

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

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

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

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

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

  return false;
}

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

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

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

  return win;
645 646
}

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

  doSetTableGroupOutputBuf(pOperator, numOfOutput, groupId);
948 949

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

1053
  int32_t numOfRows = getNumOfTotalRes(pGroupResInfo);
1054

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

  if (!hasCountFunc) {
    return TSDB_CODE_SUCCESS;
  }

  SSDataBlock* pBlock = createDataBlock();
G
Ganlin Zhao 已提交
1444
  pBlock->info.rows = 1;
1445 1446 1447
  pBlock->info.capacity = 0;

  for (int32_t i = 0; i < pOperator->exprSupp.numOfExprs; ++i) {
G
Ganlin Zhao 已提交
1448 1449 1450 1451
    SColumnInfoData colInfo = {0};
    colInfo.hasNull = true;
    colInfo.info.type = TSDB_DATA_TYPE_NULL;
    colInfo.info.bytes = 1;
1452 1453 1454 1455 1456 1457

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

G
Ganlin Zhao 已提交
1471
  blockDataEnsureCapacity(pBlock, pBlock->info.rows);
1472 1473 1474 1475 1476 1477 1478 1479 1480 1481 1482 1483 1484 1485 1486
  *ppBlock = pBlock;

  return TSDB_CODE_SUCCESS;
}

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

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


1487
// this is a blocking operator
L
Liu Jicong 已提交
1488
static int32_t doOpenAggregateOptr(SOperatorInfo* pOperator) {
H
Haojun Liao 已提交
1489 1490
  if (OPTR_IS_OPENED(pOperator)) {
    return TSDB_CODE_SUCCESS;
1491 1492
  }

H
Haojun Liao 已提交
1493
  SExecTaskInfo*    pTaskInfo = pOperator->pTaskInfo;
1494
  SAggOperatorInfo* pAggInfo = pOperator->info;
H
Haojun Liao 已提交
1495

1496 1497
  SExprSupp*     pSup = &pOperator->exprSupp;
  SOperatorInfo* downstream = pOperator->pDownstream[0];
1498

1499 1500
  int64_t st = taosGetTimestampUs();

1501 1502 1503
  int32_t order = TSDB_ORDER_ASC;
  int32_t scanFlag = MAIN_SCAN;

1504 1505 1506
  bool    hasValidBlock = false;
  bool    blockAllocated = false;

H
Haojun Liao 已提交
1507
  while (1) {
1508
    SSDataBlock* pBlock = downstream->fpSet.getNextFn(downstream);
1509
    if (pBlock == NULL) {
G
Ganlin Zhao 已提交
1510 1511
      if (!hasValidBlock &&
          downstream->operatorType != QUERY_NODE_PHYSICAL_PLAN_PARTITION) {
1512 1513 1514 1515 1516 1517 1518 1519
        createDataBlockForEmptyInput(pOperator, &pBlock);
        if (pBlock == NULL) {
          break;
        }
        blockAllocated = true;
      } else {
        break;
      }
1520
    }
1521
    hasValidBlock = true;
1522

1523 1524
    int32_t code = getTableScanInfo(pOperator, &order, &scanFlag);
    if (code != TSDB_CODE_SUCCESS) {
1525
      destroyDataBlockForEmptyInput(blockAllocated, &pBlock);
1526
      T_LONG_JMP(pTaskInfo->env, code);
1527
    }
1528

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

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

    destroyDataBlockForEmptyInput(blockAllocated, &pBlock);

1550 1551
  }

1552 1553 1554 1555 1556
  // 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);
  }

1557
  initGroupedResultInfo(&pAggInfo->groupResInfo, pAggInfo->aggSup.pResultRowHashTable, 0);
H
Haojun Liao 已提交
1558
  OPTR_SET_OPENED(pOperator);
1559

1560
  pOperator->cost.openCost = (taosGetTimestampUs() - st) / 1000.0;
1561
  return pTaskInfo->code;
H
Haojun Liao 已提交
1562 1563
}

1564
static SSDataBlock* getAggregateResult(SOperatorInfo* pOperator) {
L
Liu Jicong 已提交
1565
  SAggOperatorInfo* pAggInfo = pOperator->info;
H
Haojun Liao 已提交
1566 1567 1568 1569 1570 1571
  SOptrBasicInfo*   pInfo = &pAggInfo->binfo;

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

L
Liu Jicong 已提交
1572
  SExecTaskInfo* pTaskInfo = pOperator->pTaskInfo;
1573
  pTaskInfo->code = pOperator->fpSet._openFn(pOperator);
H
Haojun Liao 已提交
1574
  if (pTaskInfo->code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
1575
    setOperatorCompleted(pOperator);
H
Haojun Liao 已提交
1576 1577 1578
    return NULL;
  }

H
Haojun Liao 已提交
1579
  blockDataEnsureCapacity(pInfo->pRes, pOperator->resultInfo.capacity);
S
slzhou 已提交
1580 1581
  while (1) {
    doBuildResultDatablock(pOperator, pInfo, &pAggInfo->groupResInfo, pAggInfo->aggSup.pResultBuf);
H
Haojun Liao 已提交
1582
    doFilter(pInfo->pRes, pOperator->exprSupp.pFilterInfo, NULL);
S
slzhou 已提交
1583

1584
    if (!hasRemainResults(&pAggInfo->groupResInfo)) {
H
Haojun Liao 已提交
1585
      setOperatorCompleted(pOperator);
S
slzhou 已提交
1586 1587
      break;
    }
1588

S
slzhou 已提交
1589 1590 1591 1592
    if (pInfo->pRes->info.rows > 0) {
      break;
    }
  }
1593

1594
  size_t rows = blockDataGetNumOfRows(pInfo->pRes);
1595 1596
  pOperator->resultInfo.totalRows += rows;

1597
  return (rows == 0) ? NULL : pInfo->pRes;
1598 1599
}

L
Liu Jicong 已提交
1600 1601
static void doHandleRemainBlockForNewGroupImpl(SOperatorInfo* pOperator, SFillOperatorInfo* pInfo,
                                               SResultInfo* pResultInfo, SExecTaskInfo* pTaskInfo) {
1602
  pInfo->totalInputRows = pInfo->existNewGroupBlock->info.rows;
1603 1604 1605 1606 1607
  SSDataBlock* pResBlock = pInfo->pFinalRes;

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

5
54liuyao 已提交
1609
  int64_t ekey = pInfo->existNewGroupBlock->info.window.ekey;
1610 1611
  taosResetFillInfo(pInfo->pFillInfo, getFillInfoStart(pInfo->pFillInfo));

1612
  blockDataCleanup(pInfo->pRes);
1613 1614 1615 1616
  doApplyScalarCalculation(pOperator, pInfo->existNewGroupBlock, order, scanFlag);

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

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

H
Haojun Liao 已提交
1621
  pInfo->curGroupId = pInfo->existNewGroupBlock->info.id.groupId;
1622 1623 1624
  pInfo->existNewGroupBlock = NULL;
}

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

  // handle the cached new group data block
  if (pInfo->existNewGroupBlock) {
1636 1637 1638 1639 1640 1641
    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 已提交
1642
  SExprSupp*         pSup = &pOperator->exprSupp;
1643
  setInputDataBlock(pSup, pBlock, order, scanFlag, false);
1644 1645
  projectApplyFunctions(pSup->pExprInfo, pInfo->pRes, pBlock, pSup->pCtx, pSup->numOfExprs, NULL);

1646 1647 1648 1649
  // 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);
1650

1651
  projectApplyFunctions(pNoFillSupp->pExprInfo, pInfo->pRes, pBlock, pNoFillSupp->pCtx, pNoFillSupp->numOfExprs, NULL);
H
Haojun Liao 已提交
1652
  pInfo->pRes->info.id.groupId = pBlock->info.id.groupId;
1653 1654
}

S
slzhou 已提交
1655
static SSDataBlock* doFillImpl(SOperatorInfo* pOperator) {
L
Liu Jicong 已提交
1656 1657
  SFillOperatorInfo* pInfo = pOperator->info;
  SExecTaskInfo*     pTaskInfo = pOperator->pTaskInfo;
1658

H
Haojun Liao 已提交
1659
  SResultInfo* pResultInfo = &pOperator->resultInfo;
H
Haojun Liao 已提交
1660
  SSDataBlock* pResBlock = pInfo->pFinalRes;
1661 1662

  blockDataCleanup(pResBlock);
1663

H
Haojun Liao 已提交
1664 1665
  int32_t order = TSDB_ORDER_ASC;
  int32_t scanFlag = MAIN_SCAN;
1666
  getTableScanInfo(pOperator, &order, &scanFlag);
1667

1668
  doHandleRemainBlockFromNewGroup(pOperator, pInfo, pResultInfo, pTaskInfo);
1669
  if (pResBlock->info.rows > 0) {
H
Haojun Liao 已提交
1670
    pResBlock->info.id.groupId = pInfo->curGroupId;
1671
    return pResBlock;
H
Haojun Liao 已提交
1672
  }
1673

H
Haojun Liao 已提交
1674
  SOperatorInfo* pDownstream = pOperator->pDownstream[0];
L
Liu Jicong 已提交
1675
  while (1) {
1676
    SSDataBlock* pBlock = pDownstream->fpSet.getNextFn(pDownstream);
1677 1678
    if (pBlock == NULL) {
      if (pInfo->totalInputRows == 0) {
H
Haojun Liao 已提交
1679
        setOperatorCompleted(pOperator);
1680 1681
        return NULL;
      }
1682

1683
      taosFillSetStartInfo(pInfo->pFillInfo, 0, pInfo->win.ekey);
1684
    } else {
1685
      blockDataUpdateTsWindow(pBlock, pInfo->primarySrcSlotId);
1686 1687

      blockDataCleanup(pInfo->pRes);
1688 1689
      blockDataEnsureCapacity(pInfo->pRes, pBlock->info.rows);
      blockDataEnsureCapacity(pInfo->pFinalRes, pBlock->info.rows);
1690
      doApplyScalarCalculation(pOperator, pBlock, order, scanFlag);
1691

H
Haojun Liao 已提交
1692 1693
      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 已提交
1694
        pInfo->totalInputRows += pInfo->pRes->info.rows;
1695

H
Haojun Liao 已提交
1696 1697 1698 1699 1700
        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 已提交
1701
        taosFillSetInputDataBlock(pInfo->pFillInfo, pInfo->pRes);
H
Haojun Liao 已提交
1702
      } else if (pInfo->curGroupId != pBlock->info.id.groupId) {  // the new group data block
1703 1704 1705 1706 1707
        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);
1708 1709 1710
      }
    }

1711 1712
    int32_t numOfResultRows = pOperator->resultInfo.capacity - pResBlock->info.rows;
    taosFillResultDataBlock(pInfo->pFillInfo, pResBlock, numOfResultRows);
1713 1714

    // current group has no more result to return
1715
    if (pResBlock->info.rows > 0) {
1716 1717
      // 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
1718
      if (pResBlock->info.rows > pResultInfo->threshold || pBlock == NULL || pInfo->existNewGroupBlock != NULL) {
H
Haojun Liao 已提交
1719
        pResBlock->info.id.groupId = pInfo->curGroupId;
1720
        return pResBlock;
1721 1722
      }

1723
      doHandleRemainBlockFromNewGroup(pOperator, pInfo, pResultInfo, pTaskInfo);
1724
      if (pResBlock->info.rows >= pOperator->resultInfo.threshold || pBlock == NULL) {
H
Haojun Liao 已提交
1725
        pResBlock->info.id.groupId = pInfo->curGroupId;
1726
        return pResBlock;
1727 1728 1729
      }
    } else if (pInfo->existNewGroupBlock) {  // try next group
      assert(pBlock != NULL);
1730 1731 1732 1733

      blockDataCleanup(pResBlock);

      doHandleRemainBlockForNewGroupImpl(pOperator, pInfo, pResultInfo, pTaskInfo);
1734
      if (pResBlock->info.rows > pResultInfo->threshold) {
H
Haojun Liao 已提交
1735
        pResBlock->info.id.groupId = pInfo->curGroupId;
1736
        return pResBlock;
1737 1738 1739 1740 1741 1742 1743
      }
    } else {
      return NULL;
    }
  }
}

S
slzhou 已提交
1744 1745 1746 1747 1748 1749 1750 1751
static SSDataBlock* doFill(SOperatorInfo* pOperator) {
  SFillOperatorInfo* pInfo = pOperator->info;
  SExecTaskInfo*     pTaskInfo = pOperator->pTaskInfo;

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

S
slzhou 已提交
1752
  SSDataBlock* fillResult = NULL;
S
slzhou 已提交
1753
  while (true) {
S
slzhou 已提交
1754
    fillResult = doFillImpl(pOperator);
S
slzhou 已提交
1755
    if (fillResult == NULL) {
H
Haojun Liao 已提交
1756
      setOperatorCompleted(pOperator);
S
slzhou 已提交
1757 1758 1759
      break;
    }

1760
    doFilter(fillResult, pOperator->exprSupp.pFilterInfo, &pInfo->matchInfo);
S
slzhou 已提交
1761 1762 1763 1764 1765
    if (fillResult->info.rows > 0) {
      break;
    }
  }

S
slzhou 已提交
1766
  if (fillResult != NULL) {
1767
    pOperator->resultInfo.totalRows += fillResult->info.rows;
S
slzhou 已提交
1768
  }
S
slzhou 已提交
1769

S
slzhou 已提交
1770
  return fillResult;
S
slzhou 已提交
1771 1772
}

1773
void destroyExprInfo(SExprInfo* pExpr, int32_t numOfExprs) {
C
Cary Xu 已提交
1774 1775 1776 1777 1778
  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 已提交
1779 1780
      } else if (pExprInfo->base.pParam[j].type == FUNC_PARAM_TYPE_VALUE) {
        taosVariantDestroy(&pExprInfo->base.pParam[j].param);
H
Haojun Liao 已提交
1781
      }
1782
    }
C
Cary Xu 已提交
1783 1784 1785

    taosMemoryFree(pExprInfo->base.pParam);
    taosMemoryFree(pExprInfo->pExpr);
H
Haojun Liao 已提交
1786 1787 1788
  }
}

5
54liuyao 已提交
1789
void destroyOperatorInfo(SOperatorInfo* pOperator) {
1790 1791 1792 1793
  if (pOperator == NULL) {
    return;
  }

1794
  if (pOperator->fpSet.closeFn != NULL) {
1795
    pOperator->fpSet.closeFn(pOperator->info);
1796 1797
  }

H
Haojun Liao 已提交
1798
  if (pOperator->pDownstream != NULL) {
L
Liu Jicong 已提交
1799
    for (int32_t i = 0; i < pOperator->numOfDownstream; ++i) {
H
Haojun Liao 已提交
1800
      destroyOperatorInfo(pOperator->pDownstream[i]);
1801 1802
    }

wafwerar's avatar
wafwerar 已提交
1803
    taosMemoryFreeClear(pOperator->pDownstream);
H
Haojun Liao 已提交
1804
    pOperator->numOfDownstream = 0;
1805 1806
  }

1807
  cleanupExprSupp(&pOperator->exprSupp);
wafwerar's avatar
wafwerar 已提交
1808
  taosMemoryFreeClear(pOperator);
1809 1810
}

1811 1812 1813 1814 1815 1816
int32_t getBufferPgSize(int32_t rowSize, uint32_t* defaultPgsz, uint32_t* defaultBufsz) {
  *defaultPgsz = 4096;
  while (*defaultPgsz < rowSize * 4) {
    *defaultPgsz <<= 1u;
  }

1817
  // The default buffer for each operator in query is 10MB.
1818
  // at least four pages need to be in buffer
1819 1820
  // TODO: make this variable to be configurable.
  *defaultBufsz = 4096 * 2560;
1821 1822 1823 1824 1825 1826 1827
  if ((*defaultBufsz) <= (*defaultPgsz)) {
    (*defaultBufsz) = (*defaultPgsz) * 4;
  }

  return 0;
}

dengyihao's avatar
dengyihao 已提交
1828 1829
int32_t doInitAggInfoSup(SAggSupporter* pAggSup, SqlFunctionCtx* pCtx, int32_t numOfOutput, size_t keyBufSize,
                         const char* pKey) {
L
Liu Jicong 已提交
1830
  int32_t    code = 0;
1831 1832
  _hash_fn_t hashFn = taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY);

1833
  pAggSup->currentPageId = -1;
dengyihao's avatar
dengyihao 已提交
1834 1835
  pAggSup->resultRowSize = getResultRowSize(pCtx, numOfOutput);
  pAggSup->keyBuf = taosMemoryCalloc(1, keyBufSize + POINTER_BYTES + sizeof(int64_t));
1836
  pAggSup->pResultRowHashTable = tSimpleHashInit(10, hashFn);
1837

H
Haojun Liao 已提交
1838
  if (pAggSup->keyBuf == NULL || pAggSup->pResultRowHashTable == NULL) {
1839 1840 1841
    return TSDB_CODE_OUT_OF_MEMORY;
  }

dengyihao's avatar
dengyihao 已提交
1842
  uint32_t defaultPgsz = 0;
1843 1844
  uint32_t defaultBufsz = 0;
  getBufferPgSize(pAggSup->resultRowSize, &defaultPgsz, &defaultBufsz);
H
Haojun Liao 已提交
1845

wafwerar's avatar
wafwerar 已提交
1846
  if (!osTempSpaceAvailable()) {
H
Haojun Liao 已提交
1847 1848 1849
    code = TSDB_CODE_NO_AVAIL_DISK;
    qError("Init stream agg supporter failed since %s, %s", terrstr(code), pKey);
    return code;
wafwerar's avatar
wafwerar 已提交
1850
  }
1851

H
Haojun Liao 已提交
1852
  code = createDiskbasedBuf(&pAggSup->pResultBuf, defaultPgsz, defaultBufsz, pKey, tsTempDir);
H
Haojun Liao 已提交
1853
  if (code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
1854
    qError("Create agg result buf failed since %s, %s", tstrerror(code), pKey);
H
Haojun Liao 已提交
1855 1856 1857
    return code;
  }

H
Haojun Liao 已提交
1858
  return code;
1859 1860
}

1861
void cleanupAggSup(SAggSupporter* pAggSup) {
wafwerar's avatar
wafwerar 已提交
1862
  taosMemoryFreeClear(pAggSup->keyBuf);
1863
  tSimpleHashCleanup(pAggSup->pResultRowHashTable);
H
Haojun Liao 已提交
1864
  destroyDiskbasedBuf(pAggSup->pResultBuf);
1865 1866
}

H
Haojun Liao 已提交
1867
int32_t initAggSup(SExprSupp* pSup, SAggSupporter* pAggSup, SExprInfo* pExprInfo, int32_t numOfCols, size_t keyBufSize,
L
Liu Jicong 已提交
1868
                    const char* pkey) {
1869 1870 1871 1872 1873
  int32_t code = initExprSupp(pSup, pExprInfo, numOfCols);
  if (code != TSDB_CODE_SUCCESS) {
    return code;
  }

1874 1875 1876 1877 1878
  code = doInitAggInfoSup(pAggSup, pSup->pCtx, numOfCols, keyBufSize, pkey);
  if (code != TSDB_CODE_SUCCESS) {
    return code;
  }

L
Liu Jicong 已提交
1879
  for (int32_t i = 0; i < numOfCols; ++i) {
1880
    pSup->pCtx[i].saveHandle.pBuf = pAggSup->pResultBuf;
1881 1882
  }

1883
  return TSDB_CODE_SUCCESS;
H
Haojun Liao 已提交
1884 1885
}

L
Liu Jicong 已提交
1886
void initResultSizeInfo(SResultInfo* pResultInfo, int32_t numOfRows) {
wmmhello's avatar
wmmhello 已提交
1887
  ASSERT(numOfRows != 0);
1888 1889
  pResultInfo->capacity = numOfRows;
  pResultInfo->threshold = numOfRows * 0.75;
1890

1891 1892
  if (pResultInfo->threshold == 0) {
    pResultInfo->threshold = numOfRows;
1893 1894 1895
  }
}

1896 1897 1898 1899 1900
void initBasicInfo(SOptrBasicInfo* pInfo, SSDataBlock* pBlock) {
  pInfo->pRes = pBlock;
  initResultRowInfo(&pInfo->resultRowInfo);
}

5
54liuyao 已提交
1901
void* destroySqlFunctionCtx(SqlFunctionCtx* pCtx, int32_t numOfOutput) {
1902 1903 1904 1905 1906 1907 1908 1909 1910 1911
  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);
1912
    taosMemoryFreeClear(pCtx[i].subsidiaries.buf);
1913 1914 1915 1916 1917 1918 1919 1920
    taosMemoryFree(pCtx[i].input.pData);
    taosMemoryFree(pCtx[i].input.pColumnDataAgg);
  }

  taosMemoryFreeClear(pCtx);
  return NULL;
}

1921
int32_t initExprSupp(SExprSupp* pSup, SExprInfo* pExprInfo, int32_t numOfExpr) {
1922 1923 1924 1925
  pSup->pExprInfo = pExprInfo;
  pSup->numOfExprs = numOfExpr;
  if (pSup->pExprInfo != NULL) {
    pSup->pCtx = createSqlFunctionCtx(pExprInfo, numOfExpr, &pSup->rowEntryInfoOffset);
1926 1927 1928
    if (pSup->pCtx == NULL) {
      return TSDB_CODE_OUT_OF_MEMORY;
    }
1929
  }
1930 1931

  return TSDB_CODE_SUCCESS;
1932 1933
}

1934 1935 1936 1937
void cleanupExprSupp(SExprSupp* pSupp) {
  destroySqlFunctionCtx(pSupp->pCtx, pSupp->numOfExprs);
  if (pSupp->pExprInfo != NULL) {
    destroyExprInfo(pSupp->pExprInfo, pSupp->numOfExprs);
C
Cary Xu 已提交
1938
    taosMemoryFreeClear(pSupp->pExprInfo);
1939
  }
H
Haojun Liao 已提交
1940 1941 1942 1943 1944 1945

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

1946 1947 1948
  taosMemoryFree(pSupp->rowEntryInfoOffset);
}

1949 1950
SOperatorInfo* createAggregateOperatorInfo(SOperatorInfo* downstream, SAggPhysiNode* pAggNode,
                                           SExecTaskInfo* pTaskInfo) {
wafwerar's avatar
wafwerar 已提交
1951
  SAggOperatorInfo* pInfo = taosMemoryCalloc(1, sizeof(SAggOperatorInfo));
L
Liu Jicong 已提交
1952
  SOperatorInfo*    pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
H
Haojun Liao 已提交
1953 1954 1955
  if (pInfo == NULL || pOperator == NULL) {
    goto _error;
  }
H
Haojun Liao 已提交
1956

H
Haojun Liao 已提交
1957
  SSDataBlock* pResBlock = createDataBlockFromDescNode(pAggNode->node.pOutputDataBlockDesc);
H
Haojun Liao 已提交
1958 1959 1960
  initBasicInfo(&pInfo->binfo, pResBlock);

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

1963 1964
  int32_t    num = 0;
  SExprInfo* pExprInfo = createExprInfo(pAggNode->pAggFuncs, pAggNode->pGroupKeys, &num);
H
Haojun Liao 已提交
1965
  int32_t    code = initAggSup(&pOperator->exprSupp, &pInfo->aggSup, pExprInfo, num, keyBufSize, pTaskInfo->id.str);
L
Liu Jicong 已提交
1966
  if (code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
1967 1968
    goto _error;
  }
H
Haojun Liao 已提交
1969

H
Haojun Liao 已提交
1970 1971 1972 1973 1974 1975
  int32_t    numOfScalarExpr = 0;
  SExprInfo* pScalarExprInfo = NULL;
  if (pAggNode->pExprs != NULL) {
    pScalarExprInfo = createExprInfo(pAggNode->pExprs, NULL, &numOfScalarExpr);
  }

1976 1977 1978 1979
  code = initExprSupp(&pInfo->scalarExprSup, pScalarExprInfo, numOfScalarExpr);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
1980

H
Haojun Liao 已提交
1981 1982 1983 1984 1985
  code = filterInitFromNode((SNode*)pAggNode->node.pConditions, &pOperator->exprSupp.pFilterInfo, 0);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }

H
Haojun Liao 已提交
1986
  pInfo->binfo.mergeResultBlock = pAggNode->mergeDataBlock;
1987
  pInfo->groupId = UINT64_MAX;
H
Haojun Liao 已提交
1988

1989 1990 1991
  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 已提交
1992

1993 1994
  if (downstream->operatorType == QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN) {
    STableScanInfo* pTableScanInfo = downstream->info;
H
Haojun Liao 已提交
1995 1996
    pTableScanInfo->base.pdInfo.pExprSup = &pOperator->exprSupp;
    pTableScanInfo->base.pdInfo.pAggSup = &pInfo->aggSup;
1997 1998
  }

H
Haojun Liao 已提交
1999 2000 2001 2002
  code = appendDownstream(pOperator, &downstream, 1);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
2003 2004

  return pOperator;
H
Haojun Liao 已提交
2005

2006
_error:
H
Haojun Liao 已提交
2007 2008 2009 2010
  if (pInfo != NULL) {
    destroyAggOperatorInfo(pInfo);
  }

2011 2012 2013
  if (pOperator != NULL) {
    cleanupExprSupp(&pOperator->exprSupp);
  }
H
Haojun Liao 已提交
2014

2015
  taosMemoryFreeClear(pOperator);
H
Haojun Liao 已提交
2016
  pTaskInfo->code = code;
H
Haojun Liao 已提交
2017
  return NULL;
2018 2019
}

2020
void cleanupBasicInfo(SOptrBasicInfo* pInfo) {
2021
  assert(pInfo != NULL);
H
Haojun Liao 已提交
2022
  pInfo->pRes = blockDataDestroy(pInfo->pRes);
2023 2024
}

H
Haojun Liao 已提交
2025 2026 2027 2028 2029 2030 2031
static void freeItem(void* pItem) {
  void** p = pItem;
  if (*p != NULL) {
    taosMemoryFreeClear(*p);
  }
}

2032
void destroyAggOperatorInfo(void* param) {
L
Liu Jicong 已提交
2033
  SAggOperatorInfo* pInfo = (SAggOperatorInfo*)param;
L
Liu Jicong 已提交
2034 2035
  cleanupBasicInfo(&pInfo->binfo);

H
Haojun Liao 已提交
2036
  cleanupAggSup(&pInfo->aggSup);
S
shenglian zhou 已提交
2037
  cleanupExprSupp(&pInfo->scalarExprSup);
H
Haojun Liao 已提交
2038
  cleanupGroupResInfo(&pInfo->groupResInfo);
D
dapan1121 已提交
2039
  taosMemoryFreeClear(param);
2040
}
2041

2042
void destroyFillOperatorInfo(void* param) {
L
Liu Jicong 已提交
2043
  SFillOperatorInfo* pInfo = (SFillOperatorInfo*)param;
2044
  pInfo->pFillInfo = taosDestroyFillInfo(pInfo->pFillInfo);
H
Haojun Liao 已提交
2045
  pInfo->pRes = blockDataDestroy(pInfo->pRes);
H
Haojun Liao 已提交
2046 2047
  pInfo->pFinalRes = blockDataDestroy(pInfo->pFinalRes);

2048
  cleanupExprSupp(&pInfo->noFillExprSupp);
H
Haojun Liao 已提交
2049

wafwerar's avatar
wafwerar 已提交
2050
  taosMemoryFreeClear(pInfo->p);
H
Haojun Liao 已提交
2051
  taosArrayDestroy(pInfo->matchInfo.pList);
D
dapan1121 已提交
2052
  taosMemoryFreeClear(param);
2053 2054
}

H
Haojun Liao 已提交
2055 2056 2057 2058
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 已提交
2059

H
Haojun Liao 已提交
2060 2061 2062
  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 已提交
2063

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

H
Haojun Liao 已提交
2067 2068 2069 2070 2071 2072 2073
  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 已提交
2074
  pInfo->p = taosMemoryCalloc(numOfCols, POINTER_BYTES);
2075

H
Haojun Liao 已提交
2076
  if (pInfo->pFillInfo == NULL || pInfo->p == NULL) {
H
Haojun Liao 已提交
2077 2078
    taosMemoryFree(pInfo->pFillInfo);
    taosMemoryFree(pInfo->p);
H
Haojun Liao 已提交
2079 2080 2081 2082 2083 2084
    return TSDB_CODE_OUT_OF_MEMORY;
  } else {
    return TSDB_CODE_SUCCESS;
  }
}

2085
static bool isWstartColumnExist(SFillOperatorInfo* pInfo) {
2086
  if (pInfo->noFillExprSupp.numOfExprs == 0) {
2087 2088
    return false;
  }
2089 2090 2091

  for (int32_t i = 0; i < pInfo->noFillExprSupp.numOfExprs; ++i) {
    SExprInfo* exprInfo = pInfo->noFillExprSupp.pExprInfo + i;
2092
    if (exprInfo->pExpr->nodeType == QUERY_NODE_COLUMN && exprInfo->base.numOfParams == 1 &&
2093 2094 2095 2096 2097 2098 2099
        exprInfo->base.pParam[0].pCol->colType == COLUMN_TYPE_WINDOW_START) {
      return true;
    }
  }
  return false;
}

2100 2101
static int32_t createPrimaryTsExprIfNeeded(SFillOperatorInfo* pInfo, SFillPhysiNode* pPhyFillNode, SExprSupp* pExprSupp,
                                           const char* idStr) {
2102
  bool wstartExist = isWstartColumnExist(pInfo);
2103

2104 2105
  if (wstartExist == false) {
    if (pPhyFillNode->pWStartTs->type != QUERY_NODE_TARGET) {
2106
      qError("pWStartTs of fill physical node is not a target node, %s", idStr);
2107 2108 2109
      return TSDB_CODE_QRY_SYS_ERROR;
    }

2110 2111
    SExprInfo* pExpr = taosMemoryRealloc(pExprSupp->pExprInfo, (pExprSupp->numOfExprs + 1) * sizeof(SExprInfo));
    if (pExpr == NULL) {
2112 2113 2114
      return TSDB_CODE_OUT_OF_MEMORY;
    }

2115 2116 2117
    createExprFromTargetNode(&pExpr[pExprSupp->numOfExprs], (STargetNode*)pPhyFillNode->pWStartTs);
    pExprSupp->numOfExprs += 1;
    pExprSupp->pExprInfo = pExpr;
2118
  }
2119

2120 2121 2122
  return TSDB_CODE_SUCCESS;
}

L
Liu Jicong 已提交
2123 2124
SOperatorInfo* createFillOperatorInfo(SOperatorInfo* downstream, SFillPhysiNode* pPhyFillNode,
                                      SExecTaskInfo* pTaskInfo) {
2125 2126 2127 2128 2129 2130
  SFillOperatorInfo* pInfo = taosMemoryCalloc(1, sizeof(SFillOperatorInfo));
  SOperatorInfo*     pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
  if (pInfo == NULL || pOperator == NULL) {
    goto _error;
  }

H
Haojun Liao 已提交
2131
  pInfo->pRes = createDataBlockFromDescNode(pPhyFillNode->node.pOutputDataBlockDesc);
H
Haojun Liao 已提交
2132 2133 2134
  SExprInfo* pExprInfo = createExprInfo(pPhyFillNode->pFillExprs, NULL, &pInfo->numOfExpr);
  pOperator->exprSupp.pExprInfo = pExprInfo;

2135 2136 2137 2138 2139 2140 2141 2142
  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);
2143 2144 2145
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
H
Haojun Liao 已提交
2146

L
Liu Jicong 已提交
2147
  SInterval* pInterval =
2148
      QUERY_NODE_PHYSICAL_PLAN_MERGE_ALIGNED_INTERVAL == downstream->operatorType
2149 2150
          ? &((SMergeAlignedIntervalAggOperatorInfo*)downstream->info)->intervalAggOperatorInfo->interval
          : &((SIntervalAggOperatorInfo*)downstream->info)->interval;
2151

2152
  int32_t order = (pPhyFillNode->inputTsOrder == ORDER_ASC) ? TSDB_ORDER_ASC : TSDB_ORDER_DESC;
2153
  int32_t type = convertFillType(pPhyFillNode->mode);
2154

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

2157
  initResultSizeInfo(&pOperator->resultInfo, 4096);
H
Haojun Liao 已提交
2158 2159 2160 2161 2162
  blockDataEnsureCapacity(pInfo->pRes, pOperator->resultInfo.capacity);
  code = initExprSupp(&pOperator->exprSupp, pExprInfo, pInfo->numOfExpr);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
H
Haojun Liao 已提交
2163

H
Haojun Liao 已提交
2164 2165
  pInfo->primaryTsCol = ((STargetNode*)pPhyFillNode->pWStartTs)->slotId;
  pInfo->primarySrcSlotId = ((SColumnNode*)((STargetNode*)pPhyFillNode->pWStartTs)->pExpr)->slotId;
2166

2167
  int32_t numOfOutputCols = 0;
H
Haojun Liao 已提交
2168 2169
  code = extractColMatchInfo(pPhyFillNode->pFillExprs, pPhyFillNode->node.pOutputDataBlockDesc, &numOfOutputCols,
                             COL_MATCH_FROM_SLOT_ID, &pInfo->matchInfo);
2170

2171
  code = initFillInfo(pInfo, pExprInfo, pInfo->numOfExpr, pNoFillSupp->pExprInfo, pNoFillSupp->numOfExprs,
2172 2173
                      (SNodeListNode*)pPhyFillNode->pValues, pPhyFillNode->timeRange, pResultInfo->capacity,
                      pTaskInfo->id.str, pInterval, type, order);
2174 2175 2176
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
2177

H
Haojun Liao 已提交
2178
  pInfo->pFinalRes = createOneDataBlock(pInfo->pRes, false);
H
Haojun Liao 已提交
2179 2180
  blockDataEnsureCapacity(pInfo->pFinalRes, pOperator->resultInfo.capacity);

H
Haojun Liao 已提交
2181 2182 2183 2184 2185
  code = filterInitFromNode((SNode*)pPhyFillNode->node.pConditions, &pOperator->exprSupp.pFilterInfo, 0);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }

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

2190
  code = appendDownstream(pOperator, &downstream, 1);
2191
  return pOperator;
H
Haojun Liao 已提交
2192

2193
_error:
H
Haojun Liao 已提交
2194 2195 2196 2197
  if (pInfo != NULL) {
    destroyFillOperatorInfo(pInfo);
  }

2198
  pTaskInfo->code = code;
wafwerar's avatar
wafwerar 已提交
2199
  taosMemoryFreeClear(pOperator);
H
Haojun Liao 已提交
2200
  return NULL;
2201 2202
}

D
dapan1121 已提交
2203
static SExecTaskInfo* createExecTaskInfo(uint64_t queryId, uint64_t taskId, EOPTR_EXEC_MODEL model, char* dbFName) {
wafwerar's avatar
wafwerar 已提交
2204
  SExecTaskInfo* pTaskInfo = taosMemoryCalloc(1, sizeof(SExecTaskInfo));
H
Haojun Liao 已提交
2205 2206 2207 2208 2209
  if (pTaskInfo == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return NULL;
  }

2210
  setTaskStatus(pTaskInfo, TASK_NOT_COMPLETED);
H
Haojun Liao 已提交
2211

2212
  pTaskInfo->schemaInfo.dbname = strdup(dbFName);
2213
  pTaskInfo->cost.created = taosGetTimestampMs();
H
Haojun Liao 已提交
2214
  pTaskInfo->id.queryId = queryId;
dengyihao's avatar
dengyihao 已提交
2215
  pTaskInfo->execModel = model;
H
Haojun Liao 已提交
2216
  pTaskInfo->pTableInfoList = tableListCreate();
D
dapan1121 已提交
2217
  pTaskInfo->stopInfo.pStopInfo = taosArrayInit(4, sizeof(SExchangeOpStopInfo));
H
Haojun Liao 已提交
2218
  pTaskInfo->pResultBlockList = taosArrayInit(128, POINTER_BYTES);
H
Haojun Liao 已提交
2219

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

2224 2225
  return pTaskInfo;
}
H
Haojun Liao 已提交
2226

H
Haojun Liao 已提交
2227 2228
SSchemaWrapper* extractQueriedColumnSchema(SScanPhysiNode* pScanNode);

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

D
dapan1121 已提交
2237
    metaReaderClear(&mr);
2238
    return terrno;
D
dapan1121 已提交
2239
  }
2240

2241 2242
  SSchemaInfo* pSchemaInfo = &pTaskInfo->schemaInfo;
  pSchemaInfo->tablename = strdup(mr.me.name);
2243 2244

  if (mr.me.type == TSDB_SUPER_TABLE) {
2245 2246
    pSchemaInfo->sw = tCloneSSchemaWrapper(&mr.me.stbEntry.schemaRow);
    pSchemaInfo->tversion = mr.me.stbEntry.schemaTag.version;
2247
  } else if (mr.me.type == TSDB_CHILD_TABLE) {
2248 2249
    tDecoderClear(&mr.coder);

2250
    tb_uid_t suid = mr.me.ctbEntry.suid;
H
Haojun Liao 已提交
2251
    metaGetTableEntryByUidCache(&mr, suid);
2252 2253
    pSchemaInfo->sw = tCloneSSchemaWrapper(&mr.me.stbEntry.schemaRow);
    pSchemaInfo->tversion = mr.me.stbEntry.schemaTag.version;
2254
  } else {
2255
    pSchemaInfo->sw = tCloneSSchemaWrapper(&mr.me.ntbEntry.schemaRow);
2256
  }
2257 2258

  metaReaderClear(&mr);
2259

H
Haojun Liao 已提交
2260 2261 2262 2263 2264
  pSchemaInfo->qsw = extractQueriedColumnSchema(pScanNode);
  return TSDB_CODE_SUCCESS;
}

SSchemaWrapper* extractQueriedColumnSchema(SScanPhysiNode* pScanNode) {
2265 2266 2267
  int32_t numOfCols = LIST_LENGTH(pScanNode->pScanCols);
  int32_t numOfTags = LIST_LENGTH(pScanNode->pScanPseudoCols);

2268
  SSchemaWrapper* pqSw = taosMemoryCalloc(1, sizeof(SSchemaWrapper));
2269
  pqSw->pSchema = taosMemoryCalloc(numOfCols + numOfTags, sizeof(SSchema));
2270

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

H
Haojun Liao 已提交
2275 2276 2277
    SSchema* pSchema = &pqSw->pSchema[pqSw->nCols++];
    pSchema->colId = pColNode->colId;
    pSchema->type = pColNode->node.resType.type;
H
Haojun Liao 已提交
2278 2279
    pSchema->bytes = pColNode->node.resType.bytes;
    tstrncpy(pSchema->name, pColNode->colName, tListLen(pSchema->name));
2280 2281
  }

2282
  // this the tags and pseudo function columns, we only keep the tag columns
2283
  for (int32_t i = 0; i < numOfTags; ++i) {
2284 2285 2286 2287 2288 2289 2290 2291 2292
    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 已提交
2293
      pSchema->bytes = pColNode->node.resType.bytes;
H
Haojun Liao 已提交
2294
      tstrncpy(pSchema->name, pColNode->colName, tListLen(pSchema->name));
2295 2296 2297
    }
  }

H
Haojun Liao 已提交
2298
  return pqSw;
2299 2300
}

2301 2302
static void cleanupTableSchemaInfo(SSchemaInfo* pSchemaInfo) {
  taosMemoryFreeClear(pSchemaInfo->dbname);
2303
  taosMemoryFreeClear(pSchemaInfo->tablename);
2304 2305
  tDeleteSSchemaWrapper(pSchemaInfo->sw);
  tDeleteSSchemaWrapper(pSchemaInfo->qsw);
2306 2307
}

2308
static void cleanupStreamInfo(SStreamTaskInfo* pStreamInfo) { tDeleteSSchemaWrapper(pStreamInfo->schema); }
2309

2310
bool groupbyTbname(SNodeList* pGroupList) {
2311
  bool bytbname = false;
2312
  if (LIST_LENGTH(pGroupList) == 1) {
2313 2314 2315 2316 2317 2318 2319 2320 2321 2322
    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;
}

2323 2324
SOperatorInfo* createOperatorTree(SPhysiNode* pPhyNode, SExecTaskInfo* pTaskInfo, SReadHandle* pHandle, SNode* pTagCond,
                                  SNode* pTagIndexCond, const char* pUser) {
2325
  int32_t         type = nodeType(pPhyNode);
2326
  STableListInfo* pTableListInfo = pTaskInfo->pTableInfoList;
2327
  const char*     idstr = GET_TASKID(pTaskInfo);
2328

X
Xiaoyu Wang 已提交
2329
  if (pPhyNode->pChildren == NULL || LIST_LENGTH(pPhyNode->pChildren) == 0) {
2330
    SOperatorInfo* pOperator = NULL;
H
Haojun Liao 已提交
2331
    if (QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN == type) {
dengyihao's avatar
dengyihao 已提交
2332
      STableScanPhysiNode* pTableScanNode = (STableScanPhysiNode*)pPhyNode;
H
Haojun Liao 已提交
2333

2334 2335 2336 2337 2338 2339
      // 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 已提交
2340 2341
      int32_t code =
          createScanTableListInfo(&pTableScanNode->scan, pTableScanNode->pGroupTags, pTableScanNode->groupSort, pHandle,
H
Haojun Liao 已提交
2342
                                  pTableListInfo, pTagCond, pTagIndexCond, pTaskInfo);
2343
      if (code) {
wmmhello's avatar
wmmhello 已提交
2344
        pTaskInfo->code = code;
2345
        qError("failed to createScanTableListInfo, code:%s, %s", tstrerror(code), idstr);
D
dapan1121 已提交
2346 2347
        return NULL;
      }
wmmhello's avatar
wmmhello 已提交
2348

2349
      code = extractTableSchemaInfo(pHandle, &pTableScanNode->scan, pTaskInfo);
S
slzhou 已提交
2350
      if (code) {
2351
        pTaskInfo->code = terrno;
wmmhello's avatar
wmmhello 已提交
2352 2353 2354
        return NULL;
      }

2355
      pOperator = createTableScanOperatorInfo(pTableScanNode, pHandle, pTaskInfo);
D
dapan1121 已提交
2356 2357 2358 2359 2360
      if (NULL == pOperator) {
        pTaskInfo->code = terrno;
        return NULL;
      }

2361
      STableScanInfo* pScanInfo = pOperator->info;
H
Haojun Liao 已提交
2362
      pTaskInfo->cost.pRecoder = &pScanInfo->base.readRecorder;
S
slzhou 已提交
2363 2364
    } else if (QUERY_NODE_PHYSICAL_PLAN_TABLE_MERGE_SCAN == type) {
      STableMergeScanPhysiNode* pTableScanNode = (STableMergeScanPhysiNode*)pPhyNode;
H
Haojun Liao 已提交
2365 2366 2367

      int32_t code = createScanTableListInfo(&pTableScanNode->scan, pTableScanNode->pGroupTags, true, pHandle,
                                             pTableListInfo, pTagCond, pTagIndexCond, pTaskInfo);
L
Liu Jicong 已提交
2368
      if (code) {
wmmhello's avatar
wmmhello 已提交
2369
        pTaskInfo->code = code;
H
Haojun Liao 已提交
2370
        qError("failed to createScanTableListInfo, code: %s", tstrerror(code));
wmmhello's avatar
wmmhello 已提交
2371 2372
        return NULL;
      }
2373

2374
      code = extractTableSchemaInfo(pHandle, &pTableScanNode->scan, pTaskInfo);
wmmhello's avatar
wmmhello 已提交
2375 2376 2377 2378
      if (code) {
        pTaskInfo->code = terrno;
        return NULL;
      }
wmmhello's avatar
wmmhello 已提交
2379

H
Haojun Liao 已提交
2380
      pOperator = createTableMergeScanOperatorInfo(pTableScanNode, pHandle, pTaskInfo);
D
dapan1121 已提交
2381 2382 2383 2384
      if (NULL == pOperator) {
        pTaskInfo->code = terrno;
        return NULL;
      }
wmmhello's avatar
wmmhello 已提交
2385

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

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

L
Liu Jicong 已提交
2407
        for (int32_t i = 0; i < sz; i++) {
H
Haojun Liao 已提交
2408
          STableKeyInfo* pKeyInfo = tableListGetInfo(pTableListInfo, i);
2409
          qDebug("add table uid:%" PRIu64 ", gid:%" PRIu64, pKeyInfo->uid, pKeyInfo->groupId);
L
Liu Jicong 已提交
2410 2411
        }
#endif
2412
      }
2413

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

      int32_t code = createScanTableListInfo(pScanPhyNode, NULL, false, pHandle, pTableListInfo, pTagCond,
H
Haojun Liao 已提交
2423
                                             pTagIndexCond, pTaskInfo);
2424
      if (code != TSDB_CODE_SUCCESS) {
2425
        pTaskInfo->code = code;
H
Haojun Liao 已提交
2426
        qError("failed to getTableList, code: %s", tstrerror(code));
2427 2428 2429
        return NULL;
      }

H
Haojun Liao 已提交
2430
      pOperator = createTagScanOperatorInfo(pHandle, pScanPhyNode, pTaskInfo);
2431
    } else if (QUERY_NODE_PHYSICAL_PLAN_BLOCK_DIST_SCAN == type) {
2432
      SBlockDistScanPhysiNode* pBlockNode = (SBlockDistScanPhysiNode*)pPhyNode;
2433 2434

      if (pBlockNode->tableType == TSDB_SUPER_TABLE) {
H
Haojun Liao 已提交
2435 2436
        SArray* pList = taosArrayInit(4, sizeof(STableKeyInfo));
        int32_t code = vnodeGetAllTableList(pHandle->vnode, pBlockNode->uid, pList);
2437 2438 2439 2440
        if (code != TSDB_CODE_SUCCESS) {
          pTaskInfo->code = terrno;
          return NULL;
        }
H
Haojun Liao 已提交
2441

2442
        for (int32_t i = 0; i < tableListGetSize(pTableListInfo); ++i) {
H
Haojun Liao 已提交
2443
          STableKeyInfo* p = taosArrayGet(pList, i);
H
Haojun Liao 已提交
2444
          tableListAddTableInfo(pTableListInfo, p->uid, 0);
H
Haojun Liao 已提交
2445 2446
        }
        taosArrayDestroy(pList);
2447
      } else {  // Create group with only one table
H
Haojun Liao 已提交
2448
        tableListAddTableInfo(pTableListInfo, pBlockNode->uid, 0);
2449 2450
      }

2451
      pOperator = createDataBlockInfoScanOperator(pHandle, pBlockNode, pTaskInfo);
H
Haojun Liao 已提交
2452 2453 2454
    } else if (QUERY_NODE_PHYSICAL_PLAN_LAST_ROW_SCAN == type) {
      SLastRowScanPhysiNode* pScanNode = (SLastRowScanPhysiNode*)pPhyNode;

L
Liu Jicong 已提交
2455
      int32_t code = createScanTableListInfo(&pScanNode->scan, pScanNode->pGroupTags, true, pHandle, pTableListInfo,
H
Haojun Liao 已提交
2456
                                             pTagCond, pTagIndexCond, pTaskInfo);
2457 2458 2459 2460
      if (code != TSDB_CODE_SUCCESS) {
        pTaskInfo->code = code;
        return NULL;
      }
2461

2462
      code = extractTableSchemaInfo(pHandle, &pScanNode->scan, pTaskInfo);
2463 2464 2465
      if (code != TSDB_CODE_SUCCESS) {
        pTaskInfo->code = code;
        return NULL;
H
Haojun Liao 已提交
2466 2467
      }

2468
      pOperator = createCacherowsScanOperator(pScanNode, pHandle, pTaskInfo);
2469
    } else if (QUERY_NODE_PHYSICAL_PLAN_PROJECT == type) {
2470
      pOperator = createProjectOperatorInfo(NULL, (SProjectPhysiNode*)pPhyNode, pTaskInfo);
H
Haojun Liao 已提交
2471 2472
    } else {
      ASSERT(0);
H
Haojun Liao 已提交
2473
    }
2474 2475 2476 2477 2478

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

2479
    return pOperator;
H
Haojun Liao 已提交
2480 2481
  }

2482
  size_t          size = LIST_LENGTH(pPhyNode->pChildren);
2483
  SOperatorInfo** ops = taosMemoryCalloc(size, POINTER_BYTES);
2484 2485 2486 2487
  if (ops == NULL) {
    return NULL;
  }

dengyihao's avatar
dengyihao 已提交
2488
  for (int32_t i = 0; i < size; ++i) {
2489
    SPhysiNode* pChildNode = (SPhysiNode*)nodesListGetNode(pPhyNode->pChildren, i);
2490
    ops[i] = createOperatorTree(pChildNode, pTaskInfo, pHandle, pTagCond, pTagIndexCond, pUser);
2491
    if (ops[i] == NULL) {
H
Haojun Liao 已提交
2492
      taosMemoryFree(ops);
2493 2494
      return NULL;
    }
2495
  }
H
Haojun Liao 已提交
2496

2497
  SOperatorInfo* pOptr = NULL;
H
Haojun Liao 已提交
2498
  if (QUERY_NODE_PHYSICAL_PLAN_PROJECT == type) {
2499
    pOptr = createProjectOperatorInfo(ops[0], (SProjectPhysiNode*)pPhyNode, pTaskInfo);
2500
  } else if (QUERY_NODE_PHYSICAL_PLAN_HASH_AGG == type) {
H
Haojun Liao 已提交
2501 2502
    SAggPhysiNode* pAggNode = (SAggPhysiNode*)pPhyNode;
    if (pAggNode->pGroupKeys != NULL) {
H
Haojun Liao 已提交
2503
      pOptr = createGroupOperatorInfo(ops[0], pAggNode, pTaskInfo);
H
Haojun Liao 已提交
2504
    } else {
H
Haojun Liao 已提交
2505
      pOptr = createAggregateOperatorInfo(ops[0], pAggNode, pTaskInfo);
H
Haojun Liao 已提交
2506
    }
2507
  } else if (QUERY_NODE_PHYSICAL_PLAN_HASH_INTERVAL == type) {
H
Haojun Liao 已提交
2508
    SIntervalPhysiNode* pIntervalPhyNode = (SIntervalPhysiNode*)pPhyNode;
H
Haojun Liao 已提交
2509

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

2567
  taosMemoryFree(ops);
2568 2569 2570 2571
  if (pOptr) {
    pOptr->resultDataBlockId = pPhyNode->pOutputDataBlockDesc->dataBlockId;
  }

2572
  return pOptr;
2573
}
H
Haojun Liao 已提交
2574

L
Liu Jicong 已提交
2575 2576 2577 2578 2579 2580 2581 2582 2583 2584 2585 2586 2587
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 {
2588 2589 2590
    SStreamScanInfo* pInfo = pOperator->info;
    ASSERT(pInfo->pTableScanOp->operatorType == QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN);
    *ppInfo = pInfo->pTableScanOp->info;
L
Liu Jicong 已提交
2591 2592 2593 2594
    return 0;
  }
}

2595 2596 2597 2598 2599 2600 2601 2602 2603 2604 2605 2606 2607 2608 2609 2610 2611 2612 2613 2614 2615 2616
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;
}

2617
#if 0
L
Liu Jicong 已提交
2618 2619 2620 2621 2622
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;
  }
2623

L
Liu Jicong 已提交
2624 2625 2626 2627
  STableScanPhysiNode* pNode = NULL;
  if (extractTableScanNode(plan->pNode, &pNode) < 0) {
    ASSERT(0);
  }
2628

H
Haojun Liao 已提交
2629
  tsdbReaderClose(pTableScanInfo->dataReader);
2630

L
Liu Jicong 已提交
2631
  STableListInfo info = {0};
H
Haojun Liao 已提交
2632
  pTableScanInfo->dataReader = doCreateDataReader(pNode, pHandle, &info, NULL);
L
Liu Jicong 已提交
2633 2634 2635 2636
  if (pTableScanInfo->dataReader == NULL) {
    ASSERT(0);
    qError("failed to create data reader");
    return TSDB_CODE_QRY_APP_ERROR;
2637
  }
L
Liu Jicong 已提交
2638
  // TODO: set uid and ts to data reader
2639 2640
  return 0;
}
2641
#endif
2642

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

D
dapan1121 已提交
2646
  switch (pNode->type) {
D
dapan1121 已提交
2647 2648 2649 2650 2651 2652
    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 已提交
2653

D
dapan1121 已提交
2654 2655 2656
      *pParam = pInserterParam;
      break;
    }
D
dapan1121 已提交
2657
    case QUERY_NODE_PHYSICAL_PLAN_DELETE: {
2658
      SDeleterParam* pDeleterParam = taosMemoryCalloc(1, sizeof(SDeleterParam));
D
dapan1121 已提交
2659 2660 2661
      if (NULL == pDeleterParam) {
        return TSDB_CODE_OUT_OF_MEMORY;
      }
H
Haojun Liao 已提交
2662 2663 2664 2665
      int32_t tbNum = tableListGetSize(pTask->pTableInfoList);
      pDeleterParam->suid = tableListGetSuid(pTask->pTableInfoList);

      // TODO extract uid list
D
dapan1121 已提交
2666 2667 2668 2669 2670
      pDeleterParam->pUidList = taosArrayInit(tbNum, sizeof(uint64_t));
      if (NULL == pDeleterParam->pUidList) {
        taosMemoryFree(pDeleterParam);
        return TSDB_CODE_OUT_OF_MEMORY;
      }
H
Haojun Liao 已提交
2671

D
dapan1121 已提交
2672
      for (int32_t i = 0; i < tbNum; ++i) {
H
Haojun Liao 已提交
2673
        STableKeyInfo* pTable = tableListGetInfo(pTask->pTableInfoList, i);
D
dapan1121 已提交
2674 2675 2676 2677 2678 2679 2680 2681 2682 2683 2684 2685 2686
        taosArrayPush(pDeleterParam->pUidList, &pTable->uid);
      }

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

  return TSDB_CODE_SUCCESS;
}

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

D
dapan1121 已提交
2691
  *pTaskInfo = createExecTaskInfo(queryId, taskId, model, pPlan->dbFName);
H
Haojun Liao 已提交
2692 2693 2694
  if (*pTaskInfo == NULL) {
    goto _complete;
  }
H
Haojun Liao 已提交
2695

2696
  if (pHandle) {
L
Liu Jicong 已提交
2697
    /*(*pTaskInfo)->streamInfo.fillHistoryVer1 = pHandle->fillHistoryVer1;*/
2698 2699 2700
    if (pHandle->pStateBackend) {
      (*pTaskInfo)->streamInfo.pState = pHandle->pStateBackend;
    }
H
Haojun Liao 已提交
2701 2702
  }

2703
  (*pTaskInfo)->sql = sql;
D
dapan1121 已提交
2704
  sql = NULL;
H
Haojun Liao 已提交
2705

2706
  (*pTaskInfo)->pSubplan = pPlan;
2707 2708
  (*pTaskInfo)->pRoot =
      createOperatorTree(pPlan->pNode, *pTaskInfo, pHandle, pPlan->pTagCond, pPlan->pTagIndexCond, pPlan->user);
L
Liu Jicong 已提交
2709

D
dapan1121 已提交
2710
  if (NULL == (*pTaskInfo)->pRoot) {
H
Haojun Liao 已提交
2711
    terrno = (*pTaskInfo)->code;
D
dapan1121 已提交
2712
    goto _complete;
2713 2714
  }

H
Haojun Liao 已提交
2715
  return TSDB_CODE_SUCCESS;
H
Haojun Liao 已提交
2716

H
Haojun Liao 已提交
2717
_complete:
D
dapan1121 已提交
2718
  taosMemoryFree(sql);
H
Haojun Liao 已提交
2719
  doDestroyTask(*pTaskInfo);
H
Haojun Liao 已提交
2720
  return terrno;
H
Haojun Liao 已提交
2721 2722
}

H
Haojun Liao 已提交
2723 2724 2725 2726 2727
static void freeBlock(void* pParam) {
  SSDataBlock* pBlock = *(SSDataBlock**)pParam;
  blockDataDestroy(pBlock);
}

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

H
Haojun Liao 已提交
2731
  pTaskInfo->pTableInfoList = tableListDestroy(pTaskInfo->pTableInfoList);
H
Haojun Liao 已提交
2732
  destroyOperatorInfo(pTaskInfo->pRoot);
2733
  cleanupTableSchemaInfo(&pTaskInfo->schemaInfo);
2734
  cleanupStreamInfo(&pTaskInfo->streamInfo);
2735

D
dapan1121 已提交
2736
  if (!pTaskInfo->localFetch.localExec) {
D
dapan1121 已提交
2737 2738
    nodesDestroyNode((SNode*)pTaskInfo->pSubplan);
  }
2739

H
Haojun Liao 已提交
2740
  taosArrayDestroyEx(pTaskInfo->pResultBlockList, freeBlock);
D
dapan1121 已提交
2741
  taosArrayDestroy(pTaskInfo->stopInfo.pStopInfo);
wafwerar's avatar
wafwerar 已提交
2742 2743 2744
  taosMemoryFreeClear(pTaskInfo->sql);
  taosMemoryFreeClear(pTaskInfo->id.str);
  taosMemoryFreeClear(pTaskInfo);
2745 2746 2747 2748
}

static int64_t getQuerySupportBufSize(size_t numOfTables) {
  size_t s1 = sizeof(STableQueryInfo);
L
Liu Jicong 已提交
2749 2750
  //  size_t s3 = sizeof(STableCheckInfo);  buffer consumption in tsdb
  return (int64_t)(s1 * 1.5 * numOfTables);
2751 2752 2753 2754 2755 2756 2757
}

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 已提交
2758
    while (1) {
2759 2760 2761 2762 2763 2764 2765 2766 2767 2768 2769 2770 2771 2772 2773 2774 2775 2776 2777 2778 2779 2780 2781 2782 2783 2784
      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 已提交
2785

H
Haojun Liao 已提交
2786
int32_t getOperatorExplainExecInfo(SOperatorInfo* operatorInfo, SArray* pExecInfoList) {
2787
  SExplainExecInfo  execInfo = {0};
H
Haojun Liao 已提交
2788
  SExplainExecInfo* pExplainInfo = taosArrayPush(pExecInfoList, &execInfo);
2789

H
Haojun Liao 已提交
2790 2791 2792 2793 2794
  pExplainInfo->numOfRows = operatorInfo->resultInfo.totalRows;
  pExplainInfo->startupCost = operatorInfo->cost.openCost;
  pExplainInfo->totalCost = operatorInfo->cost.totalCost;
  pExplainInfo->verboseLen = 0;
  pExplainInfo->verboseInfo = NULL;
D
dapan1121 已提交
2795

2796
  if (operatorInfo->fpSet.getExplainFn) {
2797 2798
    int32_t code =
        operatorInfo->fpSet.getExplainFn(operatorInfo, &pExplainInfo->verboseInfo, &pExplainInfo->verboseLen);
D
dapan1121 已提交
2799
    if (code) {
2800
      qError("%s operator getExplainFn failed, code:%s", GET_TASKID(operatorInfo->pTaskInfo), tstrerror(code));
D
dapan1121 已提交
2801 2802 2803
      return code;
    }
  }
dengyihao's avatar
dengyihao 已提交
2804

D
dapan1121 已提交
2805
  int32_t code = 0;
D
dapan1121 已提交
2806
  for (int32_t i = 0; i < operatorInfo->numOfDownstream; ++i) {
H
Haojun Liao 已提交
2807 2808
    code = getOperatorExplainExecInfo(operatorInfo->pDownstream[i], pExecInfoList);
    if (code != TSDB_CODE_SUCCESS) {
2809
      //      taosMemoryFreeClear(*pRes);
D
dapan1121 已提交
2810 2811 2812 2813 2814
      return TSDB_CODE_QRY_OUT_OF_MEMORY;
    }
  }

  return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
2815
}
5
54liuyao 已提交
2816

2817 2818
int32_t setOutputBuf(SStreamState* pState, STimeWindow* win, SResultRow** pResult, int64_t tableGroupId,
                     SqlFunctionCtx* pCtx, int32_t numOfOutput, int32_t* rowEntryInfoOffset, SAggSupporter* pAggSup) {
2819 2820 2821 2822 2823 2824
  SWinKey key = {
      .ts = win->skey,
      .groupId = tableGroupId,
  };
  char*   value = NULL;
  int32_t size = pAggSup->resultRowSize;
5
54liuyao 已提交
2825

2826
  if (streamStateAddIfNotExist(pState, &key, (void**)&value, &size) < 0) {
2827 2828 2829 2830 2831 2832 2833 2834 2835 2836
    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;
}

2837 2838
int32_t releaseOutputBuf(SStreamState* pState, SWinKey* pKey, SResultRow* pResult) {
  streamStateReleaseBuf(pState, pKey, pResult);
2839 2840 2841
  return TSDB_CODE_SUCCESS;
}

2842 2843
int32_t saveOutputBuf(SStreamState* pState, SWinKey* pKey, SResultRow* pResult, int32_t resSize) {
  streamStatePut(pState, pKey, pResult, resSize);
2844 2845 2846
  return TSDB_CODE_SUCCESS;
}

2847
int32_t buildDataBlockFromGroupRes(SOperatorInfo* pOperator, SStreamState* pState, SSDataBlock* pBlock, SExprSupp* pSup,
2848
                                   SGroupResInfo* pGroupResInfo) {
2849
  SExecTaskInfo*  pTaskInfo = pOperator->pTaskInfo;
2850 2851 2852 2853 2854 2855 2856 2857 2858 2859 2860 2861
  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 = {
2862 2863
            .ts = *(TSKEY*)pPos->key,
            .groupId = pPos->groupId,
2864
    };
2865
    int32_t code = streamStateGet(pState, &key, &pVal, &size);
2866 2867 2868 2869 2870 2871
    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;
2872
      releaseOutputBuf(pState, &key, pRow);
2873 2874 2875
      continue;
    }

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

    if (pBlock->info.rows + pRow->numOfRows > pBlock->info.capacity) {
      ASSERT(pBlock->info.rows > 0);
2895
      releaseOutputBuf(pState, &key, pRow);
2896 2897 2898 2899 2900 2901 2902 2903 2904 2905
      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) {
2906 2907 2908 2909
        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);
2910 2911 2912 2913 2914 2915 2916 2917 2918 2919 2920 2921 2922
        }
      } 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 已提交
2923

2924
    pBlock->info.rows += pRow->numOfRows;
2925
    releaseOutputBuf(pState, &key, pRow);
2926 2927 2928 2929
  }
  blockDataUpdateTsWindow(pBlock, 0);
  return TSDB_CODE_SUCCESS;
}
5
54liuyao 已提交
2930 2931 2932 2933 2934 2935 2936

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

2937
int32_t buildSessionResultDataBlock(SOperatorInfo* pOperator, SStreamState* pState, SSDataBlock* pBlock,
5
54liuyao 已提交
2938
                                    SExprSupp* pSup, SGroupResInfo* pGroupResInfo) {
2939
  SExecTaskInfo*  pTaskInfo = pOperator->pTaskInfo;
5
54liuyao 已提交
2940 2941 2942 2943 2944 2945 2946 2947 2948 2949 2950 2951 2952
  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);
2953 2954
    if (code == -1) {
      // coverity scan
5
54liuyao 已提交
2955
      pGroupResInfo->index += 1;
2956 2957
      continue;
    }
5
54liuyao 已提交
2958 2959 2960 2961 2962 2963 2964 2965 2966
    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 已提交
2967 2968
    if (pBlock->info.id.groupId == 0) {
      pBlock->info.id.groupId = pKey->groupId;
2969

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