executorimpl.c 126.3 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"
X
Xiaoyu Wang 已提交
23
#include "tref.h"
24

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

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

H
Haojun Liao 已提交
39
#define IS_MAIN_SCAN(runtime)          ((runtime)->scanFlag == MAIN_SCAN)
40 41 42 43 44 45
#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 已提交
46
  uint32_t v = taosRand();
47 48 49 50

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

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

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

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

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

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

83
static void setBlockSMAInfo(SqlFunctionCtx* pCtx, SExprInfo* pExpr, SSDataBlock* pBlock);
84

X
Xiaoyu Wang 已提交
85
static void releaseQueryBuf(size_t numOfTables);
86

87 88
static void destroyFillOperatorInfo(void* param);
static void destroyProjectOperatorInfo(void* param);
H
Haojun Liao 已提交
89
static void destroySortOperatorInfo(void* param);
90
static void destroyAggOperatorInfo(void* param);
X
Xiaoyu Wang 已提交
91

92
static void destroyIntervalOperatorInfo(void* param);
H
Haojun Liao 已提交
93

H
Haojun Liao 已提交
94

95 96
static void destroyOperatorInfo(SOperatorInfo* pOperator);

H
Haojun Liao 已提交
97
void setOperatorCompleted(SOperatorInfo* pOperator) {
98
  pOperator->status = OP_EXEC_DONE;
H
Haojun Liao 已提交
99
  ASSERT(pOperator->pTaskInfo != NULL);
100

101
  pOperator->cost.totalCost = (taosGetTimestampUs() - pOperator->pTaskInfo->cost.start * 1000) / 1000.0;
H
Haojun Liao 已提交
102
  setTaskStatus(pOperator->pTaskInfo, TASK_COMPLETED);
103
}
104

H
Haojun Liao 已提交
105 106 107 108 109 110 111 112 113 114
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 已提交
115
int32_t operatorDummyOpenFn(SOperatorInfo* pOperator) {
116
  OPTR_SET_OPENED(pOperator);
117
  pOperator->cost.openCost = 0;
H
Haojun Liao 已提交
118
  return TSDB_CODE_SUCCESS;
H
Haojun Liao 已提交
119 120
}

H
Haojun Liao 已提交
121 122
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) {
123 124 125 126 127 128 129 130 131 132 133
  SOperatorFpSet fpSet = {
      ._openFn = openFn,
      .getNextFn = nextFn,
      .cleanupFn = cleanup,
      .closeFn = closeFn,
      .getExplainFn = explain,
  };

  return fpSet;
}

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

137
static void initCtxOutputBuffer(SqlFunctionCtx* pCtx, int32_t size);
138
static void doSetTableGroupOutputBuf(SOperatorInfo* pOperator, int32_t numOfOutput, uint64_t groupId);
139

140
#if 0
L
Liu Jicong 已提交
141 142
static bool chkResultRowFromKey(STaskRuntimeEnv* pRuntimeEnv, SResultRowInfo* pResultRowInfo, char* pData,
                                int16_t bytes, bool masterscan, uint64_t uid) {
143 144 145
  bool existed = false;
  SET_RES_WINDOW_KEY(pRuntimeEnv->keyBuf, pData, bytes, uid);

L
Liu Jicong 已提交
146 147
  SResultRow** p1 =
      (SResultRow**)taosHashGet(pRuntimeEnv->pResultRowHashTable, pRuntimeEnv->keyBuf, GET_RES_WINDOW_KEY_LEN(bytes));
148 149 150 151 152 153 154 155 156 157 158

  // in case of repeat scan/reverse scan, no new time window added.
  if (QUERY_IS_INTERVAL_QUERY(pRuntimeEnv->pQueryAttr)) {
    if (!masterscan) {  // the *p1 may be NULL in case of sliding+offset exists.
      return p1 != NULL;
    }

    if (p1 != NULL) {
      if (pResultRowInfo->size == 0) {
        existed = false;
      } else if (pResultRowInfo->size == 1) {
dengyihao's avatar
dengyihao 已提交
159
        //        existed = (pResultRowInfo->pResult[0] == (*p1));
160 161
      } else {  // check if current pResultRowInfo contains the existed pResultRow
        SET_RES_EXT_WINDOW_KEY(pRuntimeEnv->keyBuf, pData, bytes, uid, pResultRowInfo);
L
Liu Jicong 已提交
162 163
        int64_t* index =
            taosHashGet(pRuntimeEnv->pResultRowListSet, pRuntimeEnv->keyBuf, GET_RES_EXT_WINDOW_KEY_LEN(bytes));
164 165 166 167 168 169 170 171 172 173 174 175 176
        if (index != NULL) {
          existed = true;
        } else {
          existed = false;
        }
      }
    }

    return existed;
  }

  return p1 != NULL;
}
177
#endif
178

179
SResultRow* getNewResultRow(SDiskbasedBuf* pResultBuf, int32_t* currentPageId, int32_t interBufSize) {
L
Liu Jicong 已提交
180
  SFilePage* pData = NULL;
181 182 183

  // in the first scan, new space needed for results
  int32_t pageId = -1;
184
  if (*currentPageId == -1) {
185
    pData = getNewBufPage(pResultBuf, &pageId);
186 187
    pData->num = sizeof(SFilePage);
  } else {
188 189
    pData = getBufPage(pResultBuf, *currentPageId);
    pageId = *currentPageId;
190

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

195
      pData = getNewBufPage(pResultBuf, &pageId);
196 197 198 199 200 201 202 203 204 205
      if (pData != NULL) {
        pData->num = sizeof(SFilePage);
      }
    }
  }

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

206 207
  setBufPageDirty(pData, true);

208 209 210 211
  // 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;
212
  *currentPageId = pageId;
213

wmmhello's avatar
wmmhello 已提交
214
  pData->num += interBufSize;
215 216 217
  return pResultRow;
}

218 219 220 221 222 223 224
/**
 * the struct of key in hash table
 * +----------+---------------+
 * | group id |   key data    |
 * | 8 bytes  | actual length |
 * +----------+---------------+
 */
225 226 227
SResultRow* doSetResultOutBufByKey(SDiskbasedBuf* pResultBuf, SResultRowInfo* pResultRowInfo, char* pData,
                                   int16_t bytes, bool masterscan, uint64_t groupId, SExecTaskInfo* pTaskInfo,
                                   bool isIntervalQuery, SAggSupporter* pSup) {
228
  SET_RES_WINDOW_KEY(pSup->keyBuf, pData, bytes, groupId);
H
Haojun Liao 已提交
229

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

233 234
  SResultRow* pResult = NULL;

H
Haojun Liao 已提交
235 236
  // in case of repeat scan/reverse scan, no new time window added.
  if (isIntervalQuery) {
237
    if (masterscan && p1 != NULL) {  // the *p1 may be NULL in case of sliding+offset exists.
238
      pResult = getResultRowByPos(pResultBuf, p1, true);
239
      ASSERT(pResult->pageId == p1->pageId && pResult->offset == p1->offset);
H
Haojun Liao 已提交
240 241
    }
  } else {
dengyihao's avatar
dengyihao 已提交
242 243
    // In case of group by column query, the required SResultRow object must be existInCurrentResusltRowInfo in the
    // pResultRowInfo object.
H
Haojun Liao 已提交
244
    if (p1 != NULL) {
245
      // todo
246
      pResult = getResultRowByPos(pResultBuf, p1, true);
247
      ASSERT(pResult->pageId == p1->pageId && pResult->offset == p1->offset);
H
Haojun Liao 已提交
248 249 250
    }
  }

L
Liu Jicong 已提交
251
  // 1. close current opened time window
252
  if (pResultRowInfo->cur.pageId != -1 && ((pResult == NULL) || (pResult->pageId != pResultRowInfo->cur.pageId))) {
253
    SResultRowPosition pos = pResultRowInfo->cur;
X
Xiaoyu Wang 已提交
254
    SFilePage*         pPage = getBufPage(pResultBuf, pos.pageId);
255 256 257 258 259
    releaseBufPage(pResultBuf, pPage);
  }

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

263 264
    // add a new result set for a new group
    SResultRowPosition pos = {.pageId = pResult->pageId, .offset = pResult->offset};
265
    tSimpleHashPut(pSup->pResultRowHashTable, pSup->keyBuf, GET_RES_WINDOW_KEY_LEN(bytes), &pos,
L
Liu Jicong 已提交
266
                   sizeof(SResultRowPosition));
H
Haojun Liao 已提交
267 268
  }

269 270 271
  // 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 已提交
272
  // too many time window in query
273
  if (pTaskInfo->execModel == OPTR_EXEC_MODEL_BATCH &&
274
      tSimpleHashGetSize(pSup->pResultRowHashTable) > MAX_INTERVAL_TIME_WINDOW) {
275
    T_LONG_JMP(pTaskInfo->env, TSDB_CODE_QRY_TOO_MANY_TIMEWINDOW);
H
Haojun Liao 已提交
276 277
  }

H
Haojun Liao 已提交
278
  return pResult;
H
Haojun Liao 已提交
279 280
}

281
// a new buffer page for each table. Needs to opt this design
L
Liu Jicong 已提交
282
static int32_t addNewWindowResultBuf(SResultRow* pWindowRes, SDiskbasedBuf* pResultBuf, int32_t tid, uint32_t size) {
283 284 285 286
  if (pWindowRes->pageId != -1) {
    return 0;
  }

L
Liu Jicong 已提交
287
  SFilePage* pData = NULL;
288 289 290

  // in the first scan, new space needed for results
  int32_t pageId = -1;
291
  SIDList list = getDataBufPagesIdList(pResultBuf);
292 293

  if (taosArrayGetSize(list) == 0) {
294
    pData = getNewBufPage(pResultBuf, &pageId);
295
    pData->num = sizeof(SFilePage);
296 297
  } else {
    SPageInfo* pi = getLastPageInfo(list);
298
    pData = getBufPage(pResultBuf, getPageId(pi));
299
    pageId = getPageId(pi);
300

301
    if (pData->num + size > getBufPageSize(pResultBuf)) {
302
      // release current page first, and prepare the next one
303
      releaseBufPageInfo(pResultBuf, pi);
304

305
      pData = getNewBufPage(pResultBuf, &pageId);
306
      if (pData != NULL) {
307
        pData->num = sizeof(SFilePage);
308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327
      }
    }
  }

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

328
//  query_range_start, query_range_end, window_duration, window_start, window_end
329
void initExecTimeWindowInfo(SColumnInfoData* pColData, STimeWindow* pQueryWindow) {
330 331 332
  pColData->info.type = TSDB_DATA_TYPE_TIMESTAMP;
  pColData->info.bytes = sizeof(int64_t);

H
Haojun Liao 已提交
333
  colInfoDataEnsureCapacity(pColData, 5, false);
334 335 336 337 338 339 340 341 342
  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);
}

L
Liu Jicong 已提交
343
void cleanupExecTimeWindowInfo(SColumnInfoData* pColData) { colDataDestroy(pColData); }
H
Haojun Liao 已提交
344

345 346 347 348 349 350 351 352 353 354 355 356 357 358
typedef struct {
  bool    hasAgg;
  int32_t numOfRows;
  int32_t startOffset;
} SFunctionCtxStatus;

static void functionCtxSave(SqlFunctionCtx* pCtx, SFunctionCtxStatus* pStatus) {
  pStatus->hasAgg = pCtx->input.colDataAggIsSet;
  pStatus->numOfRows = pCtx->input.numOfRows;
  pStatus->startOffset = pCtx->input.startRowIndex;
}

static void functionCtxRestore(SqlFunctionCtx* pCtx, SFunctionCtxStatus* pStatus) {
  pCtx->input.colDataAggIsSet = pStatus->hasAgg;
H
Haojun Liao 已提交
359
  pCtx->input.numOfRows = pStatus->numOfRows;
360 361 362 363 364
  pCtx->input.startRowIndex = pStatus->startOffset;
}

void doApplyFunctions(SExecTaskInfo* taskInfo, SqlFunctionCtx* pCtx, SColumnInfoData* pTimeWindowData, int32_t offset,
                      int32_t forwardStep, int32_t numOfTotal, int32_t numOfOutput) {
365
  for (int32_t k = 0; k < numOfOutput; ++k) {
H
Haojun Liao 已提交
366
    // keep it temporarily
367 368
    SFunctionCtxStatus status = {0};
    functionCtxSave(&pCtx[k], &status);
369

370
    pCtx[k].input.startRowIndex = offset;
371
    pCtx[k].input.numOfRows = forwardStep;
372 373 374

    // not a whole block involved in query processing, statistics data can not be used
    // NOTE: the original value of isSet have been changed here
375 376
    if (pCtx[k].input.colDataAggIsSet && forwardStep < numOfTotal) {
      pCtx[k].input.colDataAggIsSet = false;
377 378
    }

379 380
    if (fmIsWindowPseudoColumnFunc(pCtx[k].functionId)) {
      SResultRowEntryInfo* pEntryInfo = GET_RES_INFO(&pCtx[k]);
381 382

      char* p = GET_ROWCELL_INTERBUF(pEntryInfo);
383

384
      SColumnInfoData idata = {0};
dengyihao's avatar
dengyihao 已提交
385
      idata.info.type = TSDB_DATA_TYPE_BIGINT;
386
      idata.info.bytes = tDataTypes[TSDB_DATA_TYPE_BIGINT].bytes;
dengyihao's avatar
dengyihao 已提交
387
      idata.pData = p;
388 389 390 391

      SScalarParam out = {.columnData = &idata};
      SScalarParam tw = {.numOfRows = 5, .columnData = pTimeWindowData};
      pCtx[k].sfp.process(&tw, 1, &out);
392
      pEntryInfo->numOfRes = 1;
393 394 395 396 397 398 399 400
    } 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;
401
          T_LONG_JMP(taskInfo->env, code);
402
        }
403
      }
404

405
      // restore it
406
      functionCtxRestore(&pCtx[k], &status);
407
    }
408 409 410
  }
}

411 412
static int32_t doSetInputDataBlock(SExprSupp* pExprSup, SSDataBlock* pBlock, int32_t order, int32_t scanFlag,
                                   bool createDummyCol);
413

414 415 416
static void doSetInputDataBlockInfo(SExprSupp* pExprSup, SSDataBlock* pBlock, int32_t order) {
  SqlFunctionCtx* pCtx = pExprSup->pCtx;
  for (int32_t i = 0; i < pExprSup->numOfExprs; ++i) {
417
    pCtx[i].order = order;
418
    pCtx[i].input.numOfRows = pBlock->info.rows;
419
    setBlockSMAInfo(&pCtx[i], &pExprSup->pExprInfo[i], pBlock);
420
    pCtx[i].pSrcBlock = pBlock;
421 422 423
  }
}

424
void setInputDataBlock(SExprSupp* pExprSup, SSDataBlock* pBlock, int32_t order, int32_t scanFlag, bool createDummyCol) {
425
  if (pBlock->pBlockAgg != NULL) {
426
    doSetInputDataBlockInfo(pExprSup, pBlock, order);
427
  } else {
428
    doSetInputDataBlock(pExprSup, pBlock, order, scanFlag, createDummyCol);
H
Haojun Liao 已提交
429
  }
430 431
}

L
Liu Jicong 已提交
432 433
static int32_t doCreateConstantValColumnInfo(SInputColumnInfoData* pInput, SFunctParam* pFuncParam, int32_t paramIndex,
                                             int32_t numOfRows) {
434 435 436 437 438 439 440 441
  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)
442 443
    pColInfo->info.type = pFuncParam->param.nType;
    pColInfo->info.bytes = pFuncParam->param.nLen;
444 445

    pInput->pData[paramIndex] = pColInfo;
446 447
  } else {
    pColInfo = pInput->pData[paramIndex];
448 449
  }

H
Haojun Liao 已提交
450
  colInfoDataEnsureCapacity(pColInfo, numOfRows, false);
451

452
  int8_t type = pFuncParam->param.nType;
453 454
  if (type == TSDB_DATA_TYPE_BIGINT || type == TSDB_DATA_TYPE_UBIGINT) {
    int64_t v = pFuncParam->param.i;
dengyihao's avatar
dengyihao 已提交
455
    for (int32_t i = 0; i < numOfRows; ++i) {
456 457 458 459
      colDataAppendInt64(pColInfo, i, &v);
    }
  } else if (type == TSDB_DATA_TYPE_DOUBLE) {
    double v = pFuncParam->param.d;
dengyihao's avatar
dengyihao 已提交
460
    for (int32_t i = 0; i < numOfRows; ++i) {
461 462
      colDataAppendDouble(pColInfo, i, &v);
    }
463
  } else if (type == TSDB_DATA_TYPE_VARCHAR) {
L
Liu Jicong 已提交
464
    char* tmp = taosMemoryMalloc(pFuncParam->param.nLen + VARSTR_HEADER_SIZE);
465
    STR_WITH_SIZE_TO_VARSTR(tmp, pFuncParam->param.pz, pFuncParam->param.nLen);
L
Liu Jicong 已提交
466
    for (int32_t i = 0; i < numOfRows; ++i) {
467 468
      colDataAppend(pColInfo, i, tmp, false);
    }
H
Haojun Liao 已提交
469
    taosMemoryFree(tmp);
470 471 472 473 474
  }

  return TSDB_CODE_SUCCESS;
}

475 476
static int32_t doSetInputDataBlock(SExprSupp* pExprSup, SSDataBlock* pBlock, int32_t order, int32_t scanFlag,
                                   bool createDummyCol) {
477
  int32_t         code = TSDB_CODE_SUCCESS;
478
  SqlFunctionCtx* pCtx = pExprSup->pCtx;
479

480
  for (int32_t i = 0; i < pExprSup->numOfExprs; ++i) {
L
Liu Jicong 已提交
481
    pCtx[i].order = order;
482 483
    pCtx[i].input.numOfRows = pBlock->info.rows;

L
Liu Jicong 已提交
484
    pCtx[i].pSrcBlock = pBlock;
X
Xiaoyu Wang 已提交
485
    pCtx[i].scanFlag = scanFlag;
H
Haojun Liao 已提交
486

487
    SInputColumnInfoData* pInput = &pCtx[i].input;
488
    pInput->uid = pBlock->info.uid;
C
Cary Xu 已提交
489
    pInput->colDataAggIsSet = false;
490

491
    SExprInfo* pOneExpr = &pExprSup->pExprInfo[i];
492
    for (int32_t j = 0; j < pOneExpr->base.numOfParams; ++j) {
dengyihao's avatar
dengyihao 已提交
493
      SFunctParam* pFuncParam = &pOneExpr->base.pParam[j];
G
Ganlin Zhao 已提交
494 495
      if (pFuncParam->type == FUNC_PARAM_TYPE_COLUMN) {
        int32_t slotId = pFuncParam->pCol->slotId;
dengyihao's avatar
dengyihao 已提交
496
        pInput->pData[j] = taosArrayGet(pBlock->pDataBlock, slotId);
497 498 499
        pInput->totalRows = pBlock->info.rows;
        pInput->numOfRows = pBlock->info.rows;
        pInput->startRowIndex = 0;
500

501
        // NOTE: the last parameter is the primary timestamp column
H
Haojun Liao 已提交
502
        // todo: refactor this
503
        if (fmIsImplicitTsFunc(pCtx[i].functionId) && (j == pOneExpr->base.numOfParams - 1)) {
L
Liu Jicong 已提交
504
          pInput->pPTS = pInput->pData[j];  // in case of merge function, this is not always the ts column data.
505
          //          ASSERT(pInput->pPTS->info.type == TSDB_DATA_TYPE_TIMESTAMP);
506
        }
507 508
        ASSERT(pInput->pData[j] != NULL);
      } else if (pFuncParam->type == FUNC_PARAM_TYPE_VALUE) {
509 510 511
        // 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) {
512 513 514 515
          pInput->totalRows = pBlock->info.rows;
          pInput->numOfRows = pBlock->info.rows;
          pInput->startRowIndex = 0;

516
          code = doCreateConstantValColumnInfo(pInput, pFuncParam, j, pBlock->info.rows);
517 518 519
          if (code != TSDB_CODE_SUCCESS) {
            return code;
          }
520
        }
G
Ganlin Zhao 已提交
521 522
      }
    }
H
Haojun Liao 已提交
523
  }
524 525

  return code;
H
Haojun Liao 已提交
526 527
}

528
static int32_t doAggregateImpl(SOperatorInfo* pOperator, SqlFunctionCtx* pCtx) {
529
  for (int32_t k = 0; k < pOperator->exprSupp.numOfExprs; ++k) {
H
Haojun Liao 已提交
530
    if (functionNeedToExecute(&pCtx[k])) {
531
      // todo add a dummy funtion to avoid process check
532 533 534
      if (pCtx[k].fpSet.process == NULL) {
        continue;
      }
H
Haojun Liao 已提交
535

536 537 538 539
      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;
540
      }
541 542
    }
  }
