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

X
Xiaoyu Wang 已提交
16
#include "parInsertUtil.h"
17 18

#include "catalog.h"
X
Xiaoyu Wang 已提交
19
#include "parInt.h"
X
Xiaoyu Wang 已提交
20
#include "parUtil.h"
X
Xiaoyu Wang 已提交
21
#include "querynodes.h"
C
Cary Xu 已提交
22
#include "tRealloc.h"
wmmhello's avatar
wmmhello 已提交
23
#include "tdatablock.h"
24

25
typedef struct SBlockKeyTuple {
X
Xiaoyu Wang 已提交
26 27
  TSKEY   skey;
  void*   payloadAddr;
wafwerar's avatar
wafwerar 已提交
28
  int16_t index;
29 30 31 32 33 34 35
} SBlockKeyTuple;

typedef struct SBlockKeyInfo {
  int32_t         maxBytesAlloc;
  SBlockKeyTuple* pKeyTuple;
} SBlockKeyInfo;

C
Cary Xu 已提交
36 37
typedef struct {
  int32_t   index;
C
Cary Xu 已提交
38
  SArray*   rowArray;  // array of merged rows(mem allocated by tRealloc/free by tFree)
C
Cary Xu 已提交
39
  STSchema* pSchema;
C
Cary Xu 已提交
40
  int64_t   tbUid;  // suid for child table, uid for normal table
C
Cary Xu 已提交
41 42
} SBlockRowMerger;

C
Cary Xu 已提交
43
static FORCE_INLINE void tdResetSBlockRowMerger(SBlockRowMerger* pMerger) {
C
Cary Xu 已提交
44 45 46 47 48 49 50 51 52 53 54 55 56 57
  if (pMerger) {
    pMerger->index = -1;
  }
}

static void tdFreeSBlockRowMerger(SBlockRowMerger* pMerger) {
  if (pMerger) {
    int32_t size = taosArrayGetSize(pMerger->rowArray);
    for (int32_t i = 0; i < size; ++i) {
      tFree(*(void**)taosArrayGet(pMerger->rowArray, i));
    }
    taosArrayDestroy(pMerger->rowArray);

    taosMemoryFreeClear(pMerger->pSchema);
C
Cary Xu 已提交
58
    taosMemoryFree(pMerger);
C
Cary Xu 已提交
59 60 61
  }
}

X
Xiaoyu Wang 已提交
62 63 64
static int32_t rowDataCompar(const void* lhs, const void* rhs) {
  TSKEY left = *(TSKEY*)lhs;
  TSKEY right = *(TSKEY*)rhs;
65 66 67 68 69 70 71
  if (left == right) {
    return 0;
  } else {
    return left > right ? 1 : -1;
  }
}

wafwerar's avatar
wafwerar 已提交
72 73 74 75 76 77 78 79 80 81
static int32_t rowDataComparStable(const void* lhs, const void* rhs) {
  TSKEY left = *(TSKEY*)lhs;
  TSKEY right = *(TSKEY*)rhs;
  if (left == right) {
    return ((SBlockKeyTuple*)lhs)->index - ((SBlockKeyTuple*)rhs)->index;
  } else {
    return left > right ? 1 : -1;
  }
}

X
Xiaoyu Wang 已提交
82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113
int32_t insGetExtendedRowSize(STableDataBlocks* pBlock) {
  STableComInfo* pTableInfo = &pBlock->pTableMeta->tableInfo;
  ASSERT(pBlock->rowSize == pTableInfo->rowSize);
  return pBlock->rowSize + TD_ROW_HEAD_LEN - sizeof(TSKEY) + pBlock->boundColumnInfo.extendedVarLen +
         (int32_t)TD_BITMAP_BYTES(pTableInfo->numOfColumns - 1);
}

void insGetSTSRowAppendInfo(uint8_t rowType, SParsedDataColInfo* spd, col_id_t idx, int32_t* toffset,
                            col_id_t* colIdx) {
  col_id_t schemaIdx = 0;
  if (IS_DATA_COL_ORDERED(spd)) {
    schemaIdx = spd->boundColumns[idx];
    if (TD_IS_TP_ROW_T(rowType)) {
      *toffset = (spd->cols + schemaIdx)->toffset;  // the offset of firstPart
      *colIdx = schemaIdx;
    } else {
      *toffset = idx * sizeof(SKvRowIdx);  // the offset of SKvRowIdx
      *colIdx = idx;
    }
  } else {
    ASSERT(idx == (spd->colIdxInfo + idx)->boundIdx);
    schemaIdx = (spd->colIdxInfo + idx)->schemaColIdx;
    if (TD_IS_TP_ROW_T(rowType)) {
      *toffset = (spd->cols + schemaIdx)->toffset;
      *colIdx = schemaIdx;
    } else {
      *toffset = ((spd->colIdxInfo + idx)->finalIdx) * sizeof(SKvRowIdx);
      *colIdx = (spd->colIdxInfo + idx)->finalIdx;
    }
  }
}

X
Xiaoyu Wang 已提交
114
int32_t insSetBlockInfo(SSubmitBlk* pBlocks, STableDataBlocks* dataBuf, int32_t numOfRows, SMsgBuf* pMsg) {
X
Xiaoyu Wang 已提交
115 116 117 118 119 120
  pBlocks->suid = (TSDB_NORMAL_TABLE == dataBuf->pTableMeta->tableType ? 0 : dataBuf->pTableMeta->suid);
  pBlocks->uid = dataBuf->pTableMeta->uid;
  pBlocks->sversion = dataBuf->pTableMeta->sversion;
  pBlocks->schemaLen = dataBuf->createTbReqLen;

  if (pBlocks->numOfRows + numOfRows >= INT32_MAX) {
X
Xiaoyu Wang 已提交
121
    return buildInvalidOperationMsg(pMsg, "too many rows in sql, total number of rows should be less than INT32_MAX");
X
Xiaoyu Wang 已提交
122
  }
X
Xiaoyu Wang 已提交
123 124
  pBlocks->numOfRows += numOfRows;
  return TSDB_CODE_SUCCESS;
X
Xiaoyu Wang 已提交
125 126 127
}

void insSetBoundColumnInfo(SParsedDataColInfo* pColList, SSchema* pSchema, col_id_t numOfCols) {
128 129 130
  pColList->numOfCols = numOfCols;
  pColList->numOfBound = numOfCols;
  pColList->orderStatus = ORDER_STATUS_ORDERED;  // default is ORDERED for non-bound mode
131
  pColList->boundColumns = taosMemoryCalloc(pColList->numOfCols, sizeof(col_id_t));
wafwerar's avatar
wafwerar 已提交
132
  pColList->cols = taosMemoryCalloc(pColList->numOfCols, sizeof(SBoundColumn));
133 134 135 136 137 138 139 140 141 142
  pColList->colIdxInfo = NULL;
  pColList->flen = 0;
  pColList->allNullLen = 0;

  int32_t nVar = 0;
  for (int32_t i = 0; i < pColList->numOfCols; ++i) {
    uint8_t type = pSchema[i].type;
    if (i > 0) {
      pColList->cols[i].offset = pColList->cols[i - 1].offset + pSchema[i - 1].bytes;
      pColList->cols[i].toffset = pColList->flen;
K
kailixu 已提交
143
      pColList->flen += TYPE_BYTES[type];
144 145 146 147 148 149 150 151 152 153 154 155 156
    }
    switch (type) {
      case TSDB_DATA_TYPE_BINARY:
        pColList->allNullLen += (VARSTR_HEADER_SIZE + CHAR_BYTES);
        ++nVar;
        break;
      case TSDB_DATA_TYPE_NCHAR:
        pColList->allNullLen += (VARSTR_HEADER_SIZE + TSDB_NCHAR_SIZE);
        ++nVar;
        break;
      default:
        break;
    }
X
Xiaoyu Wang 已提交
157
    pColList->boundColumns[i] = i;
158 159
  }
  pColList->allNullLen += pColList->flen;
C
Cary Xu 已提交
160
  pColList->boundNullLen = pColList->allNullLen;  // default set allNullLen
161 162 163
  pColList->extendedVarLen = (uint16_t)(nVar * sizeof(VarDataOffsetT));
}

X
Xiaoyu Wang 已提交
164
int32_t insSchemaIdxCompar(const void* lhs, const void* rhs) {
X
Xiaoyu Wang 已提交
165 166
  uint16_t left = *(uint16_t*)lhs;
  uint16_t right = *(uint16_t*)rhs;
167 168 169 170 171 172 173 174

  if (left == right) {
    return 0;
  } else {
    return left > right ? 1 : -1;
  }
}

X
Xiaoyu Wang 已提交
175
int32_t insBoundIdxCompar(const void* lhs, const void* rhs) {
X
Xiaoyu Wang 已提交
176 177
  uint16_t left = *(uint16_t*)POINTER_SHIFT(lhs, sizeof(uint16_t));
  uint16_t right = *(uint16_t*)POINTER_SHIFT(rhs, sizeof(uint16_t));
178 179 180 181 182 183 184 185

  if (left == right) {
    return 0;
  } else {
    return left > right ? 1 : -1;
  }
}

D
stmt  
dapan1121 已提交
186 187
void destroyBoundColumnInfo(void* pBoundInfo) {
  if (NULL == pBoundInfo) {
D
stmt  
dapan1121 已提交
188 189
    return;
  }
D
stmt  
dapan1121 已提交
190 191

  SParsedDataColInfo* pColList = (SParsedDataColInfo*)pBoundInfo;
X
Xiaoyu Wang 已提交
192

193
  taosMemoryFreeClear(pColList->boundColumns);
wafwerar's avatar
wafwerar 已提交
194 195
  taosMemoryFreeClear(pColList->cols);
  taosMemoryFreeClear(pColList->colIdxInfo);
196 197
}

198 199 200 201 202 203 204 205 206 207 208
void qDestroyBoundColInfo(void* pInfo) {
  if (NULL == pInfo) {
    return;
  }

  SBoundColInfo* pBoundInfo = (SBoundColInfo*)pInfo;

  taosMemoryFreeClear(pBoundInfo->pColIndex);
}


