scanoperator.c 123.9 KB
Newer Older
H
Haojun Liao 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15
/*
 * 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/>.
 */

16
#include "executorimpl.h"
H
Haojun Liao 已提交
17
#include "filter.h"
18
#include "function.h"
19
#include "functionMgt.h"
L
Liu Jicong 已提交
20
#include "os.h"
H
Haojun Liao 已提交
21
#include "querynodes.h"
22
#include "systable.h"
H
Haojun Liao 已提交
23
#include "tname.h"
24
#include "ttime.h"
H
Haojun Liao 已提交
25 26 27 28 29 30 31 32 33 34

#include "tdatablock.h"
#include "tmsg.h"

#include "query.h"
#include "tcompare.h"
#include "thash.h"
#include "ttypes.h"

#define SET_REVERSE_SCAN_FLAG(_info) ((_info)->scanFlag = REVERSE_SCAN)
35
#define SWITCH_ORDER(n)              (((n) = ((n) == TSDB_ORDER_ASC) ? TSDB_ORDER_DESC : TSDB_ORDER_ASC))
36

H
Haojun Liao 已提交
37 38 39 40 41 42 43 44 45 46 47 48
typedef struct STableMergeScanExecInfo {
  SFileBlockLoadRecorder blockRecorder;
  SSortExecInfo          sortExecInfo;
} STableMergeScanExecInfo;

typedef struct STableMergeScanSortSourceParam {
  SOperatorInfo* pOperator;
  int32_t        readerIdx;
  uint64_t       uid;
  SSDataBlock*   inputBlock;
} STableMergeScanSortSourceParam;

L
Liu Jicong 已提交
49
static bool processBlockWithProbability(const SSampleExecInfo* pInfo);
50

H
Haojun Liao 已提交
51
bool processBlockWithProbability(const SSampleExecInfo* pInfo) {
52 53 54 55 56 57 58 59 60 61 62 63
#if 0
  if (pInfo->sampleRatio == 1) {
    return true;
  }

  uint32_t val = taosRandR((uint32_t*) &pInfo->seed);
  return (val % ((uint32_t)(1/pInfo->sampleRatio))) == 0;
#else
  return true;
#endif
}

64
static void switchCtxOrder(SqlFunctionCtx* pCtx, int32_t numOfOutput) {
H
Haojun Liao 已提交
65 66 67 68 69
  for (int32_t i = 0; i < numOfOutput; ++i) {
    SWITCH_ORDER(pCtx[i].order);
  }
}

70 71 72 73 74 75 76 77 78
static void getNextTimeWindow(SInterval* pInterval, STimeWindow* tw, int32_t order) {
  int32_t factor = GET_FORWARD_DIRECTION_FACTOR(order);
  if (pInterval->intervalUnit != 'n' && pInterval->intervalUnit != 'y') {
    tw->skey += pInterval->sliding * factor;
    tw->ekey = tw->skey + pInterval->interval - 1;
    return;
  }

  int64_t key = tw->skey, interval = pInterval->interval;
79
  // convert key to second
80 81 82 83 84 85 86
  key = convertTimePrecision(key, pInterval->precision, TSDB_TIME_PRECISION_MILLI) / 1000;

  if (pInterval->intervalUnit == 'y') {
    interval *= 12;
  }

  struct tm tm;
87
  time_t    t = (time_t)key;
88 89 90 91 92
  taosLocalTime(&t, &tm);

  int mon = (int)(tm.tm_year * 12 + tm.tm_mon + interval * factor);
  tm.tm_year = mon / 12;
  tm.tm_mon = mon % 12;
wafwerar's avatar
wafwerar 已提交
93
  tw->skey = convertTimePrecision((int64_t)taosMktime(&tm) * 1000LL, TSDB_TIME_PRECISION_MILLI, pInterval->precision);
94 95 96 97

  mon = (int)(mon + interval);
  tm.tm_year = mon / 12;
  tm.tm_mon = mon % 12;
wafwerar's avatar
wafwerar 已提交
98
  tw->ekey = convertTimePrecision((int64_t)taosMktime(&tm) * 1000LL, TSDB_TIME_PRECISION_MILLI, pInterval->precision);
99 100 101 102

  tw->ekey -= 1;
}

103
static bool overlapWithTimeWindow(SInterval* pInterval, SDataBlockInfo* pBlockInfo, int32_t order) {
104 105 106 107 108 109 110
  STimeWindow w = {0};

  // 0 by default, which means it is not a interval operator of the upstream operator.
  if (pInterval->interval == 0) {
    return false;
  }

111
  if (order == TSDB_ORDER_ASC) {
112
    w = getAlignQueryTimeWindow(pInterval, pInterval->precision, pBlockInfo->window.skey);
113
    ASSERT(w.ekey >= pBlockInfo->window.skey);
114

115
    if (w.ekey < pBlockInfo->window.ekey) {
116 117 118
      return true;
    }

119 120
    while (1) {
      getNextTimeWindow(pInterval, &w, order);
121 122 123 124
      if (w.skey > pBlockInfo->window.ekey) {
        break;
      }

125
      ASSERT(w.ekey > pBlockInfo->window.ekey);
126
      if (TMAX(w.skey, pBlockInfo->window.skey) <= pBlockInfo->window.ekey) {
127 128 129 130
        return true;
      }
    }
  } else {
131
    w = getAlignQueryTimeWindow(pInterval, pInterval->precision, pBlockInfo->window.ekey);
132
    ASSERT(w.skey <= pBlockInfo->window.ekey);
133

134
    if (w.skey > pBlockInfo->window.skey) {
135 136 137
      return true;
    }

138
    while (1) {
139 140 141 142 143 144
      getNextTimeWindow(pInterval, &w, order);
      if (w.ekey < pBlockInfo->window.skey) {
        break;
      }

      assert(w.skey < pBlockInfo->window.skey);
145
      if (pBlockInfo->window.skey <= TMIN(w.ekey, pBlockInfo->window.ekey)) {
146 147 148
        return true;
      }
    }
149 150 151 152 153
  }

  return false;
}

154 155 156 157 158 159 160 161 162 163 164
// this function is for table scanner to extract temporary results of upstream aggregate results.
static SResultRow* getTableGroupOutputBuf(SOperatorInfo* pOperator, uint64_t groupId, SFilePage** pPage) {
  if (pOperator->operatorType != QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN) {
    return NULL;
  }

  int64_t buf[2] = {0};
  SET_RES_WINDOW_KEY((char*)buf, &groupId, sizeof(groupId), groupId);

  STableScanInfo* pTableScanInfo = pOperator->info;

S
slzhou 已提交
165 166
  SResultRowPosition* p1 = (SResultRowPosition*)tSimpleHashGet(pTableScanInfo->base.pdInfo.pAggSup->pResultRowHashTable,
                                                               buf, GET_RES_WINDOW_KEY_LEN(sizeof(groupId)));
167 168 169 170 171

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

H
Haojun Liao 已提交
172
  *pPage = getBufPage(pTableScanInfo->base.pdInfo.pAggSup->pResultBuf, p1->pageId);
173 174 175
  if (NULL == *pPage) {
    return NULL;
  }
L
Liu Jicong 已提交
176

177 178 179 180 181 182
  return (SResultRow*)((char*)(*pPage) + p1->offset);
}

static int32_t doDynamicPruneDataBlock(SOperatorInfo* pOperator, SDataBlockInfo* pBlockInfo, uint32_t* status) {
  STableScanInfo* pTableScanInfo = pOperator->info;

H
Haojun Liao 已提交
183
  if (pTableScanInfo->base.pdInfo.pExprSup == NULL) {
184 185 186
    return TSDB_CODE_SUCCESS;
  }

H
Haojun Liao 已提交
187
  SExprSupp* pSup1 = pTableScanInfo->base.pdInfo.pExprSup;
188 189

  SFilePage*  pPage = NULL;
H
Haojun Liao 已提交
190
  SResultRow* pRow = getTableGroupOutputBuf(pOperator, pBlockInfo->id.groupId, &pPage);
191 192 193 194 195 196 197 198 199

  if (pRow == NULL) {
    return TSDB_CODE_SUCCESS;
  }

  bool notLoadBlock = true;
  for (int32_t i = 0; i < pSup1->numOfExprs; ++i) {
    int32_t functionId = pSup1->pCtx[i].functionId;

H
Haojun Liao 已提交
200
    SResultRowEntryInfo* pEntry = getResultEntryInfo(pRow, i, pTableScanInfo->base.pdInfo.pExprSup->rowEntryInfoOffset);
201 202 203 204 205 206 207 208 209

    int32_t reqStatus = fmFuncDynDataRequired(functionId, pEntry, &pBlockInfo->window);
    if (reqStatus != FUNC_DATA_REQUIRED_NOT_LOAD) {
      notLoadBlock = false;
      break;
    }
  }

  // release buffer pages
H
Haojun Liao 已提交
210
  releaseBufPage(pTableScanInfo->base.pdInfo.pAggSup->pResultBuf, pPage);
211 212 213 214 215 216 217 218

  if (notLoadBlock) {
    *status = FUNC_DATA_REQUIRED_NOT_LOAD;
  }

  return TSDB_CODE_SUCCESS;
}

H
Haojun Liao 已提交
219
static bool doFilterByBlockSMA(SFilterInfo* pFilterInfo, SColumnDataAgg** pColsAgg, int32_t numOfCols,
220
                               int32_t numOfRows) {
H
Haojun Liao 已提交
221
  if (pColsAgg == NULL || pFilterInfo == NULL) {
H
Haojun Liao 已提交
222 223 224
    return true;
  }

H
Haojun Liao 已提交
225
  bool keep = filterRangeExecute(pFilterInfo, pColsAgg, numOfCols, numOfRows);
H
Haojun Liao 已提交
226 227 228
  return keep;
}

H
Haojun Liao 已提交
229
static bool doLoadBlockSMA(STableScanBase* pTableScanInfo, SSDataBlock* pBlock, SExecTaskInfo* pTaskInfo) {
H
Haojun Liao 已提交
230
  bool    allColumnsHaveAgg = true;
231
  int32_t code = tsdbRetrieveDatablockSMA(pTableScanInfo->dataReader, pBlock, &allColumnsHaveAgg);
H
Haojun Liao 已提交
232
  if (code != TSDB_CODE_SUCCESS) {
233
    T_LONG_JMP(pTaskInfo->env, code);
H
Haojun Liao 已提交
234 235 236 237 238 239 240 241
  }

  if (!allColumnsHaveAgg) {
    return false;
  }
  return true;
}

H
Haojun Liao 已提交
242
static void doSetTagColumnData(STableScanBase* pTableScanInfo, SSDataBlock* pBlock, SExecTaskInfo* pTaskInfo,
243
                               int32_t rows) {
H
Haojun Liao 已提交
244 245 246
  if (pTableScanInfo->pseudoSup.numOfExprs > 0) {
    SExprSupp* pSup = &pTableScanInfo->pseudoSup;

247
    int32_t code = addTagPseudoColumnData(&pTableScanInfo->readHandle, pSup->pExprInfo, pSup->numOfExprs, pBlock, rows,
248
                                          GET_TASKID(pTaskInfo), &pTableScanInfo->metaCache);
H
Haojun Liao 已提交
249
    // ignore the table not exists error, since this table may have been dropped during the scan procedure.
H
Haojun Liao 已提交
250
    if (code != TSDB_CODE_SUCCESS && code != TSDB_CODE_PAR_TABLE_NOT_EXIST) {
H
Haojun Liao 已提交
251 252
      T_LONG_JMP(pTaskInfo->env, code);
    }
H
Haojun Liao 已提交
253 254 255

    // reset the error code.
    terrno = 0;
H
Haojun Liao 已提交
256 257 258
  }
}

259
// todo handle the slimit info
260
bool applyLimitOffset(SLimitInfo* pLimitInfo, SSDataBlock* pBlock, SExecTaskInfo* pTaskInfo, SOperatorInfo* pOperator) {
261
  SLimit*     pLimit = &pLimitInfo->limit;
H
Haojun Liao 已提交
262
  const char* id = GET_TASKID(pTaskInfo);
263 264 265 266

  if (pLimit->offset > 0 && pLimitInfo->remainOffset > 0) {
    if (pLimitInfo->remainOffset >= pBlock->info.rows) {
      pLimitInfo->remainOffset -= pBlock->info.rows;
H
Haojun Liao 已提交
267
      blockDataEmpty(pBlock);
H
Haojun Liao 已提交
268
      qDebug("current block ignore due to offset, current:%" PRId64 ", %s", pLimitInfo->remainOffset, id);
269
      return false;
270 271 272 273 274 275 276 277
    } else {
      blockDataTrimFirstNRows(pBlock, pLimitInfo->remainOffset);
      pLimitInfo->remainOffset = 0;
    }
  }

  if (pLimit->limit != -1 && pLimit->limit <= (pLimitInfo->numOfOutputRows + pBlock->info.rows)) {
    // limit the output rows
278
    int32_t keep = (int32_t)(pLimit->limit - pLimitInfo->numOfOutputRows);
279 280

    blockDataKeepFirstNRows(pBlock, keep);
H
Haojun Liao 已提交
281
    qDebug("output limit %" PRId64 " has reached, %s", pLimit->limit, id);
282
    return true;
283
  }
284 285

  return false;
286 287
}

H
Haojun Liao 已提交
288
static int32_t loadDataBlock(SOperatorInfo* pOperator, STableScanBase* pTableScanInfo, SSDataBlock* pBlock,
L
Liu Jicong 已提交
289
                             uint32_t* status) {
S
slzhou 已提交
290
  SExecTaskInfo*          pTaskInfo = pOperator->pTaskInfo;
291
  SFileBlockLoadRecorder* pCost = &pTableScanInfo->readRecorder;
H
Haojun Liao 已提交
292 293

  pCost->totalBlocks += 1;
294
  pCost->totalRows += pBlock->info.rows;
295

H
Haojun Liao 已提交
296
  bool loadSMA = false;
H
Haojun Liao 已提交
297
  *status = pTableScanInfo->dataBlockLoadFlag;
H
Haojun Liao 已提交
298
  if (pOperator->exprSupp.pFilterInfo != NULL ||
299
      overlapWithTimeWindow(&pTableScanInfo->pdInfo.interval, &pBlock->info, pTableScanInfo->cond.order)) {
300 301 302 303
    (*status) = FUNC_DATA_REQUIRED_DATA_LOAD;
  }

  SDataBlockInfo* pBlockInfo = &pBlock->info;
304
  taosMemoryFreeClear(pBlock->pBlockAgg);
305 306

  if (*status == FUNC_DATA_REQUIRED_FILTEROUT) {
307 308
    qDebug("%s data block filter out, brange:%" PRId64 "-%" PRId64 ", rows:%d", GET_TASKID(pTaskInfo),
           pBlockInfo->window.skey, pBlockInfo->window.ekey, pBlockInfo->rows);
309
    pCost->filterOutBlocks += 1;
310
    pCost->totalRows += pBlock->info.rows;
311
    tsdbReleaseDataBlock(pTableScanInfo->dataReader);
312 313
    return TSDB_CODE_SUCCESS;
  } else if (*status == FUNC_DATA_REQUIRED_NOT_LOAD) {
314 315
    qDebug("%s data block skipped, brange:%" PRId64 "-%" PRId64 ", rows:%d", GET_TASKID(pTaskInfo),
           pBlockInfo->window.skey, pBlockInfo->window.ekey, pBlockInfo->rows);
316
    doSetTagColumnData(pTableScanInfo, pBlock, pTaskInfo, 1);
317
    pCost->skipBlocks += 1;
318
    tsdbReleaseDataBlock(pTableScanInfo->dataReader);
319
    return TSDB_CODE_SUCCESS;
H
Haojun Liao 已提交
320
  } else if (*status == FUNC_DATA_REQUIRED_SMA_LOAD) {
321
    pCost->loadBlockStatis += 1;
L
Liu Jicong 已提交
322
    loadSMA = true;  // mark the operation of load sma;
H
Haojun Liao 已提交
323
    bool success = doLoadBlockSMA(pTableScanInfo, pBlock, pTaskInfo);
L
Liu Jicong 已提交
324
    if (success) {  // failed to load the block sma data, data block statistics does not exist, load data block instead
325 326
      qDebug("%s data block SMA loaded, brange:%" PRId64 "-%" PRId64 ", rows:%d", GET_TASKID(pTaskInfo),
             pBlockInfo->window.skey, pBlockInfo->window.ekey, pBlockInfo->rows);
327
      doSetTagColumnData(pTableScanInfo, pBlock, pTaskInfo, 1);
328
      tsdbReleaseDataBlock(pTableScanInfo->dataReader);
329 330
      return TSDB_CODE_SUCCESS;
    } else {
331
      qDebug("%s failed to load SMA, since not all columns have SMA", GET_TASKID(pTaskInfo));
H
Haojun Liao 已提交
332
      *status = FUNC_DATA_REQUIRED_DATA_LOAD;
333
    }
H
Haojun Liao 已提交
334
  }
335

H
Haojun Liao 已提交
336
  ASSERT(*status == FUNC_DATA_REQUIRED_DATA_LOAD);
337

H
Haojun Liao 已提交
338
  // try to filter data block according to sma info
H
Haojun Liao 已提交
339
  if (pOperator->exprSupp.pFilterInfo != NULL && (!loadSMA)) {
340 341 342
    bool success = doLoadBlockSMA(pTableScanInfo, pBlock, pTaskInfo);
    if (success) {
      size_t size = taosArrayGetSize(pBlock->pDataBlock);
H
Haojun Liao 已提交
343
      bool   keep = doFilterByBlockSMA(pOperator->exprSupp.pFilterInfo, pBlock->pBlockAgg, size, pBlockInfo->rows);
344 345 346 347 348 349
      if (!keep) {
        qDebug("%s data block filter out by block SMA, brange:%" PRId64 "-%" PRId64 ", rows:%d", GET_TASKID(pTaskInfo),
               pBlockInfo->window.skey, pBlockInfo->window.ekey, pBlockInfo->rows);
        pCost->filterOutBlocks += 1;
        (*status) = FUNC_DATA_REQUIRED_FILTEROUT;

350
        tsdbReleaseDataBlock(pTableScanInfo->dataReader);
351 352
        return TSDB_CODE_SUCCESS;
      }
353
    }
H
Haojun Liao 已提交
354
  }
355

356 357 358
  // free the sma info, since it should not be involved in later computing process.
  taosMemoryFreeClear(pBlock->pBlockAgg);

359
  // try to filter data block according to current results
360 361
  doDynamicPruneDataBlock(pOperator, pBlockInfo, status);
  if (*status == FUNC_DATA_REQUIRED_NOT_LOAD) {
362
    qDebug("%s data block skipped due to dynamic prune, brange:%" PRId64 "-%" PRId64 ", rows:%d", GET_TASKID(pTaskInfo),
363 364
           pBlockInfo->window.skey, pBlockInfo->window.ekey, pBlockInfo->rows);
    pCost->skipBlocks += 1;
365
    tsdbReleaseDataBlock(pTableScanInfo->dataReader);
366
    *status = FUNC_DATA_REQUIRED_FILTEROUT;
367 368 369
    return TSDB_CODE_SUCCESS;
  }

H
Haojun Liao 已提交
370 371
  pCost->totalCheckedRows += pBlock->info.rows;
  pCost->loadBlocks += 1;
372

H
Haojun Liao 已提交
373 374
  SSDataBlock* p = tsdbRetrieveDataBlock(pTableScanInfo->dataReader, NULL);
  if (p == NULL) {
H
Haojun Liao 已提交
375
    return terrno;
H
Haojun Liao 已提交
376 377
  }

H
Haojun Liao 已提交
378
  ASSERT(p == pBlock);
379
  doSetTagColumnData(pTableScanInfo, pBlock, pTaskInfo, pBlock->info.rows);
380

H
Haojun Liao 已提交
381 382
  // restore the previous value
  pCost->totalRows -= pBlock->info.rows;
383

H
Haojun Liao 已提交
384
  if (pOperator->exprSupp.pFilterInfo != NULL) {
385
    int64_t st = taosGetTimestampUs();
H
Haojun Liao 已提交
386
    doFilter(pBlock, pOperator->exprSupp.pFilterInfo, &pTableScanInfo->matchInfo);
387

388 389
    double el = (taosGetTimestampUs() - st) / 1000.0;
    pTableScanInfo->readRecorder.filterTime += el;
390

391 392 393 394 395 396 397
    if (pBlock->info.rows == 0) {
      pCost->filterOutBlocks += 1;
      qDebug("%s data block filter out, brange:%" PRId64 "-%" PRId64 ", rows:%d, elapsed time:%.2f ms",
             GET_TASKID(pTaskInfo), pBlockInfo->window.skey, pBlockInfo->window.ekey, pBlockInfo->rows, el);
    } else {
      qDebug("%s data block filter applied, elapsed time:%.2f ms", GET_TASKID(pTaskInfo), el);
    }
398 399
  }

400 401 402 403
  bool limitReached = applyLimitOffset(&pTableScanInfo->limitInfo, pBlock, pTaskInfo, pOperator);
  if (limitReached) { // set operator flag is done
    setOperatorCompleted(pOperator);
  }
404

H
Haojun Liao 已提交
405
  pCost->totalRows += pBlock->info.rows;
H
Haojun Liao 已提交
406
  pTableScanInfo->limitInfo.numOfOutputRows = pCost->totalRows;
H
Haojun Liao 已提交
407 408 409
  return TSDB_CODE_SUCCESS;
}

H
Haojun Liao 已提交
410
static void prepareForDescendingScan(STableScanBase* pTableScanInfo, SqlFunctionCtx* pCtx, int32_t numOfOutput) {
H
Haojun Liao 已提交
411 412 413
  SET_REVERSE_SCAN_FLAG(pTableScanInfo);

  switchCtxOrder(pCtx, numOfOutput);
414
  pTableScanInfo->cond.order = TSDB_ORDER_DESC;
H
Haojun Liao 已提交
415 416
  STimeWindow* pTWindow = &pTableScanInfo->cond.twindows;
  TSWAP(pTWindow->skey, pTWindow->ekey);
H
Haojun Liao 已提交
417 418
}

419 420
typedef struct STableCachedVal {
  const char* pName;
421
  STag*       pTags;
422 423
} STableCachedVal;

424 425 426 427 428 429 430 431 432 433 434
static void freeTableCachedVal(void* param) {
  if (param == NULL) {
    return;
  }

  STableCachedVal* pVal = param;
  taosMemoryFree((void*)pVal->pName);
  taosMemoryFree(pVal->pTags);
  taosMemoryFree(pVal);
}

H
Haojun Liao 已提交
435 436 437 438 439 440 441
static STableCachedVal* createTableCacheVal(const SMetaReader* pMetaReader) {
  STableCachedVal* pVal = taosMemoryMalloc(sizeof(STableCachedVal));
  pVal->pName = strdup(pMetaReader->me.name);
  pVal->pTags = NULL;

  // only child table has tag value
  if (pMetaReader->me.type == TSDB_CHILD_TABLE) {
442
    STag* pTag = (STag*)pMetaReader->me.ctbEntry.pTags;
H
Haojun Liao 已提交
443 444 445 446 447 448 449
    pVal->pTags = taosMemoryMalloc(pTag->len);
    memcpy(pVal->pTags, pTag, pTag->len);
  }

  return pVal;
}

450 451
// const void *key, size_t keyLen, void *value
static void freeCachedMetaItem(const void* key, size_t keyLen, void* value) { freeTableCachedVal(value); }
452