543 544

  return TSDB_CODE_SUCCESS;
545 546
}

H
Haojun Liao 已提交
547
static void setPseudoOutputColInfo(SSDataBlock* pResult, SqlFunctionCtx* pCtx, SArray* pPseudoList) {
dengyihao's avatar
dengyihao 已提交
548
  size_t num = (pPseudoList != NULL) ? taosArrayGetSize(pPseudoList) : 0;
H
Haojun Liao 已提交
549 550 551 552 553
  for (int32_t i = 0; i < num; ++i) {
    pCtx[i].pOutput = taosArrayGet(pResult->pDataBlock, i);
  }
}

554
int32_t projectApplyFunctions(SExprInfo* pExpr, SSDataBlock* pResult, SSDataBlock* pSrcBlock, SqlFunctionCtx* pCtx,
X
Xiaoyu Wang 已提交
555
                              int32_t numOfOutput, SArray* pPseudoList) {
H
Haojun Liao 已提交
556
  setPseudoOutputColInfo(pResult, pCtx, pPseudoList);
557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575 576

  if (pSrcBlock == NULL) {
    for (int32_t k = 0; k < numOfOutput; ++k) {
      int32_t outputSlotId = pExpr[k].base.resSchema.slotId;

      ASSERT(pExpr[k].pExpr->nodeType == QUERY_NODE_VALUE);
      SColumnInfoData* pColInfoData = taosArrayGet(pResult->pDataBlock, outputSlotId);

      int32_t type = pExpr[k].base.pParam[0].param.nType;
      if (TSDB_DATA_TYPE_NULL == type) {
        colDataAppendNNULL(pColInfoData, 0, 1);
      } else {
        colDataAppend(pColInfoData, 0, taosVariantGet(&pExpr[k].base.pParam[0].param, type), false);
      }
    }

    pResult->info.rows = 1;
    return TSDB_CODE_SUCCESS;
  }

L
Liu Jicong 已提交
577 578 579 580
  if (pResult != pSrcBlock) {
    pResult->info.groupId = pSrcBlock->info.groupId;
    memcpy(pResult->info.parTbName, pSrcBlock->info.parTbName, TSDB_TABLE_NAME_LEN);
  }
H
Haojun Liao 已提交
581

582 583
  // if the source equals to the destination, it is to create a new column as the result of scalar
  // function or some operators.
584
  bool createNewColModel = (pResult == pSrcBlock);
585 586 587
  if (createNewColModel) {
    blockDataEnsureCapacity(pResult, pResult->info.rows);
  }
588

589 590
  int32_t numOfRows = 0;

591
  for (int32_t k = 0; k < numOfOutput; ++k) {
592 593
    int32_t               outputSlotId = pExpr[k].base.resSchema.slotId;
    SqlFunctionCtx*       pfCtx = &pCtx[k];
594
    SInputColumnInfoData* pInputData = &pfCtx->input;
595

L
Liu Jicong 已提交
596
    if (pExpr[k].pExpr->nodeType == QUERY_NODE_COLUMN) {  // it is a project query
597
      SColumnInfoData* pColInfoData = taosArrayGet(pResult->pDataBlock, outputSlotId);
598
      if (pResult->info.rows > 0 && !createNewColModel) {
599
        colDataMergeCol(pColInfoData, pResult->info.rows, (int32_t*)&pResult->info.capacity, pInputData->pData[0],
600
                        pInputData->numOfRows);
601
      } else {
602
        colDataAssign(pColInfoData, pInputData->pData[0], pInputData->numOfRows, &pResult->info);
603
      }
604

605
      numOfRows = pInputData->numOfRows;
606
    } else if (pExpr[k].pExpr->nodeType == QUERY_NODE_VALUE) {
607
      SColumnInfoData* pColInfoData = taosArrayGet(pResult->pDataBlock, outputSlotId);
608

dengyihao's avatar
dengyihao 已提交
609
      int32_t offset = createNewColModel ? 0 : pResult->info.rows;
610 611 612 613 614 615 616 617

      int32_t type = pExpr[k].base.pParam[0].param.nType;
      if (TSDB_DATA_TYPE_NULL == type) {
        colDataAppendNNULL(pColInfoData, offset, pSrcBlock->info.rows);
      } else {
        for (int32_t i = 0; i < pSrcBlock->info.rows; ++i) {
          colDataAppend(pColInfoData, i + offset, taosVariantGet(&pExpr[k].base.pParam[0].param, type), false);
        }
618
      }
619 620

      numOfRows = pSrcBlock->info.rows;
H
Haojun Liao 已提交
621
    } else if (pExpr[k].pExpr->nodeType == QUERY_NODE_OPERATOR) {
622 623 624
      SArray* pBlockList = taosArrayInit(4, POINTER_BYTES);
      taosArrayPush(pBlockList, &pSrcBlock);

625
      SColumnInfoData* pResColData = taosArrayGet(pResult->pDataBlock, outputSlotId);
626
      SColumnInfoData  idata = {.info = pResColData->info, .hasNull = true};
627

628
      SScalarParam dest = {.columnData = &idata};
X
Xiaoyu Wang 已提交
629
      int32_t      code = scalarCalculate(pExpr[k].pExpr->_optrRoot.pRootNode, pBlockList, &dest);
630 631 632 633
      if (code != TSDB_CODE_SUCCESS) {
        taosArrayDestroy(pBlockList);
        return code;
      }
634

dengyihao's avatar
dengyihao 已提交
635
      int32_t startOffset = createNewColModel ? 0 : pResult->info.rows;
636
      ASSERT(pResult->info.capacity > 0);
637

638
      colDataMergeCol(pResColData, startOffset, (int32_t*)&pResult->info.capacity, &idata, dest.numOfRows);
D
dapan1121 已提交
639
      colDataDestroy(&idata);
L
Liu Jicong 已提交
640

641
      numOfRows = dest.numOfRows;
642 643
      taosArrayDestroy(pBlockList);
    } else if (pExpr[k].pExpr->nodeType == QUERY_NODE_FUNCTION) {
644 645
      // _rowts/_c0, not tbname column
      if (fmIsPseudoColumnFunc(pfCtx->functionId) && (!fmIsScanPseudoColumnFunc(pfCtx->functionId))) {
H
Haojun Liao 已提交
646
        // do nothing
647
      } else if (fmIsIndefiniteRowsFunc(pfCtx->functionId)) {
648 649
        SResultRowEntryInfo* pResInfo = GET_RES_INFO(pfCtx);
        pfCtx->fpSet.init(pfCtx, pResInfo);
650 651 652 653 654 655 656 657 658 659

        pfCtx->pOutput = taosArrayGet(pResult->pDataBlock, outputSlotId);
        pfCtx->offset = createNewColModel ? 0 : pResult->info.rows;  // set the start offset

        // set the timestamp(_rowts) output buffer
        if (taosArrayGetSize(pPseudoList) > 0) {
          int32_t* outputColIndex = taosArrayGet(pPseudoList, 0);
          pfCtx->pTsOutput = (SColumnInfoData*)pCtx[*outputColIndex].pOutput;
        }

660 661 662 663 664
        // link pDstBlock to set selectivity value
        if (pfCtx->subsidiaries.num > 0) {
          pfCtx->pDstBlock = pResult;
        }

665
        numOfRows = pfCtx->fpSet.process(pfCtx);
H
Haojun Liao 已提交
666
      } else if (fmIsAggFunc(pfCtx->functionId)) {
G
Ganlin Zhao 已提交
667
        // selective value output should be set during corresponding function execution
668 669 670
        if (fmIsSelectValueFunc(pfCtx->functionId)) {
          continue;
        }
671 672
        // _group_key function for "partition by tbname" + csum(col_name) query
        SColumnInfoData* pOutput = taosArrayGet(pResult->pDataBlock, outputSlotId);
673
        int32_t          slotId = pfCtx->param[0].pCol->slotId;
674 675 676

        // todo handle the json tag
        SColumnInfoData* pInput = taosArrayGet(pSrcBlock->pDataBlock, slotId);
677
        for (int32_t f = 0; f < pSrcBlock->info.rows; ++f) {
678 679 680 681 682 683 684 685 686
          bool isNull = colDataIsNull_s(pInput, f);
          if (isNull) {
            colDataAppendNULL(pOutput, pResult->info.rows + f);
          } else {
            char* data = colDataGetData(pInput, f);
            colDataAppend(pOutput, pResult->info.rows + f, data, isNull);
          }
        }

H
Haojun Liao 已提交
687 688 689
      } else {
        SArray* pBlockList = taosArrayInit(4, POINTER_BYTES);
        taosArrayPush(pBlockList, &pSrcBlock);
G
Ganlin Zhao 已提交
690

691
        SColumnInfoData* pResColData = taosArrayGet(pResult->pDataBlock, outputSlotId);
692
        SColumnInfoData  idata = {.info = pResColData->info, .hasNull = true};
H
Haojun Liao 已提交
693

694
        SScalarParam dest = {.columnData = &idata};
X
Xiaoyu Wang 已提交
695
        int32_t      code = scalarCalculate((SNode*)pExpr[k].pExpr->_function.pFunctNode, pBlockList, &dest);
696 697 698 699
        if (code != TSDB_CODE_SUCCESS) {
          taosArrayDestroy(pBlockList);
          return code;
        }
700

dengyihao's avatar
dengyihao 已提交
701
        int32_t startOffset = createNewColModel ? 0 : pResult->info.rows;
702
        ASSERT(pResult->info.capacity > 0);
703
        colDataMergeCol(pResColData, startOffset, (int32_t*)&pResult->info.capacity, &idata, dest.numOfRows);
D
dapan1121 已提交
704
        colDataDestroy(&idata);
705 706

        numOfRows = dest.numOfRows;
H
Haojun Liao 已提交
707 708
        taosArrayDestroy(pBlockList);
      }
709
    } else {
710
      return TSDB_CODE_OPS_NOT_SUPPORT;
711 712
    }
  }
713

714 715 716
  if (!createNewColModel) {
    pResult->info.rows += numOfRows;
  }
717 718

  return TSDB_CODE_SUCCESS;
719 720
}

5
54liuyao 已提交
721
bool functionNeedToExecute(SqlFunctionCtx* pCtx) {
722
  struct SResultRowEntryInfo* pResInfo = GET_RES_INFO(pCtx);
723

724 725 726 727 728
  // in case of timestamp column, always generated results.
  int32_t functionId = pCtx->functionId;
  if (functionId == -1) {
    return false;
  }
729

730 731
  if (pCtx->scanFlag == REPEAT_SCAN) {
    return fmIsRepeatScanFunc(pCtx->functionId);
732 733
  }

734 735
  if (isRowEntryCompleted(pResInfo)) {
    return false;
736 737
  }

738 739 740
  return true;
}

741 742 743 744 745 746 747
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;
    }
748

749 750 751
    // Set the correct column info (data type and bytes)
    pInput->pData[paramIndex]->info.type = type;
    pInput->pData[paramIndex]->info.bytes = tDataTypes[type].bytes;
752
  }
H
Haojun Liao 已提交
753

754 755 756 757 758 759
  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;
760 761
    }
  } else {
762
    da = pInput->pColumnDataAgg[paramIndex];
763 764
  }

765
  ASSERT(!IS_VAR_DATA_TYPE(type));
766

767 768
  if (type == TSDB_DATA_TYPE_BIGINT) {
    int64_t v = pFuncParam->param.i;
769
    *da = (SColumnDataAgg){.numOfNull = 0, .min = v, .max = v, .sum = v * numOfRows};
770 771
  } else if (type == TSDB_DATA_TYPE_DOUBLE) {
    double v = pFuncParam->param.d;
772
    *da = (SColumnDataAgg){.numOfNull = 0};
773

774 775 776 777 778 779
    *(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;

780
    *da = (SColumnDataAgg){.numOfNull = 0};
781 782 783 784 785
    *(bool*)&da->min = 0;
    *(bool*)&da->max = v;
    *(bool*)&da->sum = v * numOfRows;
  } else if (type == TSDB_DATA_TYPE_TIMESTAMP) {
    // do nothing
786
  } else {
787
    ASSERT(0);
788 789
  }

790 791
  return TSDB_CODE_SUCCESS;
}
792

793
void setBlockSMAInfo(SqlFunctionCtx* pCtx, SExprInfo* pExprInfo, SSDataBlock* pBlock) {
794 795 796 797 798 799 800 801 802
  int32_t numOfRows = pBlock->info.rows;

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

  if (pBlock->pBlockAgg != NULL) {
    pInput->colDataAggIsSet = true;

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

806 807
      if (pFuncParam->type == FUNC_PARAM_TYPE_COLUMN) {
        int32_t slotId = pFuncParam->pCol->slotId;
808 809 810 811
        pInput->pColumnDataAgg[j] = pBlock->pBlockAgg[slotId];
        if (pInput->pColumnDataAgg[j] == NULL) {
          pInput->colDataAggIsSet = false;
        }
812 813 814 815

        // 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);
816 817
      } else if (pFuncParam->type == FUNC_PARAM_TYPE_VALUE) {
        doCreateConstantValColumnAggInfo(pInput, pFuncParam, pFuncParam->param.nType, j, pBlock->info.rows);
818 819
      }
    }
820
  } else {
821
    pInput->colDataAggIsSet = false;
822 823 824
  }
}

L
Liu Jicong 已提交
825
bool isTaskKilled(SExecTaskInfo* pTaskInfo) {
826 827
  // query has been executed more than tsShellActivityTimer, and the retrieve has not arrived
  // abort current query execution.
L
Liu Jicong 已提交
828 829
  if (pTaskInfo->owner != 0 &&
      ((taosGetTimestampSec() - pTaskInfo->cost.start / 1000) > 10 * getMaximumIdleDurationSec())
830
      /*(!needBuildResAfterQueryComplete(pTaskInfo))*/) {
831
    assert(pTaskInfo->cost.start != 0);
L
Liu Jicong 已提交
832 833 834
    //    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;
835 836 837 838 839
  }

  return false;
}

L
Liu Jicong 已提交
840
void setTaskKilled(SExecTaskInfo* pTaskInfo) { pTaskInfo->code = TSDB_CODE_TSC_QUERY_CANCELLED; }
841 842

/////////////////////////////////////////////////////////////////////////////////////////////
843
STimeWindow getAlignQueryTimeWindow(SInterval* pInterval, int32_t precision, int64_t key) {
L
Liu Jicong 已提交
844
  STimeWindow win = {0};
845
  win.skey = taosTimeTruncate(key, pInterval, precision);
846 847

  /*
H
Haojun Liao 已提交
848
   * if the realSkey > INT64_MAX - pInterval->interval, the query duration between
849 850
   * realSkey and realEkey must be less than one interval.Therefore, no need to adjust the query ranges.
   */
851 852 853
  win.ekey = taosTimeAdd(win.skey, pInterval->interval, pInterval->intervalUnit, precision) - 1;
  if (win.ekey < win.skey) {
    win.ekey = INT64_MAX;
854
  }
855 856

  return win;
857 858
}

L
Liu Jicong 已提交
859 860
int32_t loadDataBlockOnDemand(SExecTaskInfo* pTaskInfo, STableScanInfo* pTableScanInfo, SSDataBlock* pBlock,
                              uint32_t* status) {
861
  *status = BLK_DATA_NOT_LOAD;
862

H
Haojun Liao 已提交
863
  pBlock->pDataBlock = NULL;
L
Liu Jicong 已提交
864
  pBlock->pBlockAgg = NULL;
H
Haojun Liao 已提交
865

L
Liu Jicong 已提交
866 867
  //  int64_t groupId = pRuntimeEnv->current->groupIndex;
  //  bool    ascQuery = QUERY_IS_ASC_QUERY(pQueryAttr);
868

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

871 872
//  pCost->totalBlocks += 1;
//  pCost->totalRows += pBlock->info.rows;
H
Haojun Liao 已提交
873
#if 0
874 875 876
  // 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 已提交
877
      (QUERY_IS_INTERVAL_QUERY(pQueryAttr) && overlapWithTimeWindow(pTaskInfo, &pBlock->info))) {
878
    (*status) = BLK_DATA_DATA_LOAD;
879 880 881
  }

  // check if this data block is required to load
882
  if ((*status) != BLK_DATA_DATA_LOAD) {
883 884 885 886 887 888 889
    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 已提交
890
      bool  masterScan = IS_MAIN_SCAN(pRuntimeEnv);
891 892 893 894 895 896
      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,
897
                                    pTableScanInfo->rowEntryInfoOffset);
898 899 900
      } else {
        if (setResultOutputBufByKey(pRuntimeEnv, pTableScanInfo->pResultRowInfo, pBlock->info.uid, &win, masterScan, &pResult, groupId,
                                    pTableScanInfo->pCtx, pTableScanInfo->numOfOutput,
901
                                    pTableScanInfo->rowEntryInfoOffset) != TSDB_CODE_SUCCESS) {
902
          T_LONG_JMP(pRuntimeEnv->env, TSDB_CODE_QRY_OUT_OF_MEMORY);
903 904 905 906
        }
      }
    } else if (pQueryAttr->stableQuery && (!pQueryAttr->tsCompQuery) && (!pQueryAttr->diffQuery)) { // stable aggregate, not interval aggregate or normal column aggregate
      doSetTableGroupOutputBuf(pRuntimeEnv, pTableScanInfo->pResultRowInfo, pTableScanInfo->pCtx,
907
                               pTableScanInfo->rowEntryInfoOffset, pTableScanInfo->numOfOutput,
908 909 910 911 912 913
                               pRuntimeEnv->current->groupIndex);
    }

    if (needFilter) {
      (*status) = doFilterByBlockTimeWindow(pTableScanInfo, pBlock);
    } else {
914
      (*status) = BLK_DATA_DATA_LOAD;
915 916 917 918
    }
  }

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

921
  if ((*status) == BLK_DATA_NOT_LOAD || (*status) == BLK_DATA_FILTEROUT) {
922 923
    //qDebug("QInfo:0x%"PRIx64" data block discard, brange:%" PRId64 "-%" PRId64 ", rows:%d", pQInfo->qId, pBlockInfo->window.skey,
//           pBlockInfo->window.ekey, pBlockInfo->rows);
924
    pCost->skipBlocks += 1;
925
  } else if ((*status) == BLK_DATA_SMA_LOAD) {
926 927
    // this function never returns error?
    pCost->loadBlockStatis += 1;
928
//    tsdbRetrieveDatablockSMA(pTableScanInfo->pTsdbReadHandle, &pBlock->pBlockAgg);
929 930

    if (pBlock->pBlockAgg == NULL) {  // data block statistics does not exist, load data block
931
//      pBlock->pDataBlock = tsdbRetrieveDataBlock(pTableScanInfo->pTsdbReadHandle, NULL);
932 933 934
      pCost->totalCheckedRows += pBlock->info.rows;
    }
  } else {
935
    assert((*status) == BLK_DATA_DATA_LOAD);
936 937 938

    // load the data block statistics to perform further filter
    pCost->loadBlockStatis += 1;
939
//    tsdbRetrieveDatablockSMA(pTableScanInfo->pTsdbReadHandle, &pBlock->pBlockAgg);
940 941 942 943 944 945

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

H
Haojun Liao 已提交
946
          bool  masterScan = IS_MAIN_SCAN(pRuntimeEnv);
947 948 949 950 951
          TSKEY k = ascQuery? pBlock->info.window.skey : pBlock->info.window.ekey;

          STimeWindow win = getActiveTimeWindow(pTableScanInfo->pResultRowInfo, k, pQueryAttr);
          if (setResultOutputBufByKey(pRuntimeEnv, pTableScanInfo->pResultRowInfo, pBlock->info.uid, &win, masterScan, &pResult, groupId,
                                      pTableScanInfo->pCtx, pTableScanInfo->numOfOutput,
952
                                      pTableScanInfo->rowEntryInfoOffset) != TSDB_CODE_SUCCESS) {
953
            T_LONG_JMP(pRuntimeEnv->env, TSDB_CODE_QRY_OUT_OF_MEMORY);
954 955 956 957 958 959 960 961 962 963
          }
        }
      }
      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
964
            pCost->skipBlocks += 1;
965 966
            //qDebug("QInfo:0x%"PRIx64" data block discard, brange:%" PRId64 "-%" PRId64 ", rows:%d", pQInfo->qId,
//                   pBlockInfo->window.skey, pBlockInfo->window.ekey, pBlockInfo->rows);
967
            (*status) = BLK_DATA_FILTEROUT;
968 969 970 971 972 973 974
            return TSDB_CODE_SUCCESS;
          }
        }
      }
    }

    // current block has been discard due to filter applied