wmmhello's avatar
wmmhello 已提交
209
static int32_t createTableDataBlock(size_t defaultSize, int32_t rowSize, int32_t startOffset, STableMeta* pTableMeta,
X
Xiaoyu Wang 已提交
210
                               STableDataBlocks** dataBlocks) {
wafwerar's avatar
wafwerar 已提交
211
  STableDataBlocks* dataBuf = (STableDataBlocks*)taosMemoryCalloc(1, sizeof(STableDataBlocks));
212 213 214 215 216 217 218 219 220 221 222 223
  if (dataBuf == NULL) {
    return TSDB_CODE_TSC_OUT_OF_MEMORY;
  }

  dataBuf->nAllocSize = (uint32_t)defaultSize;
  dataBuf->headerSize = startOffset;

  // the header size will always be the startOffset value, reserved for the subumit block header
  if (dataBuf->nAllocSize <= dataBuf->headerSize) {
    dataBuf->nAllocSize = dataBuf->headerSize * 2;
  }

wafwerar's avatar
wafwerar 已提交
224
  dataBuf->pData = taosMemoryMalloc(dataBuf->nAllocSize);
225
  if (dataBuf->pData == NULL) {
wafwerar's avatar
wafwerar 已提交
226
    taosMemoryFreeClear(dataBuf);
227 228 229 230
    return TSDB_CODE_TSC_OUT_OF_MEMORY;
  }
  memset(dataBuf->pData, 0, sizeof(SSubmitBlk));

X
Xiaoyu Wang 已提交
231
  dataBuf->pTableMeta = tableMetaDup(pTableMeta);
232 233

  SParsedDataColInfo* pColInfo = &dataBuf->boundColumnInfo;
X
Xiaoyu Wang 已提交
234
  SSchema*            pSchema = getTableColumnSchema(dataBuf->pTableMeta);
X
Xiaoyu Wang 已提交
235
  insSetBoundColumnInfo(pColInfo, pSchema, dataBuf->pTableMeta->tableInfo.numOfColumns);
236

X
Xiaoyu Wang 已提交
237 238 239 240 241
  dataBuf->ordered = true;
  dataBuf->prevTS = INT64_MIN;
  dataBuf->rowSize = rowSize;
  dataBuf->size = startOffset;
  dataBuf->vgId = dataBuf->pTableMeta->vgId;
242 243 244 245 246 247 248

  assert(defaultSize > 0 && pTableMeta != NULL && dataBuf->pTableMeta != NULL);

  *dataBlocks = dataBuf;
  return TSDB_CODE_SUCCESS;
}

X
Xiaoyu Wang 已提交
249
int32_t insBuildCreateTbMsg(STableDataBlocks* pBlocks, SVCreateTbReq* pCreateTbReq) {
H
Hongze Cheng 已提交
250
  SEncoder coder = {0};
X
Xiaoyu Wang 已提交
251 252
  char*    pBuf;
  int32_t  len;
H
Hongze Cheng 已提交
253

wafwerar's avatar
wafwerar 已提交
254 255
  int32_t ret = 0;
  tEncodeSize(tEncodeSVCreateTbReq, pCreateTbReq, len, ret);
X
Xiaoyu Wang 已提交
256 257 258 259 260 261 262 263 264 265 266
  if (pBlocks->nAllocSize - pBlocks->size < len) {
    pBlocks->nAllocSize += len + pBlocks->rowSize;
    char* pTmp = taosMemoryRealloc(pBlocks->pData, pBlocks->nAllocSize);
    if (pTmp != NULL) {
      pBlocks->pData = pTmp;
      memset(pBlocks->pData + pBlocks->size, 0, pBlocks->nAllocSize - pBlocks->size);
    } else {
      pBlocks->nAllocSize -= len + pBlocks->rowSize;
      return TSDB_CODE_TSC_OUT_OF_MEMORY;
    }
  }
H
Hongze Cheng 已提交
267

X
Xiaoyu Wang 已提交
268
  pBuf = pBlocks->pData + pBlocks->size;
H
Hongze Cheng 已提交
269

H
Hongze Cheng 已提交
270
  tEncoderInit(&coder, pBuf, len);
X
Xiaoyu Wang 已提交
271
  int32_t code = tEncodeSVCreateTbReq(&coder, pCreateTbReq);
H
Hongze Cheng 已提交
272
  tEncoderClear(&coder);
X
Xiaoyu Wang 已提交
273 274
  pBlocks->size += len;
  pBlocks->createTbReqLen = len;
X
Xiaoyu Wang 已提交
275 276

  return code;
X
Xiaoyu Wang 已提交
277 278
}

X
Xiaoyu Wang 已提交
279
void insDestroyDataBlock(STableDataBlocks* pDataBlock) {
X
Xiaoyu Wang 已提交
280 281 282 283 284 285 286 287 288 289
  if (pDataBlock == NULL) {
    return;
  }

  taosMemoryFreeClear(pDataBlock->pData);
  taosMemoryFreeClear(pDataBlock->pTableMeta);
  destroyBoundColumnInfo(&pDataBlock->boundColumnInfo);
  taosMemoryFreeClear(pDataBlock);
}

X
Xiaoyu Wang 已提交
290 291 292
int32_t insGetDataBlockFromList(SHashObj* pHashList, void* id, int32_t idLen, int32_t size, int32_t startOffset,
                                int32_t rowSize, STableMeta* pTableMeta, STableDataBlocks** dataBlocks,
                                SArray* pBlockList, SVCreateTbReq* pCreateTbReq) {
293
  *dataBlocks = NULL;
D
dapan 已提交
294
  STableDataBlocks** t1 = (STableDataBlocks**)taosHashGet(pHashList, (const char*)id, idLen);
295 296 297 298 299
  if (t1 != NULL) {
    *dataBlocks = *t1;
  }

  if (*dataBlocks == NULL) {
wmmhello's avatar
wmmhello 已提交
300
    int32_t ret = createTableDataBlock((size_t)size, rowSize, startOffset, pTableMeta, dataBlocks);
301 302 303 304
    if (ret != TSDB_CODE_SUCCESS) {
      return ret;
    }

H
Hongze Cheng 已提交
305
    if (NULL != pCreateTbReq && NULL != pCreateTbReq->ctb.pTag) {
X
Xiaoyu Wang 已提交
306
      ret = insBuildCreateTbMsg(*dataBlocks, pCreateTbReq);
X
Xiaoyu Wang 已提交
307
      if (ret != TSDB_CODE_SUCCESS) {
X
Xiaoyu Wang 已提交
308
        insDestroyDataBlock(*dataBlocks);
X
Xiaoyu Wang 已提交
309 310 311 312
        return ret;
      }
    }

X
Xiaoyu Wang 已提交
313 314
    // converting to 'const char*' is to handle coverity scan errors
    taosHashPut(pHashList, (const char*)id, idLen, (const char*)dataBlocks, POINTER_BYTES);
315 316 317 318 319 320 321 322
    if (pBlockList) {
      taosArrayPush(pBlockList, dataBlocks);
    }
  }

  return TSDB_CODE_SUCCESS;
}

X
Xiaoyu Wang 已提交
323
void insDestroyBlockArrayList(SArray* pDataBlockList) {
324
  if (pDataBlockList == NULL) {
325
    return;
326 327 328 329
  }

  size_t size = taosArrayGetSize(pDataBlockList);
  for (int32_t i = 0; i < size; i++) {
330
    void* p = taosArrayGetP(pDataBlockList, i);
X
Xiaoyu Wang 已提交
331
    insDestroyDataBlock(p);
332 333 334 335 336
  }

  taosArrayDestroy(pDataBlockList);
}

X
Xiaoyu Wang 已提交
337
void insDestroyBlockHashmap(SHashObj* pDataBlockHash) {
338 339 340 341 342 343
  if (pDataBlockHash == NULL) {
    return;
  }

  void** p1 = taosHashIterate(pDataBlockHash, NULL);
  while (p1) {
wmmhello's avatar
wmmhello 已提交
344 345
    SBoundColInfo* pBlocks = *p1;
    destroyBoundColInfo(pBlocks);
346 347 348 349 350 351 352

    p1 = taosHashIterate(pDataBlockHash, p1);
  }

  taosHashCleanup(pDataBlockHash);
}

353
// data block is disordered, sort it in ascending order
C
Cary Xu 已提交
354
static int sortRemoveDataBlockDupRows(STableDataBlocks* dataBuf, SBlockKeyInfo* pBlkKeyInfo) {
X
Xiaoyu Wang 已提交
355
  SSubmitBlk* pBlocks = (SSubmitBlk*)dataBuf->pData;
356 357 358 359 360 361 362
  int16_t     nRows = pBlocks->numOfRows;

  // size is less than the total size, since duplicated rows may be removed yet.

  // allocate memory
  size_t nAlloc = nRows * sizeof(SBlockKeyTuple);
  if (pBlkKeyInfo->pKeyTuple == NULL || pBlkKeyInfo->maxBytesAlloc < nAlloc) {
X
Xiaoyu Wang 已提交
363
    char* tmp = taosMemoryRealloc(pBlkKeyInfo->pKeyTuple, nAlloc);
364 365 366
    if (tmp == NULL) {
      return TSDB_CODE_TSC_OUT_OF_MEMORY;
    }
X
Xiaoyu Wang 已提交
367
    pBlkKeyInfo->pKeyTuple = (SBlockKeyTuple*)tmp;
368 369 370 371
    pBlkKeyInfo->maxBytesAlloc = (int32_t)nAlloc;
  }
  memset(pBlkKeyInfo->pKeyTuple, 0, nAlloc);

X
Xiaoyu Wang 已提交
372
  int32_t         extendedRowSize = insGetExtendedRowSize(dataBuf);
X
Xiaoyu Wang 已提交
373 374
  SBlockKeyTuple* pBlkKeyTuple = pBlkKeyInfo->pKeyTuple;
  char*           pBlockData = pBlocks->data + pBlocks->schemaLen;
375 376
  int             n = 0;
  while (n < nRows) {
X
Xiaoyu Wang 已提交
377
    pBlkKeyTuple->skey = TD_ROW_KEY((STSRow*)pBlockData);
378
    pBlkKeyTuple->payloadAddr = pBlockData;
wafwerar's avatar
wafwerar 已提交
379
    pBlkKeyTuple->index = n;
380 381 382 383 384 385 386 387 388

    // next loop
    pBlockData += extendedRowSize;
    ++pBlkKeyTuple;
    ++n;
  }

  if (!dataBuf->ordered) {
    pBlkKeyTuple = pBlkKeyInfo->pKeyTuple;
389 390

    // todo. qsort is unstable, if timestamp is same, should get the last one
wafwerar's avatar
wafwerar 已提交
391
    taosSort(pBlkKeyTuple, nRows, sizeof(SBlockKeyTuple), rowDataComparStable);
392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421

    pBlkKeyTuple = pBlkKeyInfo->pKeyTuple;
    int32_t i = 0;
    int32_t j = 1;
    while (j < nRows) {
      TSKEY ti = (pBlkKeyTuple + i)->skey;
      TSKEY tj = (pBlkKeyTuple + j)->skey;

      if (ti == tj) {
        ++j;
        continue;
      }

      int32_t nextPos = (++i);
      if (nextPos != j) {
        memmove(pBlkKeyTuple + nextPos, pBlkKeyTuple + j, sizeof(SBlockKeyTuple));
      }
      ++j;
    }

    dataBuf->ordered = true;
    pBlocks->numOfRows = i + 1;
  }

  dataBuf->size = sizeof(SSubmitBlk) + pBlocks->numOfRows * extendedRowSize;
  dataBuf->prevTS = INT64_MIN;

  return 0;
}