453 454
int32_t addTagPseudoColumnData(SReadHandle* pHandle, const SExprInfo* pExpr, int32_t numOfExpr, SSDataBlock* pBlock,
                               int32_t rows, const char* idStr, STableMetaCacheInfo* pCache) {
455
  // currently only the tbname pseudo column
456
  if (numOfExpr <= 0) {
H
Haojun Liao 已提交
457
    return TSDB_CODE_SUCCESS;
458 459
  }

460 461
  int32_t code = 0;

462 463 464 465
  // backup the rows
  int32_t backupRows = pBlock->info.rows;
  pBlock->info.rows = rows;

466
  bool            freeReader = false;
467
  STableCachedVal val = {0};
468 469

  SMetaReader mr = {0};
470
  LRUHandle*  h = NULL;
471

472
  // 1. check if it is existed in meta cache
473
  if (pCache == NULL) {
474
    metaReaderInit(&mr, pHandle->meta, 0);
H
Haojun Liao 已提交
475
    code = metaGetTableEntryByUidCache(&mr, pBlock->info.id.uid);
476
    if (code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
477
      if (terrno == TSDB_CODE_PAR_TABLE_NOT_EXIST) {
S
slzhou 已提交
478 479
        qWarn("failed to get table meta, table may have been dropped, uid:0x%" PRIx64 ", code:%s, %s",
              pBlock->info.id.uid, tstrerror(terrno), idStr);
H
Haojun Liao 已提交
480
      } else {
S
slzhou 已提交
481 482
        qError("failed to get table meta, uid:0x%" PRIx64 ", code:%s, %s", pBlock->info.id.uid, tstrerror(terrno),
               idStr);
H
Haojun Liao 已提交
483
      }
484 485 486 487 488
      metaReaderClear(&mr);
      return terrno;
    }

    metaReaderReleaseLock(&mr);
489

490 491
    val.pName = mr.me.name;
    val.pTags = (STag*)mr.me.ctbEntry.pTags;
492 493

    freeReader = true;
494
  } else {
495 496
    pCache->metaFetch += 1;

H
Haojun Liao 已提交
497
    h = taosLRUCacheLookup(pCache->pTableMetaEntryCache, &pBlock->info.id.uid, sizeof(pBlock->info.id.uid));
498 499
    if (h == NULL) {
      metaReaderInit(&mr, pHandle->meta, 0);
H
Haojun Liao 已提交
500
      code = metaGetTableEntryByUidCache(&mr, pBlock->info.id.uid);
501
      if (code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
502
        if (terrno == TSDB_CODE_PAR_TABLE_NOT_EXIST) {
503
          qWarn("failed to get table meta, table may have been dropped, uid:0x%" PRIx64 ", code:%s, %s",
H
Haojun Liao 已提交
504
                pBlock->info.id.uid, tstrerror(terrno), idStr);
H
Haojun Liao 已提交
505
        } else {
H
Haojun Liao 已提交
506
          qError("failed to get table meta, uid:0x%" PRIx64 ", code:%s, %s", pBlock->info.id.uid, tstrerror(terrno),
507
                 idStr);
H
Haojun Liao 已提交
508
        }
509 510 511 512 513 514
        metaReaderClear(&mr);
        return terrno;
      }

      metaReaderReleaseLock(&mr);

H
Haojun Liao 已提交
515
      STableCachedVal* pVal = createTableCacheVal(&mr);
516

H
Haojun Liao 已提交
517
      val = *pVal;
518
      freeReader = true;
H
Haojun Liao 已提交
519

H
Haojun Liao 已提交
520
      int32_t ret = taosLRUCacheInsert(pCache->pTableMetaEntryCache, &pBlock->info.id.uid, sizeof(uint64_t), pVal,
521
                                       sizeof(STableCachedVal), freeCachedMetaItem, NULL, TAOS_LRU_PRIORITY_LOW);
522 523 524 525 526 527 528 529
      if (ret != TAOS_LRU_STATUS_OK) {
        qError("failed to put meta into lru cache, code:%d, %s", ret, idStr);
        freeTableCachedVal(pVal);
      }
    } else {
      pCache->cacheHit += 1;
      STableCachedVal* pVal = taosLRUCacheValue(pCache->pTableMetaEntryCache, h);
      val = *pVal;
H
Haojun Liao 已提交
530

H
Haojun Liao 已提交
531
      taosLRUCacheRelease(pCache->pTableMetaEntryCache, h, false);
532
    }
H
Haojun Liao 已提交
533

534 535
    qDebug("retrieve table meta from cache:%" PRIu64 ", hit:%" PRIu64 " miss:%" PRIu64 ", %s", pCache->metaFetch,
           pCache->cacheHit, (pCache->metaFetch - pCache->cacheHit), idStr);
H
Haojun Liao 已提交
536
  }
537

538 539
  for (int32_t j = 0; j < numOfExpr; ++j) {
    const SExprInfo* pExpr1 = &pExpr[j];
540
    int32_t          dstSlotId = pExpr1->base.resSchema.slotId;
541 542

    SColumnInfoData* pColInfoData = taosArrayGet(pBlock->pDataBlock, dstSlotId);
D
dapan1121 已提交
543
    colInfoDataCleanup(pColInfoData, pBlock->info.rows);
544

545
    int32_t functionId = pExpr1->pExpr->_function.functionId;
546 547 548

    // this is to handle the tbname
    if (fmIsScanPseudoColumnFunc(functionId)) {
549
      setTbNameColData(pBlock, pColInfoData, functionId, val.pName);
550
    } else {  // these are tags
wmmhello's avatar
wmmhello 已提交
551
      STagVal tagVal = {0};
552 553
      tagVal.cid = pExpr1->base.pParam[0].pCol->colId;
      const char* p = metaGetTableTagVal(val.pTags, pColInfoData->info.type, &tagVal);
wmmhello's avatar
wmmhello 已提交
554

555 556 557 558
      char* data = NULL;
      if (pColInfoData->info.type != TSDB_DATA_TYPE_JSON && p != NULL) {
        data = tTagValToData((const STagVal*)p, false);
      } else {
wmmhello's avatar
wmmhello 已提交
559
        data = (char*)p;
wmmhello's avatar
wmmhello 已提交
560
      }
561

H
Haojun Liao 已提交
562 563 564
      bool isNullVal = (data == NULL) || (pColInfoData->info.type == TSDB_DATA_TYPE_JSON && tTagIsJsonNull(data));
      if (isNullVal) {
        colDataAppendNNULL(pColInfoData, 0, pBlock->info.rows);
H
Haojun Liao 已提交
565
      } else if (pColInfoData->info.type != TSDB_DATA_TYPE_JSON) {
H
Haojun Liao 已提交
566
        colDataAppendNItems(pColInfoData, 0, data, pBlock->info.rows);
H
Haojun Liao 已提交
567 568 569
        if (IS_VAR_DATA_TYPE(((const STagVal*)p)->type)) {
          taosMemoryFree(data);
        }
L
Liu Jicong 已提交
570
      } else {  // todo opt for json tag
H
Haojun Liao 已提交
571
        for (int32_t i = 0; i < pBlock->info.rows; ++i) {
H
Haojun Liao 已提交
572
          colDataAppend(pColInfoData, i, data, false);
H
Haojun Liao 已提交
573
        }
574 575 576 577
      }
    }
  }

578 579
  // restore the rows
  pBlock->info.rows = backupRows;
580 581 582 583
  if (freeReader) {
    metaReaderClear(&mr);
  }

H
Haojun Liao 已提交
584
  return TSDB_CODE_SUCCESS;
585 586
}

H
Haojun Liao 已提交
587
void setTbNameColData(const SSDataBlock* pBlock, SColumnInfoData* pColInfoData, int32_t functionId, const char* name) {
588 589 590
  struct SScalarFuncExecFuncs fpSet = {0};
  fmGetScalarFuncExecFuncs(functionId, &fpSet);

H
Haojun Liao 已提交
591
  size_t len = TSDB_TABLE_FNAME_LEN + VARSTR_HEADER_SIZE;
592
  char   buf[TSDB_TABLE_FNAME_LEN + VARSTR_HEADER_SIZE] = {0};
H
Haojun Liao 已提交
593 594 595
  STR_TO_VARSTR(buf, name)

  SColumnInfoData infoData = createColumnInfoData(TSDB_DATA_TYPE_VARCHAR, len, 1);
596

H
Haojun Liao 已提交
597 598
  colInfoDataEnsureCapacity(&infoData, 1, false);
  colDataAppend(&infoData, 0, buf, false);
599

H
Haojun Liao 已提交
600
  SScalarParam srcParam = {.numOfRows = pBlock->info.rows, .columnData = &infoData};
601
  SScalarParam param = {.columnData = pColInfoData};
H
Haojun Liao 已提交
602 603 604 605 606 607 608

  if (fpSet.process != NULL) {
    fpSet.process(&srcParam, 1, &param);
  } else {
    qError("failed to get the corresponding callback function, functionId:%d", functionId);
  }

D
dapan1121 已提交
609
  colDataDestroy(&infoData);
610 611
}

612
static SSDataBlock* doTableScanImpl(SOperatorInfo* pOperator) {
H
Haojun Liao 已提交
613
  STableScanInfo* pTableScanInfo = pOperator->info;
614
  SExecTaskInfo*  pTaskInfo = pOperator->pTaskInfo;
L
Liu Jicong 已提交
615
  SSDataBlock*    pBlock = pTableScanInfo->pResBlock;
H
Haojun Liao 已提交
616

617 618
  int64_t st = taosGetTimestampUs();

H
Haojun Liao 已提交
619
  while (tsdbNextDataBlock(pTableScanInfo->base.dataReader)) {
620
    if (isTaskKilled(pTaskInfo)) {
621
      T_LONG_JMP(pTaskInfo->env, pTaskInfo->code);
622
    }
H
Haojun Liao 已提交
623

624 625 626 627 628 629
    // process this data block based on the probabilities
    bool processThisBlock = processBlockWithProbability(&pTableScanInfo->sample);
    if (!processThisBlock) {
      continue;
    }

H
Haojun Liao 已提交
630
    ASSERT(pBlock->info.id.uid != 0);
H
Haojun Liao 已提交
631
    pBlock->info.id.groupId = getTableGroupId(pTaskInfo->pTableInfoList, pBlock->info.id.uid);
632

633
    uint32_t status = 0;
H
Haojun Liao 已提交
634
    int32_t  code = loadDataBlock(pOperator, &pTableScanInfo->base, pBlock, &status);
635 636
    //    int32_t  code = loadDataBlockOnDemand(pOperator->pRuntimeEnv, pTableScanInfo, pBlock, &status);
    if (code != TSDB_CODE_SUCCESS) {
637
      T_LONG_JMP(pOperator->pTaskInfo->env, code);
638
    }
639

640 641 642
    // current block is filter out according to filter condition, continue load the next block
    if (status == FUNC_DATA_REQUIRED_FILTEROUT || pBlock->info.rows == 0) {
      continue;
643
    }
644

H
Haojun Liao 已提交
645 646
    pOperator->resultInfo.totalRows = pTableScanInfo->base.readRecorder.totalRows;
    pTableScanInfo->base.readRecorder.elapsedTime += (taosGetTimestampUs() - st) / 1000.0;
647

H
Haojun Liao 已提交
648
    pOperator->cost.totalCost = pTableScanInfo->base.readRecorder.elapsedTime;
649 650

    // todo refactor
H
Haojun Liao 已提交
651
    /*pTableScanInfo->lastStatus.uid = pBlock->info.id.uid;*/
L
Liu Jicong 已提交
652 653
    /*pTableScanInfo->lastStatus.ts = pBlock->info.window.ekey;*/
    pTaskInfo->streamInfo.lastStatus.type = TMQ_OFFSET__SNAPSHOT_DATA;
H
Haojun Liao 已提交
654
    pTaskInfo->streamInfo.lastStatus.uid = pBlock->info.id.uid;
L
Liu Jicong 已提交
655
    pTaskInfo->streamInfo.lastStatus.ts = pBlock->info.window.ekey;
656

H
Haojun Liao 已提交
657
    ASSERT(pBlock->info.id.uid != 0);
658
    return pBlock;
H
Haojun Liao 已提交
659 660 661 662
  }
  return NULL;
}

H
Haojun Liao 已提交
663
static SSDataBlock* doGroupedTableScan(SOperatorInfo* pOperator) {
H
Haojun Liao 已提交
664 665 666 667
  STableScanInfo* pTableScanInfo = pOperator->info;
  SExecTaskInfo*  pTaskInfo = pOperator->pTaskInfo;

  // The read handle is not initialized yet, since no qualified tables exists
H
Haojun Liao 已提交
668
  if (pTableScanInfo->base.dataReader == NULL || pOperator->status == OP_EXEC_DONE) {
H
Haojun Liao 已提交
669 670 671
    return NULL;
  }

672 673
  // do the ascending order traverse in the first place.
  while (pTableScanInfo->scanTimes < pTableScanInfo->scanInfo.numOfAsc) {
H
Haojun Liao 已提交
674 675 676
    SSDataBlock* p = doTableScanImpl(pOperator);
    if (p != NULL) {
      return p;
H
Haojun Liao 已提交
677 678
    }

679
    pTableScanInfo->scanTimes += 1;
680

681
    if (pTableScanInfo->scanTimes < pTableScanInfo->scanInfo.numOfAsc) {
682
      setTaskStatus(pTaskInfo, TASK_NOT_COMPLETED);
H
Haojun Liao 已提交
683
      pTableScanInfo->base.scanFlag = REPEAT_SCAN;
684
      qDebug("start to repeat ascending order scan data blocks due to query func required, %s", GET_TASKID(pTaskInfo));
H
Haojun Liao 已提交
685

686
      // do prepare for the next round table scan operation
H
Haojun Liao 已提交
687
      tsdbReaderReset(pTableScanInfo->base.dataReader, &pTableScanInfo->base.cond);
H
Haojun Liao 已提交
688
    }
689
  }
H
Haojun Liao 已提交
690

691
  int32_t total = pTableScanInfo->scanInfo.numOfAsc + pTableScanInfo->scanInfo.numOfDesc;
692
  if (pTableScanInfo->scanTimes < total) {
H
Haojun Liao 已提交
693 694 695
    if (pTableScanInfo->base.cond.order == TSDB_ORDER_ASC) {
      prepareForDescendingScan(&pTableScanInfo->base, pOperator->exprSupp.pCtx, 0);
      tsdbReaderReset(pTableScanInfo->base.dataReader, &pTableScanInfo->base.cond);
696
      qDebug("%s start to descending order scan data blocks due to query func required", GET_TASKID(pTaskInfo));
697
    }
H
Haojun Liao 已提交
698

699
    while (pTableScanInfo->scanTimes < total) {
H
Haojun Liao 已提交
700 701 702
      SSDataBlock* p = doTableScanImpl(pOperator);
      if (p != NULL) {
        return p;
703
      }
H
Haojun Liao 已提交
704

705
      pTableScanInfo->scanTimes += 1;
H
Haojun Liao 已提交
706

707
      if (pTableScanInfo->scanTimes < total) {
708
        setTaskStatus(pTaskInfo, TASK_NOT_COMPLETED);
H
Haojun Liao 已提交
709
        pTableScanInfo->base.scanFlag = REPEAT_SCAN;
H
Haojun Liao 已提交
710

711
        qDebug("%s start to repeat descending order scan data blocks", GET_TASKID(pTaskInfo));
H
Haojun Liao 已提交
712
        tsdbReaderReset(pTableScanInfo->base.dataReader, &pTableScanInfo->base.cond);
713
      }
H
Haojun Liao 已提交
714 715 716
    }
  }

wmmhello's avatar
wmmhello 已提交
717 718 719 720 721 722 723
  return NULL;
}

static SSDataBlock* doTableScan(SOperatorInfo* pOperator) {
  STableScanInfo* pInfo = pOperator->info;
  SExecTaskInfo*  pTaskInfo = pOperator->pTaskInfo;

724
  // scan table one by one sequentially
L
Liu Jicong 已提交
725
  if (pInfo->scanMode == TABLE_SCAN__TABLE_ORDER) {
H
Haojun Liao 已提交
726
    int32_t numOfTables = tableListGetSize(pTaskInfo->pTableInfoList);
H
Haojun Liao 已提交
727

L
Liu Jicong 已提交
728
    while (1) {
H
Haojun Liao 已提交
729
      SSDataBlock* result = doGroupedTableScan(pOperator);
L
Liu Jicong 已提交
730 731 732
      if (result) {
        return result;
      }
H
Haojun Liao 已提交
733

L
Liu Jicong 已提交
734 735
      // if no data, switch to next table and continue scan
      pInfo->currentTable++;
H
Haojun Liao 已提交
736
      if (pInfo->currentTable >= numOfTables) {
L
Liu Jicong 已提交
737 738
        return NULL;
      }
H
Haojun Liao 已提交
739

H
Haojun Liao 已提交
740
      STableKeyInfo* pTableInfo = tableListGetInfo(pTaskInfo->pTableInfoList, pInfo->currentTable);
H
Haojun Liao 已提交
741
      tsdbSetTableList(pInfo->base.dataReader, pTableInfo, 1);
L
Liu Jicong 已提交
742 743
      qDebug("set uid:%" PRIu64 " into scanner, total tables:%d, index:%d %s", pTableInfo->uid, numOfTables,
             pInfo->currentTable, pTaskInfo->id.str);
H
Haojun Liao 已提交
744

H
Haojun Liao 已提交
745
      tsdbReaderReset(pInfo->base.dataReader, &pInfo->base.cond);
L
Liu Jicong 已提交
746 747
      pInfo->scanTimes = 0;
    }
748 749
  } else {  // scan table group by group sequentially
    if (pInfo->currentGroupId == -1) {
H
Haojun Liao 已提交
750
      if ((++pInfo->currentGroupId) >= tableListGetOutputGroups(pTaskInfo->pTableInfoList)) {
H
Haojun Liao 已提交
751
        setOperatorCompleted(pOperator);
752 753
        return NULL;
      }
754

5
54liuyao 已提交
755
      int32_t        num = 0;
756
      STableKeyInfo* pList = NULL;
H
Haojun Liao 已提交
757
      tableListGetGroupList(pTaskInfo->pTableInfoList, pInfo->currentGroupId, &pList, &num);
H
Haojun Liao 已提交
758
      ASSERT(pInfo->base.dataReader == NULL);
759

L
Liu Jicong 已提交
760 761
      int32_t code = tsdbReaderOpen(pInfo->base.readHandle.vnode, &pInfo->base.cond, pList, num, pInfo->pResBlock,
                                    (STsdbReader**)&pInfo->base.dataReader, GET_TASKID(pTaskInfo));
762 763 764
      if (code != TSDB_CODE_SUCCESS) {
        T_LONG_JMP(pTaskInfo->env, code);
      }
wmmhello's avatar
wmmhello 已提交
765
    }
H
Haojun Liao 已提交
766

H
Haojun Liao 已提交
767
    SSDataBlock* result = doGroupedTableScan(pOperator);
768
    if (result != NULL) {
H
Haojun Liao 已提交
769
      ASSERT(result->info.id.uid != 0);
770 771
      return result;
    }
H
Haojun Liao 已提交
772

H
Haojun Liao 已提交
773
    if ((++pInfo->currentGroupId) >= tableListGetOutputGroups(pTaskInfo->pTableInfoList)) {
H
Haojun Liao 已提交
774
      setOperatorCompleted(pOperator);
775 776
      return NULL;
    }
wmmhello's avatar
wmmhello 已提交
777

778 779
    // reset value for the next group data output
    pOperator->status = OP_OPENED;
780
    resetLimitInfoForNextGroup(&pInfo->base.limitInfo);
wmmhello's avatar
wmmhello 已提交
781

5
54liuyao 已提交
782
    int32_t        num = 0;
783
    STableKeyInfo* pList = NULL;
H
Haojun Liao 已提交
784
    tableListGetGroupList(pTaskInfo->pTableInfoList, pInfo->currentGroupId, &pList, &num);
wmmhello's avatar
wmmhello 已提交
785

H
Haojun Liao 已提交
786 787
    tsdbSetTableList(pInfo->base.dataReader, pList, num);
    tsdbReaderReset(pInfo->base.dataReader, &pInfo->base.cond);
788
    pInfo->scanTimes = 0;
wmmhello's avatar
wmmhello 已提交
789

H
Haojun Liao 已提交
790
    result = doGroupedTableScan(pOperator);
791 792 793
    if (result != NULL) {
      return result;
    }
794

H
Haojun Liao 已提交
795
    setOperatorCompleted(pOperator);
796 797
    return NULL;
  }
H
Haojun Liao 已提交
798 799
}

800 801
static int32_t getTableScannerExecInfo(struct SOperatorInfo* pOptr, void** pOptrExplain, uint32_t* len) {
  SFileBlockLoadRecorder* pRecorder = taosMemoryCalloc(1, sizeof(SFileBlockLoadRecorder));
802
  STableScanInfo*         pTableScanInfo = pOptr->info;
H
Haojun Liao 已提交
803
  *pRecorder = pTableScanInfo->base.readRecorder;
804 805 806 807 808
  *pOptrExplain = pRecorder;
  *len = sizeof(SFileBlockLoadRecorder);
  return 0;
}

809
static void destroyTableScanOperatorInfo(void* param) {
810
  STableScanInfo* pTableScanInfo = (STableScanInfo*)param;
H
Haojun Liao 已提交
811
  blockDataDestroy(pTableScanInfo->pResBlock);
H
Haojun Liao 已提交
812
  cleanupQueryTableDataCond(&pTableScanInfo->base.cond);
H
Haojun Liao 已提交
813

H
Haojun Liao 已提交
814 815
  tsdbReaderClose(pTableScanInfo->base.dataReader);
  pTableScanInfo->base.dataReader = NULL;
816

H
Haojun Liao 已提交
817 818
  if (pTableScanInfo->base.matchInfo.pList != NULL) {
    taosArrayDestroy(pTableScanInfo->base.matchInfo.pList);
819
  }
L
Liu Jicong 已提交
820

H
Haojun Liao 已提交
821 822
  taosLRUCacheCleanup(pTableScanInfo->base.metaCache.pTableMetaEntryCache);
  cleanupExprSupp(&pTableScanInfo->base.pseudoSup);
D
dapan1121 已提交
823
  taosMemoryFreeClear(param);
824 825
}

826
SOperatorInfo* createTableScanOperatorInfo(STableScanPhysiNode* pTableScanNode, SReadHandle* readHandle,
827
                                           SExecTaskInfo* pTaskInfo) {
H
Haojun Liao 已提交
828 829 830
  STableScanInfo* pInfo = taosMemoryCalloc(1, sizeof(STableScanInfo));
  SOperatorInfo*  pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
  if (pInfo == NULL || pOperator == NULL) {
831
    goto _error;
H
Haojun Liao 已提交
832 833
  }

834
  SScanPhysiNode*     pScanNode = &pTableScanNode->scan;
H
Haojun Liao 已提交
835
  SDataBlockDescNode* pDescNode = pScanNode->node.pOutputDataBlockDesc;
836 837

  int32_t numOfCols = 0;
838
  int32_t code =
H
Haojun Liao 已提交
839
      extractColMatchInfo(pScanNode->pScanCols, pDescNode, &numOfCols, COL_MATCH_FROM_COL_ID, &pInfo->base.matchInfo);
840 841 842 843
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }

H
Haojun Liao 已提交
844
  initLimitInfo(pScanNode->node.pLimit, pScanNode->node.pSlimit, &pInfo->base.limitInfo);
H
Haojun Liao 已提交
845
  code = initQueryTableDataCond(&pInfo->base.cond, pTableScanNode);
846
  if (code != TSDB_CODE_SUCCESS) {
847
    goto _error;
848 849
  }

H
Haojun Liao 已提交
850
  if (pScanNode->pScanPseudoCols != NULL) {
H
Haojun Liao 已提交
851
    SExprSupp* pSup = &pInfo->base.pseudoSup;
H
Haojun Liao 已提交
852
    pSup->pExprInfo = createExprInfo(pScanNode->pScanPseudoCols, NULL, &pSup->numOfExprs);
853
    pSup->pCtx = createSqlFunctionCtx(pSup->pExprInfo, pSup->numOfExprs, &pSup->rowEntryInfoOffset);
854 855
  }

856
  pInfo->scanInfo = (SScanInfo){.numOfAsc = pTableScanNode->scanSeq[0], .numOfDesc = pTableScanNode->scanSeq[1]};
H
Haojun Liao 已提交
857 858

  pInfo->base.scanFlag = MAIN_SCAN;
H
Haojun Liao 已提交
859 860
  pInfo->base.pdInfo.interval = extractIntervalInfo(pTableScanNode);
  pInfo->base.readHandle = *readHandle;
H
Haojun Liao 已提交
861 862
  pInfo->base.dataBlockLoadFlag = pTableScanNode->dataRequired;

863 864
  pInfo->sample.sampleRatio = pTableScanNode->ratio;
  pInfo->sample.seed = taosGetTimestampSec();
865

H
Haojun Liao 已提交
866
  initResultSizeInfo(&pOperator->resultInfo, 4096);
H
Haojun Liao 已提交
867
  pInfo->pResBlock = createDataBlockFromDescNode(pDescNode);
H
Haojun Liao 已提交
868
  blockDataEnsureCapacity(pInfo->pResBlock, pOperator->resultInfo.capacity);
H
Haojun Liao 已提交
869

H
Haojun Liao 已提交
870 871 872
  code = filterInitFromNode((SNode*)pTableScanNode->scan.node.pConditions, &pOperator->exprSupp.pFilterInfo, 0);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
H
Haojun Liao 已提交
873 874
  }

wmmhello's avatar
wmmhello 已提交
875
  pInfo->currentGroupId = -1;
876
  pInfo->assignBlockUid = pTableScanNode->assignBlockUid;
877
  pInfo->hasGroupByTag = pTableScanNode->pGroupTags ? true : false;
878

L
Liu Jicong 已提交
879 880
  setOperatorInfo(pOperator, "TableScanOperator", QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN, false, OP_NOT_OPENED, pInfo,
                  pTaskInfo);
881
  pOperator->exprSupp.numOfExprs = numOfCols;
882

H
Haojun Liao 已提交
883 884
  pInfo->base.metaCache.pTableMetaEntryCache = taosLRUCacheInit(1024 * 128, -1, .5);
  if (pInfo->base.metaCache.pTableMetaEntryCache == NULL) {
885 886 887
    code = terrno;
    goto _error;
  }
888

H
Haojun Liao 已提交
889
  taosLRUCacheSetStrictCapacity(pInfo->base.metaCache.pTableMetaEntryCache, false);
890 891
  pOperator->fpSet = createOperatorFpSet(optrDummyOpenFn, doTableScan, NULL, destroyTableScanOperatorInfo,
                                         optrDefaultBufFn, getTableScannerExecInfo);
892 893 894

  // for non-blocking operator, the open cost is always 0
  pOperator->cost.openCost = 0;
H
Haojun Liao 已提交
895
  return pOperator;
896

897
_error:
898 899 900
  if (pInfo != NULL) {
    destroyTableScanOperatorInfo(pInfo);
  }
901

902 903
  taosMemoryFreeClear(pOperator);
  pTaskInfo->code = code;
904
  return NULL;
H
Haojun Liao 已提交
905 906
}

907
SOperatorInfo* createTableSeqScanOperatorInfo(void* pReadHandle, SExecTaskInfo* pTaskInfo) {
H
Haojun Liao 已提交
908
  STableScanInfo* pInfo = taosMemoryCalloc(1, sizeof(STableScanInfo));
L
Liu Jicong 已提交
909
  SOperatorInfo*  pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
H
Haojun Liao 已提交
910

H
Haojun Liao 已提交
911
  pInfo->base.dataReader = pReadHandle;
L
Liu Jicong 已提交
912
  //  pInfo->prevGroupId       = -1;
H
Haojun Liao 已提交
913

L
Liu Jicong 已提交
914 915
  setOperatorInfo(pOperator, "TableSeqScanOperator", QUERY_NODE_PHYSICAL_PLAN_TABLE_SEQ_SCAN, false, OP_NOT_OPENED,
                  pInfo, pTaskInfo);
916
  pOperator->fpSet = createOperatorFpSet(optrDummyOpenFn, doTableScanImpl, NULL, NULL, optrDefaultBufFn, NULL);
H
Haojun Liao 已提交
917 918 919
  return pOperator;
}

920
static FORCE_INLINE void doClearBufferedBlocks(SStreamScanInfo* pInfo) {
L
Liu Jicong 已提交
921 922
  taosArrayClear(pInfo->pBlockLists);
  pInfo->validBlockIndex = 0;
H
Haojun Liao 已提交
923 924
}

925
static bool isSessionWindow(SStreamScanInfo* pInfo) {
H
Haojun Liao 已提交
926
  return pInfo->windowSup.parentType == QUERY_NODE_PHYSICAL_PLAN_STREAM_SESSION;
5
54liuyao 已提交
927 928
}

929
static bool isStateWindow(SStreamScanInfo* pInfo) {
930
  return pInfo->windowSup.parentType == QUERY_NODE_PHYSICAL_PLAN_STREAM_STATE;
5
54liuyao 已提交
931
}
5
54liuyao 已提交
932

L
Liu Jicong 已提交
933
static bool isIntervalWindow(SStreamScanInfo* pInfo) {
934 935 936
  return pInfo->windowSup.parentType == QUERY_NODE_PHYSICAL_PLAN_STREAM_INTERVAL ||
         pInfo->windowSup.parentType == QUERY_NODE_PHYSICAL_PLAN_STREAM_SEMI_INTERVAL ||
         pInfo->windowSup.parentType == QUERY_NODE_PHYSICAL_PLAN_STREAM_FINAL_INTERVAL;
5
54liuyao 已提交
937 938 939
}