H
Haojun Liao 已提交
975
//    if (!doFilterByBlockSMA(pRuntimeEnv, pBlock->pBlockAgg, pTableScanInfo->pCtx, pBlockInfo->rows)) {
976
//      pCost->skipBlocks += 1;
977 978
//      qDebug("QInfo:0x%"PRIx64" data block discard, brange:%" PRId64 "-%" PRId64 ", rows:%d", pQInfo->qId, pBlockInfo->window.skey,
//             pBlockInfo->window.ekey, pBlockInfo->rows);
979
//      (*status) = BLK_DATA_FILTEROUT;
980 981 982 983 984
//      return TSDB_CODE_SUCCESS;
//    }

    pCost->totalCheckedRows += pBlockInfo->rows;
    pCost->loadBlocks += 1;
985
//    pBlock->pDataBlock = tsdbRetrieveDataBlock(pTableScanInfo->pTsdbReadHandle, NULL);
986 987 988 989 990
//    if (pBlock->pDataBlock == NULL) {
//      return terrno;
//    }

//    if (pQueryAttr->pFilters != NULL) {
991
//      filterSetColFieldData(pQueryAttr->pFilters, taosArrayGetSize(pBlock->pDataBlock), pBlock->pDataBlock);
992
//    }
993

994 995 996 997
//    if (pQueryAttr->pFilters != NULL || pRuntimeEnv->pTsBuf != NULL) {
//      filterColRowsInDataBlock(pRuntimeEnv, pBlock, ascQuery);
//    }
  }
H
Haojun Liao 已提交
998
#endif
999 1000 1001
  return TSDB_CODE_SUCCESS;
}

L
Liu Jicong 已提交
1002
static void updateTableQueryInfoForReverseScan(STableQueryInfo* pTableQueryInfo) {
1003 1004 1005 1006 1007
  if (pTableQueryInfo == NULL) {
    return;
  }
}

L
Liu Jicong 已提交
1008
void setTaskStatus(SExecTaskInfo* pTaskInfo, int8_t status) {
1009
  if (status == TASK_NOT_COMPLETED) {
H
Haojun Liao 已提交
1010
    pTaskInfo->status = status;
1011 1012
  } else {
    // QUERY_NOT_COMPLETED is not compatible with any other status, so clear its position first
1013
    CLEAR_QUERY_STATUS(pTaskInfo, TASK_NOT_COMPLETED);
H
Haojun Liao 已提交
1014
    pTaskInfo->status |= status;
1015 1016 1017
  }
}

1018
void setResultRowInitCtx(SResultRow* pResult, SqlFunctionCtx* pCtx, int32_t numOfOutput, int32_t* rowEntryInfoOffset) {
5
54liuyao 已提交
1019
  bool init = false;
1020
  for (int32_t i = 0; i < numOfOutput; ++i) {
1021
    pCtx[i].resultInfo = getResultEntryInfo(pResult, i, rowEntryInfoOffset);
5
54liuyao 已提交
1022 1023 1024
    if (init) {
      continue;
    }
1025 1026 1027 1028 1029

    struct SResultRowEntryInfo* pResInfo = pCtx[i].resultInfo;
    if (isRowEntryCompleted(pResInfo) && isRowEntryInitialized(pResInfo)) {
      continue;
    }
1030 1031 1032 1033 1034

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

1035 1036 1037 1038 1039 1040
    if (!pResInfo->initialized) {
      if (pCtx[i].functionId != -1) {
        pCtx[i].fpSet.init(&pCtx[i], pResInfo);
      } else {
        pResInfo->initialized = true;
      }
5
54liuyao 已提交
1041 1042
    } else {
      init = true;
1043 1044 1045 1046
    }
  }
}

1047 1048
static void extractQualifiedTupleByFilterResult(SSDataBlock* pBlock, const SColumnInfoData* p, bool keep,
                                                int32_t status);
1049

H
Haojun Liao 已提交
1050 1051
void doFilter(SSDataBlock* pBlock, SFilterInfo* pFilterInfo, SColMatchInfo* pColMatchInfo) {
  if (pFilterInfo == NULL || pBlock->info.rows == 0) {
S
shenglian zhou 已提交
1052 1053
    return;
  }
1054

1055
  SFilterColumnParam param1 = {.numOfCols = taosArrayGetSize(pBlock->pDataBlock), .pDataBlock = pBlock->pDataBlock};
H
Haojun Liao 已提交
1056
  int32_t code = filterSetDataFromSlotId(pFilterInfo, &param1);
1057

1058
  SColumnInfoData* p = NULL;
1059
  int32_t          status = 0;
H
Haojun Liao 已提交
1060

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

1065
  if (pColMatchInfo != NULL) {
H
Haojun Liao 已提交
1066 1067
    size_t  size = taosArrayGetSize(pColMatchInfo->pList);
    for (int32_t i = 0; i < size; ++i) {
H
Haojun Liao 已提交
1068
      SColMatchItem* pInfo = taosArrayGet(pColMatchInfo->pList, i);
1069
      if (pInfo->colId == PRIMARYKEY_TIMESTAMP_COL_ID) {
H
Haojun Liao 已提交
1070
        SColumnInfoData* pColData = taosArrayGet(pBlock->pDataBlock, pInfo->dstSlotId);
1071
        if (pColData->info.type == TSDB_DATA_TYPE_TIMESTAMP) {
H
Haojun Liao 已提交
1072
          blockDataUpdateTsWindow(pBlock, pInfo->dstSlotId);
1073 1074 1075 1076 1077 1078
          break;
        }
      }
    }
  }

1079 1080
  colDataDestroy(p);
  taosMemoryFree(p);
1081 1082
}

1083
void extractQualifiedTupleByFilterResult(SSDataBlock* pBlock, const SColumnInfoData* p, bool keep, int32_t status) {
1084 1085 1086 1087
  if (keep) {
    return;
  }

H
Haojun Liao 已提交
1088 1089 1090
  int32_t totalRows = pBlock->info.rows;

  if (status == FILTER_RESULT_ALL_QUALIFIED) {
1091
    // here nothing needs to be done
H
Haojun Liao 已提交
1092
  } else if (status == FILTER_RESULT_NONE_QUALIFIED) {
1093
    pBlock->info.rows = 0;
H
Haojun Liao 已提交
1094
  } else {
1095
    SSDataBlock* px = createOneDataBlock(pBlock, true);
1096

1097 1098
    size_t numOfCols = taosArrayGetSize(pBlock->pDataBlock);
    for (int32_t i = 0; i < numOfCols; ++i) {
1099 1100
      SColumnInfoData* pSrc = taosArrayGet(px->pDataBlock, i);
      SColumnInfoData* pDst = taosArrayGet(pBlock->pDataBlock, i);
1101
      // it is a reserved column for scalar function, and no data in this column yet.
1102
      if (pDst->pData == NULL || pSrc->pData == NULL) {
1103 1104 1105
        continue;
      }

1106 1107
      colInfoDataCleanup(pDst, pBlock->info.rows);

1108
      int32_t numOfRows = 0;
1109
      for (int32_t j = 0; j < totalRows; ++j) {
1110
        if (((int8_t*)p->pData)[j] == 0) {
D
dapan1121 已提交
1111 1112
          continue;
        }
1113

D
dapan1121 已提交
1114
        if (colDataIsNull_s(pSrc, j)) {
1115
          colDataAppendNULL(pDst, numOfRows);
D
dapan1121 已提交
1116
        } else {
1117
          colDataAppend(pDst, numOfRows, colDataGetData(pSrc, j), false);
D
dapan1121 已提交
1118
        }
1119
        numOfRows += 1;
H
Haojun Liao 已提交
1120
      }
1121

1122
      // todo this value can be assigned directly
1123 1124 1125 1126 1127
      if (pBlock->info.rows == totalRows) {
        pBlock->info.rows = numOfRows;
      } else {
        ASSERT(pBlock->info.rows == numOfRows);
      }
1128
    }
1129

dengyihao's avatar
dengyihao 已提交
1130
    blockDataDestroy(px);  // fix memory leak
1131 1132 1133
  }
}