C
Cary Xu 已提交
422 423 424 425 426 427 428 429
static void* tdGetCurRowFromBlockMerger(SBlockRowMerger* pBlkRowMerger) {
  if (pBlkRowMerger && (pBlkRowMerger->index >= 0)) {
    ASSERT(pBlkRowMerger->index < taosArrayGetSize(pBlkRowMerger->rowArray));
    return *(void**)taosArrayGet(pBlkRowMerger->rowArray, pBlkRowMerger->index);
  }
  return NULL;
}

C
Cary Xu 已提交
430
static int32_t tdBlockRowMerge(STableMeta* pTableMeta, SBlockKeyTuple* pEndKeyTp, int32_t nDupRows,
C
Cary Xu 已提交
431 432 433 434 435 436 437
                               SBlockRowMerger** pBlkRowMerger, int32_t rowSize) {
  ASSERT(nDupRows > 1);
  SBlockKeyTuple* pStartKeyTp = pEndKeyTp - (nDupRows - 1);
  ASSERT(pStartKeyTp->skey == pEndKeyTp->skey);

  // TODO: optimization if end row is all normal
#if 0
C
Cary Xu 已提交
438
  STSRow* pEndRow = (STSRow*)pEndKeyTp->payloadAddr;
C
Cary Xu 已提交
439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460
  if(isNormal(pEndRow)) { // set the end row if it is normal and return directly
    pStartKeyTp->payloadAddr = pEndKeyTp->payloadAddr;
    return TSDB_CODE_SUCCESS;
  }
#endif

  if (!(*pBlkRowMerger)) {
    (*pBlkRowMerger) = taosMemoryCalloc(1, sizeof(**pBlkRowMerger));
    if (!(*pBlkRowMerger)) {
      terrno = TSDB_CODE_OUT_OF_MEMORY;
      return TSDB_CODE_FAILED;
    }
    (*pBlkRowMerger)->index = -1;
    if (!(*pBlkRowMerger)->rowArray) {
      (*pBlkRowMerger)->rowArray = taosArrayInit(1, sizeof(void*));
      if (!(*pBlkRowMerger)->rowArray) {
        terrno = TSDB_CODE_OUT_OF_MEMORY;
        return TSDB_CODE_FAILED;
      }
    }
  }

C
Cary Xu 已提交
461 462 463 464 465 466 467 468 469 470
  if ((*pBlkRowMerger)->pSchema) {
    if ((*pBlkRowMerger)->pSchema->version != pTableMeta->sversion) {
      taosMemoryFreeClear((*pBlkRowMerger)->pSchema);
    } else {
      if ((*pBlkRowMerger)->tbUid != (pTableMeta->suid > 0 ? pTableMeta->suid : pTableMeta->uid)) {
        taosMemoryFreeClear((*pBlkRowMerger)->pSchema);
      }
    }
  }

C
Cary Xu 已提交
471
  if (!(*pBlkRowMerger)->pSchema) {
C
Cary Xu 已提交
472
    (*pBlkRowMerger)->pSchema =
H
Hongze Cheng 已提交
473
        tBuildTSchema(pTableMeta->schema, pTableMeta->tableInfo.numOfColumns, pTableMeta->sversion);
C
Cary Xu 已提交
474 475 476 477 478

    if (!(*pBlkRowMerger)->pSchema) {
      terrno = TSDB_CODE_OUT_OF_MEMORY;
      return TSDB_CODE_FAILED;
    }
C
Cary Xu 已提交
479
    (*pBlkRowMerger)->tbUid = pTableMeta->suid > 0 ? pTableMeta->suid : pTableMeta->uid;
C
Cary Xu 已提交
480 481 482 483 484
  }

  void* pDestRow = NULL;
  ++((*pBlkRowMerger)->index);
  if ((*pBlkRowMerger)->index < taosArrayGetSize((*pBlkRowMerger)->rowArray)) {
485 486
    void** pAlloc = (void**)taosArrayGet((*pBlkRowMerger)->rowArray, (*pBlkRowMerger)->index);
    if (tRealloc((uint8_t**)pAlloc, rowSize) != 0) {
C
Cary Xu 已提交
487 488
      return TSDB_CODE_FAILED;
    }
489
    pDestRow = *pAlloc;
C
Cary Xu 已提交
490 491 492 493 494 495 496 497 498 499 500 501
  } else {
    if (tRealloc((uint8_t**)&pDestRow, rowSize) != 0) {
      return TSDB_CODE_FAILED;
    }
    taosArrayPush((*pBlkRowMerger)->rowArray, &pDestRow);
  }

  // merge rows to pDestRow
  STSchema* pSchema = (*pBlkRowMerger)->pSchema;
  SArray*   pArray = taosArrayInit(pSchema->numOfCols, sizeof(SColVal));
  for (int32_t i = 0; i < pSchema->numOfCols; ++i) {
    SColVal colVal = {0};
C
Cary Xu 已提交
502
    for (int32_t j = 0; j < nDupRows; ++j) {
C
Cary Xu 已提交
503
      tTSRowGetVal((pEndKeyTp - j)->payloadAddr, pSchema, i, &colVal);
H
Hongze Cheng 已提交
504
      if (!COL_VAL_IS_NONE(&colVal)) {
C
Cary Xu 已提交
505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522
        break;
      }
    }
    taosArrayPush(pArray, &colVal);
  }
  if (tdSTSRowNew(pArray, pSchema, (STSRow**)&pDestRow) < 0) {
    taosArrayDestroy(pArray);
    return TSDB_CODE_FAILED;
  }

  taosArrayDestroy(pArray);
  return TSDB_CODE_SUCCESS;
}

// data block is disordered, sort it in ascending order, and merge dup rows if exists
static int sortMergeDataBlockDupRows(STableDataBlocks* dataBuf, SBlockKeyInfo* pBlkKeyInfo,
                                     SBlockRowMerger** ppBlkRowMerger) {
  SSubmitBlk* pBlocks = (SSubmitBlk*)dataBuf->pData;
C
Cary Xu 已提交
523
  STableMeta* pTableMeta = dataBuf->pTableMeta;
524
  int32_t     nRows = pBlocks->numOfRows;
C
Cary Xu 已提交
525 526 527 528 529 530 531 532 533 534 535 536 537 538 539 540 541

  // size is less than the total size, since duplicated rows may be removed.

  // allocate memory
  size_t nAlloc = nRows * sizeof(SBlockKeyTuple);
  if (pBlkKeyInfo->pKeyTuple == NULL || pBlkKeyInfo->maxBytesAlloc < nAlloc) {
    char* tmp = taosMemoryRealloc(pBlkKeyInfo->pKeyTuple, nAlloc);
    if (tmp == NULL) {
      return TSDB_CODE_TSC_OUT_OF_MEMORY;
    }
    pBlkKeyInfo->pKeyTuple = (SBlockKeyTuple*)tmp;
    pBlkKeyInfo->maxBytesAlloc = (int32_t)nAlloc;
  }
  memset(pBlkKeyInfo->pKeyTuple, 0, nAlloc);

  tdResetSBlockRowMerger(*ppBlkRowMerger);

X
Xiaoyu Wang 已提交
542
  int32_t         extendedRowSize = insGetExtendedRowSize(dataBuf);
C
Cary Xu 已提交
543 544
  SBlockKeyTuple* pBlkKeyTuple = pBlkKeyInfo->pKeyTuple;
  char*           pBlockData = pBlocks->data + pBlocks->schemaLen;
545
  int32_t         n = 0;
C
Cary Xu 已提交
546 547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575 576 577
  while (n < nRows) {
    pBlkKeyTuple->skey = TD_ROW_KEY((STSRow*)pBlockData);
    pBlkKeyTuple->payloadAddr = pBlockData;
    pBlkKeyTuple->index = n;

    // next loop
    pBlockData += extendedRowSize;
    ++pBlkKeyTuple;
    ++n;
  }

  if (!dataBuf->ordered) {
    pBlkKeyTuple = pBlkKeyInfo->pKeyTuple;

    taosSort(pBlkKeyTuple, nRows, sizeof(SBlockKeyTuple), rowDataComparStable);

    pBlkKeyTuple = pBlkKeyInfo->pKeyTuple;
    bool    hasDup = false;
    int32_t nextPos = 0;
    int32_t i = 0;
    int32_t j = 1;

    while (j < nRows) {
      TSKEY ti = (pBlkKeyTuple + i)->skey;
      TSKEY tj = (pBlkKeyTuple + j)->skey;

      if (ti == tj) {
        ++j;
        continue;
      }

      if ((j - i) > 1) {
C
Cary Xu 已提交
578
        if (tdBlockRowMerge(pTableMeta, (pBlkKeyTuple + j - 1), j - i, ppBlkRowMerger, extendedRowSize) < 0) {
C
Cary Xu 已提交
579 580 581
          return TSDB_CODE_FAILED;
        }
        (pBlkKeyTuple + nextPos)->payloadAddr = tdGetCurRowFromBlockMerger(*ppBlkRowMerger);
C
Cary Xu 已提交
582 583 584
        if (!hasDup) {
          hasDup = true;
        }
C
Cary Xu 已提交
585 586 587 588 589 590 591 592 593 594 595 596 597 598
        i = j;
      } else {
        if (hasDup) {
          memmove(pBlkKeyTuple + nextPos, pBlkKeyTuple + i, sizeof(SBlockKeyTuple));
        }
        ++i;
      }

      ++nextPos;
      ++j;
    }

    if ((j - i) > 1) {
      ASSERT((pBlkKeyTuple + i)->skey == (pBlkKeyTuple + j - 1)->skey);
C
Cary Xu 已提交
599
      if (tdBlockRowMerge(pTableMeta, (pBlkKeyTuple + j - 1), j - i, ppBlkRowMerger, extendedRowSize) < 0) {
C
Cary Xu 已提交
600 601 602 603 604 605 606 607
        return TSDB_CODE_FAILED;
      }
      (pBlkKeyTuple + nextPos)->payloadAddr = tdGetCurRowFromBlockMerger(*ppBlkRowMerger);
    } else if (hasDup) {
      memmove(pBlkKeyTuple + nextPos, pBlkKeyTuple + i, sizeof(SBlockKeyTuple));
    }

    dataBuf->ordered = true;
C
Cary Xu 已提交
608
    pBlocks->numOfRows = nextPos + 1;
C
Cary Xu 已提交
609 610 611 612 613 614 615 616
  }

  dataBuf->size = sizeof(SSubmitBlk) + pBlocks->numOfRows * extendedRowSize;
  dataBuf->prevTS = INT64_MIN;

  return TSDB_CODE_SUCCESS;
}