static bool isSignleIntervalWindow(SStreamScanInfo* pInfo) {
940
  return pInfo->windowSup.parentType == QUERY_NODE_PHYSICAL_PLAN_STREAM_INTERVAL;
L
Liu Jicong 已提交
941 942
}

943 944 945 946
static bool isSlidingWindow(SStreamScanInfo* pInfo) {
  return isIntervalWindow(pInfo) && pInfo->interval.interval != pInfo->interval.sliding;
}

947
static void setGroupId(SStreamScanInfo* pInfo, SSDataBlock* pBlock, int32_t groupColIndex, int32_t rowIndex) {
948 949
  SColumnInfoData* pColInfo = taosArrayGet(pBlock->pDataBlock, groupColIndex);
  uint64_t*        groupCol = (uint64_t*)pColInfo->pData;
950
  ASSERT(rowIndex < pBlock->info.rows);
951
  pInfo->groupId = groupCol[rowIndex];
952 953
}

L
Liu Jicong 已提交
954
void resetTableScanInfo(STableScanInfo* pTableScanInfo, STimeWindow* pWin) {
H
Haojun Liao 已提交
955
  pTableScanInfo->base.cond.twindows = *pWin;
L
Liu Jicong 已提交
956 957
  pTableScanInfo->scanTimes = 0;
  pTableScanInfo->currentGroupId = -1;
H
Haojun Liao 已提交
958 959
  tsdbReaderClose(pTableScanInfo->base.dataReader);
  pTableScanInfo->base.dataReader = NULL;
960 961
}

L
Liu Jicong 已提交
962 963
static SSDataBlock* readPreVersionData(SOperatorInfo* pTableScanOp, uint64_t tbUid, TSKEY startTs, TSKEY endTs,
                                       int64_t maxVersion) {
964
  STableKeyInfo tblInfo = {.uid = tbUid, .groupId = 0};
965

966
  STableScanInfo*     pTableScanInfo = pTableScanOp->info;
H
Haojun Liao 已提交
967
  SQueryTableDataCond cond = pTableScanInfo->base.cond;
968 969 970 971 972 973 974 975 976

  cond.startVersion = -1;
  cond.endVersion = maxVersion;
  cond.twindows = (STimeWindow){.skey = startTs, .ekey = endTs};

  SExecTaskInfo* pTaskInfo = pTableScanOp->pTaskInfo;

  SSDataBlock* pBlock = pTableScanInfo->pResBlock;
  STsdbReader* pReader = NULL;
L
Liu Jicong 已提交
977 978
  int32_t      code = tsdbReaderOpen(pTableScanInfo->base.readHandle.vnode, &cond, &tblInfo, 1, pBlock,
                                     (STsdbReader**)&pReader, GET_TASKID(pTaskInfo));
979 980
  if (code != TSDB_CODE_SUCCESS) {
    terrno = code;
dengyihao's avatar
dengyihao 已提交
981
    T_LONG_JMP(pTaskInfo->env, code);
982 983 984
    return NULL;
  }

H
Haojun Liao 已提交
985
  if (tsdbNextDataBlock(pReader)) {
L
Liu Jicong 已提交
986
    /*SSDataBlock* p = */ tsdbRetrieveDataBlock(pReader, NULL);
H
Haojun Liao 已提交
987
    doSetTagColumnData(&pTableScanInfo->base, pBlock, pTaskInfo, pBlock->info.rows);
H
Haojun Liao 已提交
988
    pBlock->info.id.groupId = getTableGroupId(pTaskInfo->pTableInfoList, pBlock->info.id.uid);
989 990 991 992
  }

  tsdbReaderClose(pReader);
  qDebug("retrieve prev rows:%d, skey:%" PRId64 ", ekey:%" PRId64 " uid:%" PRIu64 ", max ver:%" PRId64
5
54liuyao 已提交
993 994
         ", suid:%" PRIu64,
         pBlock->info.rows, startTs, endTs, tbUid, maxVersion, cond.suid);
995 996

  return pBlock->info.rows > 0 ? pBlock : NULL;
997 998 999 1000 1001 1002 1003 1004 1005 1006 1007
}

static uint64_t getGroupIdByCol(SStreamScanInfo* pInfo, uint64_t uid, TSKEY ts, int64_t maxVersion) {
  SSDataBlock* pPreRes = readPreVersionData(pInfo->pTableScanOp, uid, ts, ts, maxVersion);
  if (!pPreRes || pPreRes->info.rows == 0) {
    return 0;
  }
  ASSERT(pPreRes->info.rows == 1);
  return calGroupIdByData(&pInfo->partitionSup, pInfo->pPartScalarSup, pPreRes, 0);
}

5
54liuyao 已提交
1008
static uint64_t getGroupIdByUid(SStreamScanInfo* pInfo, uint64_t uid) {
H
Haojun Liao 已提交
1009
  return getTableGroupId(pInfo->pTableScanOp->pTaskInfo->pTableInfoList, uid);
1010 1011
}

5
54liuyao 已提交
1012 1013 1014 1015 1016 1017 1018 1019
static uint64_t getGroupIdByData(SStreamScanInfo* pInfo, uint64_t uid, TSKEY ts, int64_t maxVersion) {
  if (pInfo->partitionSup.needCalc) {
    return getGroupIdByCol(pInfo, uid, ts, maxVersion);
  }

  return getGroupIdByUid(pInfo, uid);
}

L
Liu Jicong 已提交
1020
static bool prepareRangeScan(SStreamScanInfo* pInfo, SSDataBlock* pBlock, int32_t* pRowIndex) {
5
54liuyao 已提交
1021 1022 1023
  if (pBlock->info.rows == 0) {
    return false;
  }
L
Liu Jicong 已提交
1024 1025 1026 1027 1028 1029 1030 1031 1032 1033
  if ((*pRowIndex) == pBlock->info.rows) {
    return false;
  }

  ASSERT(taosArrayGetSize(pBlock->pDataBlock) >= 3);
  SColumnInfoData* pStartTsCol = taosArrayGet(pBlock->pDataBlock, START_TS_COLUMN_INDEX);
  TSKEY*           startData = (TSKEY*)pStartTsCol->pData;
  SColumnInfoData* pEndTsCol = taosArrayGet(pBlock->pDataBlock, END_TS_COLUMN_INDEX);
  TSKEY*           endData = (TSKEY*)pEndTsCol->pData;
  STimeWindow      win = {.skey = startData[*pRowIndex], .ekey = endData[*pRowIndex]};
1034 1035 1036
  SColumnInfoData* pGpCol = taosArrayGet(pBlock->pDataBlock, GROUPID_COLUMN_INDEX);
  uint64_t*        gpData = (uint64_t*)pGpCol->pData;
  uint64_t         groupId = gpData[*pRowIndex];
1037 1038 1039 1040 1041 1042

  SColumnInfoData* pCalStartTsCol = taosArrayGet(pBlock->pDataBlock, CALCULATE_START_TS_COLUMN_INDEX);
  TSKEY*           calStartData = (TSKEY*)pCalStartTsCol->pData;
  SColumnInfoData* pCalEndTsCol = taosArrayGet(pBlock->pDataBlock, CALCULATE_END_TS_COLUMN_INDEX);
  TSKEY*           calEndData = (TSKEY*)pCalEndTsCol->pData;

L
Liu Jicong 已提交
1043
  setGroupId(pInfo, pBlock, GROUPID_COLUMN_INDEX, *pRowIndex);
1044 1045 1046 1047
  if (isSlidingWindow(pInfo)) {
    pInfo->updateWin.skey = calStartData[*pRowIndex];
    pInfo->updateWin.ekey = calEndData[*pRowIndex];
  }
L
Liu Jicong 已提交
1048 1049 1050
  (*pRowIndex)++;

  for (; *pRowIndex < pBlock->info.rows; (*pRowIndex)++) {
1051
    if (win.skey == startData[*pRowIndex] && groupId == gpData[*pRowIndex]) {
L
Liu Jicong 已提交
1052 1053 1054
      win.ekey = TMAX(win.ekey, endData[*pRowIndex]);
      continue;
    }
1055
    if (win.skey == endData[*pRowIndex] && groupId == gpData[*pRowIndex]) {
L
Liu Jicong 已提交
1056 1057 1058
      win.skey = TMIN(win.skey, startData[*pRowIndex]);
      continue;
    }
1059 1060
    ASSERT(!(win.skey > startData[*pRowIndex] && win.ekey < endData[*pRowIndex]) ||
           !(isInTimeWindow(&win, startData[*pRowIndex], 0) || isInTimeWindow(&win, endData[*pRowIndex], 0)));
L
Liu Jicong 已提交
1061 1062 1063 1064
    break;
  }

  resetTableScanInfo(pInfo->pTableScanOp->info, &win);
1065
  pInfo->pTableScanOp->status = OP_OPENED;
L
Liu Jicong 已提交
1066 1067 1068
  return true;
}

5
54liuyao 已提交
1069
static STimeWindow getSlidingWindow(TSKEY* startTsCol, TSKEY* endTsCol, uint64_t* gpIdCol, SInterval* pInterval,
1070
                                    SDataBlockInfo* pDataBlockInfo, int32_t* pRowIndex, bool hasGroup) {
H
Haojun Liao 已提交
1071
  SResultRowInfo dumyInfo = {0};
5
54liuyao 已提交
1072
  dumyInfo.cur.pageId = -1;
1073
  STimeWindow win = getActiveTimeWindow(NULL, &dumyInfo, startTsCol[*pRowIndex], pInterval, TSDB_ORDER_ASC);
5
54liuyao 已提交
1074 1075
  STimeWindow endWin = win;
  STimeWindow preWin = win;
5
54liuyao 已提交
1076
  uint64_t    groupId = gpIdCol[*pRowIndex];
H
Haojun Liao 已提交
1077

5
54liuyao 已提交
1078
  while (1) {
1079 1080 1081
    if (hasGroup) {
      (*pRowIndex) += 1;
    } else {
5
54liuyao 已提交
1082
      while ((groupId == gpIdCol[(*pRowIndex)] && startTsCol[*pRowIndex] <= endWin.ekey)) {
5
54liuyao 已提交
1083 1084 1085 1086 1087
        (*pRowIndex) += 1;
        if ((*pRowIndex) == pDataBlockInfo->rows) {
          break;
        }
      }
1088
    }
5
54liuyao 已提交
1089

5
54liuyao 已提交
1090 1091 1092
    do {
      preWin = endWin;
      getNextTimeWindow(pInterval, &endWin, TSDB_ORDER_ASC);
1093
    } while (endTsCol[(*pRowIndex) - 1] >= endWin.skey);
5
54liuyao 已提交
1094
    endWin = preWin;
5
54liuyao 已提交
1095
    if (win.ekey == endWin.ekey || (*pRowIndex) == pDataBlockInfo->rows || groupId != gpIdCol[*pRowIndex]) {
5
54liuyao 已提交
1096 1097 1098 1099 1100 1101
      win.ekey = endWin.ekey;
      return win;
    }
    win.ekey = endWin.ekey;
  }
}
5
54liuyao 已提交
1102

L
Liu Jicong 已提交
1103 1104 1105 1106 1107 1108 1109 1110 1111 1112 1113
static SSDataBlock* doRangeScan(SStreamScanInfo* pInfo, SSDataBlock* pSDB, int32_t tsColIndex, int32_t* pRowIndex) {
  while (1) {
    SSDataBlock* pResult = NULL;
    pResult = doTableScan(pInfo->pTableScanOp);
    if (!pResult && prepareRangeScan(pInfo, pSDB, pRowIndex)) {
      // scan next window data
      pResult = doTableScan(pInfo->pTableScanOp);
    }
    if (!pResult) {
      blockDataCleanup(pSDB);
      *pRowIndex = 0;
5
54liuyao 已提交
1114
      pInfo->updateWin = (STimeWindow){.skey = INT64_MIN, .ekey = INT64_MAX};
H
Hongze Cheng 已提交
1115
      STableScanInfo* pTableScanInfo = pInfo->pTableScanOp->info;
H
Haojun Liao 已提交
1116 1117
      tsdbReaderClose(pTableScanInfo->base.dataReader);
      pTableScanInfo->base.dataReader = NULL;
1118 1119
      return NULL;
    }
L
Liu Jicong 已提交
1120

H
Haojun Liao 已提交
1121
    doFilter(pResult, pInfo->pTableScanOp->exprSupp.pFilterInfo, NULL);
1122 1123 1124 1125
    if (pResult->info.rows == 0) {
      continue;
    }

1126 1127 1128 1129 1130 1131 1132 1133
    if (pInfo->partitionSup.needCalc) {
      SSDataBlock* tmpBlock = createOneDataBlock(pResult, true);
      blockDataCleanup(pResult);
      for (int32_t i = 0; i < tmpBlock->info.rows; i++) {
        if (calGroupIdByData(&pInfo->partitionSup, pInfo->pPartScalarSup, tmpBlock, i) == pInfo->groupId) {
          for (int32_t j = 0; j < pInfo->pTableScanOp->exprSupp.numOfExprs; j++) {
            SColumnInfoData* pSrcCol = taosArrayGet(tmpBlock->pDataBlock, j);
            SColumnInfoData* pDestCol = taosArrayGet(pResult->pDataBlock, j);
L
Liu Jicong 已提交
1134 1135
            bool             isNull = colDataIsNull(pSrcCol, tmpBlock->info.rows, i, NULL);
            char*            pSrcData = colDataGetData(pSrcCol, i);
1136 1137 1138 1139 1140
            colDataAppend(pDestCol, pResult->info.rows, pSrcData, isNull);
          }
          pResult->info.rows++;
        }
      }
H
Haojun Liao 已提交
1141 1142 1143

      blockDataDestroy(tmpBlock);

1144 1145 1146 1147
      if (pResult->info.rows > 0) {
        pResult->info.calWin = pInfo->updateWin;
        return pResult;
      }
H
Haojun Liao 已提交
1148
    } else if (pResult->info.id.groupId == pInfo->groupId) {
5
54liuyao 已提交
1149
      pResult->info.calWin = pInfo->updateWin;
1150
      return pResult;
5
54liuyao 已提交
1151 1152
    }
  }
1153
}
1154

1155
static int32_t generateSessionScanRange(SStreamScanInfo* pInfo, SSDataBlock* pSrcBlock, SSDataBlock* pDestBlock) {
5
54liuyao 已提交
1156
  blockDataCleanup(pDestBlock);
1157 1158
  if (pSrcBlock->info.rows == 0) {
    return TSDB_CODE_SUCCESS;
1159
  }
1160
  int32_t code = blockDataEnsureCapacity(pDestBlock, pSrcBlock->info.rows);
1161
  if (code != TSDB_CODE_SUCCESS) {
1162
    return code;
L
Liu Jicong 已提交
1163
  }
1164 1165
  ASSERT(taosArrayGetSize(pSrcBlock->pDataBlock) >= 3);
  SColumnInfoData* pStartTsCol = taosArrayGet(pSrcBlock->pDataBlock, START_TS_COLUMN_INDEX);
L
Liu Jicong 已提交
1166
  TSKEY*           startData = (TSKEY*)pStartTsCol->pData;
1167
  SColumnInfoData* pEndTsCol = taosArrayGet(pSrcBlock->pDataBlock, END_TS_COLUMN_INDEX);
L
Liu Jicong 已提交
1168
  TSKEY*           endData = (TSKEY*)pEndTsCol->pData;
1169 1170
  SColumnInfoData* pUidCol = taosArrayGet(pSrcBlock->pDataBlock, UID_COLUMN_INDEX);
  uint64_t*        uidCol = (uint64_t*)pUidCol->pData;
L
Liu Jicong 已提交
1171

1172 1173
  SColumnInfoData* pDestStartCol = taosArrayGet(pDestBlock->pDataBlock, START_TS_COLUMN_INDEX);
  SColumnInfoData* pDestEndCol = taosArrayGet(pDestBlock->pDataBlock, END_TS_COLUMN_INDEX);
5
54liuyao 已提交
1174
  SColumnInfoData* pDestUidCol = taosArrayGet(pDestBlock->pDataBlock, UID_COLUMN_INDEX);
1175
  SColumnInfoData* pDestGpCol = taosArrayGet(pDestBlock->pDataBlock, GROUPID_COLUMN_INDEX);
5
54liuyao 已提交
1176 1177
  SColumnInfoData* pDestCalStartTsCol = taosArrayGet(pDestBlock->pDataBlock, CALCULATE_START_TS_COLUMN_INDEX);
  SColumnInfoData* pDestCalEndTsCol = taosArrayGet(pDestBlock->pDataBlock, CALCULATE_END_TS_COLUMN_INDEX);
L
Liu Jicong 已提交
1178
  int64_t          version = pSrcBlock->info.version - 1;
1179
  for (int32_t i = 0; i < pSrcBlock->info.rows; i++) {
1180
    uint64_t groupId = getGroupIdByData(pInfo, uidCol[i], startData[i], version);
L
Liu Jicong 已提交
1181
    // gap must be 0.
5
54liuyao 已提交
1182
    SSessionKey startWin = {0};
1183
    getCurSessionWindow(pInfo->windowSup.pStreamAggSup, startData[i], startData[i], groupId, &startWin);
5
54liuyao 已提交
1184
    if (IS_INVALID_SESSION_WIN_KEY(startWin)) {
L
Liu Jicong 已提交
1185 1186 1187
      // window has been closed.
      continue;
    }
5
54liuyao 已提交
1188 1189 1190 1191 1192 1193
    SSessionKey endWin = {0};
    getCurSessionWindow(pInfo->windowSup.pStreamAggSup, endData[i], endData[i], groupId, &endWin);
    ASSERT(!IS_INVALID_SESSION_WIN_KEY(endWin));
    colDataAppend(pDestStartCol, i, (const char*)&startWin.win.skey, false);
    colDataAppend(pDestEndCol, i, (const char*)&endWin.win.ekey, false);

5
54liuyao 已提交
1194
    colDataAppendNULL(pDestUidCol, i);
L
Liu Jicong 已提交
1195
    colDataAppend(pDestGpCol, i, (const char*)&groupId, false);
5
54liuyao 已提交
1196 1197
    colDataAppendNULL(pDestCalStartTsCol, i);
    colDataAppendNULL(pDestCalEndTsCol, i);
1198
    pDestBlock->info.rows++;
L
Liu Jicong 已提交
1199
  }
1200
  return TSDB_CODE_SUCCESS;
L
Liu Jicong 已提交
1201
}
1202 1203 1204 1205 1206 1207

static int32_t generateIntervalScanRange(SStreamScanInfo* pInfo, SSDataBlock* pSrcBlock, SSDataBlock* pDestBlock) {
  blockDataCleanup(pDestBlock);
  int32_t rows = pSrcBlock->info.rows;
  if (rows == 0) {
    return TSDB_CODE_SUCCESS;
1208
  }
1209

1210 1211
  SColumnInfoData* pSrcStartTsCol = (SColumnInfoData*)taosArrayGet(pSrcBlock->pDataBlock, START_TS_COLUMN_INDEX);
  SColumnInfoData* pSrcEndTsCol = (SColumnInfoData*)taosArrayGet(pSrcBlock->pDataBlock, END_TS_COLUMN_INDEX);
1212 1213
  SColumnInfoData* pSrcUidCol = taosArrayGet(pSrcBlock->pDataBlock, UID_COLUMN_INDEX);
  SColumnInfoData* pSrcGpCol = taosArrayGet(pSrcBlock->pDataBlock, GROUPID_COLUMN_INDEX);
5
54liuyao 已提交
1214

L
Liu Jicong 已提交
1215
  uint64_t* srcUidData = (uint64_t*)pSrcUidCol->pData;
1216
  ASSERT(pSrcStartTsCol->info.type == TSDB_DATA_TYPE_TIMESTAMP);
5
54liuyao 已提交
1217 1218 1219 1220 1221 1222 1223 1224 1225 1226 1227 1228 1229 1230 1231 1232 1233 1234 1235 1236 1237 1238 1239 1240 1241 1242 1243 1244 1245 1246 1247 1248 1249 1250 1251 1252
  TSKEY*  srcStartTsCol = (TSKEY*)pSrcStartTsCol->pData;
  TSKEY*  srcEndTsCol = (TSKEY*)pSrcEndTsCol->pData;
  int64_t version = pSrcBlock->info.version - 1;

  if (pInfo->partitionSup.needCalc && srcStartTsCol[0] != srcEndTsCol[0]) {
    uint64_t     srcUid = srcUidData[0];
    TSKEY        startTs = srcStartTsCol[0];
    TSKEY        endTs = srcEndTsCol[0];
    SSDataBlock* pPreRes = readPreVersionData(pInfo->pTableScanOp, srcUid, startTs, endTs, version);
    printDataBlock(pPreRes, "pre res");
    blockDataCleanup(pSrcBlock);
    int32_t code = blockDataEnsureCapacity(pSrcBlock, pPreRes->info.rows);
    if (code != TSDB_CODE_SUCCESS) {
      return code;
    }

    SColumnInfoData* pTsCol = (SColumnInfoData*)taosArrayGet(pPreRes->pDataBlock, pInfo->primaryTsIndex);
    rows = pPreRes->info.rows;

    for (int32_t i = 0; i < rows; i++) {
      uint64_t groupId = calGroupIdByData(&pInfo->partitionSup, pInfo->pPartScalarSup, pPreRes, i);
      appendOneRowToStreamSpecialBlock(pSrcBlock, ((TSKEY*)pTsCol->pData) + i, ((TSKEY*)pTsCol->pData) + i, &srcUid,
                                       &groupId, NULL);
    }
    printDataBlock(pSrcBlock, "new delete");
  }
  uint64_t* srcGp = (uint64_t*)pSrcGpCol->pData;
  srcStartTsCol = (TSKEY*)pSrcStartTsCol->pData;
  srcEndTsCol = (TSKEY*)pSrcEndTsCol->pData;
  srcUidData = (uint64_t*)pSrcUidCol->pData;

  int32_t code = blockDataEnsureCapacity(pDestBlock, rows);
  if (code != TSDB_CODE_SUCCESS) {
    return code;
  }

1253 1254
  SColumnInfoData* pStartTsCol = taosArrayGet(pDestBlock->pDataBlock, START_TS_COLUMN_INDEX);
  SColumnInfoData* pEndTsCol = taosArrayGet(pDestBlock->pDataBlock, END_TS_COLUMN_INDEX);
1255
  SColumnInfoData* pDeUidCol = taosArrayGet(pDestBlock->pDataBlock, UID_COLUMN_INDEX);
1256 1257 1258
  SColumnInfoData* pGpCol = taosArrayGet(pDestBlock->pDataBlock, GROUPID_COLUMN_INDEX);
  SColumnInfoData* pCalStartTsCol = taosArrayGet(pDestBlock->pDataBlock, CALCULATE_START_TS_COLUMN_INDEX);
  SColumnInfoData* pCalEndTsCol = taosArrayGet(pDestBlock->pDataBlock, CALCULATE_END_TS_COLUMN_INDEX);
1259
  for (int32_t i = 0; i < rows;) {
1260
    uint64_t srcUid = srcUidData[i];
5
54liuyao 已提交
1261 1262 1263 1264 1265
    uint64_t groupId = srcGp[i];
    if (groupId == 0) {
      groupId = getGroupIdByData(pInfo, srcUid, srcStartTsCol[i], version);
    }
    TSKEY calStartTs = srcStartTsCol[i];
1266
    colDataAppend(pCalStartTsCol, pDestBlock->info.rows, (const char*)(&calStartTs), false);
5
54liuyao 已提交
1267
    STimeWindow win = getSlidingWindow(srcStartTsCol, srcEndTsCol, srcGp, &pInfo->interval, &pSrcBlock->info, &i,
1268 1269
                                       pInfo->partitionSup.needCalc);
    TSKEY       calEndTs = srcStartTsCol[i - 1];
1270 1271
    colDataAppend(pCalEndTsCol, pDestBlock->info.rows, (const char*)(&calEndTs), false);
    colDataAppend(pDeUidCol, pDestBlock->info.rows, (const char*)(&srcUid), false);
1272 1273 1274 1275
    colDataAppend(pStartTsCol, pDestBlock->info.rows, (const char*)(&win.skey), false);
    colDataAppend(pEndTsCol, pDestBlock->info.rows, (const char*)(&win.ekey), false);
    colDataAppend(pGpCol, pDestBlock->info.rows, (const char*)(&groupId), false);
    pDestBlock->info.rows++;
5
54liuyao 已提交
1276
  }
1277 1278
  return TSDB_CODE_SUCCESS;
}
1279

1280
static int32_t generateDeleteResultBlock(SStreamScanInfo* pInfo, SSDataBlock* pSrcBlock, SSDataBlock* pDestBlock) {
5
54liuyao 已提交
1281 1282 1283
  blockDataCleanup(pDestBlock);
  int32_t rows = pSrcBlock->info.rows;
  if (rows == 0) {
1284 1285
    return TSDB_CODE_SUCCESS;
  }
5
54liuyao 已提交
1286
  int32_t code = blockDataEnsureCapacity(pDestBlock, rows);
1287 1288 1289 1290
  if (code != TSDB_CODE_SUCCESS) {
    return code;
  }

5
54liuyao 已提交
1291 1292 1293 1294 1295 1296 1297 1298 1299 1300
  SColumnInfoData* pSrcStartTsCol = (SColumnInfoData*)taosArrayGet(pSrcBlock->pDataBlock, START_TS_COLUMN_INDEX);
  SColumnInfoData* pSrcEndTsCol = (SColumnInfoData*)taosArrayGet(pSrcBlock->pDataBlock, END_TS_COLUMN_INDEX);
  SColumnInfoData* pSrcUidCol = taosArrayGet(pSrcBlock->pDataBlock, UID_COLUMN_INDEX);
  uint64_t*        srcUidData = (uint64_t*)pSrcUidCol->pData;
  SColumnInfoData* pSrcGpCol = taosArrayGet(pSrcBlock->pDataBlock, GROUPID_COLUMN_INDEX);
  uint64_t*        srcGp = (uint64_t*)pSrcGpCol->pData;
  ASSERT(pSrcStartTsCol->info.type == TSDB_DATA_TYPE_TIMESTAMP);
  TSKEY*  srcStartTsCol = (TSKEY*)pSrcStartTsCol->pData;
  TSKEY*  srcEndTsCol = (TSKEY*)pSrcEndTsCol->pData;
  int64_t version = pSrcBlock->info.version - 1;
1301
  for (int32_t i = 0; i < pSrcBlock->info.rows; i++) {
5
54liuyao 已提交
1302 1303
    uint64_t srcUid = srcUidData[i];
    uint64_t groupId = srcGp[i];
L
Liu Jicong 已提交
1304
    char*    tbname[VARSTR_HEADER_SIZE + TSDB_TABLE_NAME_LEN] = {0};
5
54liuyao 已提交
1305 1306 1307
    if (groupId == 0) {
      groupId = getGroupIdByData(pInfo, srcUid, srcStartTsCol[i], version);
    }
L
Liu Jicong 已提交
1308
    if (pInfo->tbnameCalSup.pExprInfo) {
1309 1310 1311
      void* parTbname = NULL;
      streamStateGetParName(pInfo->pStreamScanOp->pTaskInfo->streamInfo.pState, groupId, &parTbname);

L
Liu Jicong 已提交
1312 1313
      memcpy(varDataVal(tbname), parTbname, TSDB_TABLE_NAME_LEN);
      varDataSetLen(tbname, strlen(varDataVal(tbname)));
L
Liu Jicong 已提交
1314
      tdbFree(parTbname);
L
Liu Jicong 已提交
1315 1316 1317
    }
    appendOneRowToStreamSpecialBlock(pDestBlock, srcStartTsCol + i, srcEndTsCol + i, srcUidData + i, &groupId,
                                     tbname[0] == 0 ? NULL : tbname);
1318 1319 1320 1321
  }
  return TSDB_CODE_SUCCESS;
}