1134
void doSetTableGroupOutputBuf(SOperatorInfo* pOperator, int32_t numOfOutput, uint64_t groupId) {
1135
  // for simple group by query without interval, all the tables belong to one group result.
1136 1137 1138
  SExecTaskInfo*    pTaskInfo = pOperator->pTaskInfo;
  SAggOperatorInfo* pAggInfo = pOperator->info;

1139
  SResultRowInfo* pResultRowInfo = &pAggInfo->binfo.resultRowInfo;
1140 1141
  SqlFunctionCtx* pCtx = pOperator->exprSupp.pCtx;
  int32_t*        rowEntryInfoOffset = pOperator->exprSupp.rowEntryInfoOffset;
1142

1143
  SResultRow* pResultRow = doSetResultOutBufByKey(pAggInfo->aggSup.pResultBuf, pResultRowInfo, (char*)&groupId,
L
Liu Jicong 已提交
1144
                                                  sizeof(groupId), true, groupId, pTaskInfo, false, &pAggInfo->aggSup);
L
Liu Jicong 已提交
1145
  assert(pResultRow != NULL);
1146 1147 1148 1149 1150 1151

  /*
   * 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 已提交
1152 1153
    int32_t ret =
        addNewWindowResultBuf(pResultRow, pAggInfo->aggSup.pResultBuf, groupId, pAggInfo->binfo.pRes->info.rowSize);
1154 1155 1156 1157 1158
    if (ret != TSDB_CODE_SUCCESS) {
      return;
    }
  }

1159
  setResultRowInitCtx(pResultRow, pCtx, numOfOutput, rowEntryInfoOffset);
1160 1161
}

1162 1163 1164
static void setExecutionContext(SOperatorInfo* pOperator, int32_t numOfOutput, uint64_t groupId) {
  SAggOperatorInfo* pAggInfo = pOperator->info;
  if (pAggInfo->groupId != UINT64_MAX && pAggInfo->groupId == groupId) {
1165 1166
    return;
  }
1167 1168

  doSetTableGroupOutputBuf(pOperator, numOfOutput, groupId);
1169 1170

  // record the current active group id
H
Haojun Liao 已提交
1171
  pAggInfo->groupId = groupId;
1172 1173
}

dengyihao's avatar
dengyihao 已提交
1174 1175
static void doUpdateNumOfRows(SqlFunctionCtx* pCtx, SResultRow* pRow, int32_t numOfExprs,
                              const int32_t* rowCellOffset) {
1176
  bool returnNotNull = false;
1177
  for (int32_t j = 0; j < numOfExprs; ++j) {
1178
    struct SResultRowEntryInfo* pResInfo = getResultEntryInfo(pRow, j, rowCellOffset);
1179 1180 1181 1182 1183 1184 1185
    if (!isRowEntryInitialized(pResInfo)) {
      continue;
    }

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

1187
    if (fmIsNotNullOutputFunc(pCtx[j].functionId)) {
1188 1189
      returnNotNull = true;
    }
1190
  }
S
shenglian zhou 已提交
1191 1192
  // 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
1193
  if (pRow->numOfRows == 0 && !returnNotNull) {
1194
    pRow->numOfRows = 1;
1195 1196 1197
  }
}

1198 1199
static void doCopyResultToDataBlock(SExprInfo* pExprInfo, int32_t numOfExprs, SResultRow* pRow, SqlFunctionCtx* pCtx,
                                    SSDataBlock* pBlock, const int32_t* rowEntryOffset, SExecTaskInfo* pTaskInfo) {
1200 1201 1202
  for (int32_t j = 0; j < numOfExprs; ++j) {
    int32_t slotId = pExprInfo[j].base.resSchema.slotId;

1203
    pCtx[j].resultInfo = getResultEntryInfo(pRow, j, rowEntryOffset);
1204
    if (pCtx[j].fpSet.finalize) {
1205
      if (strcmp(pCtx[j].pExpr->pExpr->_function.functionName, "_group_key") == 0) {
1206 1207
        // for groupkey along with functions that output multiple lines(e.g. Histogram)
        // need to match groupkey result for each output row of that function.
1208 1209 1210 1211 1212
        if (pCtx[j].resultInfo->numOfRes != 0) {
          pCtx[j].resultInfo->numOfRes = pRow->numOfRows;
        }
      }

1213 1214 1215
      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));
1216
        T_LONG_JMP(pTaskInfo->env, code);
1217 1218
      }
    } else if (strcmp(pCtx[j].pExpr->pExpr->_function.functionName, "_select_value") == 0) {
1219
      // do nothing
1220
    } else {
1221 1222
      // 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.
1223 1224 1225 1226 1227 1228 1229
      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);
      }
    }
  }
1230 1231
}

1232 1233 1234
// todo refactor. SResultRow has direct pointer in miainfo
int32_t finalizeResultRows(SDiskbasedBuf* pBuf, SResultRowPosition* resultRowPosition, SExprSupp* pSup,
                           SSDataBlock* pBlock, SExecTaskInfo* pTaskInfo) {
1235 1236 1237 1238 1239 1240 1241 1242 1243 1244 1245 1246 1247 1248 1249 1250 1251 1252 1253 1254 1255 1256 1257 1258 1259 1260
  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);
1261 1262

  releaseBufPage(pBuf, page);
1263
  pBlock->info.rows += pRow->numOfRows;
1264 1265 1266
  return 0;
}

1267 1268 1269 1270 1271 1272 1273
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;

1274
  int32_t numOfRows = getNumOfTotalRes(pGroupResInfo);
1275

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

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

H
Haojun Liao 已提交
1282
    doUpdateNumOfRows(pCtx, pRow, numOfExprs, rowEntryOffset);
1283 1284

    // no results, continue to check the next one
1285 1286
    if (pRow->numOfRows == 0) {
      pGroupResInfo->index += 1;
1287
      releaseBufPage(pBuf, page);
1288 1289 1290
      continue;
    }

1291 1292 1293 1294 1295
    if (pBlock->info.groupId == 0) {
      pBlock->info.groupId = pPos->groupId;
    } else {
      // current value belongs to different group, it can't be packed into one datablock
      if (pBlock->info.groupId != pPos->groupId) {
1296
        releaseBufPage(pBuf, page);
1297 1298 1299 1300
        break;
      }
    }

1301
    if (pBlock->info.rows + pRow->numOfRows > pBlock->info.capacity) {
1302
      ASSERT(pBlock->info.rows > 0);
1303
      releaseBufPage(pBuf, page);
1304 1305 1306 1307
      break;
    }

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

1310
    releaseBufPage(pBuf, page);
1311
    pBlock->info.rows += pRow->numOfRows;
1312 1313
  }

X
Xiaoyu Wang 已提交
1314 1315
  qDebug("%s result generated, rows:%d, groupId:%" PRIu64, GET_TASKID(pTaskInfo), pBlock->info.rows,
         pBlock->info.groupId);
1316

1317
  blockDataUpdateTsWindow(pBlock, 0);
1318 1319 1320
  return 0;
}

1321 1322 1323 1324 1325 1326 1327 1328 1329 1330 1331 1332 1333 1334 1335 1336 1337 1338 1339 1340 1341 1342 1343 1344 1345 1346 1347 1348 1349 1350 1351 1352 1353 1354 1355 1356 1357 1358 1359 1360
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
  pBlock->info.groupId = 0;
  ASSERT(!pbInfo->mergeResultBlock);
  doCopyToSDataBlock(pTaskInfo, pBlock, &pOperator->exprSupp, pBuf, pGroupResInfo);
  if (pOperator->operatorType == QUERY_NODE_PHYSICAL_PLAN_STREAM_STATE) {
    SStreamStateAggOperatorInfo* pInfo = pOperator->info;

    char* tbname = taosHashGet(pInfo->pGroupIdTbNameMap, &pBlock->info.groupId, sizeof(int64_t));
    if (tbname != NULL) {
      memcpy(pBlock->info.parTbName, tbname, TSDB_TABLE_NAME_LEN);
    } else {
      pBlock->info.parTbName[0] = 0;
    }
  } else if (pOperator->operatorType == QUERY_NODE_PHYSICAL_PLAN_STREAM_SESSION ||
             pOperator->operatorType == QUERY_NODE_PHYSICAL_PLAN_STREAM_SEMI_SESSION ||
             pOperator->operatorType == QUERY_NODE_PHYSICAL_PLAN_STREAM_FINAL_SESSION) {
    SStreamSessionAggOperatorInfo* pInfo = pOperator->info;

    char* tbname = taosHashGet(pInfo->pGroupIdTbNameMap, &pBlock->info.groupId, sizeof(int64_t));
    if (tbname != NULL) {
      memcpy(pBlock->info.parTbName, tbname, TSDB_TABLE_NAME_LEN);
    } else {
      pBlock->info.parTbName[0] = 0;
    }
  }
}

X
Xiaoyu Wang 已提交
1361 1362
void doBuildResultDatablock(SOperatorInfo* pOperator, SOptrBasicInfo* pbInfo, SGroupResInfo* pGroupResInfo,
                            SDiskbasedBuf* pBuf) {
1363
  SExecTaskInfo* pTaskInfo = pOperator->pTaskInfo;
1364
  SSDataBlock*   pBlock = pbInfo->pRes;
1365

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

1369
  blockDataCleanup(pBlock);
1370
  if (!hasRemainResults(pGroupResInfo)) {
1371 1372 1373
    return;
  }

1374 1375
  // clear the existed group id
  pBlock->info.groupId = 0;
1376 1377 1378
  if (!pbInfo->mergeResultBlock) {
    doCopyToSDataBlock(pTaskInfo, pBlock, &pOperator->exprSupp, pBuf, pGroupResInfo);
  } else {
dengyihao's avatar
dengyihao 已提交
1379
    while (hasRemainResults(pGroupResInfo)) {
1380 1381 1382
      doCopyToSDataBlock(pTaskInfo, pBlock, &pOperator->exprSupp, pBuf, pGroupResInfo);
      if (pBlock->info.rows >= pOperator->resultInfo.threshold) {
        break;
1383 1384
      }

1385 1386
      // clearing group id to continue to merge data that belong to different groups
      pBlock->info.groupId = 0;
1387
    }
1388 1389 1390

    // clear the group id info in SSDataBlock, since the client does not need it
    pBlock->info.groupId = 0;
1391 1392 1393
  }
}

L
Liu Jicong 已提交
1394 1395
void queryCostStatis(SExecTaskInfo* pTaskInfo) {
  STaskCostInfo* pSummary = &pTaskInfo->cost;
1396

1397 1398
  SFileBlockLoadRecorder* pRecorder = pSummary->pRecoder;
  if (pSummary->pRecoder != NULL) {
1399
    qDebug(
H
Haojun Liao 已提交
1400 1401 1402 1403 1404
        "%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);
1405
  }
1406 1407
}

L
Liu Jicong 已提交
1408 1409 1410
// static void updateOffsetVal(STaskRuntimeEnv *pRuntimeEnv, SDataBlockInfo *pBlockInfo) {
//   STaskAttr *pQueryAttr = pRuntimeEnv->pQueryAttr;
//   STableQueryInfo* pTableQueryInfo = pRuntimeEnv->current;
1411
//
L
Liu Jicong 已提交
1412
//   int32_t step = GET_FORWARD_DIRECTION_FACTOR(pQueryAttr->order.order);
1413
//
L
Liu Jicong 已提交
1414 1415 1416 1417
//   if (pQueryAttr->limit.offset == pBlockInfo->rows) {  // current block will ignore completed
//     pTableQueryInfo->lastKey = QUERY_IS_ASC_QUERY(pQueryAttr) ? pBlockInfo->window.ekey + step :
//     pBlockInfo->window.skey + step; pQueryAttr->limit.offset = 0; return;
//   }
1418
//
L
Liu Jicong 已提交
1419 1420 1421 1422 1423
//   if (QUERY_IS_ASC_QUERY(pQueryAttr)) {
//     pQueryAttr->pos = (int32_t)pQueryAttr->limit.offset;
//   } else {
//     pQueryAttr->pos = pBlockInfo->rows - (int32_t)pQueryAttr->limit.offset - 1;
//   }
1424
//
L
Liu Jicong 已提交
1425
//   assert(pQueryAttr->pos >= 0 && pQueryAttr->pos <= pBlockInfo->rows - 1);
1426
//
L
Liu Jicong 已提交
1427 1428
//   SArray *         pDataBlock = tsdbRetrieveDataBlock(pRuntimeEnv->pTsdbReadHandle, NULL);
//   SColumnInfoData *pColInfoData = taosArrayGet(pDataBlock, 0);
1429
//
L
Liu Jicong 已提交
1430 1431
//   // update the pQueryAttr->limit.offset value, and pQueryAttr->pos value
//   TSKEY *keys = (TSKEY *) pColInfoData->pData;
1432
//
L
Liu Jicong 已提交
1433 1434 1435
//   // update the offset value
//   pTableQueryInfo->lastKey = keys[pQueryAttr->pos];
//   pQueryAttr->limit.offset = 0;
1436
//
L
Liu Jicong 已提交
1437
//   int32_t numOfRes = tableApplyFunctionsOnBlock(pRuntimeEnv, pBlockInfo, NULL, binarySearchForKey, pDataBlock);
1438
//
L
Liu Jicong 已提交
1439 1440 1441 1442
//   //qDebug("QInfo:0x%"PRIx64" check data block, brange:%" PRId64 "-%" PRId64 ", numBlocksOfStep:%d, numOfRes:%d,
//   lastKey:%"PRId64, GET_TASKID(pRuntimeEnv),
//          pBlockInfo->window.skey, pBlockInfo->window.ekey, pBlockInfo->rows, numOfRes, pQuery->current->lastKey);
// }
1443

L
Liu Jicong 已提交
1444 1445
// void skipBlocks(STaskRuntimeEnv *pRuntimeEnv) {
//   STaskAttr *pQueryAttr = pRuntimeEnv->pQueryAttr;
1446
//
L
Liu Jicong 已提交
1447 1448 1449
//   if (pQueryAttr->limit.offset <= 0 || pQueryAttr->numOfFilterCols > 0) {
//     return;
//   }
1450
//
L
Liu Jicong 已提交
1451 1452
//   pQueryAttr->pos = 0;
//   int32_t step = GET_FORWARD_DIRECTION_FACTOR(pQueryAttr->order.order);
1453
//
L
Liu Jicong 已提交
1454 1455
//   STableQueryInfo* pTableQueryInfo = pRuntimeEnv->current;
//   TsdbQueryHandleT pTsdbReadHandle = pRuntimeEnv->pTsdbReadHandle;
1456
//
L
Liu Jicong 已提交
1457 1458 1459
//   SDataBlockInfo blockInfo = SDATA_BLOCK_INITIALIZER;
//   while (tsdbNextDataBlock(pTsdbReadHandle)) {
//     if (isTaskKilled(pRuntimeEnv->qinfo)) {
1460
//       T_LONG_JMP(pRuntimeEnv->env, TSDB_CODE_TSC_QUERY_CANCELLED);
L
Liu Jicong 已提交
1461
//     }
1462
//
L
Liu Jicong 已提交
1463
//     tsdbRetrieveDataBlockInfo(pTsdbReadHandle, &blockInfo);
1464
//
L
Liu Jicong 已提交
1465 1466 1467 1468
//     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;
1469
//
L
Liu Jicong 已提交
1470 1471 1472 1473 1474 1475 1476
//       //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;
//     }
//   }
1477
//
L
Liu Jicong 已提交
1478
//   if (terrno != TSDB_CODE_SUCCESS) {
1479
//     T_LONG_JMP(pRuntimeEnv->env, terrno);
L
Liu Jicong 已提交
1480 1481 1482 1483 1484 1485 1486
//   }
// }

// static TSKEY doSkipIntervalProcess(STaskRuntimeEnv* pRuntimeEnv, STimeWindow* win, SDataBlockInfo* pBlockInfo,
// STableQueryInfo* pTableQueryInfo) {
//   STaskAttr *pQueryAttr = pRuntimeEnv->pQueryAttr;
//   SResultRowInfo *pWindowResInfo = &pRuntimeEnv->resultRowInfo;
1487
//
L
Liu Jicong 已提交
1488 1489 1490
//   assert(pQueryAttr->limit.offset == 0);
//   STimeWindow tw = *win;
//   getNextTimeWindow(pQueryAttr, &tw);
1491
//
L
Liu Jicong 已提交
1492 1493
//   if ((tw.skey <= pBlockInfo->window.ekey && QUERY_IS_ASC_QUERY(pQueryAttr)) ||
//       (tw.ekey >= pBlockInfo->window.skey && !QUERY_IS_ASC_QUERY(pQueryAttr))) {
1494
//
L
Liu Jicong 已提交
1495 1496 1497 1498
//     // 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);
1499
//
L
Liu Jicong 已提交
1500 1501 1502 1503
//     tw = *win;
//     int32_t startPos =
//         getNextQualifiedWindow(pQueryAttr, &tw, pBlockInfo, pColInfoData->pData, binarySearchForKey, -1);
//     assert(startPos >= 0);
1504
//
L
Liu Jicong 已提交
1505 1506
//     // set the abort info
//     pQueryAttr->pos = startPos;
1507
//
L
Liu Jicong 已提交
1508 1509 1510 1511
//     // reset the query start timestamp
//     pTableQueryInfo->win.skey = ((TSKEY *)pColInfoData->pData)[startPos];
//     pQueryAttr->window.skey = pTableQueryInfo->win.skey;
//     TSKEY key = pTableQueryInfo->win.skey;
1512
//
L
Liu Jicong 已提交
1513 1514
//     pWindowResInfo->prevSKey = tw.skey;
//     int32_t index = pRuntimeEnv->resultRowInfo.curIndex;
1515
//
L
Liu Jicong 已提交
1516 1517
//     int32_t numOfRes = tableApplyFunctionsOnBlock(pRuntimeEnv, pBlockInfo, NULL, binarySearchForKey, pDataBlock);
//     pRuntimeEnv->resultRowInfo.curIndex = index;  // restore the window index
1518
//
L
Liu Jicong 已提交
1519 1520 1521 1522
//     //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);
1523
//
L
Liu Jicong 已提交
1524 1525 1526 1527 1528
//     return key;
//   } else {  // do nothing
//     pQueryAttr->window.skey      = tw.skey;
//     pWindowResInfo->prevSKey = tw.skey;
//     pTableQueryInfo->lastKey = tw.skey;
1529
//
L
Liu Jicong 已提交
1530 1531
//     return tw.skey;
//   }
1532
//
L
Liu Jicong 已提交
1533 1534 1535 1536 1537 1538 1539 1540 1541 1542
//   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);
//   }
1543
//
L
Liu Jicong 已提交
1544 1545 1546 1547 1548
//   // 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;
//   }
1549
//
L
Liu Jicong 已提交
1550 1551 1552 1553 1554 1555 1556
//   /*
//    * 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);
1557
//
L
Liu Jicong 已提交
1558 1559
//   STimeWindow w = TSWINDOW_INITIALIZER;
//   bool ascQuery = QUERY_IS_ASC_QUERY(pQueryAttr);
1560
//
L
Liu Jicong 已提交
1561 1562
//   SResultRowInfo *pWindowResInfo = &pRuntimeEnv->resultRowInfo;
//   STableQueryInfo *pTableQueryInfo = pRuntimeEnv->current;
1563
//
L
Liu Jicong 已提交
1564 1565 1566
//   SDataBlockInfo blockInfo = SDATA_BLOCK_INITIALIZER;
//   while (tsdbNextDataBlock(pRuntimeEnv->pTsdbReadHandle)) {
//     tsdbRetrieveDataBlockInfo(pRuntimeEnv->pTsdbReadHandle, &blockInfo);
1567
//
L
Liu Jicong 已提交
1568 1569 1570 1571 1572 1573 1574 1575 1576
//     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;
//     }
1577
//
L
Liu Jicong 已提交
1578 1579
//     // the first time window
//     STimeWindow win = getActiveTimeWindow(pWindowResInfo, pWindowResInfo->prevSKey, pQueryAttr);
1580
//
L
Liu Jicong 已提交
1581 1582
//     while (pQueryAttr->limit.offset > 0) {
//       STimeWindow tw = win;
1583
//
L
Liu Jicong 已提交
1584 1585 1586
//       if ((win.ekey <= blockInfo.window.ekey && ascQuery) || (win.ekey >= blockInfo.window.skey && !ascQuery)) {
//         pQueryAttr->limit.offset -= 1;
//         pWindowResInfo->prevSKey = win.skey;
1587
//
L
Liu Jicong 已提交
1588 1589 1590 1591 1592 1593
//         // 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;
//         }
//       }
1594
//
L
Liu Jicong 已提交
1595 1596 1597 1598
//       if (pQueryAttr->limit.offset == 0) {
//         *start = doSkipIntervalProcess(pRuntimeEnv, &win, &blockInfo, pTableQueryInfo);
//         return true;
//       }
1599
//
L
Liu Jicong 已提交
1600 1601
//       // current window does not ended in current data block, try next data block
//       getNextTimeWindow(pQueryAttr, &tw);
1602
//
L
Liu Jicong 已提交
1603 1604 1605 1606 1607 1608 1609 1610 1611
//       /*
//        * 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)) {
1612
//
L
Liu Jicong 已提交
1613 1614
//         SArray *pDataBlock = tsdbRetrieveDataBlock(pRuntimeEnv->pTsdbReadHandle, NULL);
//         SColumnInfoData *pColInfoData = taosArrayGet(pDataBlock, 0);
1615
//
L
Liu Jicong 已提交
1616 1617 1618
//         if ((win.ekey > blockInfo.window.ekey && ascQuery) || (win.ekey < blockInfo.window.skey && !ascQuery)) {
//           pQueryAttr->limit.offset -= 1;
//         }
1619
//
L
Liu Jicong 已提交
1620 1621 1622 1623 1624 1625 1626 1627
//         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);
1628
//
L
Liu Jicong 已提交
1629 1630 1631 1632 1633 1634 1635 1636 1637 1638 1639
//           // 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.
//       }
//     }
//   }
1640
//
L
Liu Jicong 已提交
1641 1642
//   // check for error
//   if (terrno != TSDB_CODE_SUCCESS) {
1643
//     T_LONG_JMP(pRuntimeEnv->env, terrno);
L
Liu Jicong 已提交
1644
//   }
1645
//
L
Liu Jicong 已提交
1646 1647
//   return true;
// }
1648

1649
int32_t appendDownstream(SOperatorInfo* p, SOperatorInfo** pDownstream, int32_t num) {
H
Haojun Liao 已提交
1650
  if (p->pDownstream == NULL) {
H
Haojun Liao 已提交
1651
    assert(p->numOfDownstream == 0);
1652 1653
  }

wafwerar's avatar
wafwerar 已提交
1654
  p->pDownstream = taosMemoryCalloc(1, num * POINTER_BYTES);
1655 1656 1657 1658 1659 1660 1661
  if (p->pDownstream == NULL) {
    return TSDB_CODE_OUT_OF_MEMORY;
  }

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

dengyihao's avatar
dengyihao 已提交
1664 1665
static int32_t doInitAggInfoSup(SAggSupporter* pAggSup, SqlFunctionCtx* pCtx, int32_t numOfOutput, size_t keyBufSize,
                                const char* pKey);
1666

L
Liu Jicong 已提交
1667
static bool needToMerge(SSDataBlock* pBlock, SArray* groupInfo, char** buf, int32_t rowIndex) {
1668 1669 1670 1671
  size_t size = taosArrayGetSize(groupInfo);
  if (size == 0) {
    return true;
  }
1672

1673 1674
  for (int32_t i = 0; i < size; ++i) {
    int32_t* index = taosArrayGet(groupInfo, i);
1675

1676
    SColumnInfoData* pColInfo = taosArrayGet(pBlock->pDataBlock, *index);
L
Liu Jicong 已提交
1677
    bool             isNull = colDataIsNull(pColInfo, rowIndex, pBlock->info.rows, NULL);
1678

1679 1680 1681
    if ((isNull && buf[i] != NULL) || (!isNull && buf[i] == NULL)) {
      return false;
    }
1682

1683 1684 1685 1686 1687 1688 1689 1690 1691 1692 1693 1694 1695
    char* pCell = colDataGetData(pColInfo, rowIndex);
    if (IS_VAR_DATA_TYPE(pColInfo->info.type)) {
      if (varDataLen(pCell) != varDataLen(buf[i])) {
        return false;
      } else {
        if (memcmp(varDataVal(pCell), varDataVal(buf[i]), varDataLen(pCell)) != 0) {
          return false;
        }
      }
    } else {
      if (memcmp(pCell, buf[i], pColInfo->info.bytes) != 0) {
        return false;
      }
1696 1697 1698
    }
  }

1699
  return 0;
1700 1701
}

1702
static bool saveCurrentTuple(char** rowColData, SArray* pColumnList, SSDataBlock* pBlock, int32_t rowIndex) {
L
Liu Jicong 已提交
1703
  int32_t size = (int32_t)taosArrayGetSize(pColumnList);
1704

L
Liu Jicong 已提交
1705 1706
  for (int32_t i = 0; i < size; ++i) {
    int32_t*         index = taosArrayGet(pColumnList, i);
1707
    SColumnInfoData* pColInfo = taosArrayGet(pBlock->pDataBlock, *index);
H
Haojun Liao 已提交
1708

1709 1710 1711
    char* data = colDataGetData(pColInfo, rowIndex);
    memcpy(rowColData[i], data, colDataGetLength(pColInfo, rowIndex));
  }
1712

1713 1714
  return true;
}
1715

X
Xiaoyu Wang 已提交
1716
int32_t getTableScanInfo(SOperatorInfo* pOperator, int32_t* order, int32_t* scanFlag) {
1717
  // todo add more information about exchange operation
1718
  int32_t type = pOperator->operatorType;
X
Xiaoyu Wang 已提交
1719
  if (type == QUERY_NODE_PHYSICAL_PLAN_EXCHANGE || type == QUERY_NODE_PHYSICAL_PLAN_SYSTABLE_SCAN ||
1720
      type == QUERY_NODE_PHYSICAL_PLAN_STREAM_SCAN || type == QUERY_NODE_PHYSICAL_PLAN_TAG_SCAN ||
1721
      type == QUERY_NODE_PHYSICAL_PLAN_BLOCK_DIST_SCAN || type == QUERY_NODE_PHYSICAL_PLAN_LAST_ROW_SCAN) {
1722 1723 1724
    *order = TSDB_ORDER_ASC;
    *scanFlag = MAIN_SCAN;
    return TSDB_CODE_SUCCESS;
1725
  } else if (type == QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN) {
1726 1727 1728 1729
    STableScanInfo* pTableScanInfo = pOperator->info;
    *order = pTableScanInfo->cond.order;
    *scanFlag = pTableScanInfo->scanFlag;
    return TSDB_CODE_SUCCESS;
1730 1731 1732 1733 1734
  } else if (type == QUERY_NODE_PHYSICAL_PLAN_TABLE_MERGE_SCAN) {
    STableMergeScanInfo* pTableScanInfo = pOperator->info;
    *order = pTableScanInfo->cond.order;
    *scanFlag = pTableScanInfo->scanFlag;
    return TSDB_CODE_SUCCESS;
1735
  } else {
H
Haojun Liao 已提交
1736
    if (pOperator->pDownstream == NULL || pOperator->pDownstream[0] == NULL) {
1737
      return TSDB_CODE_INVALID_PARA;
H
Haojun Liao 已提交
1738
    } else {
1739
      return getTableScanInfo(pOperator->pDownstream[0], order, scanFlag);
1740 1741 1742
    }
  }
}
1743

1744
// this is a blocking operator
L
Liu Jicong 已提交
1745
static int32_t doOpenAggregateOptr(SOperatorInfo* pOperator) {
H
Haojun Liao 已提交
1746 1747
  if (OPTR_IS_OPENED(pOperator)) {
    return TSDB_CODE_SUCCESS;
1748 1749
  }

H
Haojun Liao 已提交
1750
  SExecTaskInfo*    pTaskInfo = pOperator->pTaskInfo;
1751
  SAggOperatorInfo* pAggInfo = pOperator->info;
H
Haojun Liao 已提交
1752

1753 1754
  SExprSupp*     pSup = &pOperator->exprSupp;
  SOperatorInfo* downstream = pOperator->pDownstream[0];
1755

1756 1757
  int64_t st = taosGetTimestampUs();

1758 1759 1760
  int32_t order = TSDB_ORDER_ASC;
  int32_t scanFlag = MAIN_SCAN;

H
Haojun Liao 已提交
1761
  while (1) {
1762
    SSDataBlock* pBlock = downstream->fpSet.getNextFn(downstream);
1763 1764 1765 1766
    if (pBlock == NULL) {
      break;
    }

1767 1768
    int32_t code = getTableScanInfo(pOperator, &order, &scanFlag);
    if (code != TSDB_CODE_SUCCESS) {
1769
      T_LONG_JMP(pTaskInfo->env, code);
1770
    }
1771

1772
    // there is an scalar expression that needs to be calculated before apply the group aggregation.
1773 1774 1775
    if (pAggInfo->scalarExprSup.pExprInfo != NULL) {
      SExprSupp* pSup1 = &pAggInfo->scalarExprSup;
      code = projectApplyFunctions(pSup1->pExprInfo, pBlock, pBlock, pSup1->pCtx, pSup1->numOfExprs, NULL);
1776
      if (code != TSDB_CODE_SUCCESS) {
1777
        T_LONG_JMP(pTaskInfo->env, code);
1778
      }
1779 1780
    }

1781
    // the pDataBlock are always the same one, no need to call this again
1782
    setExecutionContext(pOperator, pOperator->exprSupp.numOfExprs, pBlock->info.groupId);
1783
    setInputDataBlock(pSup, pBlock, order, scanFlag, true);
1784
    code = doAggregateImpl(pOperator, pSup->pCtx);
1785
    if (code != 0) {
1786
      T_LONG_JMP(pTaskInfo->env, code);
1787
    }
1788 1789
  }

1790
  initGroupedResultInfo(&pAggInfo->groupResInfo, pAggInfo->aggSup.pResultRowHashTable, 0);
H
Haojun Liao 已提交
1791
  OPTR_SET_OPENED(pOperator);
1792

1793
  pOperator->cost.openCost = (taosGetTimestampUs() - st) / 1000.0;
H
Haojun Liao 已提交
1794 1795 1796
  return TSDB_CODE_SUCCESS;
}

1797
static SSDataBlock* getAggregateResult(SOperatorInfo* pOperator) {
L
Liu Jicong 已提交
1798
  SAggOperatorInfo* pAggInfo = pOperator->info;
H
Haojun Liao 已提交
1799 1800 1801 1802 1803 1804
  SOptrBasicInfo*   pInfo = &pAggInfo->binfo;

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

L
Liu Jicong 已提交
1805
  SExecTaskInfo* pTaskInfo = pOperator->pTaskInfo;
1806
  pTaskInfo->code = pOperator->fpSet._openFn(pOperator);
H
Haojun Liao 已提交
1807
  if (pTaskInfo->code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
1808
    setOperatorCompleted(pOperator);
H
Haojun Liao 已提交
1809 1810 1811
    return NULL;
  }

H
Haojun Liao 已提交
1812
  blockDataEnsureCapacity(pInfo->pRes, pOperator->resultInfo.capacity);
S
slzhou 已提交
1813 1814
  while (1) {
    doBuildResultDatablock(pOperator, pInfo, &pAggInfo->groupResInfo, pAggInfo->aggSup.pResultBuf);
H
Haojun Liao 已提交
1815
    doFilter(pInfo->pRes, pOperator->exprSupp.pFilterInfo, NULL);
S
slzhou 已提交
1816

1817
    if (!hasRemainResults(&pAggInfo->groupResInfo)) {
H
Haojun Liao 已提交
1818
      setOperatorCompleted(pOperator);
S
slzhou 已提交
1819 1820
      break;
    }
1821

S
slzhou 已提交
1822 1823 1824 1825
    if (pInfo->pRes->info.rows > 0) {
      break;
    }
  }
1826

1827
  size_t rows = blockDataGetNumOfRows(pInfo->pRes);
1828 1829
  pOperator->resultInfo.totalRows += rows;

1830
  return (rows == 0) ? NULL : pInfo->pRes;
1831 1832
}

wmmhello's avatar
wmmhello 已提交
1833
int32_t aggEncodeResultRow(SOperatorInfo* pOperator, char** result, int32_t* length) {
1834
  if (result == NULL || length == NULL) {
wmmhello's avatar
wmmhello 已提交
1835 1836 1837
    return TSDB_CODE_TSC_INVALID_INPUT;
  }
  SOptrBasicInfo* pInfo = (SOptrBasicInfo*)(pOperator->info);
1838
  SAggSupporter*  pSup = (SAggSupporter*)POINTER_SHIFT(pOperator->info, sizeof(SOptrBasicInfo));
1839
  int32_t         size = tSimpleHashGetSize(pSup->pResultRowHashTable);
1840 1841 1842
  size_t          keyLen = sizeof(uint64_t) * 2;  // estimate the key length
  int32_t         totalSize =
      sizeof(int32_t) + sizeof(int32_t) + size * (sizeof(int32_t) + keyLen + sizeof(int32_t) + pSup->resultRowSize);
wmmhello's avatar
wmmhello 已提交
1843

C
Cary Xu 已提交
1844 1845 1846 1847 1848 1849
  // no result
  if (getTotalBufSize(pSup->pResultBuf) == 0) {
    *result = NULL;
    *length = 0;
    return TSDB_CODE_SUCCESS;
  }
1850

wmmhello's avatar
wmmhello 已提交
1851
  *result = (char*)taosMemoryCalloc(1, totalSize);
L
Liu Jicong 已提交
1852
  if (*result == NULL) {
wmmhello's avatar
wmmhello 已提交
1853
    return TSDB_CODE_OUT_OF_MEMORY;
wmmhello's avatar
wmmhello 已提交
1854
  }
wmmhello's avatar
wmmhello 已提交
1855

wmmhello's avatar
wmmhello 已提交
1856
  int32_t offset = sizeof(int32_t);
wmmhello's avatar
wmmhello 已提交
1857 1858
  *(int32_t*)(*result + offset) = size;
  offset += sizeof(int32_t);
1859 1860

  // prepare memory
1861
  SResultRowPosition* pos = &pInfo->resultRowInfo.cur;
dengyihao's avatar
dengyihao 已提交
1862 1863
  void*               pPage = getBufPage(pSup->pResultBuf, pos->pageId);
  SResultRow*         pRow = (SResultRow*)((char*)pPage + pos->offset);
1864 1865
  setBufPageDirty(pPage, true);
  releaseBufPage(pSup->pResultBuf, pPage);
L
Liu Jicong 已提交
1866

1867 1868 1869 1870
  int32_t iter = 0;
  void*   pIter = NULL;
  while ((pIter = tSimpleHashIterate(pSup->pResultRowHashTable, pIter, &iter))) {
    void*               key = tSimpleHashGetKey(pIter, &keyLen);
1871
    SResultRowPosition* p1 = (SResultRowPosition*)pIter;
1872

dengyihao's avatar
dengyihao 已提交
1873
    pPage = (SFilePage*)getBufPage(pSup->pResultBuf, p1->pageId);
1874
    pRow = (SResultRow*)((char*)pPage + p1->offset);
1875 1876
    setBufPageDirty(pPage, true);
    releaseBufPage(pSup->pResultBuf, pPage);
wmmhello's avatar
wmmhello 已提交
1877 1878 1879

    // recalculate the result size
    int32_t realTotalSize = offset + sizeof(int32_t) + keyLen + sizeof(int32_t) + pSup->resultRowSize;
L
Liu Jicong 已提交
1880
    if (realTotalSize > totalSize) {
wmmhello's avatar
wmmhello 已提交
1881
      char* tmp = (char*)taosMemoryRealloc(*result, realTotalSize);
L
Liu Jicong 已提交
1882
      if (tmp == NULL) {
wafwerar's avatar
wafwerar 已提交
1883
        taosMemoryFree(*result);
wmmhello's avatar
wmmhello 已提交
1884
        *result = NULL;
wmmhello's avatar
wmmhello 已提交
1885
        return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
1886
      } else {
wmmhello's avatar
wmmhello 已提交
1887 1888 1889 1890 1891 1892 1893 1894 1895 1896 1897 1898
        *result = tmp;
      }
    }
    // save key
    *(int32_t*)(*result + offset) = keyLen;
    offset += sizeof(int32_t);
    memcpy(*result + offset, key, keyLen);
    offset += keyLen;

    // save value
    *(int32_t*)(*result + offset) = pSup->resultRowSize;
    offset += sizeof(int32_t);
1899
    memcpy(*result + offset, pRow, pSup->resultRowSize);
wmmhello's avatar
wmmhello 已提交
1900 1901 1902
    offset += pSup->resultRowSize;
  }

wmmhello's avatar
wmmhello 已提交
1903 1904 1905 1906
  *(int32_t*)(*result) = offset;
  *length = offset;

  return TDB_CODE_SUCCESS;
wmmhello's avatar
wmmhello 已提交
1907 1908
}

1909 1910 1911 1912 1913
int32_t handleLimitOffset(SOperatorInfo* pOperator, SLimitInfo* pLimitInfo, SSDataBlock* pBlock, bool holdDataInBuf) {
  if (pLimitInfo->remainGroupOffset > 0) {
    if (pLimitInfo->currentGroupId == 0) {  // it is the first group
      pLimitInfo->currentGroupId = pBlock->info.groupId;
      blockDataCleanup(pBlock);
1914
      return PROJECT_RETRIEVE_CONTINUE;
1915 1916 1917
    } else if (pLimitInfo->currentGroupId != pBlock->info.groupId) {
      // now it is the data from a new group
      pLimitInfo->remainGroupOffset -= 1;
1918 1919

      // ignore data block in current group
1920 1921
      if (pLimitInfo->remainGroupOffset > 0) {
        blockDataCleanup(pBlock);
1922 1923 1924 1925 1926
        return PROJECT_RETRIEVE_CONTINUE;
      }
    }

    // set current group id of the project operator
1927
    pLimitInfo->currentGroupId = pBlock->info.groupId;
1928 1929
  }

1930
  // here check for a new group data, we need to handle the data of the previous group.
1931 1932 1933
  if (pLimitInfo->currentGroupId != 0 && pLimitInfo->currentGroupId != pBlock->info.groupId) {
    pLimitInfo->numOfOutputGroups += 1;
    if ((pLimitInfo->slimit.limit > 0) && (pLimitInfo->slimit.limit <= pLimitInfo->numOfOutputGroups)) {
1934
      pOperator->status = OP_EXEC_DONE;
1935
      blockDataCleanup(pBlock);
1936 1937 1938 1939 1940

      return PROJECT_RETRIEVE_DONE;
    }

    // reset the value for a new group data
1941 1942
    pLimitInfo->numOfOutputRows = 0;
    pLimitInfo->remainOffset = pLimitInfo->limit.offset;
1943 1944 1945 1946 1947

    // existing rows that belongs to previous group.
    if (pBlock->info.rows > 0) {
      return PROJECT_RETRIEVE_DONE;
    }
1948 1949 1950 1951 1952
  }

  // here we reach the start position, according to the limit/offset requirements.

  // set current group id
1953
  pLimitInfo->currentGroupId = pBlock->info.groupId;
1954

1955 1956 1957
  if (pLimitInfo->remainOffset >= pBlock->info.rows) {
    pLimitInfo->remainOffset -= pBlock->info.rows;
    blockDataCleanup(pBlock);
1958
    return PROJECT_RETRIEVE_CONTINUE;
1959 1960 1961
  } else if (pLimitInfo->remainOffset < pBlock->info.rows && pLimitInfo->remainOffset > 0) {
    blockDataTrimFirstNRows(pBlock, pLimitInfo->remainOffset);
    pLimitInfo->remainOffset = 0;
1962 1963
  }

1964
  // check for the limitation in each group
1965 1966 1967 1968
  if (pLimitInfo->limit.limit >= 0 && pLimitInfo->numOfOutputRows + pBlock->info.rows >= pLimitInfo->limit.limit) {
    int32_t keepRows = (int32_t)(pLimitInfo->limit.limit - pLimitInfo->numOfOutputRows);
    blockDataKeepFirstNRows(pBlock, keepRows);
    if (pLimitInfo->slimit.limit > 0 && pLimitInfo->slimit.limit <= pLimitInfo->numOfOutputGroups) {
1969 1970 1971
      pOperator->status = OP_EXEC_DONE;
    }

1972
    return PROJECT_RETRIEVE_DONE;
1973
  }
1974

1975
  // todo optimize performance
1976 1977
  // If there are slimit/soffset value exists, multi-round result can not be packed into one group, since the
  // they may not belong to the same group the limit/offset value is not valid in this case.
1978 1979
  if ((!holdDataInBuf) || (pBlock->info.rows >= pOperator->resultInfo.threshold) || pLimitInfo->slimit.offset != -1 ||
      pLimitInfo->slimit.limit != -1) {
1980
    return PROJECT_RETRIEVE_DONE;
L
Liu Jicong 已提交
1981
  } else {  // not full enough, continue to accumulate the output data in the buffer.
1982 1983 1984 1985
    return PROJECT_RETRIEVE_CONTINUE;
  }
}

1986
static void doApplyScalarCalculation(SOperatorInfo* pOperator, SSDataBlock* pBlock, int32_t order, int32_t scanFlag);
L
Liu Jicong 已提交
1987 1988
static void doHandleRemainBlockForNewGroupImpl(SOperatorInfo* pOperator, SFillOperatorInfo* pInfo,
                                               SResultInfo* pResultInfo, SExecTaskInfo* pTaskInfo) {
1989
  pInfo->totalInputRows = pInfo->existNewGroupBlock->info.rows;
1990 1991 1992 1993 1994
  SSDataBlock* pResBlock = pInfo->pFinalRes;

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

L
Liu Jicong 已提交
1996 1997
  int64_t ekey =
      Q_STATUS_EQUAL(pTaskInfo->status, TASK_COMPLETED) ? pInfo->win.ekey : pInfo->existNewGroupBlock->info.window.ekey;
1998 1999
  taosResetFillInfo(pInfo->pFillInfo, getFillInfoStart(pInfo->pFillInfo));

2000
  blockDataCleanup(pInfo->pRes);
2001 2002 2003 2004
  doApplyScalarCalculation(pOperator, pInfo->existNewGroupBlock, order, scanFlag);

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

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

2009
  pInfo->curGroupId = pInfo->existNewGroupBlock->info.groupId;
2010 2011 2012
  pInfo->existNewGroupBlock = NULL;
}

L
Liu Jicong 已提交
2013 2014
static void doHandleRemainBlockFromNewGroup(SOperatorInfo* pOperator, SFillOperatorInfo* pInfo,
                                            SResultInfo* pResultInfo, SExecTaskInfo* pTaskInfo) {
2015
  if (taosFillHasMoreResults(pInfo->pFillInfo)) {
H
Haojun Liao 已提交
2016 2017
    int32_t numOfResultRows = pResultInfo->capacity - pInfo->pFinalRes->info.rows;
    taosFillResultDataBlock(pInfo->pFillInfo, pInfo->pFinalRes, numOfResultRows);
2018 2019
    pInfo->pRes->info.groupId = pInfo->curGroupId;
    return;
2020 2021 2022 2023
  }

  // handle the cached new group data block
  if (pInfo->existNewGroupBlock) {
2024 2025 2026 2027 2028 2029
    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 已提交
2030
  SExprSupp*         pSup = &pOperator->exprSupp;
2031
  setInputDataBlock(pSup, pBlock, order, scanFlag, false);
2032 2033
  projectApplyFunctions(pSup->pExprInfo, pInfo->pRes, pBlock, pSup->pCtx, pSup->numOfExprs, NULL);

2034 2035 2036 2037
  // 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);
2038

2039 2040
  projectApplyFunctions(pNoFillSupp->pExprInfo, pInfo->pRes, pBlock, pNoFillSupp->pCtx, pNoFillSupp->numOfExprs, NULL);
  pInfo->pRes->info.groupId = pBlock->info.groupId;
2041 2042
}

S
slzhou 已提交
2043
static SSDataBlock* doFillImpl(SOperatorInfo* pOperator) {
L
Liu Jicong 已提交
2044 2045
  SFillOperatorInfo* pInfo = pOperator->info;
  SExecTaskInfo*     pTaskInfo = pOperator->pTaskInfo;
2046

H
Haojun Liao 已提交
2047
  SResultInfo* pResultInfo = &pOperator->resultInfo;
H
Haojun Liao 已提交
2048
  SSDataBlock* pResBlock = pInfo->pFinalRes;
2049 2050

  blockDataCleanup(pResBlock);
2051

H
Haojun Liao 已提交
2052 2053
  int32_t order = TSDB_ORDER_ASC;
  int32_t scanFlag = MAIN_SCAN;
2054
  getTableScanInfo(pOperator, &order, &scanFlag);
2055

2056
  doHandleRemainBlockFromNewGroup(pOperator, pInfo, pResultInfo, pTaskInfo);
2057
  if (pResBlock->info.rows > 0) {
2058
    pResBlock->info.groupId = pInfo->curGroupId;
2059
    return pResBlock;
H
Haojun Liao 已提交
2060
  }
2061

H
Haojun Liao 已提交
2062
  SOperatorInfo* pDownstream = pOperator->pDownstream[0];
L
Liu Jicong 已提交
2063
  while (1) {
2064
    SSDataBlock* pBlock = pDownstream->fpSet.getNextFn(pDownstream);
2065 2066
    if (pBlock == NULL) {
      if (pInfo->totalInputRows == 0) {
H
Haojun Liao 已提交
2067
        setOperatorCompleted(pOperator);
2068 2069
        return NULL;
      }
2070

2071
      taosFillSetStartInfo(pInfo->pFillInfo, 0, pInfo->win.ekey);
2072
    } else {
2073
      blockDataUpdateTsWindow(pBlock, pInfo->primarySrcSlotId);
2074 2075

      blockDataCleanup(pInfo->pRes);
2076 2077
      blockDataEnsureCapacity(pInfo->pRes, pBlock->info.rows);
      blockDataEnsureCapacity(pInfo->pFinalRes, pBlock->info.rows);
2078
      doApplyScalarCalculation(pOperator, pBlock, order, scanFlag);
2079

H
Haojun Liao 已提交
2080 2081 2082
      if (pInfo->curGroupId == 0 || pInfo->curGroupId == pInfo->pRes->info.groupId) {
        pInfo->curGroupId = pInfo->pRes->info.groupId;  // the first data block
        pInfo->totalInputRows += pInfo->pRes->info.rows;
2083

H
Haojun Liao 已提交
2084 2085 2086 2087 2088
        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 已提交
2089
        taosFillSetInputDataBlock(pInfo->pFillInfo, pInfo->pRes);
L
Liu Jicong 已提交
2090
      } else if (pInfo->curGroupId != pBlock->info.groupId) {  // the new group data block
2091 2092 2093 2094 2095
        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);
2096 2097 2098
      }
    }

2099 2100
    int32_t numOfResultRows = pOperator->resultInfo.capacity - pResBlock->info.rows;
    taosFillResultDataBlock(pInfo->pFillInfo, pResBlock, numOfResultRows);
2101 2102

    // current group has no more result to return
2103
    if (pResBlock->info.rows > 0) {
2104 2105
      // 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
2106
      if (pResBlock->info.rows > pResultInfo->threshold || pBlock == NULL || pInfo->existNewGroupBlock != NULL) {
2107
        pResBlock->info.groupId = pInfo->curGroupId;
2108
        return pResBlock;
2109 2110
      }

2111
      doHandleRemainBlockFromNewGroup(pOperator, pInfo, pResultInfo, pTaskInfo);
2112
      if (pResBlock->info.rows >= pOperator->resultInfo.threshold || pBlock == NULL) {
2113
        pResBlock->info.groupId = pInfo->curGroupId;
2114
        return pResBlock;
2115 2116 2117
      }
    } else if (pInfo->existNewGroupBlock) {  // try next group
      assert(pBlock != NULL);
2118 2119 2120 2121

      blockDataCleanup(pResBlock);

      doHandleRemainBlockForNewGroupImpl(pOperator, pInfo, pResultInfo, pTaskInfo);
2122
      if (pResBlock->info.rows > pResultInfo->threshold) {
2123
        pResBlock->info.groupId = pInfo->curGroupId;
2124
        return pResBlock;
2125 2126 2127 2128 2129 2130 2131
      }
    } else {
      return NULL;
    }
  }
}

S
slzhou 已提交
2132 2133 2134 2135 2136 2137 2138 2139
static SSDataBlock* doFill(SOperatorInfo* pOperator) {
  SFillOperatorInfo* pInfo = pOperator->info;
  SExecTaskInfo*     pTaskInfo = pOperator->pTaskInfo;

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

S
slzhou 已提交
2140
  SSDataBlock* fillResult = NULL;
S
slzhou 已提交
2141
  while (true) {
S
slzhou 已提交
2142
    fillResult = doFillImpl(pOperator);
S
slzhou 已提交
2143
    if (fillResult == NULL) {
H
Haojun Liao 已提交
2144
      setOperatorCompleted(pOperator);
S
slzhou 已提交
2145 2146 2147
      break;
    }

H
Haojun Liao 已提交
2148
    doFilter(fillResult, pOperator->exprSupp.pFilterInfo, &pInfo->matchInfo );
S
slzhou 已提交
2149 2150 2151 2152 2153
    if (fillResult->info.rows > 0) {
      break;
    }
  }

S
slzhou 已提交
2154
  if (fillResult != NULL) {
2155
    pOperator->resultInfo.totalRows += fillResult->info.rows;
S
slzhou 已提交
2156
  }
S
slzhou 已提交
2157

S
slzhou 已提交
2158
  return fillResult;
S
slzhou 已提交
2159 2160
}

2161
void destroyExprInfo(SExprInfo* pExpr, int32_t numOfExprs) {
C
Cary Xu 已提交
2162 2163 2164 2165 2166
  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);
H
Haojun Liao 已提交
2167
      }
2168
    }
C
Cary Xu 已提交
2169 2170 2171

    taosMemoryFree(pExprInfo->base.pParam);
    taosMemoryFree(pExprInfo->pExpr);
H
Haojun Liao 已提交
2172 2173 2174
  }
}

2175 2176 2177 2178 2179
static void destroyOperatorInfo(SOperatorInfo* pOperator) {
  if (pOperator == NULL) {
    return;
  }

2180
  if (pOperator->fpSet.closeFn != NULL) {
2181
    pOperator->fpSet.closeFn(pOperator->info);
2182 2183
  }

H
Haojun Liao 已提交
2184
  if (pOperator->pDownstream != NULL) {
L
Liu Jicong 已提交
2185
    for (int32_t i = 0; i < pOperator->numOfDownstream; ++i) {
H
Haojun Liao 已提交
2186
      destroyOperatorInfo(pOperator->pDownstream[i]);
2187 2188
    }

wafwerar's avatar
wafwerar 已提交
2189
    taosMemoryFreeClear(pOperator->pDownstream);
H
Haojun Liao 已提交
2190
    pOperator->numOfDownstream = 0;
2191 2192
  }

2193
  cleanupExprSupp(&pOperator->exprSupp);
wafwerar's avatar
wafwerar 已提交
2194
  taosMemoryFreeClear(pOperator);
2195 2196
}

2197 2198 2199 2200 2201 2202
int32_t getBufferPgSize(int32_t rowSize, uint32_t* defaultPgsz, uint32_t* defaultBufsz) {
  *defaultPgsz = 4096;
  while (*defaultPgsz < rowSize * 4) {
    *defaultPgsz <<= 1u;
  }

2203
  // The default buffer for each operator in query is 10MB.
2204
  // at least four pages need to be in buffer
2205 2206
  // TODO: make this variable to be configurable.
  *defaultBufsz = 4096 * 2560;
2207 2208 2209 2210 2211 2212 2213
  if ((*defaultBufsz) <= (*defaultPgsz)) {
    (*defaultBufsz) = (*defaultPgsz) * 4;
  }

  return 0;
}

dengyihao's avatar
dengyihao 已提交
2214 2215
int32_t doInitAggInfoSup(SAggSupporter* pAggSup, SqlFunctionCtx* pCtx, int32_t numOfOutput, size_t keyBufSize,
                         const char* pKey) {
L
Liu Jicong 已提交
2216
  int32_t    code = 0;
2217 2218
  _hash_fn_t hashFn = taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY);

2219
  pAggSup->currentPageId = -1;
dengyihao's avatar
dengyihao 已提交
2220 2221
  pAggSup->resultRowSize = getResultRowSize(pCtx, numOfOutput);
  pAggSup->keyBuf = taosMemoryCalloc(1, keyBufSize + POINTER_BYTES + sizeof(int64_t));
2222
  pAggSup->pResultRowHashTable = tSimpleHashInit(10, hashFn);
2223

H
Haojun Liao 已提交
2224
  if (pAggSup->keyBuf == NULL || pAggSup->pResultRowHashTable == NULL) {
2225 2226 2227
    return TSDB_CODE_OUT_OF_MEMORY;
  }

dengyihao's avatar
dengyihao 已提交
2228
  uint32_t defaultPgsz = 0;
2229 2230
  uint32_t defaultBufsz = 0;
  getBufferPgSize(pAggSup->resultRowSize, &defaultPgsz, &defaultBufsz);
H
Haojun Liao 已提交
2231

wafwerar's avatar
wafwerar 已提交
2232
  if (!osTempSpaceAvailable()) {
H
Haojun Liao 已提交
2233 2234 2235
    code = TSDB_CODE_NO_AVAIL_DISK;
    qError("Init stream agg supporter failed since %s, %s", terrstr(code), pKey);
    return code;
wafwerar's avatar
wafwerar 已提交
2236
  }
2237

H
Haojun Liao 已提交
2238
  code = createDiskbasedBuf(&pAggSup->pResultBuf, defaultPgsz, defaultBufsz, pKey, tsTempDir);
H
Haojun Liao 已提交
2239
  if (code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
2240
    qError("Create agg result buf failed since %s, %s", tstrerror(code), pKey);
H
Haojun Liao 已提交
2241 2242 2243
    return code;
  }

H
Haojun Liao 已提交
2244
  return code;
2245 2246
}

2247
void cleanupAggSup(SAggSupporter* pAggSup) {
wafwerar's avatar
wafwerar 已提交
2248
  taosMemoryFreeClear(pAggSup->keyBuf);
2249
  tSimpleHashCleanup(pAggSup->pResultRowHashTable);
H
Haojun Liao 已提交
2250
  destroyDiskbasedBuf(pAggSup->pResultBuf);
2251 2252
}

L
Liu Jicong 已提交
2253 2254
int32_t initAggInfo(SExprSupp* pSup, SAggSupporter* pAggSup, SExprInfo* pExprInfo, int32_t numOfCols, size_t keyBufSize,
                    const char* pkey) {
2255 2256 2257 2258 2259
  int32_t code = initExprSupp(pSup, pExprInfo, numOfCols);
  if (code != TSDB_CODE_SUCCESS) {
    return code;
  }

2260 2261 2262 2263 2264
  code = doInitAggInfoSup(pAggSup, pSup->pCtx, numOfCols, keyBufSize, pkey);
  if (code != TSDB_CODE_SUCCESS) {
    return code;
  }

L
Liu Jicong 已提交
2265
  for (int32_t i = 0; i < numOfCols; ++i) {
2266
    pSup->pCtx[i].saveHandle.pBuf = pAggSup->pResultBuf;
2267 2268
  }

2269
  return TSDB_CODE_SUCCESS;
H
Haojun Liao 已提交
2270 2271
}

L
Liu Jicong 已提交
2272
void initResultSizeInfo(SResultInfo* pResultInfo, int32_t numOfRows) {
wmmhello's avatar
wmmhello 已提交
2273
  ASSERT(numOfRows != 0);
2274 2275
  pResultInfo->capacity = numOfRows;
  pResultInfo->threshold = numOfRows * 0.75;
2276

2277 2278
  if (pResultInfo->threshold == 0) {
    pResultInfo->threshold = numOfRows;
2279 2280 2281
  }
}

2282 2283 2284 2285 2286
void initBasicInfo(SOptrBasicInfo* pInfo, SSDataBlock* pBlock) {
  pInfo->pRes = pBlock;
  initResultRowInfo(&pInfo->resultRowInfo);
}

5
54liuyao 已提交
2287
void* destroySqlFunctionCtx(SqlFunctionCtx* pCtx, int32_t numOfOutput) {
2288 2289 2290 2291 2292 2293 2294 2295 2296 2297
  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);
2298
    taosMemoryFreeClear(pCtx[i].subsidiaries.buf);
2299 2300 2301 2302 2303 2304 2305 2306
    taosMemoryFree(pCtx[i].input.pData);
    taosMemoryFree(pCtx[i].input.pColumnDataAgg);
  }

  taosMemoryFreeClear(pCtx);
  return NULL;
}

2307
int32_t initExprSupp(SExprSupp* pSup, SExprInfo* pExprInfo, int32_t numOfExpr) {
2308 2309 2310 2311
  pSup->pExprInfo = pExprInfo;
  pSup->numOfExprs = numOfExpr;
  if (pSup->pExprInfo != NULL) {
    pSup->pCtx = createSqlFunctionCtx(pExprInfo, numOfExpr, &pSup->rowEntryInfoOffset);
2312 2313 2314
    if (pSup->pCtx == NULL) {
      return TSDB_CODE_OUT_OF_MEMORY;
    }
2315
  }
2316 2317

  return TSDB_CODE_SUCCESS;
2318 2319
}

2320 2321 2322 2323
void cleanupExprSupp(SExprSupp* pSupp) {
  destroySqlFunctionCtx(pSupp->pCtx, pSupp->numOfExprs);
  if (pSupp->pExprInfo != NULL) {
    destroyExprInfo(pSupp->pExprInfo, pSupp->numOfExprs);
C
Cary Xu 已提交
2324
    taosMemoryFreeClear(pSupp->pExprInfo);
2325
  }
H
Haojun Liao 已提交
2326 2327 2328 2329 2330 2331

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

2332 2333 2334
  taosMemoryFree(pSupp->rowEntryInfoOffset);
}

2335 2336
SOperatorInfo* createAggregateOperatorInfo(SOperatorInfo* downstream, SAggPhysiNode* pAggNode,
                                           SExecTaskInfo* pTaskInfo) {
wafwerar's avatar
wafwerar 已提交
2337
  SAggOperatorInfo* pInfo = taosMemoryCalloc(1, sizeof(SAggOperatorInfo));
L
Liu Jicong 已提交
2338
  SOperatorInfo*    pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
H
Haojun Liao 已提交
2339 2340 2341
  if (pInfo == NULL || pOperator == NULL) {
    goto _error;
  }
H
Haojun Liao 已提交
2342

H
Haojun Liao 已提交
2343 2344 2345 2346
  SSDataBlock* pResBlock = createResDataBlock(pAggNode->node.pOutputDataBlockDesc);
  initBasicInfo(&pInfo->binfo, pResBlock);

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

2349 2350 2351
  int32_t    num = 0;
  SExprInfo* pExprInfo = createExprInfo(pAggNode->pAggFuncs, pAggNode->pGroupKeys, &num);
  int32_t    code = initAggInfo(&pOperator->exprSupp, &pInfo->aggSup, pExprInfo, num, keyBufSize, pTaskInfo->id.str);
L
Liu Jicong 已提交
2352
  if (code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
2353 2354
    goto _error;
  }
H
Haojun Liao 已提交
2355

H
Haojun Liao 已提交
2356 2357 2358 2359 2360 2361
  int32_t    numOfScalarExpr = 0;
  SExprInfo* pScalarExprInfo = NULL;
  if (pAggNode->pExprs != NULL) {
    pScalarExprInfo = createExprInfo(pAggNode->pExprs, NULL, &numOfScalarExpr);
  }

2362 2363 2364 2365
  code = initExprSupp(&pInfo->scalarExprSup, pScalarExprInfo, numOfScalarExpr);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
2366

H
Haojun Liao 已提交
2367 2368 2369 2370 2371
  code = filterInitFromNode((SNode*)pAggNode->node.pConditions, &pOperator->exprSupp.pFilterInfo, 0);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }

H
Haojun Liao 已提交
2372
  pInfo->binfo.mergeResultBlock = pAggNode->mergeDataBlock;
2373
  pInfo->groupId = UINT64_MAX;
H
Haojun Liao 已提交
2374

H
Haojun Liao 已提交
2375
  setOperatorInfo(pOperator, "TableAggregate", QUERY_NODE_PHYSICAL_PLAN_HASH_AGG, true, OP_NOT_OPENED, pInfo, pTaskInfo);
5
54liuyao 已提交
2376
  pOperator->fpSet =
H
Haojun Liao 已提交
2377
      createOperatorFpSet(doOpenAggregateOptr, getAggregateResult, NULL, destroyAggOperatorInfo, NULL);
H
Haojun Liao 已提交
2378

2379 2380 2381 2382 2383 2384
  if (downstream->operatorType == QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN) {
    STableScanInfo* pTableScanInfo = downstream->info;
    pTableScanInfo->pdInfo.pExprSup = &pOperator->exprSupp;
    pTableScanInfo->pdInfo.pAggSup = &pInfo->aggSup;
  }

H
Haojun Liao 已提交
2385 2386 2387 2388
  code = appendDownstream(pOperator, &downstream, 1);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
2389 2390

  return pOperator;
H
Haojun Liao 已提交
2391

2392
_error:
H
Haojun Liao 已提交
2393 2394 2395 2396
  if (pInfo != NULL) {
    destroyAggOperatorInfo(pInfo);
  }

2397 2398 2399
  if (pOperator != NULL) {
    cleanupExprSupp(&pOperator->exprSupp);
  }
H
Haojun Liao 已提交
2400

2401
  taosMemoryFreeClear(pOperator);
H
Haojun Liao 已提交
2402
  pTaskInfo->code = code;
H
Haojun Liao 已提交
2403
  return NULL;
2404 2405
}

2406
void cleanupBasicInfo(SOptrBasicInfo* pInfo) {
2407
  assert(pInfo != NULL);
H
Haojun Liao 已提交
2408
  pInfo->pRes = blockDataDestroy(pInfo->pRes);
2409 2410
}

H
Haojun Liao 已提交
2411 2412 2413 2414 2415 2416 2417
static void freeItem(void* pItem) {
  void** p = pItem;
  if (*p != NULL) {
    taosMemoryFreeClear(*p);
  }
}

2418
void destroyAggOperatorInfo(void* param) {
L
Liu Jicong 已提交
2419
  SAggOperatorInfo* pInfo = (SAggOperatorInfo*)param;
L
Liu Jicong 已提交
2420 2421
  cleanupBasicInfo(&pInfo->binfo);

H
Haojun Liao 已提交
2422
  cleanupAggSup(&pInfo->aggSup);
S
shenglian zhou 已提交
2423
  cleanupExprSupp(&pInfo->scalarExprSup);
H
Haojun Liao 已提交
2424
  cleanupGroupResInfo(&pInfo->groupResInfo);
D
dapan1121 已提交
2425
  taosMemoryFreeClear(param);
2426
}
2427

2428
void destroyFillOperatorInfo(void* param) {
L
Liu Jicong 已提交
2429
  SFillOperatorInfo* pInfo = (SFillOperatorInfo*)param;
2430
  pInfo->pFillInfo = taosDestroyFillInfo(pInfo->pFillInfo);
H
Haojun Liao 已提交
2431
  pInfo->pRes = blockDataDestroy(pInfo->pRes);
H
Haojun Liao 已提交
2432 2433
  pInfo->pFinalRes = blockDataDestroy(pInfo->pFinalRes);

2434
  cleanupExprSupp(&pInfo->noFillExprSupp);
H
Haojun Liao 已提交
2435

wafwerar's avatar
wafwerar 已提交
2436
  taosMemoryFreeClear(pInfo->p);
H
Haojun Liao 已提交
2437
  taosArrayDestroy(pInfo->matchInfo.pList);
D
dapan1121 已提交
2438
  taosMemoryFreeClear(param);
2439 2440
}

H
Haojun Liao 已提交
2441 2442 2443 2444
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 已提交
2445

H
Haojun Liao 已提交
2446 2447 2448
  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 已提交
2449

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

H
Haojun Liao 已提交
2453 2454 2455 2456 2457 2458 2459
  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 已提交
2460
  pInfo->p = taosMemoryCalloc(numOfCols, POINTER_BYTES);
2461

H
Haojun Liao 已提交
2462
  if (pInfo->pFillInfo == NULL || pInfo->p == NULL) {
H
Haojun Liao 已提交
2463 2464
    taosMemoryFree(pInfo->pFillInfo);
    taosMemoryFree(pInfo->p);
H
Haojun Liao 已提交
2465 2466 2467 2468 2469 2470
    return TSDB_CODE_OUT_OF_MEMORY;
  } else {
    return TSDB_CODE_SUCCESS;
  }
}

2471
static bool isWstartColumnExist(SFillOperatorInfo* pInfo) {
2472
  if (pInfo->noFillExprSupp.numOfExprs == 0) {
2473 2474
    return false;
  }
2475 2476 2477

  for (int32_t i = 0; i < pInfo->noFillExprSupp.numOfExprs; ++i) {
    SExprInfo* exprInfo = pInfo->noFillExprSupp.pExprInfo + i;
2478
    if (exprInfo->pExpr->nodeType == QUERY_NODE_COLUMN && exprInfo->base.numOfParams == 1 &&
2479 2480 2481 2482 2483 2484 2485
        exprInfo->base.pParam[0].pCol->colType == COLUMN_TYPE_WINDOW_START) {
      return true;
    }
  }
  return false;
}

2486 2487
static int32_t createPrimaryTsExprIfNeeded(SFillOperatorInfo* pInfo, SFillPhysiNode* pPhyFillNode, SExprSupp* pExprSupp,
                                           const char* idStr) {
2488
  bool wstartExist = isWstartColumnExist(pInfo);
2489

2490 2491
  if (wstartExist == false) {
    if (pPhyFillNode->pWStartTs->type != QUERY_NODE_TARGET) {
2492
      qError("pWStartTs of fill physical node is not a target node, %s", idStr);
2493 2494 2495
      return TSDB_CODE_QRY_SYS_ERROR;
    }

2496 2497
    SExprInfo* pExpr = taosMemoryRealloc(pExprSupp->pExprInfo, (pExprSupp->numOfExprs + 1) * sizeof(SExprInfo));
    if (pExpr == NULL) {
2498 2499 2500
      return TSDB_CODE_OUT_OF_MEMORY;
    }

2501 2502 2503
    createExprFromTargetNode(&pExpr[pExprSupp->numOfExprs], (STargetNode*)pPhyFillNode->pWStartTs);
    pExprSupp->numOfExprs += 1;
    pExprSupp->pExprInfo = pExpr;
2504
  }
2505

2506 2507 2508
  return TSDB_CODE_SUCCESS;
}

L
Liu Jicong 已提交
2509 2510
SOperatorInfo* createFillOperatorInfo(SOperatorInfo* downstream, SFillPhysiNode* pPhyFillNode,
                                      SExecTaskInfo* pTaskInfo) {
2511 2512 2513 2514 2515 2516
  SFillOperatorInfo* pInfo = taosMemoryCalloc(1, sizeof(SFillOperatorInfo));
  SOperatorInfo*     pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
  if (pInfo == NULL || pOperator == NULL) {
    goto _error;
  }

H
Haojun Liao 已提交
2517 2518 2519 2520
  pInfo->pRes = createResDataBlock(pPhyFillNode->node.pOutputDataBlockDesc);
  SExprInfo* pExprInfo = createExprInfo(pPhyFillNode->pFillExprs, NULL, &pInfo->numOfExpr);
  pOperator->exprSupp.pExprInfo = pExprInfo;

2521 2522 2523 2524 2525 2526 2527 2528
  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);
2529 2530 2531
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
H
Haojun Liao 已提交
2532

L
Liu Jicong 已提交
2533
  SInterval* pInterval =
2534
      QUERY_NODE_PHYSICAL_PLAN_MERGE_ALIGNED_INTERVAL == downstream->operatorType
2535 2536
          ? &((SMergeAlignedIntervalAggOperatorInfo*)downstream->info)->intervalAggOperatorInfo->interval
          : &((SIntervalAggOperatorInfo*)downstream->info)->interval;
2537

2538
  int32_t order = (pPhyFillNode->inputTsOrder == ORDER_ASC) ? TSDB_ORDER_ASC : TSDB_ORDER_DESC;
2539
  int32_t type = convertFillType(pPhyFillNode->mode);
2540

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

2543
  initResultSizeInfo(&pOperator->resultInfo, 4096);
H
Haojun Liao 已提交
2544 2545 2546 2547 2548
  blockDataEnsureCapacity(pInfo->pRes, pOperator->resultInfo.capacity);
  code = initExprSupp(&pOperator->exprSupp, pExprInfo, pInfo->numOfExpr);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
H
Haojun Liao 已提交
2549

H
Haojun Liao 已提交
2550 2551
  pInfo->primaryTsCol = ((STargetNode*)pPhyFillNode->pWStartTs)->slotId;
  pInfo->primarySrcSlotId = ((SColumnNode*)((STargetNode*)pPhyFillNode->pWStartTs)->pExpr)->slotId;
2552

2553
  int32_t numOfOutputCols = 0;
H
Haojun Liao 已提交
2554 2555
  code = extractColMatchInfo(pPhyFillNode->pFillExprs, pPhyFillNode->node.pOutputDataBlockDesc, &numOfOutputCols,
                             COL_MATCH_FROM_SLOT_ID, &pInfo->matchInfo);
2556

2557
  code = initFillInfo(pInfo, pExprInfo, pInfo->numOfExpr, pNoFillSupp->pExprInfo, pNoFillSupp->numOfExprs,
2558 2559
                      (SNodeListNode*)pPhyFillNode->pValues, pPhyFillNode->timeRange, pResultInfo->capacity,
                      pTaskInfo->id.str, pInterval, type, order);
2560 2561 2562
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
2563

H
Haojun Liao 已提交
2564
  pInfo->pFinalRes = createOneDataBlock(pInfo->pRes, false);
H
Haojun Liao 已提交
2565 2566
  blockDataEnsureCapacity(pInfo->pFinalRes, pOperator->resultInfo.capacity);

H
Haojun Liao 已提交
2567 2568 2569 2570 2571
  code = filterInitFromNode((SNode*)pPhyFillNode->node.pConditions, &pOperator->exprSupp.pFilterInfo, 0);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }

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

2576
  code = appendDownstream(pOperator, &downstream, 1);
2577
  return pOperator;
H
Haojun Liao 已提交
2578

2579
_error:
H
Haojun Liao 已提交
2580 2581 2582 2583
  if (pInfo != NULL) {
    destroyFillOperatorInfo(pInfo);
  }

2584
  pTaskInfo->code = code;
wafwerar's avatar
wafwerar 已提交
2585
  taosMemoryFreeClear(pOperator);
H
Haojun Liao 已提交
2586
  return NULL;
2587 2588
}

D
dapan1121 已提交
2589
static SExecTaskInfo* createExecTaskInfo(uint64_t queryId, uint64_t taskId, EOPTR_EXEC_MODEL model, char* dbFName) {
wafwerar's avatar
wafwerar 已提交
2590
  SExecTaskInfo* pTaskInfo = taosMemoryCalloc(1, sizeof(SExecTaskInfo));
H
Haojun Liao 已提交
2591 2592 2593 2594 2595
  if (pTaskInfo == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return NULL;
  }

2596
  setTaskStatus(pTaskInfo, TASK_NOT_COMPLETED);
H
Haojun Liao 已提交
2597

2598
  pTaskInfo->schemaInfo.dbname = strdup(dbFName);
2599
  pTaskInfo->cost.created = taosGetTimestampMs();
H
Haojun Liao 已提交
2600
  pTaskInfo->id.queryId = queryId;
dengyihao's avatar
dengyihao 已提交
2601
  pTaskInfo->execModel = model;
H
Haojun Liao 已提交
2602
  pTaskInfo->pTableInfoList = tableListCreate();
H
Haojun Liao 已提交
2603

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

2608 2609
  return pTaskInfo;
}
H
Haojun Liao 已提交
2610

H
Haojun Liao 已提交
2611 2612
SSchemaWrapper* extractQueriedColumnSchema(SScanPhysiNode* pScanNode);

2613
int32_t extractTableSchemaInfo(SReadHandle* pHandle, SScanPhysiNode* pScanNode, SExecTaskInfo* pTaskInfo) {
2614 2615
  SMetaReader mr = {0};
  metaReaderInit(&mr, pHandle->meta, 0);
2616
  int32_t code = metaGetTableEntryByUid(&mr, pScanNode->uid);
2617
  if (code != TSDB_CODE_SUCCESS) {
L
Liu Jicong 已提交
2618 2619
    qError("failed to get the table meta, uid:0x%" PRIx64 ", suid:0x%" PRIx64 ", %s", pScanNode->uid, pScanNode->suid,
           GET_TASKID(pTaskInfo));
H
Haojun Liao 已提交
2620

D
dapan1121 已提交
2621
    metaReaderClear(&mr);
2622
    return terrno;
D
dapan1121 已提交
2623
  }
2624

2625 2626
  SSchemaInfo* pSchemaInfo = &pTaskInfo->schemaInfo;
  pSchemaInfo->tablename = strdup(mr.me.name);
2627 2628

  if (mr.me.type == TSDB_SUPER_TABLE) {
2629 2630
    pSchemaInfo->sw = tCloneSSchemaWrapper(&mr.me.stbEntry.schemaRow);
    pSchemaInfo->tversion = mr.me.stbEntry.schemaTag.version;
2631
  } else if (mr.me.type == TSDB_CHILD_TABLE) {
2632 2633
    tDecoderClear(&mr.coder);

2634 2635
    tb_uid_t suid = mr.me.ctbEntry.suid;
    metaGetTableEntryByUid(&mr, suid);
2636 2637
    pSchemaInfo->sw = tCloneSSchemaWrapper(&mr.me.stbEntry.schemaRow);
    pSchemaInfo->tversion = mr.me.stbEntry.schemaTag.version;
2638
  } else {
2639
    pSchemaInfo->sw = tCloneSSchemaWrapper(&mr.me.ntbEntry.schemaRow);
2640
  }
2641 2642

  metaReaderClear(&mr);
2643

H
Haojun Liao 已提交
2644 2645 2646 2647 2648
  pSchemaInfo->qsw = extractQueriedColumnSchema(pScanNode);
  return TSDB_CODE_SUCCESS;
}

SSchemaWrapper* extractQueriedColumnSchema(SScanPhysiNode* pScanNode) {
2649 2650 2651
  int32_t numOfCols = LIST_LENGTH(pScanNode->pScanCols);
  int32_t numOfTags = LIST_LENGTH(pScanNode->pScanPseudoCols);

2652
  SSchemaWrapper* pqSw = taosMemoryCalloc(1, sizeof(SSchemaWrapper));
2653
  pqSw->pSchema = taosMemoryCalloc(numOfCols + numOfTags, sizeof(SSchema));
2654

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

H
Haojun Liao 已提交
2659 2660 2661
    SSchema* pSchema = &pqSw->pSchema[pqSw->nCols++];
    pSchema->colId = pColNode->colId;
    pSchema->type = pColNode->node.resType.type;
H
Haojun Liao 已提交
2662 2663
    pSchema->bytes = pColNode->node.resType.bytes;
    tstrncpy(pSchema->name, pColNode->colName, tListLen(pSchema->name));
2664 2665
  }

2666
  // this the tags and pseudo function columns, we only keep the tag columns
2667
  for (int32_t i = 0; i < numOfTags; ++i) {
2668 2669 2670 2671 2672 2673 2674 2675 2676
    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 已提交
2677
      pSchema->bytes = pColNode->node.resType.bytes;
H
Haojun Liao 已提交
2678
      tstrncpy(pSchema->name, pColNode->colName, tListLen(pSchema->name));
2679 2680 2681
    }
  }

H
Haojun Liao 已提交
2682
  return pqSw;
2683 2684
}

2685 2686
static void cleanupTableSchemaInfo(SSchemaInfo* pSchemaInfo) {
  taosMemoryFreeClear(pSchemaInfo->dbname);
2687
  taosMemoryFreeClear(pSchemaInfo->tablename);
2688 2689
  tDeleteSSchemaWrapper(pSchemaInfo->sw);
  tDeleteSSchemaWrapper(pSchemaInfo->qsw);
2690 2691
}

2692
static void cleanupStreamInfo(SStreamTaskInfo* pStreamInfo) { tDeleteSSchemaWrapper(pStreamInfo->schema); }
2693

2694
bool groupbyTbname(SNodeList* pGroupList) {
2695
  bool bytbname = false;
2696
  if (LIST_LENGTH(pGroupList) == 1) {
2697 2698 2699 2700 2701 2702 2703 2704 2705 2706
    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;
}

2707 2708
SOperatorInfo* createOperatorTree(SPhysiNode* pPhyNode, SExecTaskInfo* pTaskInfo, SReadHandle* pHandle, SNode* pTagCond,
                                  SNode* pTagIndexCond, const char* pUser) {
2709
  int32_t         type = nodeType(pPhyNode);
2710
  STableListInfo* pTableListInfo = pTaskInfo->pTableInfoList;
2711
  const char*     idstr = GET_TASKID(pTaskInfo);
2712

X
Xiaoyu Wang 已提交
2713
  if (pPhyNode->pChildren == NULL || LIST_LENGTH(pPhyNode->pChildren) == 0) {
2714
    SOperatorInfo* pOperator = NULL;
H
Haojun Liao 已提交
2715
    if (QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN == type) {
dengyihao's avatar
dengyihao 已提交
2716
      STableScanPhysiNode* pTableScanNode = (STableScanPhysiNode*)pPhyNode;
H
Haojun Liao 已提交
2717

2718 2719 2720 2721 2722 2723
      // 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 已提交
2724 2725
      int32_t code =
          createScanTableListInfo(&pTableScanNode->scan, pTableScanNode->pGroupTags, pTableScanNode->groupSort, pHandle,
H
Haojun Liao 已提交
2726
                                  pTableListInfo, pTagCond, pTagIndexCond, pTaskInfo);
2727
      if (code) {
wmmhello's avatar
wmmhello 已提交
2728
        pTaskInfo->code = code;
2729
        qError("failed to createScanTableListInfo, code:%s, %s", tstrerror(code), idstr);
D
dapan1121 已提交
2730 2731
        return NULL;
      }
wmmhello's avatar
wmmhello 已提交
2732

2733
      code = extractTableSchemaInfo(pHandle, &pTableScanNode->scan, pTaskInfo);
S
slzhou 已提交
2734
      if (code) {
2735
        pTaskInfo->code = terrno;
wmmhello's avatar
wmmhello 已提交
2736 2737 2738
        return NULL;
      }

2739
      pOperator = createTableScanOperatorInfo(pTableScanNode, pHandle, pTaskInfo);
D
dapan1121 已提交
2740 2741 2742 2743 2744
      if (NULL == pOperator) {
        pTaskInfo->code = terrno;
        return NULL;
      }

2745 2746
      STableScanInfo* pScanInfo = pOperator->info;
      pTaskInfo->cost.pRecoder = &pScanInfo->readRecorder;
S
slzhou 已提交
2747 2748
    } else if (QUERY_NODE_PHYSICAL_PLAN_TABLE_MERGE_SCAN == type) {
      STableMergeScanPhysiNode* pTableScanNode = (STableMergeScanPhysiNode*)pPhyNode;
H
Haojun Liao 已提交
2749 2750 2751

      int32_t code = createScanTableListInfo(&pTableScanNode->scan, pTableScanNode->pGroupTags, true, pHandle,
                                             pTableListInfo, pTagCond, pTagIndexCond, pTaskInfo);
L
Liu Jicong 已提交
2752
      if (code) {
wmmhello's avatar
wmmhello 已提交
2753
        pTaskInfo->code = code;
H
Haojun Liao 已提交
2754
        qError("failed to createScanTableListInfo, code: %s", tstrerror(code));
wmmhello's avatar
wmmhello 已提交
2755 2756
        return NULL;
      }
2757

2758
      code = extractTableSchemaInfo(pHandle, &pTableScanNode->scan, pTaskInfo);
wmmhello's avatar
wmmhello 已提交
2759 2760 2761 2762
      if (code) {
        pTaskInfo->code = terrno;
        return NULL;
      }
wmmhello's avatar
wmmhello 已提交
2763

2764
      pOperator = createTableMergeScanOperatorInfo(pTableScanNode, pTableListInfo, pHandle, pTaskInfo);
D
dapan1121 已提交
2765 2766 2767 2768
      if (NULL == pOperator) {
        pTaskInfo->code = terrno;
        return NULL;
      }
wmmhello's avatar
wmmhello 已提交
2769

2770 2771
      STableScanInfo* pScanInfo = pOperator->info;
      pTaskInfo->cost.pRecoder = &pScanInfo->readRecorder;
H
Haojun Liao 已提交
2772
    } else if (QUERY_NODE_PHYSICAL_PLAN_EXCHANGE == type) {
2773 2774
      pOperator = createExchangeOperatorInfo(pHandle ? pHandle->pMsgCb->clientRpc : NULL, (SExchangePhysiNode*)pPhyNode,
                                             pTaskInfo);
H
Haojun Liao 已提交
2775
    } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_SCAN == type) {
5
54liuyao 已提交
2776
      STableScanPhysiNode* pTableScanNode = (STableScanPhysiNode*)pPhyNode;
5
54liuyao 已提交
2777
      if (pHandle->vnode) {
L
Liu Jicong 已提交
2778 2779
        int32_t code =
            createScanTableListInfo(&pTableScanNode->scan, pTableScanNode->pGroupTags, pTableScanNode->groupSort,
H
Haojun Liao 已提交
2780
                                    pHandle, pTableListInfo, pTagCond, pTagIndexCond, pTaskInfo);
L
Liu Jicong 已提交
2781
        if (code) {
wmmhello's avatar
wmmhello 已提交
2782
          pTaskInfo->code = code;
H
Haojun Liao 已提交
2783
          qError("failed to createScanTableListInfo, code: %s", tstrerror(code));
wmmhello's avatar
wmmhello 已提交
2784 2785
          return NULL;
        }
L
Liu Jicong 已提交
2786 2787

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

L
Liu Jicong 已提交
2791
        for (int32_t i = 0; i < sz; i++) {
H
Haojun Liao 已提交
2792
          STableKeyInfo* pKeyInfo = tableListGetInfo(pTableListInfo, i);
2793
          qDebug("add table uid:%" PRIu64 ", gid:%" PRIu64, pKeyInfo->uid, pKeyInfo->groupId);
L
Liu Jicong 已提交
2794 2795
        }
#endif
2796
      }
2797

H
Haojun Liao 已提交
2798
      pTaskInfo->schemaInfo.qsw = extractQueriedColumnSchema(&pTableScanNode->scan);
2799
      pOperator = createStreamScanOperatorInfo(pHandle, pTableScanNode, pTagCond, pTaskInfo);
H
Haojun Liao 已提交
2800
    } else if (QUERY_NODE_PHYSICAL_PLAN_SYSTABLE_SCAN == type) {
L
Liu Jicong 已提交
2801
      SSystemTableScanPhysiNode* pSysScanPhyNode = (SSystemTableScanPhysiNode*)pPhyNode;
2802
      pOperator = createSysTableScanOperatorInfo(pHandle, pSysScanPhyNode, pUser, pTaskInfo);
2803
    } else if (QUERY_NODE_PHYSICAL_PLAN_TAG_SCAN == type) {
X
Xiaoyu Wang 已提交
2804
      STagScanPhysiNode* pScanPhyNode = (STagScanPhysiNode*)pPhyNode;
2805 2806

      int32_t code = createScanTableListInfo(pScanPhyNode, NULL, false, pHandle, pTableListInfo, pTagCond,
H
Haojun Liao 已提交
2807
                                             pTagIndexCond, pTaskInfo);
2808
      if (code != TSDB_CODE_SUCCESS) {
2809
        pTaskInfo->code = code;
H
Haojun Liao 已提交
2810
        qError("failed to getTableList, code: %s", tstrerror(code));
2811 2812 2813
        return NULL;
      }

2814
      pOperator = createTagScanOperatorInfo(pHandle, pScanPhyNode, pTableListInfo, pTaskInfo);
2815
    } else if (QUERY_NODE_PHYSICAL_PLAN_BLOCK_DIST_SCAN == type) {
2816
      SBlockDistScanPhysiNode* pBlockNode = (SBlockDistScanPhysiNode*)pPhyNode;
2817 2818

      if (pBlockNode->tableType == TSDB_SUPER_TABLE) {
H
Haojun Liao 已提交
2819 2820
        SArray* pList = taosArrayInit(4, sizeof(STableKeyInfo));
        int32_t code = vnodeGetAllTableList(pHandle->vnode, pBlockNode->uid, pList);
2821 2822 2823 2824
        if (code != TSDB_CODE_SUCCESS) {
          pTaskInfo->code = terrno;
          return NULL;
        }
H
Haojun Liao 已提交
2825

2826
        for (int32_t i = 0; i < tableListGetSize(pTableListInfo); ++i) {
H
Haojun Liao 已提交
2827
          STableKeyInfo* p = taosArrayGet(pList, i);
H
Haojun Liao 已提交
2828
          tableListAddTableInfo(pTableListInfo, p->uid, 0);
H
Haojun Liao 已提交
2829 2830
        }
        taosArrayDestroy(pList);
2831
      } else {  // Create group with only one table
H
Haojun Liao 已提交
2832
        tableListAddTableInfo(pTableListInfo, pBlockNode->uid, 0);
2833 2834
      }

2835
      pOperator = createDataBlockInfoScanOperator(pHandle, pBlockNode, pTaskInfo);
H
Haojun Liao 已提交
2836 2837 2838
    } else if (QUERY_NODE_PHYSICAL_PLAN_LAST_ROW_SCAN == type) {
      SLastRowScanPhysiNode* pScanNode = (SLastRowScanPhysiNode*)pPhyNode;

L
Liu Jicong 已提交
2839
      int32_t code = createScanTableListInfo(&pScanNode->scan, pScanNode->pGroupTags, true, pHandle, pTableListInfo,
H
Haojun Liao 已提交
2840
                                             pTagCond, pTagIndexCond, pTaskInfo);
2841 2842 2843 2844
      if (code != TSDB_CODE_SUCCESS) {
        pTaskInfo->code = code;
        return NULL;
      }
2845

2846
      code = extractTableSchemaInfo(pHandle, &pScanNode->scan, pTaskInfo);
2847 2848 2849
      if (code != TSDB_CODE_SUCCESS) {
        pTaskInfo->code = code;
        return NULL;
H
Haojun Liao 已提交
2850 2851
      }

2852
      pOperator = createCacherowsScanOperator(pScanNode, pHandle, pTaskInfo);
2853
    } else if (QUERY_NODE_PHYSICAL_PLAN_PROJECT == type) {
2854
      pOperator = createProjectOperatorInfo(NULL, (SProjectPhysiNode*)pPhyNode, pTaskInfo);
H
Haojun Liao 已提交
2855 2856
    } else {
      ASSERT(0);
H
Haojun Liao 已提交
2857
    }
2858 2859 2860 2861 2862

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

2863
    return pOperator;
H
Haojun Liao 已提交
2864 2865
  }

2866
  size_t          size = LIST_LENGTH(pPhyNode->pChildren);
2867
  SOperatorInfo** ops = taosMemoryCalloc(size, POINTER_BYTES);
2868 2869 2870 2871
  if (ops == NULL) {
    return NULL;
  }

dengyihao's avatar
dengyihao 已提交
2872
  for (int32_t i = 0; i < size; ++i) {
2873
    SPhysiNode* pChildNode = (SPhysiNode*)nodesListGetNode(pPhyNode->pChildren, i);
2874
    ops[i] = createOperatorTree(pChildNode, pTaskInfo, pHandle, pTagCond, pTagIndexCond, pUser);
2875
    if (ops[i] == NULL) {
H
Haojun Liao 已提交
2876
      taosMemoryFree(ops);
2877 2878
      return NULL;
    }
2879
  }
H
Haojun Liao 已提交
2880

2881
  SOperatorInfo* pOptr = NULL;
H
Haojun Liao 已提交
2882
  if (QUERY_NODE_PHYSICAL_PLAN_PROJECT == type) {
2883
    pOptr = createProjectOperatorInfo(ops[0], (SProjectPhysiNode*)pPhyNode, pTaskInfo);
2884
  } else if (QUERY_NODE_PHYSICAL_PLAN_HASH_AGG == type) {
H
Haojun Liao 已提交
2885 2886
    SAggPhysiNode* pAggNode = (SAggPhysiNode*)pPhyNode;
    if (pAggNode->pGroupKeys != NULL) {
H
Haojun Liao 已提交
2887
      pOptr = createGroupOperatorInfo(ops[0], pAggNode, pTaskInfo);
H
Haojun Liao 已提交
2888
    } else {
H
Haojun Liao 已提交
2889
      pOptr = createAggregateOperatorInfo(ops[0], pAggNode, pTaskInfo);
H
Haojun Liao 已提交
2890
    }
2891
  } else if (QUERY_NODE_PHYSICAL_PLAN_HASH_INTERVAL == type) {
H
Haojun Liao 已提交
2892
    SIntervalPhysiNode* pIntervalPhyNode = (SIntervalPhysiNode*)pPhyNode;
H
Haojun Liao 已提交
2893

H
Haojun Liao 已提交
2894 2895
    bool isStream = (QUERY_NODE_PHYSICAL_PLAN_STREAM_INTERVAL == type);
    pOptr = createIntervalOperatorInfo(ops[0], pIntervalPhyNode, pTaskInfo, isStream);
2896
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_INTERVAL == type) {
2897
    pOptr = createStreamIntervalOperatorInfo(ops[0], pPhyNode, pTaskInfo);
2898 2899
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE_ALIGNED_INTERVAL == type) {
    SMergeAlignedIntervalPhysiNode* pIntervalPhyNode = (SMergeAlignedIntervalPhysiNode*)pPhyNode;
2900
    pOptr = createMergeAlignedIntervalOperatorInfo(ops[0], pIntervalPhyNode, pTaskInfo);
S
shenglian zhou 已提交
2901
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE_INTERVAL == type) {
X
Xiaoyu Wang 已提交
2902
    SMergeIntervalPhysiNode* pIntervalPhyNode = (SMergeIntervalPhysiNode*)pPhyNode;
2903
    pOptr = createMergeIntervalOperatorInfo(ops[0], pIntervalPhyNode, pTaskInfo);
5
54liuyao 已提交
2904
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_SEMI_INTERVAL == type) {
2905
    int32_t children = 0;
5
54liuyao 已提交
2906 2907
    pOptr = createStreamFinalIntervalOperatorInfo(ops[0], pPhyNode, pTaskInfo, children);
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_FINAL_INTERVAL == type) {
5
54liuyao 已提交
2908
    int32_t children = pHandle->numOfVgroups;
5
54liuyao 已提交
2909
    pOptr = createStreamFinalIntervalOperatorInfo(ops[0], pPhyNode, pTaskInfo, children);
H
Haojun Liao 已提交
2910
  } else if (QUERY_NODE_PHYSICAL_PLAN_SORT == type) {
2911
    pOptr = createSortOperatorInfo(ops[0], (SSortPhysiNode*)pPhyNode, pTaskInfo);
S
shenglian zhou 已提交
2912 2913
  } else if (QUERY_NODE_PHYSICAL_PLAN_GROUP_SORT == type) {
    pOptr = createGroupSortOperatorInfo(ops[0], (SGroupSortPhysiNode*)pPhyNode, pTaskInfo);
X
Xiaoyu Wang 已提交
2914
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE == type) {
2915
    SMergePhysiNode* pMergePhyNode = (SMergePhysiNode*)pPhyNode;
2916
    pOptr = createMultiwayMergeOperatorInfo(ops, size, pMergePhyNode, pTaskInfo);
2917
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE_SESSION == type) {
H
Haojun Liao 已提交
2918
    SSessionWinodwPhysiNode* pSessionNode = (SSessionWinodwPhysiNode*)pPhyNode;
H
Haojun Liao 已提交
2919
    pOptr = createSessionAggOperatorInfo(ops[0], pSessionNode, pTaskInfo);
2920
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_SESSION == type) {
2921 2922 2923 2924 2925
    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) {
2926
    int32_t children = pHandle->numOfVgroups;
2927
    pOptr = createStreamFinalSessionAggOperatorInfo(ops[0], pPhyNode, pTaskInfo, children);
H
Haojun Liao 已提交
2928
  } else if (QUERY_NODE_PHYSICAL_PLAN_PARTITION == type) {
2929
    pOptr = createPartitionOperatorInfo(ops[0], (SPartitionPhysiNode*)pPhyNode, pTaskInfo);
2930
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_PARTITION == type) {
2931
    pOptr = createStreamPartitionOperatorInfo(ops[0], (SStreamPartitionPhysiNode*)pPhyNode, pTaskInfo);
2932
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE_STATE == type) {
dengyihao's avatar
dengyihao 已提交
2933
    SStateWinodwPhysiNode* pStateNode = (SStateWinodwPhysiNode*)pPhyNode;
2934
    pOptr = createStatewindowOperatorInfo(ops[0], pStateNode, pTaskInfo);
2935
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_STATE == type) {
5
54liuyao 已提交
2936
    pOptr = createStreamStateAggOperatorInfo(ops[0], pPhyNode, pTaskInfo);
2937
  } else if (QUERY_NODE_PHYSICAL_PLAN_MERGE_JOIN == type) {
2938
    pOptr = createMergeJoinOperatorInfo(ops, size, (SSortMergeJoinPhysiNode*)pPhyNode, pTaskInfo);
2939
  } else if (QUERY_NODE_PHYSICAL_PLAN_FILL == type) {
H
Haojun Liao 已提交
2940
    pOptr = createFillOperatorInfo(ops[0], (SFillPhysiNode*)pPhyNode, pTaskInfo);
5
54liuyao 已提交
2941 2942
  } else if (QUERY_NODE_PHYSICAL_PLAN_STREAM_FILL == type) {
    pOptr = createStreamFillOperatorInfo(ops[0], (SStreamFillPhysiNode*)pPhyNode, pTaskInfo);
H
Haojun Liao 已提交
2943 2944
  } else if (QUERY_NODE_PHYSICAL_PLAN_INDEF_ROWS_FUNC == type) {
    pOptr = createIndefinitOutputOperatorInfo(ops[0], pPhyNode, pTaskInfo);
2945 2946
  } else if (QUERY_NODE_PHYSICAL_PLAN_INTERP_FUNC == type) {
    pOptr = createTimeSliceOperatorInfo(ops[0], pPhyNode, pTaskInfo);
H
Haojun Liao 已提交
2947 2948
  } else {
    ASSERT(0);
H
Haojun Liao 已提交
2949
  }
2950

2951
  taosMemoryFree(ops);
2952 2953 2954 2955
  if (pOptr) {
    pOptr->resultDataBlockId = pPhyNode->pOutputDataBlockDesc->dataBlockId;
  }

2956
  return pOptr;
2957
}
H
Haojun Liao 已提交
2958

L
Liu Jicong 已提交
2959 2960 2961 2962 2963 2964 2965 2966 2967 2968 2969 2970 2971
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 {
2972 2973 2974
    SStreamScanInfo* pInfo = pOperator->info;
    ASSERT(pInfo->pTableScanOp->operatorType == QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN);
    *ppInfo = pInfo->pTableScanOp->info;
L
Liu Jicong 已提交
2975 2976 2977 2978
    return 0;
  }
}

2979 2980 2981 2982 2983 2984 2985 2986 2987 2988 2989 2990 2991 2992 2993 2994 2995 2996 2997 2998 2999 3000
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;
}

3001
#if 0
L
Liu Jicong 已提交
3002 3003 3004 3005 3006
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;
  }
3007

L
Liu Jicong 已提交
3008 3009 3010 3011
  STableScanPhysiNode* pNode = NULL;
  if (extractTableScanNode(plan->pNode, &pNode) < 0) {
    ASSERT(0);
  }
3012

H
Haojun Liao 已提交
3013
  tsdbReaderClose(pTableScanInfo->dataReader);
3014

L
Liu Jicong 已提交
3015
  STableListInfo info = {0};
H
Haojun Liao 已提交
3016
  pTableScanInfo->dataReader = doCreateDataReader(pNode, pHandle, &info, NULL);
L
Liu Jicong 已提交
3017 3018 3019 3020
  if (pTableScanInfo->dataReader == NULL) {
    ASSERT(0);
    qError("failed to create data reader");
    return TSDB_CODE_QRY_APP_ERROR;
3021
  }
L
Liu Jicong 已提交
3022
  // TODO: set uid and ts to data reader
3023 3024
  return 0;
}
3025
#endif
3026

C
Cary Xu 已提交
3027
int32_t encodeOperator(SOperatorInfo* ops, char** result, int32_t* length, int32_t* nOptrWithVal) {
wmmhello's avatar
wmmhello 已提交
3028
  int32_t code = TDB_CODE_SUCCESS;
3029
  char*   pCurrent = NULL;
wmmhello's avatar
wmmhello 已提交
3030
  int32_t currLength = 0;
3031
  if (ops->fpSet.encodeResultRow) {
C
Cary Xu 已提交
3032
    if (result == NULL || length == NULL || nOptrWithVal == NULL) {
wmmhello's avatar
wmmhello 已提交
3033 3034 3035
      return TSDB_CODE_TSC_INVALID_INPUT;
    }
    code = ops->fpSet.encodeResultRow(ops, &pCurrent, &currLength);
wmmhello's avatar
wmmhello 已提交
3036

3037 3038
    if (code != TDB_CODE_SUCCESS) {
      if (*result != NULL) {
wmmhello's avatar
wmmhello 已提交
3039 3040 3041 3042
        taosMemoryFree(*result);
        *result = NULL;
      }
      return code;
C
Cary Xu 已提交
3043 3044 3045
    } else if (currLength == 0) {
      ASSERT(!pCurrent);
      goto _downstream;
wmmhello's avatar
wmmhello 已提交
3046
    }
wmmhello's avatar
wmmhello 已提交
3047

C
Cary Xu 已提交
3048 3049
    ++(*nOptrWithVal);

C
Cary Xu 已提交
3050
    ASSERT(currLength >= 0);
wmmhello's avatar
wmmhello 已提交
3051

3052
    if (*result == NULL) {
wmmhello's avatar
wmmhello 已提交
3053
      *result = (char*)taosMemoryCalloc(1, currLength + sizeof(int32_t));
wmmhello's avatar
wmmhello 已提交
3054 3055 3056 3057 3058 3059
      if (*result == NULL) {
        taosMemoryFree(pCurrent);
        return TSDB_CODE_OUT_OF_MEMORY;
      }
      memcpy(*result + sizeof(int32_t), pCurrent, currLength);
      *(int32_t*)(*result) = currLength + sizeof(int32_t);
3060
    } else {
wmmhello's avatar
wmmhello 已提交
3061
      int32_t sizePre = *(int32_t*)(*result);
3062
      char*   tmp = (char*)taosMemoryRealloc(*result, sizePre + currLength);
wmmhello's avatar
wmmhello 已提交
3063 3064 3065 3066 3067 3068 3069 3070 3071 3072 3073 3074
      if (tmp == NULL) {
        taosMemoryFree(pCurrent);
        taosMemoryFree(*result);
        *result = NULL;
        return TSDB_CODE_OUT_OF_MEMORY;
      }
      *result = tmp;
      memcpy(*result + sizePre, pCurrent, currLength);
      *(int32_t*)(*result) += currLength;
    }
    taosMemoryFree(pCurrent);
    *length = *(int32_t*)(*result);
wmmhello's avatar
wmmhello 已提交
3075 3076
  }

3077
_downstream:
wmmhello's avatar
wmmhello 已提交
3078
  for (int32_t i = 0; i < ops->numOfDownstream; ++i) {
C
Cary Xu 已提交
3079
    code = encodeOperator(ops->pDownstream[i], result, length, nOptrWithVal);
3080
    if (code != TDB_CODE_SUCCESS) {
wmmhello's avatar
wmmhello 已提交
3081
      return code;
wmmhello's avatar
wmmhello 已提交
3082 3083
    }
  }
wmmhello's avatar
wmmhello 已提交
3084
  return TDB_CODE_SUCCESS;
wmmhello's avatar
wmmhello 已提交
3085 3086
}

H
Haojun Liao 已提交
3087
int32_t decodeOperator(SOperatorInfo* ops, const char* result, int32_t length) {
wmmhello's avatar
wmmhello 已提交
3088
  int32_t code = TDB_CODE_SUCCESS;
3089 3090
  if (ops->fpSet.decodeResultRow) {
    if (result == NULL) {
wmmhello's avatar
wmmhello 已提交
3091 3092
      return TSDB_CODE_TSC_INVALID_INPUT;
    }
H
Haojun Liao 已提交
3093

3094
    ASSERT(length == *(int32_t*)result);
H
Haojun Liao 已提交
3095 3096

    const char* data = result + sizeof(int32_t);
L
Liu Jicong 已提交
3097
    code = ops->fpSet.decodeResultRow(ops, (char*)data);
3098
    if (code != TDB_CODE_SUCCESS) {
wmmhello's avatar
wmmhello 已提交
3099 3100
      return code;
    }
wmmhello's avatar
wmmhello 已提交
3101

wmmhello's avatar
wmmhello 已提交
3102
    int32_t totalLength = *(int32_t*)result;
3103 3104
    int32_t dataLength = *(int32_t*)data;

3105
    if (totalLength == dataLength + sizeof(int32_t)) {  // the last data
wmmhello's avatar
wmmhello 已提交
3106 3107
      result = NULL;
      length = 0;
3108
    } else {
wmmhello's avatar
wmmhello 已提交
3109 3110 3111 3112
      result += dataLength;
      *(int32_t*)(result) = totalLength - dataLength;
      length = totalLength - dataLength;
    }
wmmhello's avatar
wmmhello 已提交
3113 3114
  }

wmmhello's avatar
wmmhello 已提交
3115 3116
  for (int32_t i = 0; i < ops->numOfDownstream; ++i) {
    code = decodeOperator(ops->pDownstream[i], result, length);
3117
    if (code != TDB_CODE_SUCCESS) {
wmmhello's avatar
wmmhello 已提交
3118
      return code;
wmmhello's avatar
wmmhello 已提交
3119 3120
    }
  }
wmmhello's avatar
wmmhello 已提交
3121
  return TDB_CODE_SUCCESS;
wmmhello's avatar
wmmhello 已提交
3122 3123
}

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

D
dapan1121 已提交
3127
  switch (pNode->type) {
D
dapan1121 已提交
3128 3129 3130 3131 3132 3133
    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 已提交
3134

D
dapan1121 已提交
3135 3136 3137
      *pParam = pInserterParam;
      break;
    }
D
dapan1121 已提交
3138
    case QUERY_NODE_PHYSICAL_PLAN_DELETE: {
3139
      SDeleterParam* pDeleterParam = taosMemoryCalloc(1, sizeof(SDeleterParam));
D
dapan1121 已提交
3140 3141 3142
      if (NULL == pDeleterParam) {
        return TSDB_CODE_OUT_OF_MEMORY;
      }
H
Haojun Liao 已提交
3143 3144 3145 3146
      int32_t tbNum = tableListGetSize(pTask->pTableInfoList);
      pDeleterParam->suid = tableListGetSuid(pTask->pTableInfoList);

      // TODO extract uid list
D
dapan1121 已提交
3147 3148 3149 3150 3151
      pDeleterParam->pUidList = taosArrayInit(tbNum, sizeof(uint64_t));
      if (NULL == pDeleterParam->pUidList) {
        taosMemoryFree(pDeleterParam);
        return TSDB_CODE_OUT_OF_MEMORY;
      }
H
Haojun Liao 已提交
3152

D
dapan1121 已提交
3153
      for (int32_t i = 0; i < tbNum; ++i) {
H
Haojun Liao 已提交
3154
        STableKeyInfo* pTable = tableListGetInfo(pTask->pTableInfoList, i);
D
dapan1121 已提交
3155 3156 3157 3158 3159 3160 3161 3162 3163 3164 3165 3166 3167
        taosArrayPush(pDeleterParam->pUidList, &pTable->uid);
      }

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

  return TSDB_CODE_SUCCESS;
}

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

D
dapan1121 已提交
3172
  *pTaskInfo = createExecTaskInfo(queryId, taskId, model, pPlan->dbFName);
H
Haojun Liao 已提交
3173 3174 3175
  if (*pTaskInfo == NULL) {
    goto _complete;
  }
H
Haojun Liao 已提交
3176

3177
  if (pHandle) {
L
Liu Jicong 已提交
3178
    /*(*pTaskInfo)->streamInfo.fillHistoryVer1 = pHandle->fillHistoryVer1;*/
3179 3180 3181
    if (pHandle->pStateBackend) {
      (*pTaskInfo)->streamInfo.pState = pHandle->pStateBackend;
    }
H
Haojun Liao 已提交
3182 3183
  }