617
// Erase the empty space reserved for binary data
618
static int trimDataBlock(void* pDataBlock, STableDataBlocks* pTableDataBlock, SBlockKeyTuple* blkKeyTuple) {
619
  // TODO: optimize this function, handle the case while binary is not presented
X
Xiaoyu Wang 已提交
620
  int32_t     nonDataLen = sizeof(SSubmitBlk) + pTableDataBlock->createTbReqLen;
621
  SSubmitBlk* pBlock = pDataBlock;
X
Xiaoyu Wang 已提交
622 623
  memcpy(pDataBlock, pTableDataBlock->pData, nonDataLen);
  pDataBlock = (char*)pDataBlock + nonDataLen;
624

X
Xiaoyu Wang 已提交
625
  pBlock->schemaLen = pTableDataBlock->createTbReqLen;
626 627
  pBlock->dataLen = 0;

628 629 630 631 632 633 634
  int32_t numOfRows = pBlock->numOfRows;
  for (int32_t i = 0; i < numOfRows; ++i) {
    void*     payload = (blkKeyTuple + i)->payloadAddr;
    TDRowLenT rowTLen = TD_ROW_LEN((STSRow*)payload);
    memcpy(pDataBlock, payload, rowTLen);
    pDataBlock = POINTER_SHIFT(pDataBlock, rowTLen);
    pBlock->dataLen += rowTLen;
635 636
  }

637
  return pBlock->dataLen + pBlock->schemaLen;
638 639
}

640
int32_t insMergeTableDataBlocks(SHashObj* pHashObj, SArray** pVgDataBlocks) {
S
Shengliang Guan 已提交
641
  const int INSERT_HEAD_SIZE = sizeof(SSubmitReq);
642
  int       code = 0;
D
dapan 已提交
643
  SHashObj* pVnodeDataBlockHashList = taosHashInit(128, taosGetDefaultHashFunction(TSDB_DATA_TYPE_INT), true, false);
644 645 646
  SArray*   pVnodeDataBlockList = taosArrayInit(8, POINTER_BYTES);

  STableDataBlocks** p = taosHashIterate(pHashObj, NULL);
X
Xiaoyu Wang 已提交
647 648
  STableDataBlocks*  pOneTableBlock = *p;
  SBlockKeyInfo      blkKeyInfo = {0};  // share by pOneTableBlock
X
Xiaoyu Wang 已提交
649
  SBlockRowMerger*   pBlkRowMerger = NULL;
C
Cary Xu 已提交
650

651
  while (pOneTableBlock) {
X
Xiaoyu Wang 已提交
652
    SSubmitBlk* pBlocks = (SSubmitBlk*)pOneTableBlock->pData;
653 654
    if (pBlocks->numOfRows > 0) {
      STableDataBlocks* dataBuf = NULL;
X
Xiaoyu Wang 已提交
655
      pOneTableBlock->pTableMeta->vgId = pOneTableBlock->vgId;  // for schemaless, restore origin vgId
X
Xiaoyu Wang 已提交
656 657 658
      int32_t ret = insGetDataBlockFromList(pVnodeDataBlockHashList, &pOneTableBlock->vgId,
                                            sizeof(pOneTableBlock->vgId), TSDB_PAYLOAD_SIZE, INSERT_HEAD_SIZE, 0,
                                            pOneTableBlock->pTableMeta, &dataBuf, pVnodeDataBlockList, NULL);
659
      if (ret != TSDB_CODE_SUCCESS) {
C
Cary Xu 已提交
660
        tdFreeSBlockRowMerger(pBlkRowMerger);
661
        taosHashCleanup(pVnodeDataBlockHashList);
X
Xiaoyu Wang 已提交
662
        insDestroyBlockArrayList(pVnodeDataBlockList);
wafwerar's avatar
wafwerar 已提交
663
        taosMemoryFreeClear(blkKeyInfo.pKeyTuple);
664 665
        return ret;
      }
X
Xiaoyu Wang 已提交
666
      ASSERT(pOneTableBlock->pTableMeta->tableInfo.rowSize > 0);
667
      // the maximum expanded size in byte when a row-wise data is converted to SDataRow format
668
      int64_t destSize = dataBuf->size + pOneTableBlock->size +
X
Xiaoyu Wang 已提交
669 670
                         sizeof(STColumn) * getNumOfColumns(pOneTableBlock->pTableMeta) +
                         pOneTableBlock->createTbReqLen;
671 672 673

      if (dataBuf->nAllocSize < destSize) {
        dataBuf->nAllocSize = (uint32_t)(destSize * 1.5);
wafwerar's avatar
wafwerar 已提交
674
        char* tmp = taosMemoryRealloc(dataBuf->pData, dataBuf->nAllocSize);
675 676 677
        if (tmp != NULL) {
          dataBuf->pData = tmp;
        } else {  // failed to allocate memory, free already allocated memory and return error code
C
Cary Xu 已提交
678
          tdFreeSBlockRowMerger(pBlkRowMerger);
679
          taosHashCleanup(pVnodeDataBlockHashList);
X
Xiaoyu Wang 已提交
680
          insDestroyBlockArrayList(pVnodeDataBlockList);
wafwerar's avatar
wafwerar 已提交
681 682
          taosMemoryFreeClear(dataBuf->pData);
          taosMemoryFreeClear(blkKeyInfo.pKeyTuple);
683 684 685 686
          return TSDB_CODE_TSC_OUT_OF_MEMORY;
        }
      }

687 688 689 690 691 692 693
      if ((code = sortMergeDataBlockDupRows(pOneTableBlock, &blkKeyInfo, &pBlkRowMerger)) != 0) {
        tdFreeSBlockRowMerger(pBlkRowMerger);
        taosHashCleanup(pVnodeDataBlockHashList);
        insDestroyBlockArrayList(pVnodeDataBlockList);
        taosMemoryFreeClear(dataBuf->pData);
        taosMemoryFreeClear(blkKeyInfo.pKeyTuple);
        return code;
694
      }
695
      ASSERT(blkKeyInfo.pKeyTuple != NULL && pBlocks->numOfRows > 0);
696 697

      // erase the empty space reserved for binary data
698
      int32_t finalLen = trimDataBlock(dataBuf->pData + dataBuf->size, pOneTableBlock, blkKeyInfo.pKeyTuple);
699 700 701 702 703 704 705 706 707 708 709 710 711 712 713

      dataBuf->size += (finalLen + sizeof(SSubmitBlk));
      assert(dataBuf->size <= dataBuf->nAllocSize);
      dataBuf->numOfTables += 1;
    }

    p = taosHashIterate(pHashObj, p);
    if (p == NULL) {
      break;
    }

    pOneTableBlock = *p;
  }

  // free the table data blocks;
C
Cary Xu 已提交
714
  tdFreeSBlockRowMerger(pBlkRowMerger);
715
  taosHashCleanup(pVnodeDataBlockHashList);
wafwerar's avatar
wafwerar 已提交
716
  taosMemoryFreeClear(blkKeyInfo.pKeyTuple);
717
  *pVgDataBlocks = pVnodeDataBlockList;
718 719 720
  return TSDB_CODE_SUCCESS;
}

X
Xiaoyu Wang 已提交
721
int32_t insAllocateMemForSize(STableDataBlocks* pDataBlock, int32_t allSize) {
X
Xiaoyu Wang 已提交
722
  size_t   remain = pDataBlock->nAllocSize - pDataBlock->size;
D
stmt  
dapan1121 已提交
723
  uint32_t nAllocSizeOld = pDataBlock->nAllocSize;
X
Xiaoyu Wang 已提交
724

D
stmt  
dapan1121 已提交
725 726 727 728
  // expand the allocated size
  if (remain < allSize) {
    pDataBlock->nAllocSize = (pDataBlock->size + allSize) * 1.5;

X
Xiaoyu Wang 已提交
729
    char* tmp = taosMemoryRealloc(pDataBlock->pData, (size_t)pDataBlock->nAllocSize);
D
stmt  
dapan1121 已提交
730 731 732 733 734 735 736 737 738 739 740 741 742
    if (tmp != NULL) {
      pDataBlock->pData = tmp;
      memset(pDataBlock->pData + pDataBlock->size, 0, pDataBlock->nAllocSize - pDataBlock->size);
    } else {
      // do nothing, if allocate more memory failed
      pDataBlock->nAllocSize = nAllocSizeOld;
      return TSDB_CODE_TSC_OUT_OF_MEMORY;
    }
  }

  return TSDB_CODE_SUCCESS;
}