1322 1323 1324 1325
static int32_t generateScanRange(SStreamScanInfo* pInfo, SSDataBlock* pSrcBlock, SSDataBlock* pDestBlock) {
  int32_t code = TSDB_CODE_SUCCESS;
  if (isIntervalWindow(pInfo)) {
    code = generateIntervalScanRange(pInfo, pSrcBlock, pDestBlock);
1326
  } else if (isSessionWindow(pInfo) || isStateWindow(pInfo)) {
1327
    code = generateSessionScanRange(pInfo, pSrcBlock, pDestBlock);
5
54liuyao 已提交
1328 1329
  } else {
    code = generateDeleteResultBlock(pInfo, pSrcBlock, pDestBlock);
1330
  }
1331
  pDestBlock->info.type = STREAM_CLEAR;
1332
  pDestBlock->info.version = pSrcBlock->info.version;
1333
  pDestBlock->info.dataLoad = 1;
1334 1335 1336 1337
  blockDataUpdateTsWindow(pDestBlock, 0);
  return code;
}

L
Liu Jicong 已提交
1338 1339 1340 1341 1342 1343 1344 1345 1346 1347 1348 1349 1350 1351 1352 1353 1354 1355 1356 1357 1358 1359 1360 1361 1362 1363 1364 1365 1366 1367
#if 0
void calBlockTag(SStreamScanInfo* pInfo, SSDataBlock* pBlock) {
  SExprSupp*    pTagCalSup = &pInfo->tagCalSup;
  SStreamState* pState = pInfo->pStreamScanOp->pTaskInfo->streamInfo.pState;
  if (pTagCalSup == NULL || pTagCalSup->numOfExprs == 0) return;
  if (pBlock == NULL || pBlock->info.rows == 0) return;

  void*   tag = NULL;
  int32_t tagLen = 0;
  if (streamStateGetParTag(pState, pBlock->info.id.groupId, &tag, &tagLen) == 0) {
    pBlock->info.tagLen = tagLen;
    void* pTag = taosMemoryRealloc(pBlock->info.pTag, tagLen);
    if (pTag == NULL) {
      tdbFree(tag);
      taosMemoryFree(pBlock->info.pTag);
      pBlock->info.pTag = NULL;
      pBlock->info.tagLen = 0;
      return;
    }
    pBlock->info.pTag = pTag;
    memcpy(pBlock->info.pTag, tag, tagLen);
    tdbFree(tag);
    return;
  } else {
    pBlock->info.pTag = NULL;
  }
  tdbFree(tag);
}
#endif

5
54liuyao 已提交
1368
static void calBlockTbName(SStreamScanInfo* pInfo, SSDataBlock* pBlock) {
1369 1370
  SExprSupp*    pTbNameCalSup = &pInfo->tbnameCalSup;
  SStreamState* pState = pInfo->pStreamScanOp->pTaskInfo->streamInfo.pState;
5
54liuyao 已提交
1371 1372
  blockDataCleanup(pInfo->pCreateTbRes);
  if (pInfo->tbnameCalSup.numOfExprs == 0 && pInfo->tagCalSup.numOfExprs == 0) {
L
Liu Jicong 已提交
1373
    pBlock->info.parTbName[0] = 0;
L
Liu Jicong 已提交
1374
  } else {
5
54liuyao 已提交
1375 1376
    appendCreateTableRow(pInfo->pStreamScanOp->pTaskInfo->streamInfo.pState, &pInfo->tbnameCalSup, &pInfo->tagCalSup,
                         pBlock->info.id.groupId, pBlock, 0, pInfo->pCreateTbRes);
L
Liu Jicong 已提交
1377 1378 1379
  }
}

1380 1381
void appendOneRowToStreamSpecialBlock(SSDataBlock* pBlock, TSKEY* pStartTs, TSKEY* pEndTs, uint64_t* pUid,
                                      uint64_t* pGp, void* pTbName) {
1382 1383
  SColumnInfoData* pStartTsCol = taosArrayGet(pBlock->pDataBlock, START_TS_COLUMN_INDEX);
  SColumnInfoData* pEndTsCol = taosArrayGet(pBlock->pDataBlock, END_TS_COLUMN_INDEX);
1384 1385
  SColumnInfoData* pUidCol = taosArrayGet(pBlock->pDataBlock, UID_COLUMN_INDEX);
  SColumnInfoData* pGpCol = taosArrayGet(pBlock->pDataBlock, GROUPID_COLUMN_INDEX);
1386 1387
  SColumnInfoData* pCalStartCol = taosArrayGet(pBlock->pDataBlock, CALCULATE_START_TS_COLUMN_INDEX);
  SColumnInfoData* pCalEndCol = taosArrayGet(pBlock->pDataBlock, CALCULATE_END_TS_COLUMN_INDEX);
1388
  SColumnInfoData* pTableCol = taosArrayGet(pBlock->pDataBlock, TABLE_NAME_COLUMN_INDEX);
1389 1390
  colDataAppend(pStartTsCol, pBlock->info.rows, (const char*)pStartTs, false);
  colDataAppend(pEndTsCol, pBlock->info.rows, (const char*)pEndTs, false);
1391 1392
  colDataAppend(pUidCol, pBlock->info.rows, (const char*)pUid, false);
  colDataAppend(pGpCol, pBlock->info.rows, (const char*)pGp, false);
1393 1394
  colDataAppend(pCalStartCol, pBlock->info.rows, (const char*)pStartTs, false);
  colDataAppend(pCalEndCol, pBlock->info.rows, (const char*)pEndTs, false);
1395
  colDataAppend(pTableCol, pBlock->info.rows, (const char*)pTbName, pTbName == NULL);
1396
  pBlock->info.rows++;
5
54liuyao 已提交
1397 1398
}

1399
static void checkUpdateData(SStreamScanInfo* pInfo, bool invertible, SSDataBlock* pBlock, bool out) {
1400 1401
  if (out) {
    blockDataCleanup(pInfo->pUpdateDataRes);
5
54liuyao 已提交
1402
    blockDataEnsureCapacity(pInfo->pUpdateDataRes, pBlock->info.rows * 2);
1403
  }
1404 1405
  SColumnInfoData* pColDataInfo = taosArrayGet(pBlock->pDataBlock, pInfo->primaryTsIndex);
  ASSERT(pColDataInfo->info.type == TSDB_DATA_TYPE_TIMESTAMP);
5
54liuyao 已提交
1406
  TSKEY* tsCol = (TSKEY*)pColDataInfo->pData;
H
Haojun Liao 已提交
1407
  bool   tableInserted = updateInfoIsTableInserted(pInfo->pUpdateInfo, pBlock->info.id.uid);
1408
  for (int32_t rowId = 0; rowId < pBlock->info.rows; rowId++) {
5
54liuyao 已提交
1409 1410
    SResultRowInfo dumyInfo;
    dumyInfo.cur.pageId = -1;
L
Liu Jicong 已提交
1411
    bool        isClosed = false;
5
54liuyao 已提交
1412
    STimeWindow win = {.skey = INT64_MIN, .ekey = INT64_MAX};
L
Liu Jicong 已提交
1413
    if (tableInserted && isOverdue(tsCol[rowId], &pInfo->twAggSup)) {
5
54liuyao 已提交
1414 1415 1416
      win = getActiveTimeWindow(NULL, &dumyInfo, tsCol[rowId], &pInfo->interval, TSDB_ORDER_ASC);
      isClosed = isCloseWindow(&win, &pInfo->twAggSup);
    }
5
54liuyao 已提交
1417
    // must check update info first.
H
Haojun Liao 已提交
1418
    bool update = updateInfoIsUpdated(pInfo->pUpdateInfo, pBlock->info.id.uid, tsCol[rowId]);
L
Liu Jicong 已提交
1419
    bool closedWin = isClosed && isSignleIntervalWindow(pInfo) &&
H
Haojun Liao 已提交
1420
                     isDeletedStreamWindow(&win, pBlock->info.id.groupId,
1421
                                           pInfo->pTableScanOp->pTaskInfo->streamInfo.pState, &pInfo->twAggSup);
L
Liu Jicong 已提交
1422
    if ((update || closedWin) && out) {
L
Liu Jicong 已提交
1423
      qDebug("stream update check not pass, update %d, closedWin %d", update, closedWin);
5
54liuyao 已提交
1424
      uint64_t gpId = 0;
H
Haojun Liao 已提交
1425
      appendOneRowToStreamSpecialBlock(pInfo->pUpdateDataRes, tsCol + rowId, tsCol + rowId, &pBlock->info.id.uid, &gpId,
1426
                                       NULL);
5
54liuyao 已提交
1427 1428
      if (closedWin && pInfo->partitionSup.needCalc) {
        gpId = calGroupIdByData(&pInfo->partitionSup, pInfo->pPartScalarSup, pBlock, rowId);
S
slzhou 已提交
1429 1430
        appendOneRowToStreamSpecialBlock(pInfo->pUpdateDataRes, tsCol + rowId, tsCol + rowId, &pBlock->info.id.uid,
                                         &gpId, NULL);
5
54liuyao 已提交
1431
      }
1432 1433
    }
  }
1434 1435
  if (out && pInfo->pUpdateDataRes->info.rows > 0) {
    pInfo->pUpdateDataRes->info.version = pBlock->info.version;
1436
    pInfo->pUpdateDataRes->info.dataLoad = 1;
1437
    blockDataUpdateTsWindow(pInfo->pUpdateDataRes, 0);
1438
    pInfo->pUpdateDataRes->info.type = pInfo->partitionSup.needCalc ? STREAM_DELETE_DATA : STREAM_CLEAR;
5
54liuyao 已提交
1439 1440
  }
}
L
Liu Jicong 已提交
1441

1442
static int32_t setBlockIntoRes(SStreamScanInfo* pInfo, const SSDataBlock* pBlock, bool filter) {
L
Liu Jicong 已提交
1443 1444
  SDataBlockInfo* pBlockInfo = &pInfo->pRes->info;
  SOperatorInfo*  pOperator = pInfo->pStreamScanOp;
L
Liu Jicong 已提交
1445
  SExecTaskInfo*  pTaskInfo = pOperator->pTaskInfo;
L
Liu Jicong 已提交
1446

1447 1448
  blockDataEnsureCapacity(pInfo->pRes, pBlock->info.rows);

L
Liu Jicong 已提交
1449
  pInfo->pRes->info.rows = pBlock->info.rows;
H
Haojun Liao 已提交
1450
  pInfo->pRes->info.id.uid = pBlock->info.id.uid;
L
Liu Jicong 已提交
1451
  pInfo->pRes->info.type = STREAM_NORMAL;
1452
  pInfo->pRes->info.version = pBlock->info.version;
L
Liu Jicong 已提交
1453

H
Haojun Liao 已提交
1454
  pInfo->pRes->info.id.groupId = getTableGroupId(pTaskInfo->pTableInfoList, pBlock->info.id.uid);
L
Liu Jicong 已提交
1455 1456

  // todo extract method
H
Haojun Liao 已提交
1457 1458 1459
  for (int32_t i = 0; i < taosArrayGetSize(pInfo->matchInfo.pList); ++i) {
    SColMatchItem* pColMatchInfo = taosArrayGet(pInfo->matchInfo.pList, i);
    if (!pColMatchInfo->needOutput) {
L
Liu Jicong 已提交
1460 1461 1462 1463 1464 1465 1466
      continue;
    }

    bool colExists = false;
    for (int32_t j = 0; j < blockDataGetNumOfCols(pBlock); ++j) {
      SColumnInfoData* pResCol = bdGetColumnInfoData(pBlock, j);
      if (pResCol->info.colId == pColMatchInfo->colId) {
H
Haojun Liao 已提交
1467
        SColumnInfoData* pDst = taosArrayGet(pInfo->pRes->pDataBlock, pColMatchInfo->dstSlotId);
1468
        colDataAssign(pDst, pResCol, pBlock->info.rows, &pInfo->pRes->info);
L
Liu Jicong 已提交
1469 1470 1471 1472 1473 1474 1475
        colExists = true;
        break;
      }
    }

    // the required column does not exists in submit block, let's set it to be all null value
    if (!colExists) {
H
Haojun Liao 已提交
1476
      SColumnInfoData* pDst = taosArrayGet(pInfo->pRes->pDataBlock, pColMatchInfo->dstSlotId);
L
Liu Jicong 已提交
1477 1478 1479 1480 1481 1482
      colDataAppendNNULL(pDst, 0, pBlockInfo->rows);
    }
  }

  // currently only the tbname pseudo column
  if (pInfo->numOfPseudoExpr > 0) {
L
Liu Jicong 已提交
1483
    int32_t code = addTagPseudoColumnData(&pInfo->readHandle, pInfo->pPseudoExpr, pInfo->numOfPseudoExpr, pInfo->pRes,
1484
                                          pInfo->pRes->info.rows, GET_TASKID(pTaskInfo), NULL);
K
kailixu 已提交
1485 1486
    // ignore the table not exists error, since this table may have been dropped during the scan procedure.
    if (code != TSDB_CODE_SUCCESS && code != TSDB_CODE_PAR_TABLE_NOT_EXIST) {
L
Liu Jicong 已提交
1487
      blockDataFreeRes((SSDataBlock*)pBlock);
1488
      T_LONG_JMP(pTaskInfo->env, code);
H
Haojun Liao 已提交
1489
    }
K
kailixu 已提交
1490 1491 1492

    // reset the error code.
    terrno = 0;
L
Liu Jicong 已提交
1493 1494
  }

1495
  if (filter) {
H
Haojun Liao 已提交
1496
    doFilter(pInfo->pRes, pOperator->exprSupp.pFilterInfo, NULL);
1497
  }
1498

1499
  pInfo->pRes->info.dataLoad = 1;
L
Liu Jicong 已提交
1500
  blockDataUpdateTsWindow(pInfo->pRes, pInfo->primaryTsIndex);
L
Liu Jicong 已提交
1501
  blockDataFreeRes((SSDataBlock*)pBlock);
L
Liu Jicong 已提交
1502

L
Liu Jicong 已提交
1503
  calBlockTbName(pInfo, pInfo->pRes);
L
Liu Jicong 已提交
1504 1505
  return 0;
}
5
54liuyao 已提交
1506

L
Liu Jicong 已提交
1507
static SSDataBlock* doQueueScan(SOperatorInfo* pOperator) {
1508 1509
  SExecTaskInfo*   pTaskInfo = pOperator->pTaskInfo;
  SStreamScanInfo* pInfo = pOperator->info;
H
Haojun Liao 已提交
1510

L
Liu Jicong 已提交
1511
  qDebug("queue scan called");
L
Liu Jicong 已提交
1512

L
Liu Jicong 已提交
1513
  if (pTaskInfo->streamInfo.submit.msgStr != NULL) {
L
Liu Jicong 已提交
1514 1515
    if (pInfo->tqReader->msg2.msgStr == NULL) {
      /*pInfo->tqReader->pMsg = pTaskInfo->streamInfo.pReq;*/
L
Liu Jicong 已提交
1516

L
Liu Jicong 已提交
1517
      /*const SSubmitReq* pSubmit = pInfo->tqReader->pMsg;*/
L
Liu Jicong 已提交
1518 1519
      /*if (tqReaderSetDataMsg(pInfo->tqReader, pSubmit, 0) < 0) {*/
      /*void* msgStr = pTaskInfo->streamInfo.*/
L
Liu Jicong 已提交
1520
      SPackedData submit = pTaskInfo->streamInfo.submit;
L
Liu Jicong 已提交
1521
      if (tqReaderSetSubmitReq2(pInfo->tqReader, submit.msgStr, submit.msgLen, submit.ver) < 0) {
L
Liu Jicong 已提交
1522
        qError("submit msg messed up when initing stream submit block %p", submit.msgStr);
L
Liu Jicong 已提交
1523
        pInfo->tqReader->msg2 = (SPackedData){0};
L
Liu Jicong 已提交
1524
        pInfo->tqReader->setMsg = 0;
L
Liu Jicong 已提交
1525 1526 1527 1528 1529 1530 1531
        ASSERT(0);
      }
    }

    blockDataCleanup(pInfo->pRes);
    SDataBlockInfo* pBlockInfo = &pInfo->pRes->info;

L
Liu Jicong 已提交
1532
    while (tqNextDataBlock2(pInfo->tqReader)) {
L
Liu Jicong 已提交
1533 1534
      SSDataBlock block = {0};

1535
      int32_t code = tqRetrieveDataBlock2(&block, pInfo->tqReader, NULL);
L
Liu Jicong 已提交
1536 1537 1538 1539 1540

      if (code != TSDB_CODE_SUCCESS || block.info.rows == 0) {
        continue;
      }

1541
      setBlockIntoRes(pInfo, &block, true);
L
Liu Jicong 已提交
1542 1543 1544 1545 1546 1547

      if (pBlockInfo->rows > 0) {
        return pInfo->pRes;
      }
    }

L
Liu Jicong 已提交
1548
    pInfo->tqReader->msg2 = (SPackedData){0};
L
Liu Jicong 已提交
1549
    pInfo->tqReader->setMsg = 0;
L
Liu Jicong 已提交
1550
    pTaskInfo->streamInfo.submit = (SPackedData){0};
L
Liu Jicong 已提交
1551
    return NULL;
L
Liu Jicong 已提交
1552 1553
  }

L
Liu Jicong 已提交
1554 1555 1556
  if (pTaskInfo->streamInfo.prepareStatus.type == TMQ_OFFSET__SNAPSHOT_DATA) {
    SSDataBlock* pResult = doTableScan(pInfo->pTableScanOp);
    if (pResult && pResult->info.rows > 0) {
L
Liu Jicong 已提交
1557 1558
      qDebug("queue scan tsdb return %d rows min:%" PRId64 " max:%" PRId64, pResult->info.rows,
             pResult->info.window.skey, pResult->info.window.ekey);
1559
      pTaskInfo->streamInfo.returned = 1;
L
Liu Jicong 已提交
1560 1561
      return pResult;
    } else {
1562 1563
      if (!pTaskInfo->streamInfo.returned) {
        STableScanInfo* pTSInfo = pInfo->pTableScanOp->info;
H
Haojun Liao 已提交
1564 1565
        tsdbReaderClose(pTSInfo->base.dataReader);
        pTSInfo->base.dataReader = NULL;
1566
        tqOffsetResetToLog(&pTaskInfo->streamInfo.prepareStatus, pTaskInfo->streamInfo.snapshotVer);
1567
        qDebug("queue scan tsdb over, switch to wal ver %" PRId64 "", pTaskInfo->streamInfo.snapshotVer + 1);
1568
        if (tqSeekVer(pInfo->tqReader, pTaskInfo->streamInfo.snapshotVer + 1) < 0) {
1569
          tqOffsetResetToLog(&pTaskInfo->streamInfo.lastStatus, pTaskInfo->streamInfo.snapshotVer);
1570 1571 1572 1573
          return NULL;
        }
        ASSERT(pInfo->tqReader->pWalReader->curVersion == pTaskInfo->streamInfo.snapshotVer + 1);
      } else {
L
Liu Jicong 已提交
1574 1575
        return NULL;
      }
1576 1577 1578
    }
  }

L
Liu Jicong 已提交
1579 1580 1581 1582 1583 1584
  if (pTaskInfo->streamInfo.prepareStatus.type == TMQ_OFFSET__LOG) {
    while (1) {
      SFetchRet ret = {0};
      tqNextBlock(pInfo->tqReader, &ret);
      if (ret.fetchType == FETCH_TYPE__DATA) {
        blockDataCleanup(pInfo->pRes);
1585
        if (setBlockIntoRes(pInfo, &ret.data, true) < 0) {
L
Liu Jicong 已提交
1586 1587 1588
          ASSERT(0);
        }
        if (pInfo->pRes->info.rows > 0) {
L
Liu Jicong 已提交
1589
          pOperator->status = OP_EXEC_RECV;
L
Liu Jicong 已提交
1590
          qDebug("queue scan log return %d rows", pInfo->pRes->info.rows);
L
Liu Jicong 已提交
1591 1592 1593 1594
          return pInfo->pRes;
        }
      } else if (ret.fetchType == FETCH_TYPE__META) {
        ASSERT(0);
L
Liu Jicong 已提交
1595 1596 1597
        //        pTaskInfo->streamInfo.lastStatus = ret.offset;
        //        pTaskInfo->streamInfo.metaBlk = ret.meta;
        //        return NULL;
L
Liu Jicong 已提交
1598 1599
      } else if (ret.fetchType == FETCH_TYPE__NONE ||
                 (ret.fetchType == FETCH_TYPE__SEP && pOperator->status == OP_EXEC_RECV)) {
L
Liu Jicong 已提交
1600
        pTaskInfo->streamInfo.lastStatus = ret.offset;
1601 1602 1603 1604
        ASSERT(pTaskInfo->streamInfo.lastStatus.version >= pTaskInfo->streamInfo.prepareStatus.version);
        ASSERT(pTaskInfo->streamInfo.lastStatus.version + 1 == pInfo->tqReader->pWalReader->curVersion);
        char formatBuf[80];
        tFormatOffset(formatBuf, 80, &ret.offset);
L
Liu Jicong 已提交
1605
        qDebug("queue scan log return null, offset %s", formatBuf);
L
Liu Jicong 已提交
1606
        pOperator->status = OP_OPENED;
L
Liu Jicong 已提交
1607 1608 1609
        return NULL;
      }
    }
L
Liu Jicong 已提交
1610
#if 0
1611
    } else if (pTaskInfo->streamInfo.prepareStatus.type == TMQ_OFFSET__SNAPSHOT_DATA) {
L
Liu Jicong 已提交
1612
    SSDataBlock* pResult = doTableScan(pInfo->pTableScanOp);
L
Liu Jicong 已提交
1613 1614 1615 1616 1617 1618
    if (pResult && pResult->info.rows > 0) {
      qDebug("stream scan tsdb return %d rows", pResult->info.rows);
      return pResult;
    }
    qDebug("stream scan tsdb return null");
    return NULL;
L
Liu Jicong 已提交
1619
#endif
L
Liu Jicong 已提交
1620 1621 1622
  } else {
    ASSERT(0);
    return NULL;
H
Haojun Liao 已提交
1623
  }
L
Liu Jicong 已提交
1624 1625
}

L
Liu Jicong 已提交
1626 1627 1628 1629 1630 1631 1632 1633 1634 1635 1636 1637 1638 1639 1640 1641 1642 1643 1644 1645 1646 1647 1648 1649 1650 1651 1652 1653
static int32_t filterDelBlockByUid(SSDataBlock* pDst, const SSDataBlock* pSrc, SStreamScanInfo* pInfo) {
  STqReader* pReader = pInfo->tqReader;
  int32_t    rows = pSrc->info.rows;
  blockDataEnsureCapacity(pDst, rows);

  SColumnInfoData* pSrcStartCol = taosArrayGet(pSrc->pDataBlock, START_TS_COLUMN_INDEX);
  uint64_t*        startCol = (uint64_t*)pSrcStartCol->pData;
  SColumnInfoData* pSrcEndCol = taosArrayGet(pSrc->pDataBlock, END_TS_COLUMN_INDEX);
  uint64_t*        endCol = (uint64_t*)pSrcEndCol->pData;
  SColumnInfoData* pSrcUidCol = taosArrayGet(pSrc->pDataBlock, UID_COLUMN_INDEX);
  uint64_t*        uidCol = (uint64_t*)pSrcUidCol->pData;

  SColumnInfoData* pDstStartCol = taosArrayGet(pDst->pDataBlock, START_TS_COLUMN_INDEX);
  SColumnInfoData* pDstEndCol = taosArrayGet(pDst->pDataBlock, END_TS_COLUMN_INDEX);
  SColumnInfoData* pDstUidCol = taosArrayGet(pDst->pDataBlock, UID_COLUMN_INDEX);
  int32_t          j = 0;
  for (int32_t i = 0; i < rows; i++) {
    if (taosHashGet(pReader->tbIdHash, &uidCol[i], sizeof(uint64_t))) {
      colDataAppend(pDstStartCol, j, (const char*)&startCol[i], false);
      colDataAppend(pDstEndCol, j, (const char*)&endCol[i], false);
      colDataAppend(pDstUidCol, j, (const char*)&uidCol[i], false);

      colDataAppendNULL(taosArrayGet(pDst->pDataBlock, GROUPID_COLUMN_INDEX), j);
      colDataAppendNULL(taosArrayGet(pDst->pDataBlock, CALCULATE_START_TS_COLUMN_INDEX), j);
      colDataAppendNULL(taosArrayGet(pDst->pDataBlock, CALCULATE_END_TS_COLUMN_INDEX), j);
      j++;
    }
  }
L
Liu Jicong 已提交
1654
  uint32_t cap = pDst->info.capacity;
L
Liu Jicong 已提交
1655 1656
  pDst->info = pSrc->info;
  pDst->info.rows = j;
L
Liu Jicong 已提交
1657
  pDst->info.capacity = cap;
L
Liu Jicong 已提交
1658 1659 1660 1661

  return 0;
}