3184
  (*pTaskInfo)->sql = sql;
D
dapan1121 已提交
3185
  sql = NULL;
H
Haojun Liao 已提交
3186

3187
  (*pTaskInfo)->pSubplan = pPlan;
3188 3189
  (*pTaskInfo)->pRoot =
      createOperatorTree(pPlan->pNode, *pTaskInfo, pHandle, pPlan->pTagCond, pPlan->pTagIndexCond, pPlan->user);
L
Liu Jicong 已提交
3190

D
dapan1121 已提交
3191
  if (NULL == (*pTaskInfo)->pRoot) {
H
Haojun Liao 已提交
3192
    terrno = (*pTaskInfo)->code;
D
dapan1121 已提交
3193
    goto _complete;
3194 3195
  }

H
Haojun Liao 已提交
3196
  return TSDB_CODE_SUCCESS;
H
Haojun Liao 已提交
3197

H
Haojun Liao 已提交
3198
_complete:
D
dapan1121 已提交
3199
  taosMemoryFree(sql);
H
Haojun Liao 已提交
3200
  doDestroyTask(*pTaskInfo);
H
Haojun Liao 已提交
3201
  return terrno;
H
Haojun Liao 已提交
3202 3203
}

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

H
Haojun Liao 已提交
3207
  pTaskInfo->pTableInfoList = tableListDestroy(pTaskInfo->pTableInfoList);