X
Xiaoyu Wang 已提交
743 744 745 746 747 748 749
int32_t insInitRowBuilder(SRowBuilder* pBuilder, int16_t schemaVer, SParsedDataColInfo* pColInfo) {
  ASSERT(pColInfo->numOfCols > 0 && (pColInfo->numOfBound <= pColInfo->numOfCols));
  tdSRowInit(pBuilder, schemaVer);
  tdSRowSetExtendedInfo(pBuilder, pColInfo->numOfCols, pColInfo->numOfBound, pColInfo->flen, pColInfo->allNullLen,
                        pColInfo->boundNullLen);
  return TSDB_CODE_SUCCESS;
}
X
Xiaoyu Wang 已提交
750

X
Xiaoyu Wang 已提交
751 752 753 754 755 756 757 758
static char* tableNameGetPosition(SToken* pToken, char target) {
  bool inEscape = false;
  bool inQuote = false;
  char quotaStr = 0;

  for (uint32_t i = 0; i < pToken->n; ++i) {
    if (*(pToken->z + i) == target && (!inEscape) && (!inQuote)) {
      return pToken->z + i;
759 760
    }

X
Xiaoyu Wang 已提交
761 762 763 764 765 766 767 768 769 770 771 772 773 774 775
    if (*(pToken->z + i) == TS_ESCAPE_CHAR) {
      if (!inQuote) {
        inEscape = !inEscape;
      }
    }

    if (*(pToken->z + i) == '\'' || *(pToken->z + i) == '"') {
      if (!inEscape) {
        if (!inQuote) {
          quotaStr = *(pToken->z + i);
          inQuote = !inQuote;
        } else if (quotaStr == *(pToken->z + i)) {
          inQuote = !inQuote;
        }
      }
776 777 778
    }
  }

X
Xiaoyu Wang 已提交
779
  return NULL;
780 781
}

X
Xiaoyu Wang 已提交
782 783 784 785 786 787 788 789
int32_t insCreateSName(SName* pName, SToken* pTableName, int32_t acctId, const char* dbName, SMsgBuf* pMsgBuf) {
  const char* msg1 = "name too long";
  const char* msg2 = "invalid database name";
  const char* msg3 = "db is not specified";
  const char* msg4 = "invalid table name";

  int32_t code = TSDB_CODE_SUCCESS;
  char*   p = tableNameGetPosition(pTableName, TS_PATH_DELIMITER[0]);
D
stmt  
dapan1121 已提交
790

X
Xiaoyu Wang 已提交
791 792
  if (p != NULL) {  // db has been specified in sql string so we ignore current db path
    assert(*p == TS_PATH_DELIMITER[0]);
D
stmt  
dapan1121 已提交
793

X
Xiaoyu Wang 已提交
794 795 796
    int32_t dbLen = p - pTableName->z;
    if (dbLen <= 0) {
      return buildInvalidOperationMsg(pMsgBuf, msg2);
D
stmt  
dapan1121 已提交
797
    }
X
Xiaoyu Wang 已提交
798 799 800
    char name[TSDB_DB_FNAME_LEN] = {0};
    strncpy(name, pTableName->z, dbLen);
    int32_t actualDbLen = strdequote(name);
X
Xiaoyu Wang 已提交
801

X
Xiaoyu Wang 已提交
802 803 804 805
    code = tNameSetDbName(pName, acctId, name, actualDbLen);
    if (code != TSDB_CODE_SUCCESS) {
      return buildInvalidOperationMsg(pMsgBuf, msg1);
    }
X
Xiaoyu Wang 已提交
806

X
Xiaoyu Wang 已提交
807 808 809 810
    int32_t tbLen = pTableName->n - dbLen - 1;
    if (tbLen <= 0) {
      return buildInvalidOperationMsg(pMsgBuf, msg4);
    }
D
stmt  
dapan1121 已提交
811

X
Xiaoyu Wang 已提交
812 813 814
    char tbname[TSDB_TABLE_FNAME_LEN] = {0};
    strncpy(tbname, p + 1, tbLen);
    /*tbLen = */ strdequote(tbname);
D
stmt  
dapan1121 已提交
815

X
Xiaoyu Wang 已提交
816 817 818 819 820 821 822 823 824 825 826 827 828 829 830 831 832 833
    code = tNameFromString(pName, tbname, T_NAME_TABLE);
    if (code != 0) {
      return buildInvalidOperationMsg(pMsgBuf, msg1);
    }
  } else {  // get current DB name first, and then set it into path
    if (pTableName->n >= TSDB_TABLE_NAME_LEN) {
      return buildInvalidOperationMsg(pMsgBuf, msg1);
    }

    assert(pTableName->n < TSDB_TABLE_FNAME_LEN);

    char name[TSDB_TABLE_FNAME_LEN] = {0};
    strncpy(name, pTableName->z, pTableName->n);
    strdequote(name);

    if (dbName == NULL) {
      return buildInvalidOperationMsg(pMsgBuf, msg3);
    }
X
Xiaoyu Wang 已提交
834

X
Xiaoyu Wang 已提交
835 836 837 838 839
    code = tNameSetDbName(pName, acctId, dbName, strlen(dbName));
    if (code != TSDB_CODE_SUCCESS) {
      code = buildInvalidOperationMsg(pMsgBuf, msg2);
      return code;
    }
X
Xiaoyu Wang 已提交
840

X
Xiaoyu Wang 已提交
841 842 843
    code = tNameFromString(pName, name, T_NAME_TABLE);
    if (code != 0) {
      code = buildInvalidOperationMsg(pMsgBuf, msg1);
D
dapan1121 已提交
844 845 846
    }
  }

X
Xiaoyu Wang 已提交
847 848 849 850
  if (NULL != strchr(pName->tname, '.')) {
    code = generateSyntaxErrMsgExt(pMsgBuf, TSDB_CODE_PAR_INVALID_IDENTIFIER_NAME, "The table name cannot contain '.'");
  }

X
Xiaoyu Wang 已提交
851
  return code;
D
stmt  
dapan1121 已提交
852 853
}

D
dapan1121 已提交
854
int16_t insFindCol(SToken* pColname, int16_t start, int16_t end, SSchema* pSchema) {
X
Xiaoyu Wang 已提交
855 856 857 858 859
  while (start < end) {
    if (strlen(pSchema[start].name) == pColname->n && strncmp(pColname->z, pSchema[start].name, pColname->n) == 0) {
      return start;
    }
    ++start;
D
stmt  
dapan1121 已提交
860
  }
X
Xiaoyu Wang 已提交
861 862
  return -1;
}
D
stmt  
dapan1121 已提交
863

X
Xiaoyu Wang 已提交
864
void insBuildCreateTbReq(SVCreateTbReq* pTbReq, const char* tname, STag* pTag, int64_t suid, const char* sname,
865
                         SArray* tagName, uint8_t tagNum, int32_t ttl) {
X
Xiaoyu Wang 已提交
866 867 868 869 870 871
  pTbReq->type = TD_CHILD_TABLE;
  pTbReq->name = strdup(tname);
  pTbReq->ctb.suid = suid;
  pTbReq->ctb.tagNum = tagNum;
  if (sname) pTbReq->ctb.stbName = strdup(sname);
  pTbReq->ctb.pTag = (uint8_t*)pTag;
H
Haojun Liao 已提交
872
  pTbReq->ctb.tagName = taosArrayDup(tagName, NULL);
873
  pTbReq->ttl = ttl;
X
Xiaoyu Wang 已提交
874 875 876 877
  pTbReq->commentLen = -1;

  return;
}
D
stmt  
dapan1121 已提交
878

X
Xiaoyu Wang 已提交
879 880 881
int32_t insMemRowAppend(SMsgBuf* pMsgBuf, const void* value, int32_t len, void* param) {
  SMemParam*   pa = (SMemParam*)param;
  SRowBuilder* rb = pa->rb;
X
Xiaoyu Wang 已提交
882

X
Xiaoyu Wang 已提交
883 884 885
  if (value == NULL) {  // it is a null data
    tdAppendColValToRow(rb, pa->schema->colId, pa->schema->type, TD_VTYPE_NULL, value, false, pa->toffset, pa->colIdx);
    return TSDB_CODE_SUCCESS;
D
dapan 已提交
886
  }
X
Xiaoyu Wang 已提交
887

X
Xiaoyu Wang 已提交
888 889 890 891 892 893 894 895 896 897 898 899 900 901 902 903 904 905 906 907 908
  if (TSDB_DATA_TYPE_BINARY == pa->schema->type) {
    const char* rowEnd = tdRowEnd(rb->pBuf);
    STR_WITH_SIZE_TO_VARSTR(rowEnd, value, len);
    tdAppendColValToRow(rb, pa->schema->colId, pa->schema->type, TD_VTYPE_NORM, rowEnd, false, pa->toffset, pa->colIdx);
  } else if (TSDB_DATA_TYPE_NCHAR == pa->schema->type) {
    // if the converted output len is over than pColumnModel->bytes, return error: 'Argument list too long'
    int32_t     output = 0;
    const char* rowEnd = tdRowEnd(rb->pBuf);
    if (!taosMbsToUcs4(value, len, (TdUcs4*)varDataVal(rowEnd), pa->schema->bytes - VARSTR_HEADER_SIZE, &output)) {
      if (errno == E2BIG) {
        return generateSyntaxErrMsg(pMsgBuf, TSDB_CODE_PAR_VALUE_TOO_LONG, pa->schema->name);
      }
      char buf[512] = {0};
      snprintf(buf, tListLen(buf), "%s", strerror(errno));
      return buildSyntaxErrMsg(pMsgBuf, buf, value);
    }
    varDataSetLen(rowEnd, output);
    tdAppendColValToRow(rb, pa->schema->colId, pa->schema->type, TD_VTYPE_NORM, rowEnd, false, pa->toffset, pa->colIdx);
  } else {
    tdAppendColValToRow(rb, pa->schema->colId, pa->schema->type, TD_VTYPE_NORM, value, false, pa->toffset, pa->colIdx);
  }
D
stmt  
dapan1121 已提交
909

D
stmt  
dapan1121 已提交
910 911 912
  return TSDB_CODE_SUCCESS;
}