5
54liuyao 已提交
1662 1663 1664 1665 1666 1667 1668 1669 1670 1671 1672 1673 1674 1675 1676 1677 1678
// for partition by tag
static void setBlockGroupIdByUid(SStreamScanInfo* pInfo, SSDataBlock* pBlock) {
  SColumnInfoData* pStartTsCol = taosArrayGet(pBlock->pDataBlock, START_TS_COLUMN_INDEX);
  TSKEY*           startTsCol = (TSKEY*)pStartTsCol->pData;
  SColumnInfoData* pGpCol = taosArrayGet(pBlock->pDataBlock, GROUPID_COLUMN_INDEX);
  uint64_t*        gpCol = (uint64_t*)pGpCol->pData;
  SColumnInfoData* pUidCol = taosArrayGet(pBlock->pDataBlock, UID_COLUMN_INDEX);
  uint64_t*        uidCol = (uint64_t*)pUidCol->pData;
  int32_t          rows = pBlock->info.rows;
  if (!pInfo->partitionSup.needCalc) {
    for (int32_t i = 0; i < rows; i++) {
      uint64_t groupId = getGroupIdByUid(pInfo, uidCol[i]);
      colDataAppend(pGpCol, i, (const char*)&groupId, false);
    }
  }
}

5
54liuyao 已提交
1679 1680 1681 1682 1683 1684 1685 1686 1687 1688 1689 1690 1691 1692 1693 1694 1695 1696
static void doCheckUpdate(SStreamScanInfo* pInfo, TSKEY endKey) {
  if (pInfo->pUpdateInfo) {
    checkUpdateData(pInfo, true, pInfo->pRes, true);
    pInfo->twAggSup.maxTs = TMAX(pInfo->twAggSup.maxTs, endKey);
    if (pInfo->pUpdateDataRes->info.rows > 0) {
      pInfo->updateResIndex = 0;
      if (pInfo->pUpdateDataRes->info.type == STREAM_CLEAR) {
        pInfo->scanMode = STREAM_SCAN_FROM_UPDATERES;
      } else if (pInfo->pUpdateDataRes->info.type == STREAM_INVERT) {
        pInfo->scanMode = STREAM_SCAN_FROM_RES;
        // return pInfo->pUpdateDataRes;
      } else if (pInfo->pUpdateDataRes->info.type == STREAM_DELETE_DATA) {
        pInfo->scanMode = STREAM_SCAN_FROM_DELETE_DATA;
      }
    }
  }
}

L
Liu Jicong 已提交
1697 1698 1699 1700 1701
static SSDataBlock* doStreamScan(SOperatorInfo* pOperator) {
  // NOTE: this operator does never check if current status is done or not
  SExecTaskInfo*   pTaskInfo = pOperator->pTaskInfo;
  SStreamScanInfo* pInfo = pOperator->info;

L
Liu Jicong 已提交
1702
  qDebug("stream scan called");
H
Haojun Liao 已提交
1703

1704 1705
  if (pTaskInfo->streamInfo.recoverStep == STREAM_RECOVER_STEP__PREPARE1 ||
      pTaskInfo->streamInfo.recoverStep == STREAM_RECOVER_STEP__PREPARE2) {
L
Liu Jicong 已提交
1706
    STableScanInfo* pTSInfo = pInfo->pTableScanOp->info;
H
Haojun Liao 已提交
1707
    memcpy(&pTSInfo->base.cond, &pTaskInfo->streamInfo.tableCond, sizeof(SQueryTableDataCond));
1708
    if (pTaskInfo->streamInfo.recoverStep == STREAM_RECOVER_STEP__PREPARE1) {
H
Haojun Liao 已提交
1709 1710 1711 1712
      pTSInfo->base.cond.startVersion = 0;
      pTSInfo->base.cond.endVersion = pTaskInfo->streamInfo.fillHistoryVer1;
      qDebug("stream recover step 1, from %" PRId64 " to %" PRId64, pTSInfo->base.cond.startVersion,
             pTSInfo->base.cond.endVersion);
1713
    } else {
H
Haojun Liao 已提交
1714 1715 1716 1717
      pTSInfo->base.cond.startVersion = pTaskInfo->streamInfo.fillHistoryVer1 + 1;
      pTSInfo->base.cond.endVersion = pTaskInfo->streamInfo.fillHistoryVer2;
      qDebug("stream recover step 2, from %" PRId64 " to %" PRId64, pTSInfo->base.cond.startVersion,
             pTSInfo->base.cond.endVersion);
1718
    }
L
Liu Jicong 已提交
1719 1720

    /*resetTableScanInfo(pTSInfo, pWin);*/
H
Haojun Liao 已提交
1721 1722
    tsdbReaderClose(pTSInfo->base.dataReader);
    pTSInfo->base.dataReader = NULL;
L
Liu Jicong 已提交
1723
    pInfo->pTableScanOp->status = OP_OPENED;
L
Liu Jicong 已提交
1724

L
Liu Jicong 已提交
1725 1726 1727
    pTSInfo->scanTimes = 0;
    pTSInfo->currentGroupId = -1;
    pTaskInfo->streamInfo.recoverStep = STREAM_RECOVER_STEP__SCAN;
L
Liu Jicong 已提交
1728
    pTaskInfo->streamInfo.recoverScanFinished = false;
L
Liu Jicong 已提交
1729 1730 1731
  }

  if (pTaskInfo->streamInfo.recoverStep == STREAM_RECOVER_STEP__SCAN) {
L
Liu Jicong 已提交
1732 1733 1734 1735 1736
    if (pInfo->blockRecoverContiCnt > 100) {
      pInfo->blockRecoverTotCnt += pInfo->blockRecoverContiCnt;
      pInfo->blockRecoverContiCnt = 0;
      return NULL;
    }
5
54liuyao 已提交
1737 1738 1739 1740 1741 1742 1743 1744 1745 1746 1747 1748 1749

    switch (pInfo->scanMode) {
      case STREAM_SCAN_FROM_RES: {
        pInfo->scanMode = STREAM_SCAN_FROM_READERHANDLE;
        printDataBlock(pInfo->pRecoverRes, "scan recover");
        return pInfo->pRecoverRes;
      } break;
      default:
        break;
    }

    pInfo->pRecoverRes = doTableScan(pInfo->pTableScanOp);
    if (pInfo->pRecoverRes != NULL) {
L
Liu Jicong 已提交
1750
      pInfo->blockRecoverContiCnt++;
5
54liuyao 已提交
1751
      calBlockTbName(pInfo, pInfo->pRecoverRes);
1752
      if (pInfo->pUpdateInfo) {
5
54liuyao 已提交
1753
        TSKEY maxTs = updateInfoFillBlockData(pInfo->pUpdateInfo, pInfo->pRecoverRes, pInfo->primaryTsIndex);
L
Liu Jicong 已提交
1754
        pInfo->twAggSup.maxTs = TMAX(pInfo->twAggSup.maxTs, maxTs);
1755
      }
5
54liuyao 已提交
1756 1757 1758 1759 1760 1761 1762
      if (pInfo->pCreateTbRes->info.rows > 0) {
        pInfo->scanMode = STREAM_SCAN_FROM_RES;
        return pInfo->pCreateTbRes;
      }
      qDebug("stream recover scan get block, rows %d", pInfo->pRecoverRes->info.rows);
      printDataBlock(pInfo->pRecoverRes, "scan recover");
      return pInfo->pRecoverRes;
L
Liu Jicong 已提交
1763 1764
    }
    pTaskInfo->streamInfo.recoverStep = STREAM_RECOVER_STEP__NONE;
L
Liu Jicong 已提交
1765
    STableScanInfo* pTSInfo = pInfo->pTableScanOp->info;
H
Haojun Liao 已提交
1766 1767
    tsdbReaderClose(pTSInfo->base.dataReader);
    pTSInfo->base.dataReader = NULL;
1768

H
Haojun Liao 已提交
1769 1770
    pTSInfo->base.cond.startVersion = -1;
    pTSInfo->base.cond.endVersion = -1;
L
Liu Jicong 已提交
1771

L
Liu Jicong 已提交
1772
    pTaskInfo->streamInfo.recoverScanFinished = true;
L
Liu Jicong 已提交
1773 1774 1775
    return NULL;
  }

5
54liuyao 已提交
1776
  size_t total = taosArrayGetSize(pInfo->pBlockLists);
5
54liuyao 已提交
1777
// TODO: refactor
L
Liu Jicong 已提交
1778
FETCH_NEXT_BLOCK:
L
Liu Jicong 已提交
1779
  if (pInfo->blockType == STREAM_INPUT__DATA_BLOCK) {
H
Haojun Liao 已提交
1780
    if (pInfo->validBlockIndex >= total) {
L
Liu Jicong 已提交
1781
      doClearBufferedBlocks(pInfo);
L
Liu Jicong 已提交
1782
      /*pOperator->status = OP_EXEC_DONE;*/
H
Haojun Liao 已提交
1783 1784 1785
      return NULL;
    }

1786
    int32_t      current = pInfo->validBlockIndex++;
L
Liu Jicong 已提交
1787 1788
    SPackedData* pPacked = taosArrayGet(pInfo->pBlockLists, current);
    SSDataBlock* pBlock = pPacked->pDataBlock;
H
Haojun Liao 已提交
1789 1790
    if (pBlock->info.id.groupId && pBlock->info.parTbName[0]) {
      streamStatePutParName(pTaskInfo->streamInfo.pState, pBlock->info.id.groupId, pBlock->info.parTbName);
1791
    }
1792
    // TODO move into scan
5
54liuyao 已提交
1793 1794
    pBlock->info.calWin.skey = INT64_MIN;
    pBlock->info.calWin.ekey = INT64_MAX;
1795
    pBlock->info.dataLoad = 1;
1796
    blockDataUpdateTsWindow(pBlock, 0);
1797
    switch (pBlock->info.type) {
L
Liu Jicong 已提交
1798 1799 1800
      case STREAM_NORMAL:
      case STREAM_GET_ALL:
        return pBlock;
1801 1802 1803
      case STREAM_RETRIEVE: {
        pInfo->blockType = STREAM_INPUT__DATA_SUBMIT;
        pInfo->scanMode = STREAM_SCAN_FROM_DATAREADER_RETRIEVE;
1804 1805
        copyDataBlock(pInfo->pUpdateRes, pBlock);
        prepareRangeScan(pInfo, pInfo->pUpdateRes, &pInfo->updateResIndex);
1806 1807 1808
        updateInfoAddCloseWindowSBF(pInfo->pUpdateInfo);
      } break;
      case STREAM_DELETE_DATA: {
1809
        printDataBlock(pBlock, "stream scan delete recv");
L
Liu Jicong 已提交
1810
        SSDataBlock* pDelBlock = NULL;
L
Liu Jicong 已提交
1811
        if (pInfo->tqReader) {
L
Liu Jicong 已提交
1812
          pDelBlock = createSpecialDataBlock(STREAM_DELETE_DATA);
L
Liu Jicong 已提交
1813
          filterDelBlockByUid(pDelBlock, pBlock, pInfo);
L
Liu Jicong 已提交
1814 1815
        } else {
          pDelBlock = pBlock;
L
Liu Jicong 已提交
1816
        }
5
54liuyao 已提交
1817 1818
        setBlockGroupIdByUid(pInfo, pDelBlock);
        printDataBlock(pDelBlock, "stream scan delete recv filtered");
5
54liuyao 已提交
1819 1820 1821 1822 1823 1824
        if (pDelBlock->info.rows == 0) {
          if (pInfo->tqReader) {
            blockDataDestroy(pDelBlock);
          }
          goto FETCH_NEXT_BLOCK;
        }
1825
        if (!isIntervalWindow(pInfo) && !isSessionWindow(pInfo) && !isStateWindow(pInfo)) {
L
Liu Jicong 已提交
1826
          generateDeleteResultBlock(pInfo, pDelBlock, pInfo->pDeleteDataRes);
1827
          pInfo->pDeleteDataRes->info.type = STREAM_DELETE_RESULT;
L
Liu Jicong 已提交
1828
          printDataBlock(pDelBlock, "stream scan delete result");
H
Haojun Liao 已提交
1829 1830
          blockDataDestroy(pDelBlock);

L
Liu Jicong 已提交
1831 1832 1833 1834 1835
          if (pInfo->pDeleteDataRes->info.rows > 0) {
            return pInfo->pDeleteDataRes;
          } else {
            goto FETCH_NEXT_BLOCK;
          }
1836 1837 1838
        } else {
          pInfo->blockType = STREAM_INPUT__DATA_SUBMIT;
          pInfo->updateResIndex = 0;
L
Liu Jicong 已提交
1839
          generateScanRange(pInfo, pDelBlock, pInfo->pUpdateRes);
1840 1841 1842
          prepareRangeScan(pInfo, pInfo->pUpdateRes, &pInfo->updateResIndex);
          copyDataBlock(pInfo->pDeleteDataRes, pInfo->pUpdateRes);
          pInfo->pDeleteDataRes->info.type = STREAM_DELETE_DATA;
L
Liu Jicong 已提交
1843 1844 1845 1846
          printDataBlock(pDelBlock, "stream scan delete data");
          if (pInfo->tqReader) {
            blockDataDestroy(pDelBlock);
          }
L
Liu Jicong 已提交
1847
          if (pInfo->pDeleteDataRes->info.rows > 0) {
5
54liuyao 已提交
1848
            pInfo->scanMode = STREAM_SCAN_FROM_DATAREADER_RANGE;
L
Liu Jicong 已提交
1849 1850 1851 1852
            return pInfo->pDeleteDataRes;
          } else {
            goto FETCH_NEXT_BLOCK;
          }
1853
        }
1854 1855 1856
      } break;
      default:
        break;
5
54liuyao 已提交
1857
    }
1858
    // printDataBlock(pBlock, "stream scan recv");
1859
    return pBlock;
L
Liu Jicong 已提交
1860
  } else if (pInfo->blockType == STREAM_INPUT__DATA_SUBMIT) {
L
Liu Jicong 已提交
1861
    qDebug("scan mode %d", pInfo->scanMode);
5
54liuyao 已提交
1862 1863 1864
    switch (pInfo->scanMode) {
      case STREAM_SCAN_FROM_RES: {
        pInfo->scanMode = STREAM_SCAN_FROM_READERHANDLE;
5
54liuyao 已提交
1865 1866 1867 1868
        doCheckUpdate(pInfo, pInfo->pRes->info.window.ekey);
        doFilter(pInfo->pRes, pOperator->exprSupp.pFilterInfo, NULL);
        pInfo->pRes->info.dataLoad = 1;
        blockDataUpdateTsWindow(pInfo->pRes, pInfo->primaryTsIndex);
5
54liuyao 已提交
1869 1870
        return pInfo->pRes;
      } break;
1871
      case STREAM_SCAN_FROM_DELETE_DATA: {
1872 1873 1874 1875 1876 1877 1878
        generateScanRange(pInfo, pInfo->pUpdateDataRes, pInfo->pUpdateRes);
        prepareRangeScan(pInfo, pInfo->pUpdateRes, &pInfo->updateResIndex);
        pInfo->scanMode = STREAM_SCAN_FROM_DATAREADER_RANGE;
        copyDataBlock(pInfo->pDeleteDataRes, pInfo->pUpdateRes);
        pInfo->pDeleteDataRes->info.type = STREAM_DELETE_DATA;
        return pInfo->pDeleteDataRes;
      } break;
5
54liuyao 已提交
1879 1880 1881 1882 1883 1884 1885 1886 1887 1888
      case STREAM_SCAN_FROM_UPDATERES: {
        generateScanRange(pInfo, pInfo->pUpdateDataRes, pInfo->pUpdateRes);
        prepareRangeScan(pInfo, pInfo->pUpdateRes, &pInfo->updateResIndex);
        pInfo->scanMode = STREAM_SCAN_FROM_DATAREADER_RANGE;
        return pInfo->pUpdateRes;
      } break;
      case STREAM_SCAN_FROM_DATAREADER_RANGE:
      case STREAM_SCAN_FROM_DATAREADER_RETRIEVE: {
        SSDataBlock* pSDB = doRangeScan(pInfo, pInfo->pUpdateRes, pInfo->primaryTsIndex, &pInfo->updateResIndex);
        if (pSDB) {
1889
          STableScanInfo* pTableScanInfo = pInfo->pTableScanOp->info;
H
Haojun Liao 已提交
1890 1891
          uint64_t        version = getReaderMaxVersion(pTableScanInfo->base.dataReader);
          updateInfoSetScanRange(pInfo->pUpdateInfo, &pTableScanInfo->base.cond.twindows, pInfo->groupId, version);
5
54liuyao 已提交
1892 1893
          pSDB->info.type = pInfo->scanMode == STREAM_SCAN_FROM_DATAREADER_RANGE ? STREAM_NORMAL : STREAM_PULL_DATA;
          checkUpdateData(pInfo, true, pSDB, false);
1894
          // printDataBlock(pSDB, "stream scan update");
L
Liu Jicong 已提交
1895
          calBlockTbName(pInfo, pSDB);
5
54liuyao 已提交
1896 1897
          return pSDB;
        }
1898
        blockDataCleanup(pInfo->pUpdateDataRes);
5
54liuyao 已提交
1899 1900 1901 1902
        pInfo->scanMode = STREAM_SCAN_FROM_READERHANDLE;
      } break;
      default:
        break;
1903
    }
1904

1905
    SStreamAggSupporter* pSup = pInfo->windowSup.pStreamAggSup;
5
54liuyao 已提交
1906
    if (isStateWindow(pInfo) && pSup->pScanBlock->info.rows > 0) {
1907 1908
      pInfo->scanMode = STREAM_SCAN_FROM_DATAREADER_RANGE;
      pInfo->updateResIndex = 0;
5
54liuyao 已提交
1909 1910
      copyDataBlock(pInfo->pUpdateRes, pSup->pScanBlock);
      blockDataCleanup(pSup->pScanBlock);
1911 1912
      prepareRangeScan(pInfo, pInfo->pUpdateRes, &pInfo->updateResIndex);
      return pInfo->pUpdateRes;
5
54liuyao 已提交
1913
    }
5
54liuyao 已提交
1914

H
Haojun Liao 已提交
1915 1916
    SDataBlockInfo* pBlockInfo = &pInfo->pRes->info;

1917
    int32_t totBlockNum = taosArrayGetSize(pInfo->pBlockLists);
1918

L
Liu Jicong 已提交
1919
  NEXT_SUBMIT_BLK:
1920
    while (1) {
L
Liu Jicong 已提交
1921
      if (pInfo->tqReader->msg2.msgStr == NULL) {
1922
        if (pInfo->validBlockIndex >= totBlockNum) {
5
54liuyao 已提交
1923
          updateInfoDestoryColseWinSBF(pInfo->pUpdateInfo);
L
Liu Jicong 已提交
1924
          doClearBufferedBlocks(pInfo);
L
Liu Jicong 已提交
1925
          qDebug("stream scan return empty, consume block %d", totBlockNum);
1926 1927
          return NULL;
        }
1928

L
Liu Jicong 已提交
1929 1930
        int32_t      current = pInfo->validBlockIndex++;
        SPackedData* pSubmit = taosArrayGet(pInfo->pBlockLists, current);
L
Liu Jicong 已提交
1931
        /*if (tqReaderSetDataMsg(pInfo->tqReader, pSubmit, 0) < 0) {*/
L
Liu Jicong 已提交
1932
        if (tqReaderSetSubmitReq2(pInfo->tqReader, pSubmit->msgStr, pSubmit->msgLen, pSubmit->ver) < 0) {
1933 1934 1935 1936
          qError("submit msg messed up when initing stream submit block %p, current %d, total %d", pSubmit, current,
                 totBlockNum);
          continue;
        }
H
Haojun Liao 已提交
1937 1938
      }

1939 1940
      blockDataCleanup(pInfo->pRes);

L
Liu Jicong 已提交
1941
      while (tqNextDataBlock2(pInfo->tqReader)) {
1942
        SSDataBlock block = {0};
1943

1944
        int32_t code = tqRetrieveDataBlock2(&block, pInfo->tqReader, NULL);
1945 1946 1947 1948 1949

        if (code != TSDB_CODE_SUCCESS || block.info.rows == 0) {
          continue;
        }

1950
        setBlockIntoRes(pInfo, &block, false);
1951

H
Haojun Liao 已提交
1952
        if (updateInfoIgnore(pInfo->pUpdateInfo, &pInfo->pRes->info.window, pInfo->pRes->info.id.groupId,
L
Liu Jicong 已提交
1953
                             pInfo->pRes->info.version)) {
1954 1955 1956 1957 1958
          printDataBlock(pInfo->pRes, "stream scan ignore");
          blockDataCleanup(pInfo->pRes);
          continue;
        }

5
54liuyao 已提交
1959 1960 1961
        if (pInfo->pCreateTbRes->info.rows > 0) {
          pInfo->scanMode = STREAM_SCAN_FROM_RES;
          return pInfo->pCreateTbRes;
1962 1963
        }

5
54liuyao 已提交
1964
        doCheckUpdate(pInfo, pBlockInfo->window.ekey);
H
Haojun Liao 已提交
1965
        doFilter(pInfo->pRes, pOperator->exprSupp.pFilterInfo, NULL);
1966
        pInfo->pRes->info.dataLoad = 1;
1967 1968 1969
        blockDataUpdateTsWindow(pInfo->pRes, pInfo->primaryTsIndex);

        if (pBlockInfo->rows > 0 || pInfo->pUpdateDataRes->info.rows > 0) {
1970 1971 1972
          break;
        }
      }
1973
      if (pBlockInfo->rows > 0 || pInfo->pUpdateDataRes->info.rows > 0) {
5
54liuyao 已提交
1974
        break;
J
jiacy-jcy 已提交
1975
      } else {
1976
        continue;
5
54liuyao 已提交
1977
      }
H
Haojun Liao 已提交
1978 1979 1980 1981
    }

    // record the scan action.
    pInfo->numOfExec++;
1982
    pOperator->resultInfo.totalRows += pBlockInfo->rows;
1983
    // printDataBlock(pInfo->pRes, "stream scan");
H
Haojun Liao 已提交
1984

L
Liu Jicong 已提交
1985
    qDebug("scan rows: %d", pBlockInfo->rows);
L
Liu Jicong 已提交
1986 1987 1988
    if (pBlockInfo->rows > 0) {
      return pInfo->pRes;
    }
1989 1990 1991 1992 1993 1994

    if (pInfo->pUpdateDataRes->info.rows > 0) {
      goto FETCH_NEXT_BLOCK;
    }

    goto NEXT_SUBMIT_BLK;
L
Liu Jicong 已提交
1995 1996 1997
  } else {
    ASSERT(0);
    return NULL;
H
Haojun Liao 已提交
1998 1999 2000
  }
}

H
Haojun Liao 已提交
2001
static SArray* extractTableIdList(const STableListInfo* pTableListInfo) {
2002 2003 2004
  SArray* tableIdList = taosArrayInit(4, sizeof(uint64_t));

  // Transfer the Array of STableKeyInfo into uid list.
H
Haojun Liao 已提交
2005 2006 2007
  size_t size = tableListGetSize(pTableListInfo);
  for (int32_t i = 0; i < size; ++i) {
    STableKeyInfo* pkeyInfo = tableListGetInfo(pTableListInfo, i);
2008 2009 2010 2011 2012 2013
    taosArrayPush(tableIdList, &pkeyInfo->uid);
  }

  return tableIdList;
}