H
Haojun Liao 已提交
3208
  destroyOperatorInfo(pTaskInfo->pRoot);
3209
  cleanupTableSchemaInfo(&pTaskInfo->schemaInfo);
3210
  cleanupStreamInfo(&pTaskInfo->streamInfo);
3211

D
dapan1121 已提交
3212
  if (!pTaskInfo->localFetch.localExec) {
D
dapan1121 已提交
3213 3214
    nodesDestroyNode((SNode*)pTaskInfo->pSubplan);
  }
3215

wafwerar's avatar
wafwerar 已提交
3216 3217 3218
  taosMemoryFreeClear(pTaskInfo->sql);
  taosMemoryFreeClear(pTaskInfo->id.str);
  taosMemoryFreeClear(pTaskInfo);
3219 3220 3221 3222
}

static int64_t getQuerySupportBufSize(size_t numOfTables) {
  size_t s1 = sizeof(STableQueryInfo);
L
Liu Jicong 已提交
3223 3224
  //  size_t s3 = sizeof(STableCheckInfo);  buffer consumption in tsdb
  return (int64_t)(s1 * 1.5 * numOfTables);
3225 3226 3227 3228 3229 3230 3231
}

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 已提交
3232
    while (1) {
3233 3234 3235 3236 3237 3238 3239 3240 3241 3242 3243 3244 3245 3246 3247 3248 3249 3250 3251 3252 3253 3254 3255 3256 3257 3258
      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 已提交
3259

H
Haojun Liao 已提交
3260
int32_t getOperatorExplainExecInfo(SOperatorInfo* operatorInfo, SArray* pExecInfoList) {
3261
  SExplainExecInfo  execInfo = {0};
H
Haojun Liao 已提交
3262
  SExplainExecInfo* pExplainInfo = taosArrayPush(pExecInfoList, &execInfo);
3263

H
Haojun Liao 已提交
3264 3265 3266 3267 3268
  pExplainInfo->numOfRows = operatorInfo->resultInfo.totalRows;
  pExplainInfo->startupCost = operatorInfo->cost.openCost;
  pExplainInfo->totalCost = operatorInfo->cost.totalCost;
  pExplainInfo->verboseLen = 0;
  pExplainInfo->verboseInfo = NULL;
D
dapan1121 已提交
3269

3270
  if (operatorInfo->fpSet.getExplainFn) {
3271 3272
    int32_t code =
        operatorInfo->fpSet.getExplainFn(operatorInfo, &pExplainInfo->verboseInfo, &pExplainInfo->verboseLen);
D
dapan1121 已提交
3273
    if (code) {
3274
      qError("%s operator getExplainFn failed, code:%s", GET_TASKID(operatorInfo->pTaskInfo), tstrerror(code));
D
dapan1121 已提交
3275 3276 3277
      return code;
    }
  }
dengyihao's avatar
dengyihao 已提交
3278

D
dapan1121 已提交
3279
  int32_t code = 0;
D
dapan1121 已提交
3280
  for (int32_t i = 0; i < operatorInfo->numOfDownstream; ++i) {
H
Haojun Liao 已提交
3281 3282
    code = getOperatorExplainExecInfo(operatorInfo->pDownstream[i], pExecInfoList);
    if (code != TSDB_CODE_SUCCESS) {
3283
      //      taosMemoryFreeClear(*pRes);
D
dapan1121 已提交
3284 3285 3286 3287 3288
      return TSDB_CODE_QRY_OUT_OF_MEMORY;
    }
  }

  return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
3289
}
5
54liuyao 已提交
3290

3291 3292
int32_t setOutputBuf(SStreamState* pState, STimeWindow* win, SResultRow** pResult, int64_t tableGroupId,
                     SqlFunctionCtx* pCtx, int32_t numOfOutput, int32_t* rowEntryInfoOffset, SAggSupporter* pAggSup) {
3293 3294 3295 3296 3297 3298
  SWinKey key = {
      .ts = win->skey,
      .groupId = tableGroupId,
  };
  char*   value = NULL;
  int32_t size = pAggSup->resultRowSize;
5
54liuyao 已提交
3299

3300
  if (streamStateAddIfNotExist(pState, &key, (void**)&value, &size) < 0) {
3301 3302 3303 3304 3305 3306 3307 3308 3309 3310
    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;
}

3311 3312
int32_t releaseOutputBuf(SStreamState* pState, SWinKey* pKey, SResultRow* pResult) {
  streamStateReleaseBuf(pState, pKey, pResult);
3313 3314 3315
  return TSDB_CODE_SUCCESS;
}

3316 3317
int32_t saveOutputBuf(SStreamState* pState, SWinKey* pKey, SResultRow* pResult, int32_t resSize) {
  streamStatePut(pState, pKey, pResult, resSize);
3318 3319 3320
  return TSDB_CODE_SUCCESS;
}

3321
int32_t buildDataBlockFromGroupRes(SOperatorInfo* pOperator, SStreamState* pState, SSDataBlock* pBlock, SExprSupp* pSup,
3322
                                   SGroupResInfo* pGroupResInfo) {
3323
  SExecTaskInfo*  pTaskInfo = pOperator->pTaskInfo;
3324 3325 3326 3327 3328 3329 3330 3331 3332 3333 3334 3335
  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 = {
3336 3337
            .ts = *(TSKEY*)pPos->key,
            .groupId = pPos->groupId,
3338
    };
3339
    int32_t code = streamStateGet(pState, &key, &pVal, &size);
3340 3341 3342 3343 3344 3345
    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;
3346
      releaseOutputBuf(pState, &key, pRow);
3347 3348 3349 3350 3351
      continue;
    }

    if (pBlock->info.groupId == 0) {
      pBlock->info.groupId = pPos->groupId;
3352 3353 3354 3355
      SStreamIntervalOperatorInfo* pInfo = pOperator->info;
      char* tbname = taosHashGet(pInfo->pGroupIdTbNameMap, &pBlock->info.groupId, sizeof(int64_t));
      if (tbname != NULL) {
        memcpy(pBlock->info.parTbName, tbname, TSDB_TABLE_NAME_LEN);
L
Liu Jicong 已提交
3356 3357
      } else {
        pBlock->info.parTbName[0] = 0;
3358
      }
3359 3360 3361
    } else {
      // current value belongs to different group, it can't be packed into one datablock
      if (pBlock->info.groupId != pPos->groupId) {
3362
        releaseOutputBuf(pState, &key, pRow);
3363 3364 3365 3366 3367 3368
        break;
      }
    }

    if (pBlock->info.rows + pRow->numOfRows > pBlock->info.capacity) {
      ASSERT(pBlock->info.rows > 0);
3369
      releaseOutputBuf(pState, &key, pRow);
3370 3371 3372 3373 3374 3375 3376 3377 3378 3379
      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) {
3380 3381 3382 3383
        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);
3384 3385 3386 3387 3388 3389 3390 3391 3392 3393 3394 3395 3396
        }
      } 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 已提交