X
Xiaoyu Wang 已提交
913 914 915 916 917
int32_t insCheckTimestamp(STableDataBlocks* pDataBlocks, const char* start) {
  // once the data block is disordered, we do NOT keep previous timestamp any more
  if (!pDataBlocks->ordered) {
    return TSDB_CODE_SUCCESS;
  }
D
dapan 已提交
918

X
Xiaoyu Wang 已提交
919 920 921
  TSKEY k = *(TSKEY*)start;
  if (k <= pDataBlocks->prevTS) {
    pDataBlocks->ordered = false;
D
stmt  
dapan1121 已提交
922 923
  }

X
Xiaoyu Wang 已提交
924 925
  pDataBlocks->prevTS = k;
  return TSDB_CODE_SUCCESS;
D
stmt  
dapan1121 已提交
926 927
}

X
Xiaoyu Wang 已提交
928 929 930 931 932 933 934 935 936 937 938 939 940 941 942 943 944 945
static void buildMsgHeader(STableDataBlocks* src, SVgDataBlocks* blocks) {
  SSubmitReq* submit = (SSubmitReq*)blocks->pData;
  submit->header.vgId = htonl(blocks->vg.vgId);
  submit->header.contLen = htonl(blocks->size);
  submit->length = submit->header.contLen;
  submit->numOfBlocks = htonl(blocks->numOfTables);
  SSubmitBlk* blk = (SSubmitBlk*)(submit + 1);
  int32_t     numOfBlocks = blocks->numOfTables;
  while (numOfBlocks--) {
    int32_t dataLen = blk->dataLen;
    int32_t schemaLen = blk->schemaLen;
    blk->uid = htobe64(blk->uid);
    blk->suid = htobe64(blk->suid);
    blk->sversion = htonl(blk->sversion);
    blk->dataLen = htonl(blk->dataLen);
    blk->schemaLen = htonl(blk->schemaLen);
    blk->numOfRows = htonl(blk->numOfRows);
    blk = (SSubmitBlk*)(blk->data + schemaLen + dataLen);
D
stmt  
dapan1121 已提交
946
  }
X
Xiaoyu Wang 已提交
947
}
D
stmt  
dapan1121 已提交
948

X
Xiaoyu Wang 已提交
949 950 951 952
int32_t insBuildOutput(SHashObj* pVgroupsHashObj, SArray* pVgDataBlocks, SArray** pDataBlocks) {
  size_t numOfVg = taosArrayGetSize(pVgDataBlocks);
  *pDataBlocks = taosArrayInit(numOfVg, POINTER_BYTES);
  if (NULL == *pDataBlocks) {
X
Xiaoyu Wang 已提交
953 954 955
    return TSDB_CODE_TSC_OUT_OF_MEMORY;
  }
  for (size_t i = 0; i < numOfVg; ++i) {
X
Xiaoyu Wang 已提交
956
    STableDataBlocks* src = taosArrayGetP(pVgDataBlocks, i);
X
Xiaoyu Wang 已提交
957 958 959 960
    SVgDataBlocks*    dst = taosMemoryCalloc(1, sizeof(SVgDataBlocks));
    if (NULL == dst) {
      return TSDB_CODE_TSC_OUT_OF_MEMORY;
    }
X
Xiaoyu Wang 已提交
961
    taosHashGetDup(pVgroupsHashObj, (const char*)&src->vgId, sizeof(src->vgId), &dst->vg);
X
Xiaoyu Wang 已提交
962 963 964 965
    dst->numOfTables = src->numOfTables;
    dst->size = src->size;
    TSWAP(dst->pData, src->pData);
    buildMsgHeader(src, dst);
X
Xiaoyu Wang 已提交
966
    taosArrayPush(*pDataBlocks, &dst);
X
Xiaoyu Wang 已提交
967 968
  }
  return TSDB_CODE_SUCCESS;
D
stmt  
dapan1121 已提交
969
}
X
Xiaoyu Wang 已提交
970

D
dapan1121 已提交
971
static void initBoundCols(int32_t ncols, int16_t* pBoundCols) {
X
Xiaoyu Wang 已提交
972 973 974 975 976 977 978 979 980 981 982 983 984 985 986 987
  for (int32_t i = 0; i < ncols; ++i) {
    pBoundCols[i] = i;
  }
}

static void initColValues(STableMeta* pTableMeta, SArray* pValues) {
  SSchema* pSchemas = getTableColumnSchema(pTableMeta);
  for (int32_t i = 0; i < pTableMeta->tableInfo.numOfColumns; ++i) {
    SColVal val = COL_VAL_NONE(pSchemas[i].colId, pSchemas[i].type);
    taosArrayPush(pValues, &val);
  }
}

int32_t insInitBoundColsInfo(int32_t numOfBound, SBoundColInfo* pInfo) {
  pInfo->numOfCols = numOfBound;
  pInfo->numOfBound = numOfBound;
D
dapan1121 已提交
988
  pInfo->pColIndex = taosMemoryCalloc(numOfBound, sizeof(int16_t));
X
Xiaoyu Wang 已提交
989 990 991 992 993 994 995
  if (NULL == pInfo->pColIndex) {
    return TSDB_CODE_OUT_OF_MEMORY;
  }
  initBoundCols(numOfBound, pInfo->pColIndex);
  return TSDB_CODE_SUCCESS;
}

X
Xiaoyu Wang 已提交
996
void insCheckTableDataOrder(STableDataCxt* pTableCxt, TSKEY tsKey) {
X
Xiaoyu Wang 已提交
997
  // once the data block is disordered, we do NOT keep last timestamp any more
X
Xiaoyu Wang 已提交
998 999 1000 1001
  if (!pTableCxt->ordered) {
    return;
  }

X
Xiaoyu Wang 已提交
1002
  if (tsKey < pTableCxt->lastTs) {
X
Xiaoyu Wang 已提交
1003 1004 1005
    pTableCxt->ordered = false;
  }

X
Xiaoyu Wang 已提交
1006 1007 1008 1009
  if (tsKey == pTableCxt->lastTs) {
    pTableCxt->duplicateTs = true;
  }

X
Xiaoyu Wang 已提交
1010 1011 1012 1013
  pTableCxt->lastTs = tsKey;
  return;
}

1014
void destroyBoundColInfo(SBoundColInfo* pInfo) { taosMemoryFreeClear(pInfo->pColIndex); }
X
Xiaoyu Wang 已提交
1015

D
dapan1121 已提交
1016
static int32_t createTableDataCxt(STableMeta* pTableMeta, SVCreateTbReq** pCreateTbReq, STableDataCxt** pOutput, bool colMode) {
X
Xiaoyu Wang 已提交
1017 1018 1019 1020 1021 1022 1023
  STableDataCxt* pTableCxt = taosMemoryCalloc(1, sizeof(STableDataCxt));
  if (NULL == pTableCxt) {
    return TSDB_CODE_OUT_OF_MEMORY;
  }

  int32_t code = TSDB_CODE_SUCCESS;

X
Xiaoyu Wang 已提交
1024
  pTableCxt->lastTs = 0;
X
Xiaoyu Wang 已提交
1025 1026
  pTableCxt->ordered = true;
  pTableCxt->duplicateTs = false;
X
Xiaoyu Wang 已提交
1027

X
Xiaoyu Wang 已提交
1028 1029 1030 1031 1032 1033 1034 1035 1036 1037 1038 1039 1040 1041 1042 1043 1044 1045 1046 1047 1048 1049 1050
  pTableCxt->pMeta = tableMetaDup(pTableMeta);
  if (NULL == pTableCxt->pMeta) {
    code = TSDB_CODE_OUT_OF_MEMORY;
  }
  if (TSDB_CODE_SUCCESS == code) {
    pTableCxt->pSchema =
        tBuildTSchema(getTableColumnSchema(pTableMeta), pTableMeta->tableInfo.numOfColumns, pTableMeta->sversion);
    if (NULL == pTableCxt->pSchema) {
      code = TSDB_CODE_OUT_OF_MEMORY;
    }
  }
  if (TSDB_CODE_SUCCESS == code) {
    code = insInitBoundColsInfo(pTableMeta->tableInfo.numOfColumns, &pTableCxt->boundColsInfo);
  }
  if (TSDB_CODE_SUCCESS == code) {
    pTableCxt->pValues = taosArrayInit(pTableMeta->tableInfo.numOfColumns, sizeof(SColVal));
    if (NULL == pTableCxt->pValues) {
      code = TSDB_CODE_OUT_OF_MEMORY;
    } else {
      initColValues(pTableMeta, pTableCxt->pValues);
    }
  }
  if (TSDB_CODE_SUCCESS == code) {
X
Xiaoyu Wang 已提交
1051 1052
    pTableCxt->pData = taosMemoryCalloc(1, sizeof(SSubmitTbData));
    if (NULL == pTableCxt->pData) {
X
Xiaoyu Wang 已提交
1053
      code = TSDB_CODE_OUT_OF_MEMORY;
X
Xiaoyu Wang 已提交
1054
    } else {
X
Xiaoyu Wang 已提交
1055
      pTableCxt->pData->flags = NULL != *pCreateTbReq ? SUBMIT_REQ_AUTO_CREATE_TABLE : 0;
D
dapan1121 已提交
1056
      pTableCxt->pData->flags |= colMode ? SUBMIT_REQ_COLUMN_DATA_FORMAT : 0;
X
Xiaoyu Wang 已提交
1057 1058 1059
      pTableCxt->pData->suid = pTableMeta->suid;
      pTableCxt->pData->uid = pTableMeta->uid;
      pTableCxt->pData->sver = pTableMeta->sversion;
X
Xiaoyu Wang 已提交
1060 1061
      pTableCxt->pData->pCreateTbReq = *pCreateTbReq;
      *pCreateTbReq = NULL;
D
dapan1121 已提交
1062 1063 1064 1065 1066 1067 1068 1069 1070 1071
      if (pTableCxt->pData->flags & SUBMIT_REQ_COLUMN_DATA_FORMAT) {
        pTableCxt->pData->aCol = taosArrayInit(128, sizeof(SColData));
        if (NULL == pTableCxt->pData->aCol) {
          code = TSDB_CODE_OUT_OF_MEMORY;
        }
      } else {
        pTableCxt->pData->aRowP = taosArrayInit(128, POINTER_BYTES);
        if (NULL == pTableCxt->pData->aRowP) {
          code = TSDB_CODE_OUT_OF_MEMORY;
        }
X
Xiaoyu Wang 已提交
1072
      }
X
Xiaoyu Wang 已提交
1073 1074 1075 1076 1077 1078 1079 1080 1081 1082 1083 1084 1085
    }
  }

  if (TSDB_CODE_SUCCESS == code) {
    *pOutput = pTableCxt;
  } else {
    taosMemoryFree(pTableCxt);
  }

  return code;
}