2014
static SSDataBlock* doRawScan(SOperatorInfo* pOperator) {
L
Liu Jicong 已提交
2015 2016
  // NOTE: this operator does never check if current status is done or not
  SExecTaskInfo*      pTaskInfo = pOperator->pTaskInfo;
2017
  SStreamRawScanInfo* pInfo = pOperator->info;
L
Liu Jicong 已提交
2018
  pTaskInfo->streamInfo.metaRsp.metaRspLen = 0;  // use metaRspLen !=0 to judge if data is meta
wmmhello's avatar
wmmhello 已提交
2019
  pTaskInfo->streamInfo.metaRsp.metaRsp = NULL;
2020

wmmhello's avatar
wmmhello 已提交
2021
  qDebug("tmqsnap doRawScan called");
L
Liu Jicong 已提交
2022
  if (pTaskInfo->streamInfo.prepareStatus.type == TMQ_OFFSET__SNAPSHOT_DATA) {
2023
    if (pInfo->dataReader && tsdbNextDataBlock(pInfo->dataReader)) {
wmmhello's avatar
wmmhello 已提交
2024
      if (isTaskKilled(pTaskInfo)) {
2025
        longjmp(pTaskInfo->env, pTaskInfo->code);
wmmhello's avatar
wmmhello 已提交
2026
      }
2027

H
Haojun Liao 已提交
2028 2029
      SSDataBlock* pBlock = tsdbRetrieveDataBlock(pInfo->dataReader, NULL);
      if (pBlock == NULL) {
wmmhello's avatar
wmmhello 已提交
2030
        longjmp(pTaskInfo->env, terrno);
wmmhello's avatar
wmmhello 已提交
2031 2032
      }

H
Haojun Liao 已提交
2033
      qDebug("tmqsnap doRawScan get data uid:%" PRId64 "", pBlock->info.id.uid);
wmmhello's avatar
wmmhello 已提交
2034
      pTaskInfo->streamInfo.lastStatus.type = TMQ_OFFSET__SNAPSHOT_DATA;
H
Haojun Liao 已提交
2035
      pTaskInfo->streamInfo.lastStatus.uid = pBlock->info.id.uid;
wmmhello's avatar
wmmhello 已提交
2036 2037 2038
      pTaskInfo->streamInfo.lastStatus.ts = pBlock->info.window.ekey;
      return pBlock;
    }
wmmhello's avatar
wmmhello 已提交
2039 2040

    SMetaTableInfo mtInfo = getUidfromSnapShot(pInfo->sContext);
L
Liu Jicong 已提交
2041
    if (mtInfo.uid == 0) {  // read snapshot done, change to get data from wal
wmmhello's avatar
wmmhello 已提交
2042 2043
      qDebug("tmqsnap read snapshot done, change to get data from wal");
      pTaskInfo->streamInfo.prepareStatus.uid = mtInfo.uid;
wmmhello's avatar
wmmhello 已提交
2044 2045
      pTaskInfo->streamInfo.lastStatus.type = TMQ_OFFSET__LOG;
      pTaskInfo->streamInfo.lastStatus.version = pInfo->sContext->snapVersion;
L
Liu Jicong 已提交
2046
    } else {
wmmhello's avatar
wmmhello 已提交
2047 2048
      pTaskInfo->streamInfo.prepareStatus.uid = mtInfo.uid;
      pTaskInfo->streamInfo.prepareStatus.ts = INT64_MIN;
2049
      qDebug("tmqsnap change get data uid:%" PRId64 "", mtInfo.uid);
wmmhello's avatar
wmmhello 已提交
2050 2051
      qStreamPrepareScan(pTaskInfo, &pTaskInfo->streamInfo.prepareStatus, pInfo->sContext->subType);
    }
2052
    tDeleteSSchemaWrapper(mtInfo.schema);
wmmhello's avatar
wmmhello 已提交
2053
    qDebug("tmqsnap stream scan tsdb return null");
wmmhello's avatar
wmmhello 已提交
2054
    return NULL;
L
Liu Jicong 已提交
2055 2056 2057 2058 2059 2060 2061
  } else if (pTaskInfo->streamInfo.prepareStatus.type == TMQ_OFFSET__SNAPSHOT_META) {
    SSnapContext* sContext = pInfo->sContext;
    void*         data = NULL;
    int32_t       dataLen = 0;
    int16_t       type = 0;
    int64_t       uid = 0;
    if (getMetafromSnapShot(sContext, &data, &dataLen, &type, &uid) < 0) {
wmmhello's avatar
wmmhello 已提交
2062
      qError("tmqsnap getMetafromSnapShot error");
wmmhello's avatar
wmmhello 已提交
2063
      taosMemoryFreeClear(data);
2064 2065 2066
      return NULL;
    }

L
Liu Jicong 已提交
2067
    if (!sContext->queryMetaOrData) {  // change to get data next poll request
wmmhello's avatar
wmmhello 已提交
2068 2069 2070 2071
      pTaskInfo->streamInfo.lastStatus.type = TMQ_OFFSET__SNAPSHOT_META;
      pTaskInfo->streamInfo.lastStatus.uid = uid;
      pTaskInfo->streamInfo.metaRsp.rspOffset.type = TMQ_OFFSET__SNAPSHOT_DATA;
      pTaskInfo->streamInfo.metaRsp.rspOffset.uid = 0;
wmmhello's avatar
wmmhello 已提交
2072
      pTaskInfo->streamInfo.metaRsp.rspOffset.ts = INT64_MIN;
L
Liu Jicong 已提交
2073
    } else {
wmmhello's avatar
wmmhello 已提交
2074 2075 2076 2077 2078 2079 2080
      pTaskInfo->streamInfo.lastStatus.type = TMQ_OFFSET__SNAPSHOT_META;
      pTaskInfo->streamInfo.lastStatus.uid = uid;
      pTaskInfo->streamInfo.metaRsp.rspOffset = pTaskInfo->streamInfo.lastStatus;
      pTaskInfo->streamInfo.metaRsp.resMsgType = type;
      pTaskInfo->streamInfo.metaRsp.metaRspLen = dataLen;
      pTaskInfo->streamInfo.metaRsp.metaRsp = data;
    }
2081

wmmhello's avatar
wmmhello 已提交
2082
    return NULL;
2083
  }
L
Liu Jicong 已提交
2084 2085 2086 2087 2088 2089 2090 2091 2092 2093 2094 2095 2096 2097 2098 2099 2100 2101 2102 2103 2104 2105 2106 2107 2108 2109 2110 2111 2112 2113 2114 2115 2116 2117 2118 2119 2120 2121
  //  else if (pTaskInfo->streamInfo.prepareStatus.type == TMQ_OFFSET__LOG) {
  //    int64_t fetchVer = pTaskInfo->streamInfo.prepareStatus.version + 1;
  //
  //    while(1){
  //      if (tqFetchLog(pInfo->tqReader->pWalReader, pInfo->sContext->withMeta, &fetchVer, &pInfo->pCkHead) < 0) {
  //        qDebug("tmqsnap tmq poll: consumer log end. offset %" PRId64, fetchVer);
  //        pTaskInfo->streamInfo.lastStatus.version = fetchVer;
  //        pTaskInfo->streamInfo.lastStatus.type = TMQ_OFFSET__LOG;
  //        return NULL;
  //      }
  //      SWalCont* pHead = &pInfo->pCkHead->head;
  //      qDebug("tmqsnap tmq poll: consumer log offset %" PRId64 " msgType %d", fetchVer, pHead->msgType);
  //
  //      if (pHead->msgType == TDMT_VND_SUBMIT) {
  //        SSubmitReq* pCont = (SSubmitReq*)&pHead->body;
  //        tqReaderSetDataMsg(pInfo->tqReader, pCont, 0);
  //        SSDataBlock* block = tqLogScanExec(pInfo->sContext->subType, pInfo->tqReader, pInfo->pFilterOutTbUid,
  //        &pInfo->pRes); if(block){
  //          pTaskInfo->streamInfo.lastStatus.type = TMQ_OFFSET__LOG;
  //          pTaskInfo->streamInfo.lastStatus.version = fetchVer;
  //          qDebug("tmqsnap fetch data msg, ver:%" PRId64 ", type:%d", pHead->version, pHead->msgType);
  //          return block;
  //        }else{
  //          fetchVer++;
  //        }
  //      } else{
  //        ASSERT(pInfo->sContext->withMeta);
  //        ASSERT(IS_META_MSG(pHead->msgType));
  //        qDebug("tmqsnap fetch meta msg, ver:%" PRId64 ", type:%d", pHead->version, pHead->msgType);
  //        pTaskInfo->streamInfo.metaRsp.rspOffset.version = fetchVer;
  //        pTaskInfo->streamInfo.metaRsp.rspOffset.type = TMQ_OFFSET__LOG;
  //        pTaskInfo->streamInfo.metaRsp.resMsgType = pHead->msgType;
  //        pTaskInfo->streamInfo.metaRsp.metaRspLen = pHead->bodyLen;
  //        pTaskInfo->streamInfo.metaRsp.metaRsp = taosMemoryMalloc(pHead->bodyLen);
  //        memcpy(pTaskInfo->streamInfo.metaRsp.metaRsp, pHead->body, pHead->bodyLen);
  //        return NULL;
  //      }
  //    }
2122 2123 2124
  return NULL;
}

wmmhello's avatar
wmmhello 已提交
2125
static void destroyRawScanOperatorInfo(void* param) {
wmmhello's avatar
wmmhello 已提交
2126 2127 2128 2129 2130 2131
  SStreamRawScanInfo* pRawScan = (SStreamRawScanInfo*)param;
  tsdbReaderClose(pRawScan->dataReader);
  destroySnapContext(pRawScan->sContext);
  taosMemoryFree(pRawScan);
}

L
Liu Jicong 已提交
2132 2133 2134
// for subscribing db or stb (not including column),
// if this scan is used, meta data can be return
// and schemas are decided when scanning
2135
SOperatorInfo* createRawScanOperatorInfo(SReadHandle* pHandle, SExecTaskInfo* pTaskInfo) {
L
Liu Jicong 已提交
2136 2137 2138 2139 2140
  // create operator
  // create tb reader
  // create meta reader
  // create tq reader

H
Haojun Liao 已提交
2141 2142
  int32_t code = TSDB_CODE_SUCCESS;

2143
  SStreamRawScanInfo* pInfo = taosMemoryCalloc(1, sizeof(SStreamRawScanInfo));
L
Liu Jicong 已提交
2144
  SOperatorInfo*      pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
2145
  if (pInfo == NULL || pOperator == NULL) {
H
Haojun Liao 已提交
2146 2147
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto _end;
2148 2149
  }

wmmhello's avatar
wmmhello 已提交
2150 2151
  pInfo->vnode = pHandle->vnode;

2152
  pInfo->sContext = pHandle->sContext;
L
Liu Jicong 已提交
2153 2154
  setOperatorInfo(pOperator, "RawScanOperator", QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN, false, OP_NOT_OPENED, pInfo,
                  pTaskInfo);
2155

2156
  pOperator->fpSet = createOperatorFpSet(NULL, doRawScan, NULL, destroyRawScanOperatorInfo, optrDefaultBufFn, NULL);
2157
  return pOperator;
H
Haojun Liao 已提交
2158

L
Liu Jicong 已提交
2159
_end:
H
Haojun Liao 已提交
2160 2161 2162 2163
  taosMemoryFree(pInfo);
  taosMemoryFree(pOperator);
  pTaskInfo->code = code;
  return NULL;
L
Liu Jicong 已提交
2164 2165
}

2166
static void destroyStreamScanOperatorInfo(void* param) {
2167 2168
  SStreamScanInfo* pStreamScan = (SStreamScanInfo*)param;
  if (pStreamScan->pTableScanOp && pStreamScan->pTableScanOp->info) {
5
54liuyao 已提交
2169
    destroyOperatorInfo(pStreamScan->pTableScanOp);
2170 2171 2172 2173
  }
  if (pStreamScan->tqReader) {
    tqCloseReader(pStreamScan->tqReader);
  }
H
Haojun Liao 已提交
2174 2175
  if (pStreamScan->matchInfo.pList) {
    taosArrayDestroy(pStreamScan->matchInfo.pList);
2176
  }
C
Cary Xu 已提交
2177 2178
  if (pStreamScan->pPseudoExpr) {
    destroyExprInfo(pStreamScan->pPseudoExpr, pStreamScan->numOfPseudoExpr);
L
Liu Jicong 已提交
2179
    taosMemoryFree(pStreamScan->pPseudoExpr);
C
Cary Xu 已提交
2180
  }
C
Cary Xu 已提交
2181

L
Liu Jicong 已提交
2182
  cleanupExprSupp(&pStreamScan->tbnameCalSup);
5
54liuyao 已提交
2183
  cleanupExprSupp(&pStreamScan->tagCalSup);
L
Liu Jicong 已提交
2184

L
Liu Jicong 已提交
2185
  updateInfoDestroy(pStreamScan->pUpdateInfo);
2186 2187 2188 2189
  blockDataDestroy(pStreamScan->pRes);
  blockDataDestroy(pStreamScan->pUpdateRes);
  blockDataDestroy(pStreamScan->pPullDataRes);
  blockDataDestroy(pStreamScan->pDeleteDataRes);
5
54liuyao 已提交
2190
  blockDataDestroy(pStreamScan->pUpdateDataRes);
5
54liuyao 已提交
2191
  blockDataDestroy(pStreamScan->pCreateTbRes);
2192 2193 2194 2195
  taosArrayDestroy(pStreamScan->pBlockLists);
  taosMemoryFree(pStreamScan);
}

2196
SOperatorInfo* createStreamScanOperatorInfo(SReadHandle* pHandle, STableScanPhysiNode* pTableScanNode, SNode* pTagCond,
2197
                                            SExecTaskInfo* pTaskInfo) {
H
Haojun Liao 已提交
2198
  SArray*          pColIds = NULL;
2199 2200
  SStreamScanInfo* pInfo = taosMemoryCalloc(1, sizeof(SStreamScanInfo));
  SOperatorInfo*   pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
2201

H
Haojun Liao 已提交
2202
  if (pInfo == NULL || pOperator == NULL) {
S
Shengliang Guan 已提交
2203
    terrno = TSDB_CODE_OUT_OF_MEMORY;
2204
    goto _error;
H
Haojun Liao 已提交
2205 2206
  }

2207
  SScanPhysiNode*     pScanPhyNode = &pTableScanNode->scan;
2208
  SDataBlockDescNode* pDescNode = pScanPhyNode->node.pOutputDataBlockDesc;
H
Haojun Liao 已提交
2209

2210
  pInfo->pTagCond = pTagCond;
2211
  pInfo->pGroupTags = pTableScanNode->pGroupTags;
2212

2213
  int32_t numOfCols = 0;
2214 2215
  int32_t code =
      extractColMatchInfo(pScanPhyNode->pScanCols, pDescNode, &numOfCols, COL_MATCH_FROM_COL_ID, &pInfo->matchInfo);
H
Haojun Liao 已提交
2216 2217 2218
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
2219

H
Haojun Liao 已提交
2220
  int32_t numOfOutput = taosArrayGetSize(pInfo->matchInfo.pList);
H
Haojun Liao 已提交
2221
  pColIds = taosArrayInit(numOfOutput, sizeof(int16_t));
2222
  for (int32_t i = 0; i < numOfOutput; ++i) {
H
Haojun Liao 已提交
2223
    SColMatchItem* id = taosArrayGet(pInfo->matchInfo.pList, i);
2224 2225

    int16_t colId = id->colId;
2226
    taosArrayPush(pColIds, &colId);
2227
    if (id->colId == PRIMARYKEY_TIMESTAMP_COL_ID) {
H
Haojun Liao 已提交
2228
      pInfo->primaryTsIndex = id->dstSlotId;
5
54liuyao 已提交
2229
    }
H
Haojun Liao 已提交
2230 2231
  }

L
Liu Jicong 已提交
2232 2233 2234 2235 2236 2237 2238 2239 2240 2241 2242 2243 2244
  if (pTableScanNode->pSubtable != NULL) {
    SExprInfo* pSubTableExpr = taosMemoryCalloc(1, sizeof(SExprInfo));
    if (pSubTableExpr == NULL) {
      terrno = TSDB_CODE_OUT_OF_MEMORY;
      goto _error;
    }
    pInfo->tbnameCalSup.pExprInfo = pSubTableExpr;
    createExprFromOneNode(pSubTableExpr, pTableScanNode->pSubtable, 0);
    if (initExprSupp(&pInfo->tbnameCalSup, pSubTableExpr, 1) != 0) {
      goto _error;
    }
  }

2245 2246
  if (pTableScanNode->pTags != NULL) {
    int32_t    numOfTags;
5
54liuyao 已提交
2247
    SExprInfo* pTagExpr = createExpr(pTableScanNode->pTags, &numOfTags);
2248 2249 2250 2251 2252 2253 2254 2255 2256 2257
    if (pTagExpr == NULL) {
      terrno = TSDB_CODE_OUT_OF_MEMORY;
      goto _error;
    }
    if (initExprSupp(&pInfo->tagCalSup, pTagExpr, numOfTags) != 0) {
      terrno = TSDB_CODE_OUT_OF_MEMORY;
      goto _error;
    }
  }

L
Liu Jicong 已提交
2258
  pInfo->pBlockLists = taosArrayInit(4, sizeof(SPackedData));
H
Haojun Liao 已提交
2259
  if (pInfo->pBlockLists == NULL) {
2260 2261
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    goto _error;
H
Haojun Liao 已提交
2262 2263
  }

5
54liuyao 已提交
2264
  if (pHandle->vnode) {
L
Liu Jicong 已提交
2265
    SOperatorInfo*  pTableScanOp = createTableScanOperatorInfo(pTableScanNode, pHandle, pTaskInfo);
L
Liu Jicong 已提交
2266
    STableScanInfo* pTSInfo = (STableScanInfo*)pTableScanOp->info;
2267
    if (pHandle->version > 0) {
H
Haojun Liao 已提交
2268
      pTSInfo->base.cond.endVersion = pHandle->version;
2269
    }
L
Liu Jicong 已提交
2270

2271
    STableKeyInfo* pList = NULL;
5
54liuyao 已提交
2272
    int32_t        num = 0;
H
Haojun Liao 已提交
2273
    tableListGetGroupList(pTaskInfo->pTableInfoList, 0, &pList, &num);
2274

2275
    if (pHandle->initTableReader) {
L
Liu Jicong 已提交
2276
      pTSInfo->scanMode = TABLE_SCAN__TABLE_ORDER;
H
Haojun Liao 已提交
2277
      pTSInfo->base.dataReader = NULL;
L
Liu Jicong 已提交
2278 2279
      code = tsdbReaderOpen(pHandle->vnode, &pTSInfo->base.cond, pList, num, pTSInfo->pResBlock,
                            &pTSInfo->base.dataReader, NULL);
dengyihao's avatar
dengyihao 已提交
2280 2281
      if (code != 0) {
        terrno = code;
H
Haojun Liao 已提交
2282
        destroyTableScanOperatorInfo(pTableScanOp);
2283
        goto _error;
L
Liu Jicong 已提交
2284
      }
L
Liu Jicong 已提交
2285 2286
    }

L
Liu Jicong 已提交
2287 2288 2289 2290
    if (pHandle->initTqReader) {
      ASSERT(pHandle->tqReader == NULL);
      pInfo->tqReader = tqOpenReader(pHandle->vnode);
      ASSERT(pInfo->tqReader);
2291
    } else {
L
Liu Jicong 已提交
2292 2293
      ASSERT(pHandle->tqReader);
      pInfo->tqReader = pHandle->tqReader;
2294 2295
    }

2296
    pInfo->pUpdateInfo = NULL;
2297
    pInfo->pTableScanOp = pTableScanOp;
2298 2299 2300
    if (pInfo->pTableScanOp->pTaskInfo->streamInfo.pState) {
      streamStateSetNumber(pInfo->pTableScanOp->pTaskInfo->streamInfo.pState, -1);
    }
L
Liu Jicong 已提交
2301

L
Liu Jicong 已提交
2302 2303
    pInfo->readHandle = *pHandle;
    pInfo->tableUid = pScanPhyNode->uid;
L
Liu Jicong 已提交
2304
    pTaskInfo->streamInfo.snapshotVer = pHandle->version;
5
54liuyao 已提交
2305 2306
    pInfo->pCreateTbRes = buildCreateTableBlock(&pInfo->tbnameCalSup, &pInfo->tagCalSup);
    blockDataEnsureCapacity(pInfo->pCreateTbRes, 8);
L
Liu Jicong 已提交
2307

L
Liu Jicong 已提交
2308
    // set the extract column id to streamHandle
L
Liu Jicong 已提交
2309
    tqReaderSetColIdList(pInfo->tqReader, pColIds);
H
Haojun Liao 已提交
2310
    SArray* tableIdList = extractTableIdList(pTaskInfo->pTableInfoList);
2311
    code = tqReaderSetTbUidList(pInfo->tqReader, tableIdList);
L
Liu Jicong 已提交
2312 2313 2314 2315 2316
    if (code != 0) {
      taosArrayDestroy(tableIdList);
      goto _error;
    }
    taosArrayDestroy(tableIdList);
H
Haojun Liao 已提交
2317
    memcpy(&pTaskInfo->streamInfo.tableCond, &pTSInfo->base.cond, sizeof(SQueryTableDataCond));
L
Liu Jicong 已提交
2318 2319
  } else {
    taosArrayDestroy(pColIds);
H
Haojun Liao 已提交
2320
    pColIds = NULL;
5
54liuyao 已提交
2321 2322
  }

2323 2324 2325 2326 2327
  // create the pseduo columns info
  if (pTableScanNode->scan.pScanPseudoCols != NULL) {
    pInfo->pPseudoExpr = createExprInfo(pTableScanNode->scan.pScanPseudoCols, NULL, &pInfo->numOfPseudoExpr);
  }

H
Haojun Liao 已提交
2328 2329 2330 2331 2332
  code = filterInitFromNode((SNode*)pScanPhyNode->node.pConditions, &pOperator->exprSupp.pFilterInfo, 0);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }

H
Haojun Liao 已提交
2333
  pInfo->pRes = createDataBlockFromDescNode(pDescNode);
2334
  pInfo->pUpdateRes = createSpecialDataBlock(STREAM_CLEAR);
2335
  pInfo->scanMode = STREAM_SCAN_FROM_READERHANDLE;
L
Liu Jicong 已提交
2336
  pInfo->windowSup = (SWindowSupporter){.pStreamAggSup = NULL, .gap = -1, .parentType = QUERY_NODE_PHYSICAL_PLAN};
2337
  pInfo->groupId = 0;
2338
  pInfo->pPullDataRes = createSpecialDataBlock(STREAM_RETRIEVE);
2339
  pInfo->pStreamScanOp = pOperator;
2340
  pInfo->deleteDataIndex = 0;
2341
  pInfo->pDeleteDataRes = createSpecialDataBlock(STREAM_DELETE_DATA);
5
54liuyao 已提交
2342
  pInfo->updateWin = (STimeWindow){.skey = INT64_MAX, .ekey = INT64_MAX};
2343
  pInfo->pUpdateDataRes = createSpecialDataBlock(STREAM_CLEAR);
X
Xiaoyu Wang 已提交
2344
  pInfo->assignBlockUid = pTableScanNode->assignBlockUid;
2345
  pInfo->partitionSup.needCalc = false;
L
Liu Jicong 已提交
2346

L
Liu Jicong 已提交
2347 2348
  setOperatorInfo(pOperator, "StreamScanOperator", QUERY_NODE_PHYSICAL_PLAN_STREAM_SCAN, false, OP_NOT_OPENED, pInfo,
                  pTaskInfo);
2349
  pOperator->exprSupp.numOfExprs = taosArrayGetSize(pInfo->pRes->pDataBlock);
H
Haojun Liao 已提交
2350

L
Liu Jicong 已提交
2351
  __optr_fn_t nextFn = pTaskInfo->execModel == OPTR_EXEC_MODEL_STREAM ? doStreamScan : doQueueScan;
L
Liu Jicong 已提交
2352 2353
  pOperator->fpSet =
      createOperatorFpSet(optrDummyOpenFn, nextFn, NULL, destroyStreamScanOperatorInfo, optrDefaultBufFn, NULL);
2354

H
Haojun Liao 已提交
2355
  return pOperator;
2356

L
Liu Jicong 已提交
2357
_error:
H
Haojun Liao 已提交
2358 2359 2360 2361 2362 2363 2364 2365
  if (pColIds != NULL) {
    taosArrayDestroy(pColIds);
  }

  if (pInfo != NULL) {
    destroyStreamScanOperatorInfo(pInfo);
  }

2366 2367
  taosMemoryFreeClear(pOperator);
  return NULL;
H
Haojun Liao 已提交
2368 2369
}

2370
static SSDataBlock* doTagScan(SOperatorInfo* pOperator) {
H
Haojun Liao 已提交
2371 2372 2373 2374
  if (pOperator->status == OP_EXEC_DONE) {
    return NULL;
  }

2375 2376 2377
  SExecTaskInfo* pTaskInfo = pOperator->pTaskInfo;

  STagScanInfo* pInfo = pOperator->info;
2378
  SExprInfo*    pExprInfo = &pOperator->exprSupp.pExprInfo[0];
2379
  SSDataBlock*  pRes = pInfo->pRes;
2380
  blockDataCleanup(pRes);
H
Haojun Liao 已提交
2381

H
Haojun Liao 已提交
2382
  int32_t size = tableListGetSize(pTaskInfo->pTableInfoList);
wmmhello's avatar
wmmhello 已提交
2383
  if (size == 0) {
H
Haojun Liao 已提交
2384 2385 2386 2387
    setTaskStatus(pTaskInfo, TASK_COMPLETED);
    return NULL;
  }

2388 2389 2390
  char        str[512] = {0};
  int32_t     count = 0;
  SMetaReader mr = {0};
2391
  metaReaderInit(&mr, pInfo->readHandle.meta, 0);
H
Haojun Liao 已提交
2392

wmmhello's avatar
wmmhello 已提交
2393
  while (pInfo->curPos < size && count < pOperator->resultInfo.capacity) {
H
Haojun Liao 已提交
2394
    STableKeyInfo* item = tableListGetInfo(pTaskInfo->pTableInfoList, pInfo->curPos);
L
Liu Jicong 已提交
2395
    int32_t        code = metaGetTableEntryByUid(&mr, item->uid);
2396
    tDecoderClear(&mr.coder);
H
Haojun Liao 已提交
2397
    if (code != TSDB_CODE_SUCCESS) {
L
Liu Jicong 已提交
2398 2399
      qError("failed to get table meta, uid:0x%" PRIx64 ", code:%s, %s", item->uid, tstrerror(terrno),
             GET_TASKID(pTaskInfo));
H
Haojun Liao 已提交
2400
      metaReaderClear(&mr);
2401
      T_LONG_JMP(pTaskInfo->env, terrno);
H
Haojun Liao 已提交
2402
    }
H
Haojun Liao 已提交
2403

2404
    for (int32_t j = 0; j < pOperator->exprSupp.numOfExprs; ++j) {
2405 2406 2407 2408 2409 2410
      SColumnInfoData* pDst = taosArrayGet(pRes->pDataBlock, pExprInfo[j].base.resSchema.slotId);

      // refactor later
      if (fmIsScanPseudoColumnFunc(pExprInfo[j].pExpr->_function.functionId)) {
        STR_TO_VARSTR(str, mr.me.name);
        colDataAppend(pDst, count, str, false);
2411
      } else {  // it is a tag value
wmmhello's avatar
wmmhello 已提交
2412 2413
        STagVal val = {0};
        val.cid = pExprInfo[j].base.pParam[0].pCol->colId;
2414
        const char* p = metaGetTableTagVal(mr.me.ctbEntry.pTags, pDst->info.type, &val);
wmmhello's avatar
wmmhello 已提交
2415

2416 2417 2418 2419
        char* data = NULL;
        if (pDst->info.type != TSDB_DATA_TYPE_JSON && p != NULL) {
          data = tTagValToData((const STagVal*)p, false);
        } else {
wmmhello's avatar
wmmhello 已提交
2420 2421
          data = (char*)p;
        }
L
Liu Jicong 已提交
2422 2423
        colDataAppend(pDst, count, data,
                      (data == NULL) || (pDst->info.type == TSDB_DATA_TYPE_JSON && tTagIsJsonNull(data)));
2424

2425 2426
        if (pDst->info.type != TSDB_DATA_TYPE_JSON && p != NULL && IS_VAR_DATA_TYPE(((const STagVal*)p)->type) &&
            data != NULL) {
wmmhello's avatar
wmmhello 已提交
2427
          taosMemoryFree(data);
wmmhello's avatar
wmmhello 已提交
2428
        }
H
Haojun Liao 已提交
2429 2430 2431
      }
    }

2432
    count += 1;
wmmhello's avatar
wmmhello 已提交
2433
    if (++pInfo->curPos >= size) {
H
Haojun Liao 已提交
2434
      setOperatorCompleted(pOperator);
H
Haojun Liao 已提交
2435 2436 2437
    }
  }

2438 2439
  metaReaderClear(&mr);

2440
  // qDebug("QInfo:0x%"PRIx64" create tag values results completed, rows:%d", GET_TASKID(pRuntimeEnv), count);
H
Haojun Liao 已提交
2441
  if (pOperator->status == OP_EXEC_DONE) {
2442
    setTaskStatus(pTaskInfo, TASK_COMPLETED);
H
Haojun Liao 已提交
2443 2444 2445
  }

  pRes->info.rows = count;
wmmhello's avatar
wmmhello 已提交
2446
  pOperator->resultInfo.totalRows += count;
2447

2448
  return (pRes->info.rows == 0) ? NULL : pInfo->pRes;
H
Haojun Liao 已提交
2449 2450
}