3397

3398
    pBlock->info.rows += pRow->numOfRows;
3399
    releaseOutputBuf(pState, &key, pRow);
3400 3401 3402 3403
  }
  blockDataUpdateTsWindow(pBlock, 0);
  return TSDB_CODE_SUCCESS;
}
5
54liuyao 已提交
3404 3405 3406 3407 3408 3409 3410

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

3411
int32_t buildSessionResultDataBlock(SOperatorInfo* pOperator, SStreamState* pState, SSDataBlock* pBlock,
5
54liuyao 已提交
3412
                                    SExprSupp* pSup, SGroupResInfo* pGroupResInfo) {
3413
  SExecTaskInfo*  pTaskInfo = pOperator->pTaskInfo;
5
54liuyao 已提交
3414 3415 3416 3417 3418 3419 3420 3421 3422 3423 3424 3425 3426
  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);
3427 3428 3429 3430
    if (code == -1) {
      // coverity scan
      continue;
    }
5
54liuyao 已提交
3431 3432 3433 3434 3435 3436 3437 3438 3439 3440 3441
    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;
    }

    if (pBlock->info.groupId == 0) {
      pBlock->info.groupId = pKey->groupId;
3442 3443 3444 3445 3446 3447 3448 3449 3450 3451 3452 3453 3454 3455 3456 3457 3458 3459 3460 3461 3462 3463 3464 3465 3466

      if (pOperator->operatorType == QUERY_NODE_PHYSICAL_PLAN_STREAM_STATE) {
        SStreamStateAggOperatorInfo* pInfo = pOperator->info;

        char* tbname = taosHashGet(pInfo->pGroupIdTbNameMap, &pBlock->info.groupId, sizeof(int64_t));
        if (tbname != NULL) {
          memcpy(pBlock->info.parTbName, tbname, TSDB_TABLE_NAME_LEN);
        } else {
          pBlock->info.parTbName[0] = 0;
        }
      } else if (pOperator->operatorType == QUERY_NODE_PHYSICAL_PLAN_STREAM_SESSION ||
                 pOperator->operatorType == QUERY_NODE_PHYSICAL_PLAN_STREAM_SEMI_SESSION ||
                 pOperator->operatorType == QUERY_NODE_PHYSICAL_PLAN_STREAM_FINAL_SESSION) {
        SStreamSessionAggOperatorInfo* pInfo = pOperator->info;

        char* tbname = taosHashGet(pInfo->pGroupIdTbNameMap, &pBlock->info.groupId, sizeof(int64_t));
        if (tbname != NULL) {
          memcpy(pBlock->info.parTbName, tbname, TSDB_TABLE_NAME_LEN);
        } else {
          pBlock->info.parTbName[0] = 0;
        }
      } else {
        ASSERT(0);
      }

5
54liuyao 已提交
3467 3468 3469 3470 3471 3472 3473 3474 3475 3476 3477 3478 3479 3480 3481 3482 3483 3484 3485 3486 3487 3488 3489 3490 3491 3492 3493 3494 3495 3496 3497 3498 3499 3500 3501 3502 3503 3504 3505 3506 3507 3508 3509 3510 3511
    } else {
      // current value belongs to different group, it can't be packed into one datablock
      if (pBlock->info.groupId != pKey->groupId) {
        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;
3512
}