int32_t insGetTableDataCxt(SHashObj* pHash, void* id, int32_t idLen, STableMeta* pTableMeta,
D
dapan1121 已提交
1086
                           SVCreateTbReq** pCreateTbReq, STableDataCxt** pTableCxt, bool colMode) {
wmmhello's avatar
wmmhello 已提交
1087 1088 1089
  STableDataCxt** tmp = (STableDataCxt**)taosHashGet(pHash, id, idLen);
  if (NULL != tmp) {
    *pTableCxt = *tmp;
X
Xiaoyu Wang 已提交
1090 1091
    return TSDB_CODE_SUCCESS;
  }
D
dapan1121 已提交
1092
  int32_t code = createTableDataCxt(pTableMeta, pCreateTbReq, pTableCxt, colMode);
X
Xiaoyu Wang 已提交
1093 1094 1095 1096 1097 1098
  if (TSDB_CODE_SUCCESS == code) {
    code = taosHashPut(pHash, id, idLen, pTableCxt, POINTER_BYTES);
  }
  return code;
}

X
Xiaoyu Wang 已提交
1099 1100 1101 1102 1103 1104 1105
static void destroyColVal(void* p) {
  SColVal* pVal = p;
  if (TSDB_DATA_TYPE_NCHAR == pVal->type) {
    taosMemoryFree(pVal->value.pData);
  }
}

X
Xiaoyu Wang 已提交
1106 1107 1108 1109 1110 1111 1112 1113
void insDestroyTableDataCxt(STableDataCxt* pTableCxt) {
  if (NULL == pTableCxt) {
    return;
  }

  taosMemoryFreeClear(pTableCxt->pMeta);
  tDestroyTSchema(pTableCxt->pSchema);
  destroyBoundColInfo(&pTableCxt->boundColsInfo);
X
Xiaoyu Wang 已提交
1114
  taosArrayDestroyEx(pTableCxt->pValues, destroyColVal);
H
Hongze Cheng 已提交
1115
  if (pTableCxt->pData) {
H
Hongze Cheng 已提交
1116
    tDestroySSubmitTbData(pTableCxt->pData, TSDB_MSG_FLG_ENCODE);
H
Hongze Cheng 已提交
1117 1118
    taosMemoryFree(pTableCxt->pData);
  }
X
Xiaoyu Wang 已提交
1119
  taosMemoryFree(pTableCxt);
X
Xiaoyu Wang 已提交
1120 1121 1122 1123 1124 1125 1126
}

void insDestroyVgroupDataCxt(SVgroupDataCxt* pVgCxt) {
  if (NULL == pVgCxt) {
    return;
  }

H
Hongze Cheng 已提交
1127
  tDestroySSubmitReq2(pVgCxt->pData, TSDB_MSG_FLG_ENCODE);
X
Xiaoyu Wang 已提交
1128
  taosMemoryFree(pVgCxt);
X
Xiaoyu Wang 已提交
1129 1130 1131 1132 1133 1134 1135 1136 1137 1138 1139 1140 1141 1142 1143 1144 1145 1146 1147 1148 1149 1150 1151 1152 1153 1154 1155 1156 1157 1158 1159 1160 1161 1162 1163 1164 1165 1166 1167 1168 1169 1170 1171 1172 1173 1174 1175
}

void insDestroyVgroupDataCxtList(SArray* pVgCxtList) {
  if (NULL == pVgCxtList) {
    return;
  }

  size_t size = taosArrayGetSize(pVgCxtList);
  for (int32_t i = 0; i < size; i++) {
    void* p = taosArrayGetP(pVgCxtList, i);
    insDestroyVgroupDataCxt(p);
  }

  taosArrayDestroy(pVgCxtList);
}

void insDestroyVgroupDataCxtHashMap(SHashObj* pVgCxtHash) {
  if (NULL == pVgCxtHash) {
    return;
  }

  void** p = taosHashIterate(pVgCxtHash, NULL);
  while (p) {
    insDestroyVgroupDataCxt(*(SVgroupDataCxt**)p);

    p = taosHashIterate(pVgCxtHash, p);
  }

  taosHashCleanup(pVgCxtHash);
}

void insDestroyTableDataCxtHashMap(SHashObj* pTableCxtHash) {
  if (NULL == pTableCxtHash) {
    return;
  }

  void** p = taosHashIterate(pTableCxtHash, NULL);
  while (p) {
    insDestroyTableDataCxt(*(STableDataCxt**)p);

    p = taosHashIterate(pTableCxtHash, p);
  }

  taosHashCleanup(pTableCxtHash);
}

static int32_t fillVgroupDataCxt(STableDataCxt* pTableCxt, SVgroupDataCxt* pVgCxt) {
X
Xiaoyu Wang 已提交
1176 1177 1178
  if (NULL == pVgCxt->pData->aSubmitTbData) {
    pVgCxt->pData->aSubmitTbData = taosArrayInit(128, sizeof(SSubmitTbData));
    if (NULL == pVgCxt->pData->aSubmitTbData) {
X
Xiaoyu Wang 已提交
1179 1180 1181
      return TSDB_CODE_OUT_OF_MEMORY;
    }
  }
X
Xiaoyu Wang 已提交
1182
  taosArrayPush(pVgCxt->pData->aSubmitTbData, pTableCxt->pData);
X
Xiaoyu Wang 已提交
1183
  taosMemoryFreeClear(pTableCxt->pData);
X
Xiaoyu Wang 已提交
1184 1185 1186 1187

  return TSDB_CODE_SUCCESS;
}

X
Xiaoyu Wang 已提交
1188 1189 1190 1191 1192 1193 1194 1195 1196 1197 1198 1199 1200 1201 1202 1203 1204 1205 1206 1207 1208 1209 1210
static int32_t createVgroupDataCxt(STableDataCxt* pTableCxt, SHashObj* pVgroupHash, SArray* pVgroupList,
                                   SVgroupDataCxt** pOutput) {
  SVgroupDataCxt* pVgCxt = taosMemoryCalloc(1, sizeof(SVgroupDataCxt));
  if (NULL == pVgCxt) {
    return TSDB_CODE_OUT_OF_MEMORY;
  }
  pVgCxt->pData = taosMemoryCalloc(1, sizeof(SSubmitReq2));
  if (NULL == pVgCxt->pData) {
    insDestroyVgroupDataCxt(pVgCxt);
    return TSDB_CODE_OUT_OF_MEMORY;
  }

  pVgCxt->vgId = pTableCxt->pMeta->vgId;
  int32_t code = taosHashPut(pVgroupHash, &pVgCxt->vgId, sizeof(pVgCxt->vgId), &pVgCxt, POINTER_BYTES);
  if (TSDB_CODE_SUCCESS == code) {
    taosArrayPush(pVgroupList, &pVgCxt);
    *pOutput = pVgCxt;
  } else {
    insDestroyVgroupDataCxt(pVgCxt);
  }
  return code;
}

D
dapan1121 已提交
1211 1212 1213 1214 1215 1216 1217 1218 1219 1220 1221 1222
int insColDataComp(const void* lp, const void* rp) {
  SColData* pLeft = (SColData*)lp;
  SColData* pRight = (SColData*)rp;
  if (pLeft->cid < pRight->cid) {
    return -1;
  } else if (pLeft->cid > pRight->cid) {
    return 1;
  }

  return 0;
}