2451
static void destroyTagScanOperatorInfo(void* param) {
H
Haojun Liao 已提交
2452 2453
  STagScanInfo* pInfo = (STagScanInfo*)param;
  pInfo->pRes = blockDataDestroy(pInfo->pRes);
H
Haojun Liao 已提交
2454
  taosArrayDestroy(pInfo->matchInfo.pList);
D
dapan1121 已提交
2455
  taosMemoryFreeClear(param);
H
Haojun Liao 已提交
2456 2457
}

S
slzhou 已提交
2458 2459
SOperatorInfo* createTagScanOperatorInfo(SReadHandle* pReadHandle, STagScanPhysiNode* pPhyNode,
                                         SExecTaskInfo* pTaskInfo) {
2460
  STagScanInfo*  pInfo = taosMemoryCalloc(1, sizeof(STagScanInfo));
H
Haojun Liao 已提交
2461 2462 2463 2464 2465
  SOperatorInfo* pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
  if (pInfo == NULL || pOperator == NULL) {
    goto _error;
  }

2466 2467 2468 2469
  SDataBlockDescNode* pDescNode = pPhyNode->node.pOutputDataBlockDesc;

  int32_t    numOfExprs = 0;
  SExprInfo* pExprInfo = createExprInfo(pPhyNode->pScanPseudoCols, NULL, &numOfExprs);
2470
  int32_t    code = initExprSupp(&pOperator->exprSupp, pExprInfo, numOfExprs);
2471 2472 2473
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
2474

H
Haojun Liao 已提交
2475 2476
  int32_t num = 0;
  code = extractColMatchInfo(pPhyNode->pScanPseudoCols, pDescNode, &num, COL_MATCH_FROM_COL_ID, &pInfo->matchInfo);
2477 2478 2479
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
2480

H
Haojun Liao 已提交
2481
  pInfo->pRes = createDataBlockFromDescNode(pDescNode);
2482 2483
  pInfo->readHandle = *pReadHandle;
  pInfo->curPos = 0;
2484

L
Liu Jicong 已提交
2485 2486
  setOperatorInfo(pOperator, "TagScanOperator", QUERY_NODE_PHYSICAL_PLAN_TAG_SCAN, false, OP_NOT_OPENED, pInfo,
                  pTaskInfo);
2487
  initResultSizeInfo(&pOperator->resultInfo, 4096);
2488 2489
  blockDataEnsureCapacity(pInfo->pRes, pOperator->resultInfo.capacity);

L
Liu Jicong 已提交
2490 2491
  pOperator->fpSet =
      createOperatorFpSet(optrDummyOpenFn, doTagScan, NULL, destroyTagScanOperatorInfo, optrDefaultBufFn, NULL);
H
Haojun Liao 已提交
2492 2493

  return pOperator;
2494

2495
_error:
H
Haojun Liao 已提交
2496 2497 2498 2499 2500
  taosMemoryFree(pInfo);
  taosMemoryFree(pOperator);
  terrno = TSDB_CODE_OUT_OF_MEMORY;
  return NULL;
}
2501

dengyihao's avatar
dengyihao 已提交
2502
static SSDataBlock* getTableDataBlockImpl(void* param) {
dengyihao's avatar
opt mem  
dengyihao 已提交
2503 2504 2505 2506 2507 2508 2509
  STableMergeScanSortSourceParam* source = param;
  SOperatorInfo*                  pOperator = source->pOperator;
  STableMergeScanInfo*            pInfo = pOperator->info;
  SExecTaskInfo*                  pTaskInfo = pOperator->pTaskInfo;
  int32_t                         readIdx = source->readerIdx;
  SSDataBlock*                    pBlock = source->inputBlock;

H
Haojun Liao 已提交
2510
  SQueryTableDataCond* pQueryCond = taosArrayGet(pInfo->queryConds, readIdx);
dengyihao's avatar
opt mem  
dengyihao 已提交
2511

L
Liu Jicong 已提交
2512 2513
  int64_t      st = taosGetTimestampUs();
  void*        p = tableListGetInfo(pTaskInfo->pTableInfoList, readIdx + pInfo->tableStartIndex);
H
Haojun Liao 已提交
2514
  SReadHandle* pHandle = &pInfo->base.readHandle;
dengyihao's avatar
dengyihao 已提交
2515

L
Liu Jicong 已提交
2516 2517
  int32_t code =
      tsdbReaderOpen(pHandle->vnode, pQueryCond, p, 1, pBlock, &pInfo->base.dataReader, GET_TASKID(pTaskInfo));
dengyihao's avatar
dengyihao 已提交
2518
  if (code != 0) {
H
Haojun Liao 已提交
2519
    T_LONG_JMP(pTaskInfo->env, code);
dengyihao's avatar
dengyihao 已提交
2520
  }
dengyihao's avatar
opt mem  
dengyihao 已提交
2521

H
Haojun Liao 已提交
2522
  STsdbReader* reader = pInfo->base.dataReader;
2523
  qTrace("tsdb/read-table-data: %p, enter next reader", reader);
dengyihao's avatar
opt mem  
dengyihao 已提交
2524
  while (tsdbNextDataBlock(reader)) {
H
Haojun Liao 已提交
2525
    if (isTaskKilled(pTaskInfo)) {
2526
      T_LONG_JMP(pTaskInfo->env, pTaskInfo->code);
dengyihao's avatar
opt mem  
dengyihao 已提交
2527 2528 2529
    }

    // process this data block based on the probabilities
H
Haojun Liao 已提交
2530
    bool processThisBlock = processBlockWithProbability(&pInfo->sample);
dengyihao's avatar
opt mem  
dengyihao 已提交
2531 2532 2533 2534
    if (!processThisBlock) {
      continue;
    }

H
Haojun Liao 已提交
2535
    if (pQueryCond->order == TSDB_ORDER_ASC) {
dengyihao's avatar
opt mem  
dengyihao 已提交
2536 2537 2538 2539
      pQueryCond->twindows.skey = pBlock->info.window.ekey + 1;
    } else {
      pQueryCond->twindows.ekey = pBlock->info.window.skey - 1;
    }
dengyihao's avatar
opt mem  
dengyihao 已提交
2540 2541

    uint32_t status = 0;
H
Haojun Liao 已提交
2542
    loadDataBlock(pOperator, &pInfo->base, pBlock, &status);
S
slzhou 已提交
2543
    //    code = loadDataBlockFromOneTable(pOperator, pTableScanInfo, pBlock, &status);
dengyihao's avatar
opt mem  
dengyihao 已提交
2544
    if (code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
2545
      T_LONG_JMP(pTaskInfo->env, code);
dengyihao's avatar
opt mem  
dengyihao 已提交
2546 2547 2548 2549 2550 2551 2552
    }

    // current block is filter out according to filter condition, continue load the next block
    if (status == FUNC_DATA_REQUIRED_FILTEROUT || pBlock->info.rows == 0) {
      continue;
    }

H
Haojun Liao 已提交
2553
    pBlock->info.id.groupId = getTableGroupId(pTaskInfo->pTableInfoList, pBlock->info.id.uid);
dengyihao's avatar
opt mem  
dengyihao 已提交
2554

H
Haojun Liao 已提交
2555
    pOperator->resultInfo.totalRows += pBlock->info.rows;
H
Haojun Liao 已提交
2556
    pInfo->base.readRecorder.elapsedTime += (taosGetTimestampUs() - st) / 1000.0;
dengyihao's avatar
opt mem  
dengyihao 已提交
2557

2558
    qTrace("tsdb/read-table-data: %p, close reader", reader);
H
Haojun Liao 已提交
2559 2560
    tsdbReaderClose(pInfo->base.dataReader);
    pInfo->base.dataReader = NULL;
dengyihao's avatar
opt mem  
dengyihao 已提交
2561 2562
    return pBlock;
  }
H
Haojun Liao 已提交
2563

H
Haojun Liao 已提交
2564 2565
  tsdbReaderClose(pInfo->base.dataReader);
  pInfo->base.dataReader = NULL;
dengyihao's avatar
opt mem  
dengyihao 已提交
2566 2567 2568
  return NULL;
}

2569 2570 2571
SArray* generateSortByTsInfo(SArray* colMatchInfo, int32_t order) {
  int32_t tsTargetSlotId = 0;
  for (int32_t i = 0; i < taosArrayGetSize(colMatchInfo); ++i) {
H
Haojun Liao 已提交
2572
    SColMatchItem* colInfo = taosArrayGet(colMatchInfo, i);
2573
    if (colInfo->colId == PRIMARYKEY_TIMESTAMP_COL_ID) {
H
Haojun Liao 已提交
2574
      tsTargetSlotId = colInfo->dstSlotId;
2575 2576 2577
    }
  }

2578 2579 2580
  SArray*         pList = taosArrayInit(1, sizeof(SBlockOrderInfo));
  SBlockOrderInfo bi = {0};
  bi.order = order;
2581
  bi.slotId = tsTargetSlotId;
2582 2583 2584 2585 2586 2587 2588
  bi.nullFirst = NULL_ORDER_FIRST;

  taosArrayPush(pList, &bi);

  return pList;
}

H
Haojun Liao 已提交
2589
int32_t dumpQueryTableCond(const SQueryTableDataCond* src, SQueryTableDataCond* dst) {
dengyihao's avatar
opt mem  
dengyihao 已提交
2590 2591 2592 2593 2594 2595 2596
  memcpy((void*)dst, (void*)src, sizeof(SQueryTableDataCond));
  dst->colList = taosMemoryCalloc(src->numOfCols, sizeof(SColumnInfo));
  for (int i = 0; i < src->numOfCols; i++) {
    dst->colList[i] = src->colList[i];
  }
  return 0;
}
H
Haojun Liao 已提交
2597

2598
int32_t startGroupTableMergeScan(SOperatorInfo* pOperator) {
2599 2600 2601
  STableMergeScanInfo* pInfo = pOperator->info;
  SExecTaskInfo*       pTaskInfo = pOperator->pTaskInfo;

S
slzhou 已提交
2602
  {
H
Haojun Liao 已提交
2603
    size_t  numOfTables = tableListGetSize(pTaskInfo->pTableInfoList);
S
slzhou 已提交
2604
    int32_t i = pInfo->tableStartIndex + 1;
H
Haojun Liao 已提交
2605
    for (; i < numOfTables; ++i) {
H
Haojun Liao 已提交
2606
      STableKeyInfo* tableKeyInfo = tableListGetInfo(pTaskInfo->pTableInfoList, i);
S
slzhou 已提交
2607 2608 2609 2610 2611 2612
      if (tableKeyInfo->groupId != pInfo->groupId) {
        break;
      }
    }
    pInfo->tableEndIndex = i - 1;
  }
2613

S
slzhou 已提交
2614 2615
  int32_t tableStartIdx = pInfo->tableStartIndex;
  int32_t tableEndIdx = pInfo->tableEndIndex;
2616

H
Haojun Liao 已提交
2617
  pInfo->base.dataReader = NULL;
2618

2619 2620
  // todo the total available buffer should be determined by total capacity of buffer of this task.
  // the additional one is reserved for merge result
S
slzhou 已提交
2621
  pInfo->sortBufSize = pInfo->bufPageSize * (tableEndIdx - tableStartIdx + 1 + 1);
2622
  int32_t numOfBufPage = pInfo->sortBufSize / pInfo->bufPageSize;
L
Liu Jicong 已提交
2623 2624
  pInfo->pSortHandle = tsortCreateSortHandle(pInfo->pSortInfo, SORT_MULTISOURCE_MERGE, pInfo->bufPageSize, numOfBufPage,
                                             pInfo->pSortInputBlock, pTaskInfo->id.str);
2625

dengyihao's avatar
dengyihao 已提交
2626
  tsortSetFetchRawDataFp(pInfo->pSortHandle, getTableDataBlockImpl, NULL, NULL);
dengyihao's avatar
opt mem  
dengyihao 已提交
2627 2628 2629 2630 2631 2632

  // one table has one data block
  int32_t numOfTable = tableEndIdx - tableStartIdx + 1;
  pInfo->queryConds = taosArrayInit(numOfTable, sizeof(SQueryTableDataCond));

  for (int32_t i = 0; i < numOfTable; ++i) {
2633 2634 2635 2636
    STableMergeScanSortSourceParam param = {0};
    param.readerIdx = i;
    param.pOperator = pOperator;
    param.inputBlock = createOneDataBlock(pInfo->pResBlock, false);
H
Haojun Liao 已提交
2637 2638
    blockDataEnsureCapacity(param.inputBlock, pOperator->resultInfo.capacity);

2639
    taosArrayPush(pInfo->sortSourceParams, &param);
dengyihao's avatar
opt mem  
dengyihao 已提交
2640 2641

    SQueryTableDataCond cond;
H
Haojun Liao 已提交
2642
    dumpQueryTableCond(&pInfo->base.cond, &cond);
dengyihao's avatar
opt mem  
dengyihao 已提交
2643
    taosArrayPush(pInfo->queryConds, &cond);
2644 2645
  }

dengyihao's avatar
opt mem  
dengyihao 已提交
2646
  for (int32_t i = 0; i < numOfTable; ++i) {
2647
    SSortSource*                    ps = taosMemoryCalloc(1, sizeof(SSortSource));
2648
    STableMergeScanSortSourceParam* param = taosArrayGet(pInfo->sortSourceParams, i);
2649
    ps->param = param;
2650
    ps->onlyRef = true;
2651 2652 2653 2654 2655 2656
    tsortAddSource(pInfo->pSortHandle, ps);
  }

  int32_t code = tsortOpen(pInfo->pSortHandle);

  if (code != TSDB_CODE_SUCCESS) {
2657
    T_LONG_JMP(pTaskInfo->env, terrno);
2658 2659
  }

2660 2661 2662 2663 2664 2665 2666
  return TSDB_CODE_SUCCESS;
}

int32_t stopGroupTableMergeScan(SOperatorInfo* pOperator) {
  STableMergeScanInfo* pInfo = pOperator->info;
  SExecTaskInfo*       pTaskInfo = pOperator->pTaskInfo;

dengyihao's avatar
dengyihao 已提交
2667
  int32_t numOfTable = taosArrayGetSize(pInfo->queryConds);
2668

2669 2670 2671 2672 2673 2674 2675
  SSortExecInfo sortExecInfo = tsortGetSortExecInfo(pInfo->pSortHandle);
  pInfo->sortExecInfo.sortMethod = sortExecInfo.sortMethod;
  pInfo->sortExecInfo.sortBuffer = sortExecInfo.sortBuffer;
  pInfo->sortExecInfo.loops += sortExecInfo.loops;
  pInfo->sortExecInfo.readBytes += sortExecInfo.readBytes;
  pInfo->sortExecInfo.writeBytes += sortExecInfo.writeBytes;

dengyihao's avatar
dengyihao 已提交
2676
  for (int32_t i = 0; i < numOfTable; ++i) {
2677 2678 2679
    STableMergeScanSortSourceParam* param = taosArrayGet(pInfo->sortSourceParams, i);
    blockDataDestroy(param->inputBlock);
  }
2680 2681
  taosArrayClear(pInfo->sortSourceParams);

2682
  tsortDestroySortHandle(pInfo->pSortHandle);
dengyihao's avatar
dengyihao 已提交
2683
  pInfo->pSortHandle = NULL;
2684

dengyihao's avatar
opt mem  
dengyihao 已提交
2685 2686 2687
  for (int32_t i = 0; i < taosArrayGetSize(pInfo->queryConds); i++) {
    SQueryTableDataCond* cond = taosArrayGet(pInfo->queryConds, i);
    taosMemoryFree(cond->colList);
2688
  }
dengyihao's avatar
opt mem  
dengyihao 已提交
2689 2690 2691
  taosArrayDestroy(pInfo->queryConds);
  pInfo->queryConds = NULL;

2692
  resetLimitInfoForNextGroup(&pInfo->limitInfo);
2693 2694 2695
  return TSDB_CODE_SUCCESS;
}

2696 2697
// all data produced by this function only belongs to one group
// slimit/soffset does not need to be concerned here, since this function only deal with data within one group.
L
Liu Jicong 已提交
2698 2699
SSDataBlock* getSortedTableMergeScanBlockData(SSortHandle* pHandle, SSDataBlock* pResBlock, int32_t capacity,
                                              SOperatorInfo* pOperator) {
2700 2701 2702
  STableMergeScanInfo* pInfo = pOperator->info;
  SExecTaskInfo*       pTaskInfo = pOperator->pTaskInfo;

2703
  blockDataCleanup(pResBlock);
2704 2705

  while (1) {
2706
    STupleHandle* pTupleHandle = tsortNextTuple(pHandle);
2707 2708 2709 2710
    if (pTupleHandle == NULL) {
      break;
    }

2711 2712
    appendOneRowToDataBlock(pResBlock, pTupleHandle);
    if (pResBlock->info.rows >= capacity) {
2713 2714 2715 2716
      break;
    }
  }

2717 2718 2719
  applyLimitOffset(&pInfo->limitInfo, pResBlock, pTaskInfo, pOperator);
  pInfo->limitInfo.numOfOutputRows += pResBlock->info.rows;

2720 2721 2722
  qDebug("%s get sorted row block, rows:%d, limit:%"PRId64, GET_TASKID(pTaskInfo), pResBlock->info.rows,
      pInfo->limitInfo.numOfOutputRows);

2723
  return (pResBlock->info.rows > 0) ? pResBlock : NULL;
2724 2725 2726 2727 2728 2729 2730 2731 2732 2733 2734 2735
}

SSDataBlock* doTableMergeScan(SOperatorInfo* pOperator) {
  if (pOperator->status == OP_EXEC_DONE) {
    return NULL;
  }

  SExecTaskInfo*       pTaskInfo = pOperator->pTaskInfo;
  STableMergeScanInfo* pInfo = pOperator->info;

  int32_t code = pOperator->fpSet._openFn(pOperator);
  if (code != TSDB_CODE_SUCCESS) {
2736
    T_LONG_JMP(pTaskInfo->env, code);
2737
  }
2738

H
Haojun Liao 已提交
2739
  size_t tableListSize = tableListGetSize(pTaskInfo->pTableInfoList);
S
slzhou 已提交
2740 2741
  if (!pInfo->hasGroupId) {
    pInfo->hasGroupId = true;
2742

S
slzhou 已提交
2743
    if (tableListSize == 0) {
H
Haojun Liao 已提交
2744
      setOperatorCompleted(pOperator);
2745 2746
      return NULL;
    }
S
slzhou 已提交
2747
    pInfo->tableStartIndex = 0;
H
Haojun Liao 已提交
2748
    pInfo->groupId = ((STableKeyInfo*)tableListGetInfo(pTaskInfo->pTableInfoList, pInfo->tableStartIndex))->groupId;
2749 2750
    startGroupTableMergeScan(pOperator);
  }
2751

S
slzhou 已提交
2752 2753
  SSDataBlock* pBlock = NULL;
  while (pInfo->tableStartIndex < tableListSize) {
L
Liu Jicong 已提交
2754 2755
    pBlock = getSortedTableMergeScanBlockData(pInfo->pSortHandle, pInfo->pResBlock, pOperator->resultInfo.capacity,
                                              pOperator);
S
slzhou 已提交
2756
    if (pBlock != NULL) {
H
Haojun Liao 已提交
2757
      pBlock->info.id.groupId = pInfo->groupId;
S
slzhou 已提交
2758 2759 2760
      pOperator->resultInfo.totalRows += pBlock->info.rows;
      return pBlock;
    } else {
2761
      // Data of this group are all dumped, let's try the next group
S
slzhou 已提交
2762 2763
      stopGroupTableMergeScan(pOperator);
      if (pInfo->tableEndIndex >= tableListSize - 1) {
H
Haojun Liao 已提交
2764
        setOperatorCompleted(pOperator);
S
slzhou 已提交
2765 2766
        break;
      }
2767

S
slzhou 已提交
2768
      pInfo->tableStartIndex = pInfo->tableEndIndex + 1;
H
Haojun Liao 已提交
2769
      pInfo->groupId = tableListGetInfo(pTaskInfo->pTableInfoList, pInfo->tableStartIndex)->groupId;
S
slzhou 已提交
2770 2771
      startGroupTableMergeScan(pOperator);
    }
wmmhello's avatar
wmmhello 已提交
2772 2773
  }

2774 2775 2776
  return pBlock;
}

2777
void destroyTableMergeScanOperatorInfo(void* param) {
2778
  STableMergeScanInfo* pTableScanInfo = (STableMergeScanInfo*)param;
H
Haojun Liao 已提交
2779
  cleanupQueryTableDataCond(&pTableScanInfo->base.cond);
2780

dengyihao's avatar
dengyihao 已提交
2781 2782 2783
  int32_t numOfTable = taosArrayGetSize(pTableScanInfo->queryConds);

  for (int32_t i = 0; i < numOfTable; i++) {
H
Haojun Liao 已提交
2784 2785
    STableMergeScanSortSourceParam* p = taosArrayGet(pTableScanInfo->sortSourceParams, i);
    blockDataDestroy(p->inputBlock);
2786
  }
H
Haojun Liao 已提交
2787

2788
  taosArrayDestroy(pTableScanInfo->sortSourceParams);
dengyihao's avatar
dengyihao 已提交
2789 2790
  tsortDestroySortHandle(pTableScanInfo->pSortHandle);
  pTableScanInfo->pSortHandle = NULL;
2791

H
Haojun Liao 已提交
2792 2793
  tsdbReaderClose(pTableScanInfo->base.dataReader);
  pTableScanInfo->base.dataReader = NULL;
dengyihao's avatar
opt mem  
dengyihao 已提交
2794

dengyihao's avatar
opt mem  
dengyihao 已提交
2795 2796 2797
  for (int i = 0; i < taosArrayGetSize(pTableScanInfo->queryConds); i++) {
    SQueryTableDataCond* pCond = taosArrayGet(pTableScanInfo->queryConds, i);
    taosMemoryFree(pCond->colList);
2798
  }
dengyihao's avatar
opt mem  
dengyihao 已提交
2799
  taosArrayDestroy(pTableScanInfo->queryConds);
2800

H
Haojun Liao 已提交
2801 2802
  if (pTableScanInfo->base.matchInfo.pList != NULL) {
    taosArrayDestroy(pTableScanInfo->base.matchInfo.pList);
2803 2804 2805 2806 2807 2808
  }

  pTableScanInfo->pResBlock = blockDataDestroy(pTableScanInfo->pResBlock);
  pTableScanInfo->pSortInputBlock = blockDataDestroy(pTableScanInfo->pSortInputBlock);

  taosArrayDestroy(pTableScanInfo->pSortInfo);
H
Haojun Liao 已提交
2809
  cleanupExprSupp(&pTableScanInfo->base.pseudoSup);
L
Liu Jicong 已提交
2810

H
Haojun Liao 已提交
2811 2812 2813 2814
  tsdbReaderClose(pTableScanInfo->base.dataReader);
  pTableScanInfo->base.dataReader = NULL;
  taosLRUCacheCleanup(pTableScanInfo->base.metaCache.pTableMetaEntryCache);

D
dapan1121 已提交
2815
  taosMemoryFreeClear(param);
2816 2817 2818 2819
}

int32_t getTableMergeScanExplainExecInfo(SOperatorInfo* pOptr, void** pOptrExplain, uint32_t* len) {
  ASSERT(pOptr != NULL);
2820 2821
  // TODO: merge these two info into one struct
  STableMergeScanExecInfo* execInfo = taosMemoryCalloc(1, sizeof(STableMergeScanExecInfo));
L
Liu Jicong 已提交
2822
  STableMergeScanInfo*     pInfo = pOptr->info;
H
Haojun Liao 已提交
2823
  execInfo->blockRecorder = pInfo->base.readRecorder;
2824
  execInfo->sortExecInfo = pInfo->sortExecInfo;
2825 2826 2827

  *pOptrExplain = execInfo;
  *len = sizeof(STableMergeScanExecInfo);
L
Liu Jicong 已提交
2828

2829 2830 2831
  return TSDB_CODE_SUCCESS;
}

H
Haojun Liao 已提交
2832 2833
SOperatorInfo* createTableMergeScanOperatorInfo(STableScanPhysiNode* pTableScanNode, SReadHandle* readHandle,
                                                SExecTaskInfo* pTaskInfo) {
2834 2835 2836 2837 2838
  STableMergeScanInfo* pInfo = taosMemoryCalloc(1, sizeof(STableMergeScanInfo));
  SOperatorInfo*       pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
  if (pInfo == NULL || pOperator == NULL) {
    goto _error;
  }
2839

2840 2841 2842
  SDataBlockDescNode* pDescNode = pTableScanNode->scan.node.pOutputDataBlockDesc;

  int32_t numOfCols = 0;
2843
  int32_t code = extractColMatchInfo(pTableScanNode->scan.pScanCols, pDescNode, &numOfCols, COL_MATCH_FROM_COL_ID,
H
Haojun Liao 已提交
2844
                                     &pInfo->base.matchInfo);
H
Haojun Liao 已提交
2845 2846 2847
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }
2848

H
Haojun Liao 已提交
2849
  code = initQueryTableDataCond(&pInfo->base.cond, pTableScanNode);
2850
  if (code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
2851
    taosArrayDestroy(pInfo->base.matchInfo.pList);
2852 2853 2854 2855
    goto _error;
  }

  if (pTableScanNode->scan.pScanPseudoCols != NULL) {
H
Haojun Liao 已提交
2856
    SExprSupp* pSup = &pInfo->base.pseudoSup;
2857 2858
    pSup->pExprInfo = createExprInfo(pTableScanNode->scan.pScanPseudoCols, NULL, &pSup->numOfExprs);
    pSup->pCtx = createSqlFunctionCtx(pSup->pExprInfo, pSup->numOfExprs, &pSup->rowEntryInfoOffset);
2859 2860 2861 2862
  }

  pInfo->scanInfo = (SScanInfo){.numOfAsc = pTableScanNode->scanSeq[0], .numOfDesc = pTableScanNode->scanSeq[1]};

H
Haojun Liao 已提交
2863 2864 2865 2866 2867 2868
  pInfo->base.metaCache.pTableMetaEntryCache = taosLRUCacheInit(1024 * 128, -1, .5);
  if (pInfo->base.metaCache.pTableMetaEntryCache == NULL) {
    code = terrno;
    goto _error;
  }

H
Haojun Liao 已提交
2869 2870
  pInfo->base.dataBlockLoadFlag = FUNC_DATA_REQUIRED_DATA_LOAD;
  pInfo->base.scanFlag = MAIN_SCAN;
H
Haojun Liao 已提交
2871
  pInfo->base.readHandle = *readHandle;
2872 2873 2874

  pInfo->base.limitInfo.limit.limit = -1;
  pInfo->base.limitInfo.slimit.limit = -1;
H
Haojun Liao 已提交
2875

2876
  pInfo->sample.sampleRatio = pTableScanNode->ratio;
L
Liu Jicong 已提交
2877
  pInfo->sample.seed = taosGetTimestampSec();
H
Haojun Liao 已提交
2878 2879 2880 2881 2882 2883

  code = filterInitFromNode((SNode*)pTableScanNode->scan.node.pConditions, &pOperator->exprSupp.pFilterInfo, 0);
  if (code != TSDB_CODE_SUCCESS) {
    goto _error;
  }

H
Haojun Liao 已提交
2884
  initResultSizeInfo(&pOperator->resultInfo, 1024);
H
Haojun Liao 已提交
2885
  pInfo->pResBlock = createDataBlockFromDescNode(pDescNode);
H
Haojun Liao 已提交
2886 2887
  blockDataEnsureCapacity(pInfo->pResBlock, pOperator->resultInfo.capacity);

2888
  pInfo->sortSourceParams = taosArrayInit(64, sizeof(STableMergeScanSortSourceParam));
2889

H
Haojun Liao 已提交
2890
  pInfo->pSortInfo = generateSortByTsInfo(pInfo->base.matchInfo.pList, pInfo->base.cond.order);
2891
  pInfo->pSortInputBlock = createOneDataBlock(pInfo->pResBlock, false);
2892
  initLimitInfo(pTableScanNode->scan.node.pLimit, pTableScanNode->scan.node.pSlimit, &pInfo->limitInfo);
2893

dengyihao's avatar
dengyihao 已提交
2894
  int32_t  rowSize = pInfo->pResBlock->info.rowSize;
A
Alex Duan 已提交
2895 2896
  uint32_t nCols = taosArrayGetSize(pInfo->pResBlock->pDataBlock);
  pInfo->bufPageSize = getProperSortPageSize(rowSize, nCols);
2897

L
Liu Jicong 已提交
2898 2899
  setOperatorInfo(pOperator, "TableMergeScanOperator", QUERY_NODE_PHYSICAL_PLAN_TABLE_MERGE_SCAN, false, OP_NOT_OPENED,
                  pInfo, pTaskInfo);
L
Liu Jicong 已提交
2900
  pOperator->exprSupp.numOfExprs = numOfCols;
2901

2902 2903
  pOperator->fpSet = createOperatorFpSet(optrDummyOpenFn, doTableMergeScan, NULL, destroyTableMergeScanOperatorInfo,
                                         optrDefaultBufFn, getTableMergeScanExplainExecInfo);
2904 2905 2906 2907 2908 2909 2910 2911 2912
  pOperator->cost.openCost = 0;
  return pOperator;

_error:
  pTaskInfo->code = TSDB_CODE_OUT_OF_MEMORY;
  taosMemoryFree(pInfo);
  taosMemoryFree(pOperator);
  return NULL;
}
S
shenglian zhou 已提交
2913 2914 2915 2916

// ====================================================================================================================
// TableCountScanOperator
static SSDataBlock* doTableCountScan(SOperatorInfo* pOperator);
S
slzhou 已提交
2917
static void         destoryTableCountScanOperator(void* param);
S
slzhou 已提交
2918 2919 2920 2921 2922 2923
static void         buildVnodeGroupedStbTableCount(STableCountScanOperatorInfo* pInfo, STableCountScanSupp* pSupp,
                                                   SSDataBlock* pRes, char* dbName, tb_uid_t stbUid);
static void         buildVnodeGroupedNtbTableCount(STableCountScanOperatorInfo* pInfo, STableCountScanSupp* pSupp,
                                                   SSDataBlock* pRes, char* dbName);
static void         buildVnodeFilteredTbCount(SOperatorInfo* pOperator, STableCountScanOperatorInfo* pInfo,
                                              STableCountScanSupp* pSupp, SSDataBlock* pRes, char* dbName);
L
Liu Jicong 已提交
2924 2925
static void         buildVnodeGroupedTableCount(SOperatorInfo* pOperator, STableCountScanOperatorInfo* pInfo,
                                                STableCountScanSupp* pSupp, SSDataBlock* pRes, int32_t vgId, char* dbName);
S
slzhou 已提交
2926 2927 2928 2929 2930 2931 2932
static SSDataBlock* buildVnodeDbTableCount(SOperatorInfo* pOperator, STableCountScanOperatorInfo* pInfo,
                                           STableCountScanSupp* pSupp, SSDataBlock* pRes);
static void         buildSysDbGroupedTableCount(SOperatorInfo* pOperator, STableCountScanOperatorInfo* pInfo,
                                                STableCountScanSupp* pSupp, SSDataBlock* pRes, size_t infodbTableNum,
                                                size_t perfdbTableNum);
static void         buildSysDbFilterTableCount(SOperatorInfo* pOperator, STableCountScanSupp* pSupp, SSDataBlock* pRes,
                                               size_t infodbTableNum, size_t perfdbTableNum);
S
slzhou 已提交
2933 2934 2935 2936 2937 2938 2939 2940 2941 2942 2943 2944 2945 2946 2947 2948 2949 2950 2951 2952 2953 2954 2955 2956 2957 2958 2959 2960 2961 2962 2963 2964 2965 2966 2967 2968 2969 2970 2971 2972 2973 2974 2975 2976 2977 2978 2979 2980 2981 2982 2983 2984 2985 2986 2987 2988 2989 2990 2991 2992 2993
static const char*  GROUP_TAG_DB_NAME = "db_name";
static const char*  GROUP_TAG_STABLE_NAME = "stable_name";

int32_t tblCountScanGetGroupTagsSlotId(const SNodeList* scanCols, STableCountScanSupp* supp) {
  if (scanCols != NULL) {
    SNode* pNode = NULL;
    FOREACH(pNode, scanCols) {
      if (nodeType(pNode) != QUERY_NODE_TARGET) {
        return TSDB_CODE_QRY_SYS_ERROR;
      }
      STargetNode* targetNode = (STargetNode*)pNode;
      if (nodeType(targetNode->pExpr) != QUERY_NODE_COLUMN) {
        return TSDB_CODE_QRY_SYS_ERROR;
      }
      SColumnNode* colNode = (SColumnNode*)(targetNode->pExpr);
      if (strcmp(colNode->colName, GROUP_TAG_DB_NAME) == 0) {
        supp->dbNameSlotId = targetNode->slotId;
      } else if (strcmp(colNode->colName, GROUP_TAG_STABLE_NAME) == 0) {
        supp->stbNameSlotId = targetNode->slotId;
      }
    }
  }
  return TSDB_CODE_SUCCESS;
}

int32_t tblCountScanGetCountSlotId(const SNodeList* pseudoCols, STableCountScanSupp* supp) {
  if (pseudoCols != NULL) {
    SNode* pNode = NULL;
    FOREACH(pNode, pseudoCols) {
      if (nodeType(pNode) != QUERY_NODE_TARGET) {
        return TSDB_CODE_QRY_SYS_ERROR;
      }
      STargetNode* targetNode = (STargetNode*)pNode;
      if (nodeType(targetNode->pExpr) != QUERY_NODE_FUNCTION) {
        return TSDB_CODE_QRY_SYS_ERROR;
      }
      SFunctionNode* funcNode = (SFunctionNode*)(targetNode->pExpr);
      if (funcNode->funcType == FUNCTION_TYPE_TABLE_COUNT) {
        supp->tbCountSlotId = targetNode->slotId;
      }
    }
  }
  return TSDB_CODE_SUCCESS;
}

int32_t tblCountScanGetInputs(SNodeList* groupTags, SName* tableName, STableCountScanSupp* supp) {
  if (groupTags != NULL) {
    SNode* pNode = NULL;
    FOREACH(pNode, groupTags) {
      if (nodeType(pNode) != QUERY_NODE_COLUMN) {
        return TSDB_CODE_QRY_SYS_ERROR;
      }
      SColumnNode* colNode = (SColumnNode*)pNode;
      if (strcmp(colNode->colName, GROUP_TAG_DB_NAME) == 0) {
        supp->groupByDbName = true;
      }
      if (strcmp(colNode->colName, GROUP_TAG_STABLE_NAME) == 0) {
        supp->groupByStbName = true;
      }
    }
  } else {
S
slzhou 已提交
2994 2995
    strncpy(supp->dbNameFilter, tNameGetDbNameP(tableName), TSDB_DB_NAME_LEN);
    strncpy(supp->stbNameFilter, tNameGetTableName(tableName), TSDB_TABLE_NAME_LEN);
S
slzhou 已提交
2996 2997 2998 2999 3000 3001 3002 3003 3004 3005 3006 3007 3008 3009 3010 3011 3012 3013 3014 3015 3016 3017 3018 3019 3020 3021 3022 3023
  }
  return TSDB_CODE_SUCCESS;
}

int32_t getTableCountScanSupp(SNodeList* groupTags, SName* tableName, SNodeList* scanCols, SNodeList* pseudoCols,
                              STableCountScanSupp* supp, SExecTaskInfo* taskInfo) {
  int32_t code = 0;
  code = tblCountScanGetInputs(groupTags, tableName, supp);
  if (code != TSDB_CODE_SUCCESS) {
    qError("%s get table count scan supp. get inputs error", GET_TASKID(taskInfo));
    return code;
  }
  supp->dbNameSlotId = -1;
  supp->stbNameSlotId = -1;
  supp->tbCountSlotId = -1;

  code = tblCountScanGetGroupTagsSlotId(scanCols, supp);
  if (code != TSDB_CODE_SUCCESS) {
    qError("%s get table count scan supp. get group tags slot id error", GET_TASKID(taskInfo));
    return code;
  }
  code = tblCountScanGetCountSlotId(pseudoCols, supp);
  if (code != TSDB_CODE_SUCCESS) {
    qError("%s get table count scan supp. get count error", GET_TASKID(taskInfo));
    return code;
  }
  return code;
}
S
shenglian zhou 已提交
3024

S
slzhou 已提交
3025
SOperatorInfo* createTableCountScanOperatorInfo(SReadHandle* readHandle, STableCountScanPhysiNode* pTblCountScanNode,
S
shenglian zhou 已提交
3026 3027 3028
                                                SExecTaskInfo* pTaskInfo) {
  int32_t code = TSDB_CODE_SUCCESS;

S
slzhou 已提交
3029
  SScanPhysiNode*              pScanNode = &pTblCountScanNode->scan;
S
slzhou 已提交
3030
  STableCountScanOperatorInfo* pInfo = taosMemoryCalloc(1, sizeof(STableCountScanOperatorInfo));
S
slzhou 已提交
3031
  SOperatorInfo*               pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
S
shenglian zhou 已提交
3032 3033 3034 3035 3036 3037 3038 3039 3040

  if (!pInfo || !pOperator) {
    goto _error;
  }

  pInfo->readHandle = *readHandle;

  SDataBlockDescNode* pDescNode = pScanNode->node.pOutputDataBlockDesc;
  initResultSizeInfo(&pOperator->resultInfo, 1);
3041
  pInfo->pRes = createDataBlockFromDescNode(pDescNode);
S
shenglian zhou 已提交
3042 3043
  blockDataEnsureCapacity(pInfo->pRes, pOperator->resultInfo.capacity);

S
slzhou 已提交
3044 3045 3046
  getTableCountScanSupp(pTblCountScanNode->pGroupTags, &pTblCountScanNode->scan.tableName,
                        pTblCountScanNode->scan.pScanCols, pTblCountScanNode->scan.pScanPseudoCols, &pInfo->supp,
                        pTaskInfo);
S
shenglian zhou 已提交
3047 3048 3049

  setOperatorInfo(pOperator, "TableCountScanOperator", QUERY_NODE_PHYSICAL_PLAN_TABLE_COUNT_SCAN, false, OP_NOT_OPENED,
                  pInfo, pTaskInfo);
L
Liu Jicong 已提交
3050 3051
  pOperator->fpSet = createOperatorFpSet(optrDummyOpenFn, doTableCountScan, NULL, destoryTableCountScanOperator,
                                         optrDefaultBufFn, NULL);
S
shenglian zhou 已提交
3052 3053 3054 3055 3056 3057 3058 3059 3060 3061 3062
  return pOperator;

_error:
  if (pInfo != NULL) {
    destoryTableCountScanOperator(pInfo);
  }
  taosMemoryFreeClear(pOperator);
  pTaskInfo->code = code;
  return NULL;
}

S
slzhou 已提交
3063 3064 3065
void fillTableCountScanDataBlock(STableCountScanSupp* pSupp, char* dbName, char* stbName, int64_t count,
                                 SSDataBlock* pRes) {
  if (pSupp->dbNameSlotId != -1) {
3066
    ASSERT(strlen(dbName));
S
slzhou 已提交
3067
    SColumnInfoData* colInfoData = taosArrayGet(pRes->pDataBlock, pSupp->dbNameSlotId);
H
Haojun Liao 已提交
3068 3069 3070 3071

    char varDbName[TSDB_DB_NAME_LEN + VARSTR_HEADER_SIZE] = {0};
    tstrncpy(varDataVal(varDbName), dbName, TSDB_DB_NAME_LEN);

S
slzhou 已提交
3072 3073 3074 3075 3076 3077
    varDataSetLen(varDbName, strlen(dbName));
    colDataAppend(colInfoData, 0, varDbName, false);
  }

  if (pSupp->stbNameSlotId != -1) {
    SColumnInfoData* colInfoData = taosArrayGet(pRes->pDataBlock, pSupp->stbNameSlotId);
3078
    if (strlen(stbName) != 0) {
S
slzhou 已提交
3079
      char varStbName[TSDB_TABLE_NAME_LEN + VARSTR_HEADER_SIZE] = {0};
H
Haojun Liao 已提交
3080
      strncpy(varDataVal(varStbName), stbName, TSDB_TABLE_NAME_LEN);
3081 3082 3083 3084 3085
      varDataSetLen(varStbName, strlen(stbName));
      colDataAppend(colInfoData, 0, varStbName, false);
    } else {
      colDataAppendNULL(colInfoData, 0);
    }
S
slzhou 已提交
3086 3087 3088
  }

  if (pSupp->tbCountSlotId != -1) {
S
slzhou 已提交
3089
    SColumnInfoData* colInfoData = taosArrayGet(pRes->pDataBlock, pSupp->tbCountSlotId);
S
slzhou 已提交
3090 3091 3092 3093 3094
    colDataAppend(colInfoData, 0, (char*)&count, false);
  }
  pRes->info.rows = 1;
}

S
slzhou 已提交
3095
static SSDataBlock* buildSysDbTableCount(SOperatorInfo* pOperator, STableCountScanOperatorInfo* pInfo) {
S
slzhou 已提交
3096 3097 3098
  STableCountScanSupp* pSupp = &pInfo->supp;
  SSDataBlock*         pRes = pInfo->pRes;

S
slzhou 已提交
3099
  size_t infodbTableNum;
S
slzhou 已提交
3100
  getInfosDbMeta(NULL, &infodbTableNum);
S
slzhou 已提交
3101
  size_t perfdbTableNum;
S
slzhou 已提交
3102 3103 3104
  getPerfDbMeta(NULL, &perfdbTableNum);

  if (pSupp->groupByDbName) {
S
slzhou 已提交
3105
    buildSysDbGroupedTableCount(pOperator, pInfo, pSupp, pRes, infodbTableNum, perfdbTableNum);
S
slzhou 已提交
3106 3107
    return (pRes->info.rows > 0) ? pRes : NULL;
  } else {
S
slzhou 已提交
3108
    buildSysDbFilterTableCount(pOperator, pSupp, pRes, infodbTableNum, perfdbTableNum);
S
slzhou 已提交
3109 3110 3111 3112
    return (pRes->info.rows > 0) ? pRes : NULL;
  }
}

S
slzhou 已提交
3113 3114 3115 3116 3117 3118 3119 3120 3121 3122 3123 3124 3125 3126 3127 3128 3129 3130 3131 3132 3133 3134 3135 3136 3137 3138 3139 3140 3141
static void buildSysDbFilterTableCount(SOperatorInfo* pOperator, STableCountScanSupp* pSupp, SSDataBlock* pRes,
                                       size_t infodbTableNum, size_t perfdbTableNum) {
  if (strcmp(pSupp->dbNameFilter, TSDB_INFORMATION_SCHEMA_DB) == 0) {
    fillTableCountScanDataBlock(pSupp, TSDB_INFORMATION_SCHEMA_DB, "", infodbTableNum, pRes);
  } else if (strcmp(pSupp->dbNameFilter, TSDB_PERFORMANCE_SCHEMA_DB) == 0) {
    fillTableCountScanDataBlock(pSupp, TSDB_PERFORMANCE_SCHEMA_DB, "", perfdbTableNum, pRes);
  } else if (strlen(pSupp->dbNameFilter) == 0) {
    fillTableCountScanDataBlock(pSupp, "", "", infodbTableNum + perfdbTableNum, pRes);
  }
  setOperatorCompleted(pOperator);
}

static void buildSysDbGroupedTableCount(SOperatorInfo* pOperator, STableCountScanOperatorInfo* pInfo,
                                        STableCountScanSupp* pSupp, SSDataBlock* pRes, size_t infodbTableNum,
                                        size_t perfdbTableNum) {
  if (pInfo->currGrpIdx == 0) {
    uint64_t groupId = calcGroupId(TSDB_INFORMATION_SCHEMA_DB, strlen(TSDB_INFORMATION_SCHEMA_DB));
    pRes->info.id.groupId = groupId;
    fillTableCountScanDataBlock(pSupp, TSDB_INFORMATION_SCHEMA_DB, "", infodbTableNum, pRes);
  } else if (pInfo->currGrpIdx == 1) {
    uint64_t groupId = calcGroupId(TSDB_PERFORMANCE_SCHEMA_DB, strlen(TSDB_PERFORMANCE_SCHEMA_DB));
    pRes->info.id.groupId = groupId;
    fillTableCountScanDataBlock(pSupp, TSDB_PERFORMANCE_SCHEMA_DB, "", perfdbTableNum, pRes);
  } else {
    setOperatorCompleted(pOperator);
  }
  pInfo->currGrpIdx++;
}

S
shenglian zhou 已提交
3142
static SSDataBlock* doTableCountScan(SOperatorInfo* pOperator) {
S
slzhou 已提交
3143 3144 3145 3146
  SExecTaskInfo*               pTaskInfo = pOperator->pTaskInfo;
  STableCountScanOperatorInfo* pInfo = pOperator->info;
  STableCountScanSupp*         pSupp = &pInfo->supp;
  SSDataBlock*                 pRes = pInfo->pRes;
S
slzhou 已提交
3147
  blockDataCleanup(pRes);
3148

S
slzhou 已提交
3149 3150 3151
  if (pOperator->status == OP_EXEC_DONE) {
    return NULL;
  }
S
slzhou 已提交
3152
  if (pInfo->readHandle.mnd != NULL) {
S
slzhou 已提交
3153
    return buildSysDbTableCount(pOperator, pInfo);
S
slzhou 已提交
3154
  }
S
slzhou 已提交
3155

S
slzhou 已提交
3156 3157 3158 3159 3160
  return buildVnodeDbTableCount(pOperator, pInfo, pSupp, pRes);
}

static SSDataBlock* buildVnodeDbTableCount(SOperatorInfo* pOperator, STableCountScanOperatorInfo* pInfo,
                                           STableCountScanSupp* pSupp, SSDataBlock* pRes) {
S
slzhou 已提交
3161 3162
  const char* db = NULL;
  int32_t     vgId = 0;
S
slzhou 已提交
3163
  char        dbName[TSDB_DB_NAME_LEN] = {0};
S
slzhou 已提交
3164

S
slzhou 已提交
3165 3166 3167 3168 3169 3170
  // get dbname
  vnodeGetInfo(pInfo->readHandle.vnode, &db, &vgId);
  SName sn = {0};
  tNameFromString(&sn, db, T_NAME_ACCT | T_NAME_DB);
  tNameGetDbName(&sn, dbName);

S
slzhou 已提交
3171
  if (pSupp->groupByDbName) {
S
slzhou 已提交
3172 3173 3174 3175 3176 3177 3178 3179 3180 3181 3182 3183 3184 3185
    buildVnodeGroupedTableCount(pOperator, pInfo, pSupp, pRes, vgId, dbName);
  } else {
    buildVnodeFilteredTbCount(pOperator, pInfo, pSupp, pRes, dbName);
  }
  return pRes->info.rows > 0 ? pRes : NULL;
}

static void buildVnodeGroupedTableCount(SOperatorInfo* pOperator, STableCountScanOperatorInfo* pInfo,
                                        STableCountScanSupp* pSupp, SSDataBlock* pRes, int32_t vgId, char* dbName) {
  if (pSupp->groupByStbName) {
    if (pInfo->stbUidList == NULL) {
      pInfo->stbUidList = taosArrayInit(16, sizeof(tb_uid_t));
      if (vnodeGetStbIdList(pInfo->readHandle.vnode, 0, pInfo->stbUidList) < 0) {
        qError("vgId:%d, failed to get stb id list error: %s", vgId, terrstr());
S
slzhou 已提交
3186
      }
S
slzhou 已提交
3187 3188 3189 3190 3191 3192 3193 3194 3195 3196
    }
    if (pInfo->currGrpIdx < taosArrayGetSize(pInfo->stbUidList)) {
      tb_uid_t stbUid = *(tb_uid_t*)taosArrayGet(pInfo->stbUidList, pInfo->currGrpIdx);
      buildVnodeGroupedStbTableCount(pInfo, pSupp, pRes, dbName, stbUid);

      pInfo->currGrpIdx++;
    } else if (pInfo->currGrpIdx == taosArrayGetSize(pInfo->stbUidList)) {
      buildVnodeGroupedNtbTableCount(pInfo, pSupp, pRes, dbName);

      pInfo->currGrpIdx++;
S
slzhou 已提交
3197
    } else {
S
slzhou 已提交
3198
      setOperatorCompleted(pOperator);
S
slzhou 已提交
3199 3200
    }
  } else {
S
slzhou 已提交
3201 3202 3203 3204 3205 3206 3207 3208 3209 3210 3211 3212 3213 3214 3215 3216 3217
    uint64_t groupId = calcGroupId(dbName, strlen(dbName));
    pRes->info.id.groupId = groupId;
    int64_t dbTableCount = metaGetTbNum(pInfo->readHandle.meta);
    fillTableCountScanDataBlock(pSupp, dbName, "", dbTableCount, pRes);
    setOperatorCompleted(pOperator);
  }
}

static void buildVnodeFilteredTbCount(SOperatorInfo* pOperator, STableCountScanOperatorInfo* pInfo,
                                      STableCountScanSupp* pSupp, SSDataBlock* pRes, char* dbName) {
  if (strlen(pSupp->dbNameFilter) != 0) {
    if (strlen(pSupp->stbNameFilter) != 0) {
      tb_uid_t      uid = metaGetTableEntryUidByName(pInfo->readHandle.meta, pSupp->stbNameFilter);
      SMetaStbStats stats = {0};
      metaGetStbStats(pInfo->readHandle.meta, uid, &stats);
      int64_t ctbNum = stats.ctbNum;
      fillTableCountScanDataBlock(pSupp, dbName, pSupp->stbNameFilter, ctbNum, pRes);
S
slzhou 已提交
3218 3219 3220
    } else {
      int64_t tbNumVnode = metaGetTbNum(pInfo->readHandle.meta);
      fillTableCountScanDataBlock(pSupp, dbName, "", tbNumVnode, pRes);
S
slzhou 已提交
3221
    }
S
slzhou 已提交
3222 3223 3224
  } else {
    int64_t tbNumVnode = metaGetTbNum(pInfo->readHandle.meta);
    fillTableCountScanDataBlock(pSupp, dbName, "", tbNumVnode, pRes);
S
slzhou 已提交
3225
  }
S
slzhou 已提交
3226 3227 3228 3229 3230 3231 3232 3233 3234 3235
  setOperatorCompleted(pOperator);
}

static void buildVnodeGroupedNtbTableCount(STableCountScanOperatorInfo* pInfo, STableCountScanSupp* pSupp,
                                           SSDataBlock* pRes, char* dbName) {
  char fullStbName[TSDB_TABLE_FNAME_LEN] = {0};
  snprintf(fullStbName, TSDB_TABLE_FNAME_LEN, "%s.%s", dbName, "");
  uint64_t groupId = calcGroupId(fullStbName, strlen(fullStbName));
  pRes->info.id.groupId = groupId;
  int64_t ntbNum = metaGetNtbNum(pInfo->readHandle.meta);
3236 3237 3238
  if (ntbNum != 0) {
    fillTableCountScanDataBlock(pSupp, dbName, "", ntbNum, pRes);
  }
S
slzhou 已提交
3239 3240 3241 3242 3243 3244 3245 3246 3247 3248 3249 3250 3251 3252 3253 3254 3255
}

static void buildVnodeGroupedStbTableCount(STableCountScanOperatorInfo* pInfo, STableCountScanSupp* pSupp,
                                           SSDataBlock* pRes, char* dbName, tb_uid_t stbUid) {
  char stbName[TSDB_TABLE_NAME_LEN] = {0};
  metaGetTableSzNameByUid(pInfo->readHandle.meta, stbUid, stbName);

  char fullStbName[TSDB_TABLE_FNAME_LEN] = {0};
  snprintf(fullStbName, TSDB_TABLE_FNAME_LEN, "%s.%s", dbName, stbName);
  uint64_t groupId = calcGroupId(fullStbName, strlen(fullStbName));
  pRes->info.id.groupId = groupId;

  SMetaStbStats stats = {0};
  metaGetStbStats(pInfo->readHandle.meta, stbUid, &stats);
  int64_t ctbNum = stats.ctbNum;

  fillTableCountScanDataBlock(pSupp, dbName, stbName, ctbNum, pRes);
S
shenglian zhou 已提交
3256 3257 3258
}

static void destoryTableCountScanOperator(void* param) {
S
slzhou 已提交
3259
  STableCountScanOperatorInfo* pTableCountScanInfo = param;
S
shenglian zhou 已提交
3260 3261
  blockDataDestroy(pTableCountScanInfo->pRes);

S
slzhou 已提交
3262
  taosArrayDestroy(pTableCountScanInfo->stbUidList);
S
shenglian zhou 已提交
3263 3264
  taosMemoryFreeClear(param);
}