X
Xiaoyu Wang 已提交
1223 1224 1225 1226 1227 1228 1229 1230 1231 1232
int32_t insMergeTableDataCxt(SHashObj* pTableHash, SArray** pVgDataBlocks) {
  SHashObj* pVgroupHash = taosHashInit(128, taosGetDefaultHashFunction(TSDB_DATA_TYPE_INT), true, false);
  SArray*   pVgroupList = taosArrayInit(8, POINTER_BYTES);
  if (NULL == pVgroupHash || NULL == pVgroupList) {
    taosHashCleanup(pVgroupHash);
    taosArrayDestroy(pVgroupList);
    return TSDB_CODE_TSC_OUT_OF_MEMORY;
  }

  int32_t code = TSDB_CODE_SUCCESS;
D
dapan1121 已提交
1233 1234
  bool colFormat = false;
  
X
Xiaoyu Wang 已提交
1235
  void* p = taosHashIterate(pTableHash, NULL);
D
dapan1121 已提交
1236 1237 1238 1239 1240
  if (p) {
    STableDataCxt* pTableCxt = *(STableDataCxt**)p;
    colFormat = (0 != (pTableCxt->pData->flags & SUBMIT_REQ_COLUMN_DATA_FORMAT));
  }
  
X
Xiaoyu Wang 已提交
1241 1242
  while (TSDB_CODE_SUCCESS == code && NULL != p) {
    STableDataCxt* pTableCxt = *(STableDataCxt**)p;
D
dapan1121 已提交
1243 1244 1245
    if (colFormat) {
      taosArraySort(pTableCxt->pData->aCol, insColDataComp);
      
D
dapan1121 已提交
1246
      tColDataSortMerge(pTableCxt->pData->aCol);
D
dapan1121 已提交
1247
    } else {
1248 1249 1250 1251 1252 1253
      if (!pTableCxt->ordered) {
        tRowSort(pTableCxt->pData->aRowP);
      }
      if (!pTableCxt->ordered || pTableCxt->duplicateTs) {
        code = tRowMerge(pTableCxt->pData->aRowP, pTableCxt->pSchema, 0);
      }
D
dapan1121 已提交
1254 1255
    }
    
X
Xiaoyu Wang 已提交
1256
    if (TSDB_CODE_SUCCESS == code) {
X
Xiaoyu Wang 已提交
1257
      SVgroupDataCxt* pVgCxt = NULL;
X
Xiaoyu Wang 已提交
1258
      int32_t         vgId = pTableCxt->pMeta->vgId;
X
Xiaoyu Wang 已提交
1259 1260
      void**          p = taosHashGet(pVgroupHash, &vgId, sizeof(vgId));
      if (NULL == p) {
X
Xiaoyu Wang 已提交
1261
        code = createVgroupDataCxt(pTableCxt, pVgroupHash, pVgroupList, &pVgCxt);
X
Xiaoyu Wang 已提交
1262 1263
      } else {
        pVgCxt = *(SVgroupDataCxt**)p;
X
Xiaoyu Wang 已提交
1264 1265
      }
      if (TSDB_CODE_SUCCESS == code) {
X
Xiaoyu Wang 已提交
1266 1267 1268 1269 1270 1271 1272 1273 1274 1275 1276 1277 1278 1279 1280 1281 1282 1283
        code = fillVgroupDataCxt(pTableCxt, pVgCxt);
      }
    }
    if (TSDB_CODE_SUCCESS == code) {
      p = taosHashIterate(pTableHash, p);
    }
  }

  taosHashCleanup(pVgroupHash);
  if (TSDB_CODE_SUCCESS == code) {
    *pVgDataBlocks = pVgroupList;
  } else {
    taosArrayDestroy(pVgroupList);
  }

  return code;
}

X
Xiaoyu Wang 已提交
1284 1285 1286 1287 1288
static int32_t buildSubmitReq(int32_t vgId, SSubmitReq2* pReq, void** pData, uint32_t* pLen) {
  int32_t  code = TSDB_CODE_SUCCESS;
  uint32_t len = 0;
  void*    pBuf = NULL;
  tEncodeSize(tEncodeSSubmitReq2, pReq, len, code);
X
Xiaoyu Wang 已提交
1289 1290
  if (TSDB_CODE_SUCCESS == code) {
    SEncoder encoder;
X
Xiaoyu Wang 已提交
1291 1292 1293
    len += sizeof(SMsgHead);
    pBuf = taosMemoryMalloc(len);
    if (NULL == pBuf) {
X
Xiaoyu Wang 已提交
1294 1295
      return TSDB_CODE_TSC_OUT_OF_MEMORY;
    }
X
Xiaoyu Wang 已提交
1296 1297 1298
    ((SMsgHead*)pBuf)->vgId = htonl(vgId);
    ((SMsgHead*)pBuf)->contLen = htonl(len);
    tEncoderInit(&encoder, POINTER_SHIFT(pBuf, sizeof(SMsgHead)), len - sizeof(SMsgHead));
X
Xiaoyu Wang 已提交
1299 1300 1301
    code = tEncodeSSubmitReq2(&encoder, pReq);
    tEncoderClear(&encoder);
  }
X
Xiaoyu Wang 已提交
1302 1303 1304 1305 1306 1307

  if (TSDB_CODE_SUCCESS == code) {
    *pData = pBuf;
    *pLen = len;
  } else {
    taosMemoryFree(pBuf);
X
Xiaoyu Wang 已提交
1308 1309 1310 1311
  }
  return code;
}

X
Xiaoyu Wang 已提交
1312 1313 1314 1315 1316 1317
static void destroyVgDataBlocks(void* p) {
  SVgDataBlocks* pVg = p;
  taosMemoryFree(pVg->pData);
  taosMemoryFree(pVg);
}

X
Xiaoyu Wang 已提交
1318 1319 1320 1321 1322 1323 1324 1325 1326 1327 1328 1329 1330 1331 1332
int32_t insBuildVgDataBlocks(SHashObj* pVgroupsHashObj, SArray* pVgDataCxtList, SArray** pVgDataBlocks) {
  size_t  numOfVg = taosArrayGetSize(pVgDataCxtList);
  SArray* pDataBlocks = taosArrayInit(numOfVg, POINTER_BYTES);
  if (NULL == pDataBlocks) {
    return TSDB_CODE_TSC_OUT_OF_MEMORY;
  }

  int32_t code = TSDB_CODE_SUCCESS;
  for (size_t i = 0; TSDB_CODE_SUCCESS == code && i < numOfVg; ++i) {
    SVgroupDataCxt* src = taosArrayGetP(pVgDataCxtList, i);
    SVgDataBlocks*  dst = taosMemoryCalloc(1, sizeof(SVgDataBlocks));
    if (NULL == dst) {
      code = TSDB_CODE_TSC_OUT_OF_MEMORY;
    }
    if (TSDB_CODE_SUCCESS == code) {
X
Xiaoyu Wang 已提交
1333
      dst->numOfTables = taosArrayGetSize(src->pData->aSubmitTbData);
X
Xiaoyu Wang 已提交
1334 1335 1336
      code = taosHashGetDup(pVgroupsHashObj, (const char*)&src->vgId, sizeof(src->vgId), &dst->vg);
    }
    if (TSDB_CODE_SUCCESS == code) {
X
Xiaoyu Wang 已提交
1337
      code = buildSubmitReq(src->vgId, src->pData, &dst->pData, &dst->size);
X
Xiaoyu Wang 已提交
1338 1339 1340 1341 1342 1343 1344 1345 1346
    }
    if (TSDB_CODE_SUCCESS == code) {
      code = (NULL == taosArrayPush(pDataBlocks, &dst) ? TSDB_CODE_TSC_OUT_OF_MEMORY : TSDB_CODE_SUCCESS);
    }
  }

  if (TSDB_CODE_SUCCESS == code) {
    *pVgDataBlocks = pDataBlocks;
  } else {
X
Xiaoyu Wang 已提交
1347
    taosArrayDestroyP(pDataBlocks, destroyVgDataBlocks);
X
Xiaoyu Wang 已提交
1348 1349 1350 1351
  }

  return code;
}
wmmhello's avatar
wmmhello 已提交
1352 1353 1354 1355 1356 1357 1358 1359 1360 1361 1362 1363 1364 1365 1366 1367 1368 1369 1370 1371 1372 1373 1374 1375 1376 1377 1378 1379 1380 1381 1382 1383 1384 1385 1386 1387 1388 1389 1390 1391 1392 1393 1394 1395 1396 1397 1398 1399 1400 1401 1402 1403 1404 1405 1406 1407 1408 1409 1410 1411 1412 1413 1414 1415 1416 1417 1418 1419 1420 1421 1422 1423 1424 1425

int rawBlockBindData(SQuery *query, STableMeta* pTableMeta, SRetrieveTableRsp* pRsp, SVCreateTbReq *pCreateTb){
  STableDataCxt* pTableCxt = NULL;
  int ret = insGetTableDataCxt(((SVnodeModifOpStmt *)(query->pRoot))->pTableBlockHashObj, &pTableMeta->uid, sizeof(pTableMeta->uid),
                           pTableMeta, &pCreateTb, &pTableCxt, true);
  if (ret != TSDB_CODE_SUCCESS) {
    uError("insGetTableDataCxt error");
    goto end;
  }

  // no need to bind, because select * get all fields
  ret = initTableColSubmitData(pTableCxt);
  if (ret != TSDB_CODE_SUCCESS) {
    uError( "initTableColSubmitData error");
    goto end;
  }

  char* p = (char*)pRsp->data;
  // | version | total length | total rows | total columns | flag seg| block group id | column schema | each column length |
  p += sizeof(int32_t);
  p += sizeof(int32_t);

  int32_t numOfRows = *(int32_t*)p;
  p += sizeof(int32_t);

  int32_t numOfCols = *(int32_t*)p;
  p += sizeof(int32_t);

  p += sizeof(int32_t);
  p += sizeof(uint64_t);

  int8_t *fields = p;
  p += numOfCols * (sizeof(int8_t) + sizeof(int32_t));

  int32_t* colLength = (int32_t*)p;
  p += sizeof(int32_t) * numOfCols;

  char* pStart = p;

  SSchema*            pSchema = getTableColumnSchema(pTableCxt->pMeta);
  SBoundColInfo*      boundInfo = &pTableCxt->boundColsInfo;

  if(boundInfo->numOfBound != numOfCols){
    uError("boundInfo->numOfBound:%d != numOfCols:%d", boundInfo->numOfBound, numOfCols);
    ret = TSDB_CODE_INVALID_PARA;
    goto end;
  }
  for (int c = 0; c < boundInfo->numOfBound; ++c) {
    SSchema* pColSchema = &pSchema[boundInfo->pColIndex[c]];
    SColData* pCol = taosArrayGet(pTableCxt->pData->aCol, c);

    if (*fields != pColSchema->type && *(int32_t*)(fields + sizeof(int8_t)) != pColSchema->bytes) {
      uError( "type or bytes not equal");
      ret = TSDB_CODE_INVALID_PARA;
      goto end;
    }

    colLength[c] = htonl(colLength[c]);
    int8_t* offset = pStart;
    if (IS_VAR_DATA_TYPE(pColSchema->type)) {
      pStart += numOfRows * sizeof(int32_t);
    } else {
      pStart += BitmapLen(numOfRows);
    }
    char *pData = pStart;

    tColDataAddValueByDataBlock(pCol, pColSchema->type, pColSchema->bytes, numOfRows, offset, pData);
    fields += sizeof(int8_t) + sizeof(int32_t);
    pStart += colLength[c];
  }

end:
  return ret;
}