taosdump.c 95.5 KB
Newer Older
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16
/*
 * 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/>.
 */

#include <iconv.h>
H
Hui Li 已提交
17 18 19
#include <sys/stat.h>
#include <sys/syscall.h>

H
Haojun Liao 已提交
20
#include "../../../include/client/taos.h"
21 22
#include "os.h"
#include "taosdef.h"
H
Hongze Cheng 已提交
23
#include "tmsg.h"
24 25 26 27 28
#include "tglobal.h"
#include "tsclient.h"
#include "tsdb.h"
#include "tutil.h"

29 30
#define TSDB_SUPPORT_NANOSECOND 1

31 32 33
#define MAX_FILE_NAME_LEN       256             // max file name length on linux is 255
#define COMMAND_SIZE            65536
#define MAX_RECORDS_PER_REQ     32766
34 35
//#define DEFAULT_DUMP_FILE "taosdump.sql"

36 37 38
// for strncpy buffer overflow
#define min(a, b) (((a) < (b)) ? (a) : (b))

39 40 41 42
static int  converStringToReadable(char *str, int size, char *buf, int bufsize);
static int  convertNCharToReadable(char *str, int size, char *buf, int bufsize);
static void taosDumpCharset(FILE *fp);
static void taosLoadFileCharset(FILE *fp, char *fcharset);
43 44 45 46 47 48

typedef struct {
  short bytes;
  int8_t type;
} SOColInfo;

49 50 51 52 53 54 55 56 57 58 59 60 61
#define debugPrint(fmt, ...) \
    do { if (g_args.debug_print || g_args.verbose_print) \
      fprintf(stderr, "DEBG: "fmt, __VA_ARGS__); } while(0)

#define verbosePrint(fmt, ...) \
    do { if (g_args.verbose_print) \
        fprintf(stderr, "VERB: "fmt, __VA_ARGS__); } while(0)

#define performancePrint(fmt, ...) \
    do { if (g_args.performance_print) \
        fprintf(stderr, "VERB: "fmt, __VA_ARGS__); } while(0)

#define errorPrint(fmt, ...) \
62
    do { fprintf(stderr, "\033[31m"); fprintf(stderr, "ERROR: "fmt, __VA_ARGS__); fprintf(stderr, "\033[0m"); } while(0)
63

64 65 66 67 68 69 70 71 72 73 74 75 76 77
static bool isStringNumber(char *input)
{
    int len = strlen(input);
    if (0 == len) {
        return false;
    }

    for (int i = 0; i < len; i++) {
        if (!isdigit(input[i]))
            return false;
    }

    return true;
}
78

79 80
// -------------------------- SHOW DATABASE INTERFACE-----------------------
enum _show_db_index {
81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100
    TSDB_SHOW_DB_NAME_INDEX,
    TSDB_SHOW_DB_CREATED_TIME_INDEX,
    TSDB_SHOW_DB_NTABLES_INDEX,
    TSDB_SHOW_DB_VGROUPS_INDEX,
    TSDB_SHOW_DB_REPLICA_INDEX,
    TSDB_SHOW_DB_QUORUM_INDEX,
    TSDB_SHOW_DB_DAYS_INDEX,
    TSDB_SHOW_DB_KEEP_INDEX,
    TSDB_SHOW_DB_CACHE_INDEX,
    TSDB_SHOW_DB_BLOCKS_INDEX,
    TSDB_SHOW_DB_MINROWS_INDEX,
    TSDB_SHOW_DB_MAXROWS_INDEX,
    TSDB_SHOW_DB_WALLEVEL_INDEX,
    TSDB_SHOW_DB_FSYNC_INDEX,
    TSDB_SHOW_DB_COMP_INDEX,
    TSDB_SHOW_DB_CACHELAST_INDEX,
    TSDB_SHOW_DB_PRECISION_INDEX,
    TSDB_SHOW_DB_UPDATE_INDEX,
    TSDB_SHOW_DB_STATUS_INDEX,
    TSDB_MAX_SHOW_DB
101 102 103 104
};

// -----------------------------------------SHOW TABLES CONFIGURE -------------------------------------
enum _show_tables_index {
105 106 107 108 109 110 111 112
    TSDB_SHOW_TABLES_NAME_INDEX,
    TSDB_SHOW_TABLES_CREATED_TIME_INDEX,
    TSDB_SHOW_TABLES_COLUMNS_INDEX,
    TSDB_SHOW_TABLES_METRIC_INDEX,
    TSDB_SHOW_TABLES_UID_INDEX,
    TSDB_SHOW_TABLES_TID_INDEX,
    TSDB_SHOW_TABLES_VGID_INDEX,
    TSDB_MAX_SHOW_TABLES
113 114 115 116
};

// ---------------------------------- DESCRIBE METRIC CONFIGURE ------------------------------
enum _describe_table_index {
117 118 119 120 121
    TSDB_DESCRIBE_METRIC_FIELD_INDEX,
    TSDB_DESCRIBE_METRIC_TYPE_INDEX,
    TSDB_DESCRIBE_METRIC_LENGTH_INDEX,
    TSDB_DESCRIBE_METRIC_NOTE_INDEX,
    TSDB_MAX_DESCRIBE_METRIC
122 123
};

124 125
#define COL_NOTE_LEN    128

126
typedef struct {
127 128 129 130
    char field[TSDB_COL_NAME_LEN + 1];
    char type[16];
    int length;
    char note[COL_NOTE_LEN];
131 132 133
} SColDes;

typedef struct {
134 135
    char name[TSDB_TABLE_NAME_LEN];
    SColDes cols[];
136 137 138 139
} STableDef;

extern char version[];

140 141 142
#define DB_PRECISION_LEN   8
#define DB_STATUS_LEN      16

143
typedef struct {
144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162
    char     name[TSDB_DB_NAME_LEN];
    char     create_time[32];
    int32_t  ntables;
    int32_t  vgroups;
    int16_t  replica;
    int16_t  quorum;
    int16_t  days;
    char     keeplist[32];
    //int16_t  daysToKeep;
    //int16_t  daysToKeep1;
    //int16_t  daysToKeep2;
    int32_t  cache; //MB
    int32_t  blocks;
    int32_t  minrows;
    int32_t  maxrows;
    int8_t   wallevel;
    int32_t  fsync;
    int8_t   comp;
    int8_t   cachelast;
163
    char     precision[DB_PRECISION_LEN];   // time resolution
164
    int8_t   update;
165
    char     status[DB_STATUS_LEN];
166 167 168
} SDbInfo;

typedef struct {
169 170
    char name[TSDB_TABLE_NAME_LEN];
    char metric[TSDB_TABLE_NAME_LEN];
171 172 173
} STableRecord;

typedef struct {
174 175
    bool isMetric;
    STableRecord tableRecord;
176 177 178
} STableRecordInfo;

typedef struct {
179 180 181 182
    pthread_t threadID;
    int32_t   threadIndex;
    int32_t   totalThreads;
    char      dbName[TSDB_DB_NAME_LEN];
183
    int         precision;
184 185 186
    void     *taosCon;
    int64_t   rowsOfDumpOut;
    int64_t   tablesOfDumpOut;
187 188
} SThreadParaObj;

189
typedef struct {
190 191 192 193
    int64_t   totalRowsOfDumpOut;
    int64_t   totalChildTblsOfDumpOut;
    int32_t   totalSuperTblsOfDumpOut;
    int32_t   totalDatabasesOfDumpOut;
194 195
} resultStatistics;

196
static int64_t g_totalDumpOutRows = 0;
197

198
SDbInfo **g_dbInfos = NULL;
199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218

const char *argp_program_version = version;
const char *argp_program_bug_address = "<support@taosdata.com>";

/* Program documentation. */
static char doc[] = "";
/* "Argp example #4 -- a program with somewhat more complicated\ */
/*         options\ */
/*         \vThis part of the documentation comes *after* the options;\ */
/*         note that the text is automatically filled, but it's possible\ */
/*         to force a line-break, e.g.\n<-- here."; */

/* A description of the arguments we accept. */
static char args_doc[] = "dbname [tbname ...]\n--databases dbname ...\n--all-databases\n-i inpath\n-o outpath";

/* Keys for options without short-options. */
#define OPT_ABORT 1 /* –abort */

/* The options we understand. */
static struct argp_option options[] = {
219 220 221 222
    // connection option
    {"host", 'h', "HOST",    0,  "Server host dumping data from. Default is localhost.", 0},
    {"user", 'u', "USER",    0,  "User name used to connect to server. Default is root.", 0},
#ifdef _TD_POWER_
223
    {"password", 'p', 0,    0,  "User password to connect to server. Default is powerdb.", 0},
224
#else
225
    {"password", 'p', 0,    0,  "User password to connect to server. Default is taosdata.", 0},
226 227 228 229 230 231 232 233
#endif
    {"port", 'P', "PORT",        0,  "Port to connect", 0},
    {"mysqlFlag",     'q', "MYSQLFLAG",   0,  "mysqlFlag, Default is 0", 0},
    // input/output file
    {"outpath", 'o', "OUTPATH",     0,  "Output file path.", 1},
    {"inpath", 'i', "INPATH",      0,  "Input file path.", 1},
    {"resultFile", 'r', "RESULTFILE",  0,  "DumpOut/In Result file path and name.", 1},
#ifdef _TD_POWER_
234
    {"config-dir", 'c', "CONFIG_DIR",  0,  "Configure directory. Default is /etc/power/taos.cfg.", 1},
235
#else
236
    {"config-dir", 'c', "CONFIG_DIR",  0,  "Configure directory. Default is /etc/taos/taos.cfg.", 1},
237 238 239 240 241 242 243 244 245
#endif
    {"encode", 'e', "ENCODE", 0,  "Input file encoding.", 1},
    // dump unit options
    {"all-databases", 'A', 0, 0,  "Dump all databases.", 2},
    {"databases", 'D', 0, 0,  "Dump assigned databases", 2},
    {"allow-sys",   'a', 0, 0,  "Allow to dump sys database", 2},
    // dump format options
    {"schemaonly", 's', 0, 0,  "Only dump schema.", 2},
    {"without-property", 'N', 0, 0,  "Dump schema without properties.", 2},
246
    {"avro", 'v', 0, 0,  "Dump apache avro format data file. By default, dump sql command sequence.", 2},
247 248
    {"start-time",    'S', "START_TIME",  0,  "Start time to dump. Either epoch or ISO8601/RFC3339 format is acceptable. ISO8601 format example: 2017-10-01T00:00:00.000+0800 or 2017-10-0100:00:00:000+0800 or '2017-10-01 00:00:00.000+0800'",  4},
    {"end-time",      'E', "END_TIME",    0,  "End time to dump. Either epoch or ISO8601/RFC3339 format is acceptable. ISO8601 format example: 2017-10-01T00:00:00.000+0800 or 2017-10-0100:00:00.000+0800 or '2017-10-01 00:00:00.000+0800'",  5},
249 250 251 252
    {"data-batch",  'B', "DATA_BATCH",  0,  "Number of data point per insert statement. Max value is 32766. Default is 1.", 3},
    {"max-sql-len", 'L', "SQL_LEN",     0,  "Max length of one sql. Default is 65480.",   3},
    {"table-batch", 't', "TABLE_BATCH", 0,  "Number of table dumpout into one output file. Default is 1.",  3},
    {"thread_num",  'T', "THREAD_NUM",  0,  "Number of thread for dump in file. Default is 5.", 3},
253
    {"debug",   'g', 0, 0,  "Print debug info.",    8},
254 255
    {0}
};
256 257

/* Used by main to communicate with parse_opt. */
258
typedef struct arguments {
259
    // connection option
260 261
    char    *host;
    char    *user;
262
    char    password[SHELL_MAX_PASSWORD_LEN];
263 264 265 266 267 268 269 270
    uint16_t port;
    uint16_t mysqlFlag;
    // output file
    char     outpath[MAX_FILE_NAME_LEN];
    char     inpath[MAX_FILE_NAME_LEN];
    // result file
    char    *resultFile;
    char    *encode;
271
    // dump unit option
272 273
    bool     all_databases;
    bool     databases;
274
    // dump format option
275 276 277 278
    bool     schemaonly;
    bool     with_property;
    bool     avro;
    int64_t  start_time;
279
    char     humanStartTime[28];
280
    int64_t  end_time;
281
    char     humanEndTime[28];
282
    char     precision[8];
283

284 285 286 287 288 289 290 291 292 293 294 295 296
    int32_t  data_batch;
    int32_t  max_sql_len;
    int32_t  table_batch; // num of table which will be dump into one output file.
    bool     allow_sys;
    // other options
    int32_t  thread_num;
    int      abort;
    char   **arg_list;
    int       arg_list_len;
    bool      isDumpIn;
    bool      debug_print;
    bool      verbose_print;
    bool      performance_print;
297 298

    int         dbCount;
299
} SArguments;
300 301

/* Our argp parser. */
302 303
static error_t parse_opt(int key, char *arg, struct argp_state *state);

304
static struct argp argp = {options, parse_opt, args_doc, doc};
305 306
static resultStatistics g_resultStatistics = {0};
static FILE *g_fpOfResult = NULL;
H
Hui Li 已提交
307
static int g_numOfCores = 1;
308

309 310 311 312 313 314 315 316 317 318 319
static int taosDumpOut();
static int taosDumpIn();
static void taosDumpCreateDbClause(SDbInfo *dbInfo, bool isDumpProperty,
        FILE *fp);
static int taosDumpDb(SDbInfo *dbInfo, FILE *fp, TAOS *taosCon);
static int32_t taosDumpStable(char *table, FILE *fp, TAOS* taosCon,
        char* dbName);
static void taosDumpCreateTableClause(STableDef *tableDes, int numOfCols,
        FILE *fp, char* dbName);
static void taosDumpCreateMTableClause(STableDef *tableDes, char *metric,
        int numOfCols, FILE *fp, char* dbName);
320
static int32_t taosDumpTable(char *tbName, char *metric,
321
        FILE *fp, TAOS* taosCon, char* dbName, int precision);
322 323
static int taosDumpTableData(FILE *fp, char *tbName,
        TAOS* taosCon, char* dbName,
324
        int precision,
325
        char *jsonAvroSchema);
326 327
static int taosCheckParam(struct arguments *arguments);
static void taosFreeDbInfos();
328 329 330 331
static void taosStartDumpOutWorkThreads(
        int32_t numOfThread,
        char *dbName,
        int precision);
332

333
struct arguments g_args = {
334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355
    // connection option
    NULL,
    "root",
#ifdef _TD_POWER_
    "powerdb",
#else
    "taosdata",
#endif
    0,
    0,
    // outpath and inpath
    "",
    "",
    "./dump_result.txt",
    NULL,
    // dump unit option
    false,
    false,
    // dump format option
    false,      // schemeonly
    true,       // with_property
    false,      // avro format
356 357
    -INT64_MAX + 1, // start_time
    {0},        // humanStartTime
358
    INT64_MAX,  // end_time
359
    {0},        // humanEndTime
360
    "ms",       // precision
361 362 363 364 365 366 367 368 369 370 371 372
    1,          // data_batch
    TSDB_MAX_SQL_LEN,   // max_sql_len
    1,          // table_batch
    false,      // allow_sys
    // other options
    5,          // thread_num
    0,          // abort
    NULL,       // arg_list
    0,          // arg_list_len
    false,      // isDumpIn
    false,      // debug_print
    false,      // verbose_print
373 374
    false,      // performance_print
        0,      // dbCount
375
};
376

377 378 379 380 381 382 383 384 385
static void errorPrintReqArg2(char *program, char *wrong_arg)
{
    fprintf(stderr,
            "%s: option requires a number argument '-%s'\n",
            program, wrong_arg);
    fprintf(stderr,
            "Try `taosdump --help' or `taosdump --usage' for more information.\n");
}

386 387 388 389 390 391 392 393 394
static void errorPrintReqArg3(char *program, char *wrong_arg)
{
    fprintf(stderr,
            "%s: option '%s' requires an argument\n",
            program, wrong_arg);
    fprintf(stderr,
            "Try `taosdump --help' or `taosdump --usage' for more information.\n");
}

395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414
/* Parse a single option. */
static error_t parse_opt(int key, char *arg, struct argp_state *state) {
    /* Get the input argument from argp_parse, which we
       know is a pointer to our arguments structure. */
    wordexp_t full_path;

    switch (key) {
        // connection option
        case 'a':
            g_args.allow_sys = true;
            break;
        case 'h':
            g_args.host = arg;
            break;
        case 'u':
            g_args.user = arg;
            break;
        case 'p':
            break;
        case 'P':
415 416 417 418
            if (!isStringNumber(arg)) {
                errorPrintReqArg2("taosdump", "P");
                exit(EXIT_FAILURE);
            }
419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449
            g_args.port = atoi(arg);
            break;
        case 'q':
            g_args.mysqlFlag = atoi(arg);
            break;
        case 'o':
            if (wordexp(arg, &full_path, 0) != 0) {
                errorPrint("Invalid path %s\n", arg);
                return -1;
            }
            tstrncpy(g_args.outpath, full_path.we_wordv[0],
                    MAX_FILE_NAME_LEN);
            wordfree(&full_path);
            break;
        case 'g':
            g_args.debug_print = true;
            break;
        case 'i':
            g_args.isDumpIn = true;
            if (wordexp(arg, &full_path, 0) != 0) {
                errorPrint("Invalid path %s\n", arg);
                return -1;
            }
            tstrncpy(g_args.inpath, full_path.we_wordv[0],
                    MAX_FILE_NAME_LEN);
            wordfree(&full_path);
            break;
        case 'r':
            g_args.resultFile = arg;
            break;
        case 'c':
450 451 452 453
            if (0 == strlen(arg)) {
                errorPrintReqArg3("taosdump", "-c or --config-dir");
                exit(EXIT_FAILURE);
            }
454 455
            if (wordexp(arg, &full_path, 0) != 0) {
                errorPrint("Invalid path %s\n", arg);
456
                exit(EXIT_FAILURE);
457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476
            }
            tstrncpy(configDir, full_path.we_wordv[0], MAX_FILE_NAME_LEN);
            wordfree(&full_path);
            break;
        case 'e':
            g_args.encode = arg;
            break;
            // dump unit option
        case 'A':
            break;
        case 'D':
            g_args.databases = true;
            break;
            // dump format option
        case 's':
            g_args.schemaonly = true;
            break;
        case 'N':
            g_args.with_property = false;
            break;
477
        case 'v':
478 479 480 481 482 483
            g_args.avro = true;
            break;
        case 'S':
            // parse time here.
            break;
        case 'E':
484
            break;
485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505
        case 'B':
            g_args.data_batch = atoi(arg);
            if (g_args.data_batch > MAX_RECORDS_PER_REQ) {
                g_args.data_batch = MAX_RECORDS_PER_REQ;
            }
            break;
        case 'L':
            {
                int32_t len = atoi(arg);
                if (len > TSDB_MAX_ALLOWED_SQL_LEN) {
                    len = TSDB_MAX_ALLOWED_SQL_LEN;
                } else if (len < TSDB_MAX_SQL_LEN) {
                    len = TSDB_MAX_SQL_LEN;
                }
                g_args.max_sql_len = len;
                break;
            }
        case 't':
            g_args.table_batch = atoi(arg);
            break;
        case 'T':
506 507 508 509
            if (!isStringNumber(arg)) {
                errorPrint("%s", "\n\t-T need a number following!\n");
                exit(EXIT_FAILURE);
            }
510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526
            g_args.thread_num = atoi(arg);
            break;
        case OPT_ABORT:
            g_args.abort = 1;
            break;
        case ARGP_KEY_ARG:
            g_args.arg_list     = &state->argv[state->next - 1];
            g_args.arg_list_len = state->argc - state->next + 1;
            state->next             = state->argc;
            break;

        default:
            return ARGP_ERR_UNKNOWN;
    }
    return 0;
}

527
static int queryDbImpl(TAOS *taos, char *command) {
528 529 530 531 532 533 534 535 536
    int i;
    TAOS_RES *res = NULL;
    int32_t   code = -1;

    for (i = 0; i < 5; i++) {
        if (NULL != res) {
            taos_free_result(res);
            res = NULL;
        }
537

538 539 540 541 542
        res = taos_query(taos, command);
        code = taos_errno(res);
        if (0 == code) {
            break;
        }
543
    }
544

545 546 547 548 549
    if (code != 0) {
        errorPrint("Failed to run <%s>, reason: %s\n", command, taos_errstr(res));
        taos_free_result(res);
        //taos_close(taos);
        return -1;
550 551 552
    }

    taos_free_result(res);
553
    return 0;
554 555
}

556
UNUSED_FUNC static void parse_precision_first(
557 558 559 560 561 562 563 564 565 566 567 568 569 570 571
        int argc, char *argv[], SArguments *arguments) {
    for (int i = 1; i < argc; i++) {
        if (strcmp(argv[i], "-C") == 0) {
            if (NULL == argv[i+1]) {
                errorPrint("%s need a valid value following!\n", argv[i]);
                exit(-1);
            }
            char *tmp = strdup(argv[i+1]);
            if (tmp == NULL) {
                errorPrint("%s() LN%d, strdup() cannot allocate memory\n",
                        __func__, __LINE__);
                exit(-1);
            }
            if ((0 != strncasecmp(tmp, "ms", strlen("ms")))
                    && (0 != strncasecmp(tmp, "us", strlen("us")))
572 573 574 575
#if TSDB_SUPPORT_NANOSECOND == 1
                    && (0 != strncasecmp(tmp, "ns", strlen("ns")))
#endif
                    ) {
576 577 578 579 580
                //
                errorPrint("input precision: %s is invalid value\n", tmp);
                free(tmp);
                exit(-1);
            }
581 582
            tstrncpy(g_args.precision, tmp,
                min(DB_PRECISION_LEN, strlen(tmp) + 1));
583 584 585 586 587
            free(tmp);
        }
    }
}

588
static void parse_args(
589
        int argc, char *argv[], SArguments *arguments) {
590

591
    for (int i = 1; i < argc; i++) {
592 593 594 595
        if ((strncmp(argv[i], "-p", 2) == 0)
              || (strncmp(argv[i], "--password", 10) == 0)) {
            if ((strlen(argv[i]) == 2)
                  || (strncmp(argv[i], "--password", 10) == 0)) {
596
                printf("Enter password: ");
597
                taosSetConsoleEcho(false);
598 599 600
                if(scanf("%20s", arguments->password) > 1) {
                    errorPrint("%s() LN%d, password read error!\n", __func__, __LINE__);
                }
601
                taosSetConsoleEcho(true);
602
            } else {
603 604 605
                tstrncpy(arguments->password, (char *)(argv[i] + 2),
                        SHELL_MAX_PASSWORD_LEN);
                strcpy(argv[i], "-p");
606
            }
607 608 609 610 611 612 613 614 615 616
        } else if (strcmp(argv[i], "-gg") == 0) {
            arguments->verbose_print = true;
            strcpy(argv[i], "");
        } else if (strcmp(argv[i], "-PP") == 0) {
            arguments->performance_print = true;
            strcpy(argv[i], "");
        } else if (strcmp(argv[i], "-A") == 0) {
            g_args.all_databases = true;
        } else {
            continue;
617
        }
618

619 620 621
    }
}

622 623 624 625 626 627 628 629 630 631 632 633 634 635 636 637 638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657 658 659 660 661 662 663 664 665 666 667 668 669 670 671 672 673 674 675 676 677 678 679 680 681 682 683 684 685 686 687 688
static void copyHumanTimeToArg(char *timeStr, bool isStartTime)
{
    if (isStartTime)
        strcpy(g_args.humanStartTime, timeStr);
    else
        strcpy(g_args.humanEndTime, timeStr);
}

static void copyTimestampToArg(char *timeStr, bool isStartTime)
{
    if (isStartTime)
        g_args.start_time = atol(timeStr);
    else
        g_args.end_time = atol(timeStr);
}

static void parse_timestamp(
        int argc, char *argv[], SArguments *arguments) {
    for (int i = 1; i < argc; i++) {
        char *tmp;
        bool isStartTime = false;
        bool isEndTime = false;

        if (strcmp(argv[i], "-S") == 0) {
            isStartTime = true;
        } else if (strcmp(argv[i], "-E") == 0) {
            isEndTime = true;
        }

        if (isStartTime || isEndTime) {
            if (NULL == argv[i+1]) {
                errorPrint("%s need a valid value following!\n", argv[i]);
                exit(-1);
            }
            tmp = strdup(argv[i+1]);

            if (strchr(tmp, ':') && strchr(tmp, '-')) {
                copyHumanTimeToArg(tmp, isStartTime);
            } else {
                copyTimestampToArg(tmp, isStartTime);
            }
        }
    }
}

static int getPrecisionByString(char *precision)
{
    if (0 == strncasecmp(precision,
                "ms", 2)) {
        return TSDB_TIME_PRECISION_MILLI;
    } else if (0 == strncasecmp(precision,
                "us", 2)) {
        return TSDB_TIME_PRECISION_MICRO;
#if TSDB_SUPPORT_NANOSECOND == 1
    } else if (0 == strncasecmp(precision,
                "ns", 2)) {
        return TSDB_TIME_PRECISION_NANO;
#endif
    } else {
        errorPrint("Invalid time precision: %s",
                precision);
    }

    return -1;
}

/*
689 690
static void parse_timestamp(
        int argc, char *argv[], SArguments *arguments) {
691 692 693
    for (int i = 1; i < argc; i++) {
        if ((strcmp(argv[i], "-S") == 0)
                || (strcmp(argv[i], "-E") == 0)) {
694 695 696 697 698 699 700 701 702 703
            if (NULL == argv[i+1]) {
                errorPrint("%s need a valid value following!\n", argv[i]);
                exit(-1);
            }
            char *tmp = strdup(argv[i+1]);
            if (NULL == tmp) {
                errorPrint("%s() LN%d, strdup() cannot allocate memory\n",
                        __func__, __LINE__);
                exit(-1);
            }
704

705 706
            int64_t tmpEpoch;
            if (strchr(tmp, ':') && strchr(tmp, '-')) {
707
                strcpy(g_args.humanStartTime, tmp)
708 709 710 711 712 713 714
                int32_t timePrec;
                if (0 == strncasecmp(arguments->precision,
                            "ms", strlen("ms"))) {
                    timePrec = TSDB_TIME_PRECISION_MILLI;
                } else if (0 == strncasecmp(arguments->precision,
                            "us", strlen("us"))) {
                    timePrec = TSDB_TIME_PRECISION_MICRO;
715
#if TSDB_SUPPORT_NANOSECOND == 1
716 717 718
                } else if (0 == strncasecmp(arguments->precision,
                            "ns", strlen("ns"))) {
                    timePrec = TSDB_TIME_PRECISION_NANO;
719
#endif
720
                } else {
721 722 723 724 725 726 727 728 729 730 731 732
                    errorPrint("Invalid time precision: %s",
                            arguments->precision);
                    free(tmp);
                    return;
                }

                if (TSDB_CODE_SUCCESS != taosParseTime(
                            tmp, &tmpEpoch, strlen(tmp),
                            timePrec, 0)) {
                    errorPrint("Input %s, end time error!\n", tmp);
                    free(tmp);
                    return;
733 734
                }
            } else {
735
                tstrncpy(arguments->precision, "n/a", strlen("n/a") + 1);
736
                tmpEpoch = atoll(tmp);
737
            }
738

739
            sprintf(argv[i+1], "%"PRId64"", tmpEpoch);
740 741 742
            debugPrint("%s() LN%d, tmp is: %s, argv[%d]: %s\n",
                    __func__, __LINE__, tmp, i, argv[i]);
            free(tmp);
743 744 745
        }
    }
}
746
*/
747

748
int main(int argc, char *argv[]) {
749 750 751
    static char verType[32] = {0};
    sprintf(verType, "version: %s\n", version);
    argp_program_version = verType;
752

753 754 755
    int ret = 0;
    /* Parse our arguments; every option seen by parse_opt will be
       reflected in arguments. */
756
    if (argc > 1) {
757
//        parse_precision_first(argc, argv, &g_args);
758
        parse_timestamp(argc, argv, &g_args);
759
        parse_args(argc, argv, &g_args);
760
    }
761

762
    argp_parse(&argp, argc, argv, 0, 0, &g_args);
763

764 765 766 767 768 769 770
    if (g_args.abort) {
#ifndef _ALPINE
        error(10, 0, "ABORTED");
#else
        abort();
#endif
    }
771

772 773 774 775 776 777 778 779 780 781 782 783 784 785 786 787 788
    printf("====== arguments config ======\n");
    {
        printf("host: %s\n", g_args.host);
        printf("user: %s\n", g_args.user);
        printf("password: %s\n", g_args.password);
        printf("port: %u\n", g_args.port);
        printf("mysqlFlag: %d\n", g_args.mysqlFlag);
        printf("outpath: %s\n", g_args.outpath);
        printf("inpath: %s\n", g_args.inpath);
        printf("resultFile: %s\n", g_args.resultFile);
        printf("encode: %s\n", g_args.encode);
        printf("all_databases: %s\n", g_args.all_databases?"true":"false");
        printf("databases: %d\n", g_args.databases);
        printf("schemaonly: %s\n", g_args.schemaonly?"true":"false");
        printf("with_property: %s\n", g_args.with_property?"true":"false");
        printf("avro format: %s\n", g_args.avro?"true":"false");
        printf("start_time: %" PRId64 "\n", g_args.start_time);
789
        printf("human readable start time: %s \n", g_args.humanStartTime);
790
        printf("end_time: %" PRId64 "\n", g_args.end_time);
791
        printf("human readable end time: %s \n", g_args.humanEndTime);
792
        printf("precision: %s\n", g_args.precision);
793 794 795 796 797 798 799 800 801 802 803 804 805
        printf("data_batch: %d\n", g_args.data_batch);
        printf("max_sql_len: %d\n", g_args.max_sql_len);
        printf("table_batch: %d\n", g_args.table_batch);
        printf("thread_num: %d\n", g_args.thread_num);
        printf("allow_sys: %d\n", g_args.allow_sys);
        printf("abort: %d\n", g_args.abort);
        printf("isDumpIn: %d\n", g_args.isDumpIn);
        printf("arg_list_len: %d\n", g_args.arg_list_len);
        printf("debug_print: %d\n", g_args.debug_print);

        for (int32_t i = 0; i < g_args.arg_list_len; i++) {
            printf("arg_list[%d]: %s\n", i, g_args.arg_list[i]);
        }
806
    }
807 808 809 810
    printf("==============================\n");
    if (taosCheckParam(&g_args) < 0) {
        exit(EXIT_FAILURE);
    }
811

812 813 814 815 816 817 818 819 820 821 822 823 824 825 826 827 828 829 830 831 832 833 834 835
    g_fpOfResult = fopen(g_args.resultFile, "a");
    if (NULL == g_fpOfResult) {
        errorPrint("Failed to open %s for save result\n", g_args.resultFile);
        exit(-1);
    };

    fprintf(g_fpOfResult, "#############################################################################\n");
    fprintf(g_fpOfResult, "============================== arguments config =============================\n");
    {
        fprintf(g_fpOfResult, "host: %s\n", g_args.host);
        fprintf(g_fpOfResult, "user: %s\n", g_args.user);
        fprintf(g_fpOfResult, "password: %s\n", g_args.password);
        fprintf(g_fpOfResult, "port: %u\n", g_args.port);
        fprintf(g_fpOfResult, "mysqlFlag: %d\n", g_args.mysqlFlag);
        fprintf(g_fpOfResult, "outpath: %s\n", g_args.outpath);
        fprintf(g_fpOfResult, "inpath: %s\n", g_args.inpath);
        fprintf(g_fpOfResult, "resultFile: %s\n", g_args.resultFile);
        fprintf(g_fpOfResult, "encode: %s\n", g_args.encode);
        fprintf(g_fpOfResult, "all_databases: %s\n", g_args.all_databases?"true":"false");
        fprintf(g_fpOfResult, "databases: %d\n", g_args.databases);
        fprintf(g_fpOfResult, "schemaonly: %s\n", g_args.schemaonly?"true":"false");
        fprintf(g_fpOfResult, "with_property: %s\n", g_args.with_property?"true":"false");
        fprintf(g_fpOfResult, "avro format: %s\n", g_args.avro?"true":"false");
        fprintf(g_fpOfResult, "start_time: %" PRId64 "\n", g_args.start_time);
836
        fprintf(g_fpOfResult, "human readable start time: %s \n", g_args.humanStartTime);
837
        fprintf(g_fpOfResult, "end_time: %" PRId64 "\n", g_args.end_time);
838
        fprintf(g_fpOfResult, "human readable end time: %s \n", g_args.humanEndTime);
839
        fprintf(g_fpOfResult, "precision: %s\n", g_args.precision);
840 841 842 843 844 845 846 847 848 849 850 851
        fprintf(g_fpOfResult, "data_batch: %d\n", g_args.data_batch);
        fprintf(g_fpOfResult, "max_sql_len: %d\n", g_args.max_sql_len);
        fprintf(g_fpOfResult, "table_batch: %d\n", g_args.table_batch);
        fprintf(g_fpOfResult, "thread_num: %d\n", g_args.thread_num);
        fprintf(g_fpOfResult, "allow_sys: %d\n", g_args.allow_sys);
        fprintf(g_fpOfResult, "abort: %d\n", g_args.abort);
        fprintf(g_fpOfResult, "isDumpIn: %d\n", g_args.isDumpIn);
        fprintf(g_fpOfResult, "arg_list_len: %d\n", g_args.arg_list_len);

        for (int32_t i = 0; i < g_args.arg_list_len; i++) {
            fprintf(g_fpOfResult, "arg_list[%d]: %s\n", i, g_args.arg_list[i]);
        }
852 853
    }

854 855 856 857 858 859 860 861 862 863 864 865 866 867 868 869 870 871 872 873 874 875 876 877 878 879 880 881 882 883 884 885
    g_numOfCores = (int32_t)sysconf(_SC_NPROCESSORS_ONLN);

    time_t tTime = time(NULL);
    struct tm tm = *localtime(&tTime);

    if (g_args.isDumpIn) {
        fprintf(g_fpOfResult, "============================== DUMP IN ============================== \n");
        fprintf(g_fpOfResult, "# DumpIn start time:                   %d-%02d-%02d %02d:%02d:%02d\n",
                tm.tm_year + 1900, tm.tm_mon + 1,
                tm.tm_mday, tm.tm_hour, tm.tm_min, tm.tm_sec);
        if (taosDumpIn() < 0) {
            ret = -1;
        }
    } else {
        fprintf(g_fpOfResult, "============================== DUMP OUT ============================== \n");
        fprintf(g_fpOfResult, "# DumpOut start time:                   %d-%02d-%02d %02d:%02d:%02d\n",
                tm.tm_year + 1900, tm.tm_mon + 1,
                tm.tm_mday, tm.tm_hour, tm.tm_min, tm.tm_sec);
        if (taosDumpOut() < 0) {
            ret = -1;
        } else {
            fprintf(g_fpOfResult, "\n============================== TOTAL STATISTICS ============================== \n");
            fprintf(g_fpOfResult, "# total database count:     %d\n",
                    g_resultStatistics.totalDatabasesOfDumpOut);
            fprintf(g_fpOfResult, "# total super table count:  %d\n",
                    g_resultStatistics.totalSuperTblsOfDumpOut);
            fprintf(g_fpOfResult, "# total child table count:  %"PRId64"\n",
                    g_resultStatistics.totalChildTblsOfDumpOut);
            fprintf(g_fpOfResult, "# total row count:          %"PRId64"\n",
                    g_resultStatistics.totalRowsOfDumpOut);
        }
    }
886

887 888
    fprintf(g_fpOfResult, "\n");
    fclose(g_fpOfResult);
889

890
    return ret;
891 892
}

893
static void taosFreeDbInfos() {
894
    if (g_dbInfos == NULL) return;
895 896
    for (int i = 0; i < g_args.dbCount; i++)
        tfree(g_dbInfos[i]);
897
    tfree(g_dbInfos);
898 899 900
}

// check table is normal table or super table
901 902 903 904 905 906 907 908 909 910 911 912 913 914
static int taosGetTableRecordInfo(
        char *table, STableRecordInfo *pTableRecordInfo, TAOS *taosCon) {
    TAOS_ROW row = NULL;
    bool isSet = false;
    TAOS_RES *result     = NULL;

    memset(pTableRecordInfo, 0, sizeof(STableRecordInfo));

    char* tempCommand = (char *)malloc(COMMAND_SIZE);
    if (tempCommand == NULL) {
        errorPrint("%s() LN%d, failed to allocate memory\n",
                __func__, __LINE__);
        return -1;
    }
915

916
    sprintf(tempCommand, "show tables like %s", table);
917

918 919
    result = taos_query(taosCon, tempCommand);
    int32_t code = taos_errno(result);
920

921 922 923 924 925 926 927
    if (code != 0) {
        errorPrint("%s() LN%d, failed to run command %s\n",
                __func__, __LINE__, tempCommand);
        free(tempCommand);
        taos_free_result(result);
        return -1;
    }
928

929 930 931 932 933
    TAOS_FIELD *fields = taos_fetch_fields(result);

    while ((row = taos_fetch_row(result)) != NULL) {
        isSet = true;
        pTableRecordInfo->isMetric = false;
934
        tstrncpy(pTableRecordInfo->tableRecord.name,
935
                (char *)row[TSDB_SHOW_TABLES_NAME_INDEX],
936
                min(TSDB_TABLE_NAME_LEN,
937
                    fields[TSDB_SHOW_TABLES_NAME_INDEX].bytes + 1));
938
        tstrncpy(pTableRecordInfo->tableRecord.metric,
939
                (char *)row[TSDB_SHOW_TABLES_METRIC_INDEX],
940
                min(TSDB_TABLE_NAME_LEN,
941
                    fields[TSDB_SHOW_TABLES_METRIC_INDEX].bytes + 1));
942 943
        break;
    }
944

945
    taos_free_result(result);
946
    result = NULL;
947

948 949 950 951
    if (isSet) {
        free(tempCommand);
        return 0;
    }
952

953
    sprintf(tempCommand, "show stables like %s", table);
954

955 956
    result = taos_query(taosCon, tempCommand);
    code = taos_errno(result);
957

958 959 960 961 962 963 964
    if (code != 0) {
        errorPrint("%s() LN%d, failed to run command %s\n",
                __func__, __LINE__, tempCommand);
        free(tempCommand);
        taos_free_result(result);
        return -1;
    }
965

966 967 968 969 970 971 972
    while ((row = taos_fetch_row(result)) != NULL) {
        isSet = true;
        pTableRecordInfo->isMetric = true;
        tstrncpy(pTableRecordInfo->tableRecord.metric, table,
                TSDB_TABLE_NAME_LEN);
        break;
    }
973

974
    taos_free_result(result);
975
    result = NULL;
976

977 978 979 980 981 982
    if (isSet) {
        free(tempCommand);
        return 0;
    }
    errorPrint("%s() LN%d, invalid table/metric %s\n",
            __func__, __LINE__, table);
983
    free(tempCommand);
984
    return -1;
985 986 987
}


988 989
static int32_t taosSaveAllNormalTableToTempFile(TAOS *taosCon, char*meter,
        char* metric, int* fd) {
990 991 992 993 994 995 996 997 998 999
    STableRecord tableRecord;

    if (-1 == *fd) {
        *fd = open(".tables.tmp.0",
                O_RDWR | O_CREAT, S_IRWXU | S_IRGRP | S_IXGRP | S_IROTH);
        if (*fd == -1) {
            errorPrint("%s() LN%d, failed to open temp file: .tables.tmp.0\n",
                    __func__, __LINE__);
            return -1;
        }
1000
    }
1001

1002 1003 1004
    memset(&tableRecord, 0, sizeof(STableRecord));
    tstrncpy(tableRecord.name, meter, TSDB_TABLE_NAME_LEN);
    tstrncpy(tableRecord.metric, metric, TSDB_TABLE_NAME_LEN);
1005

1006 1007
    taosWrite(*fd, &tableRecord, sizeof(STableRecord));
    return 0;
1008 1009
}

1010 1011 1012 1013 1014 1015
static int32_t taosSaveTableOfMetricToTempFile(
        TAOS *taosCon, char* metric,
        int32_t*  totalNumOfThread) {
    TAOS_ROW row;
    int fd = -1;
    STableRecord tableRecord;
1016

1017 1018 1019 1020 1021
    char* tmpCommand = (char *)malloc(COMMAND_SIZE);
    if (tmpCommand == NULL) {
        errorPrint("%s() LN%d, failed to allocate memory\n", __func__, __LINE__);
        return -1;
    }
1022

1023
    sprintf(tmpCommand, "select tbname from %s", metric);
1024

1025 1026 1027 1028 1029 1030 1031 1032 1033
    TAOS_RES *res = taos_query(taosCon, tmpCommand);
    int32_t code = taos_errno(res);
    if (code != 0) {
        errorPrint("%s() LN%d, failed to run command %s\n",
                __func__, __LINE__, tmpCommand);
        free(tmpCommand);
        taos_free_result(res);
        return -1;
    }
1034 1035
    free(tmpCommand);

1036 1037 1038 1039 1040 1041 1042 1043 1044 1045 1046 1047
    char     tmpBuf[MAX_FILE_NAME_LEN];
    memset(tmpBuf, 0, MAX_FILE_NAME_LEN);
    sprintf(tmpBuf, ".select-tbname.tmp");
    fd = open(tmpBuf, O_RDWR | O_CREAT | O_TRUNC, S_IRWXU | S_IRGRP | S_IXGRP | S_IROTH);
    if (fd == -1) {
        errorPrint("%s() LN%d, failed to open temp file: %s\n",
                __func__, __LINE__, tmpBuf);
        taos_free_result(res);
        return -1;
    }

    TAOS_FIELD *fields = taos_fetch_fields(res);
1048

1049 1050
    int32_t  numOfTable  = 0;
    while ((row = taos_fetch_row(res)) != NULL) {
H
Hui Li 已提交
1051

1052 1053 1054
        memset(&tableRecord, 0, sizeof(STableRecord));
        tstrncpy(tableRecord.name, (char *)row[0], fields[0].bytes);
        tstrncpy(tableRecord.metric, metric, TSDB_TABLE_NAME_LEN);
1055

1056 1057
        taosWrite(fd, &tableRecord, sizeof(STableRecord));
        numOfTable++;
1058
    }
1059 1060
    taos_free_result(res);
    lseek(fd, 0, SEEK_SET);
1061

1062 1063 1064 1065 1066 1067 1068 1069 1070 1071 1072
    int maxThreads = g_args.thread_num;
    int tableOfPerFile ;
    if (numOfTable <= g_args.thread_num) {
        tableOfPerFile = 1;
        maxThreads = numOfTable;
    } else {
        tableOfPerFile = numOfTable / g_args.thread_num;
        if (0 != numOfTable % g_args.thread_num) {
            tableOfPerFile += 1;
        }
    }
1073

1074 1075 1076 1077 1078 1079
    char* tblBuf = (char*)calloc(1, tableOfPerFile * sizeof(STableRecord));
    if (NULL == tblBuf){
        errorPrint("%s() LN%d, failed to calloc %" PRIzu "\n",
                __func__, __LINE__, tableOfPerFile * sizeof(STableRecord));
        close(fd);
        return -1;
1080
    }
H
Hui Li 已提交
1081

1082 1083
    int32_t  numOfThread = *totalNumOfThread;
    int      subFd = -1;
1084
    for (; numOfThread <= maxThreads; numOfThread++) {
1085 1086 1087 1088 1089 1090 1091 1092 1093 1094 1095 1096 1097 1098 1099 1100 1101 1102 1103 1104 1105 1106 1107 1108 1109
        memset(tmpBuf, 0, MAX_FILE_NAME_LEN);
        sprintf(tmpBuf, ".tables.tmp.%d", numOfThread);
        subFd = open(tmpBuf, O_RDWR | O_CREAT | O_TRUNC, S_IRWXU | S_IRGRP | S_IXGRP | S_IROTH);
        if (subFd == -1) {
            errorPrint("%s() LN%d, failed to open temp file: %s\n",
                    __func__, __LINE__, tmpBuf);
            for (int32_t loopCnt = 0; loopCnt < numOfThread; loopCnt++) {
                sprintf(tmpBuf, ".tables.tmp.%d", loopCnt);
                (void)remove(tmpBuf);
            }
            sprintf(tmpBuf, ".select-tbname.tmp");
            (void)remove(tmpBuf);
            free(tblBuf);
            close(fd);
            return -1;
        }

        // read tableOfPerFile for fd, write to subFd
        ssize_t readLen = read(fd, tblBuf, tableOfPerFile * sizeof(STableRecord));
        if (readLen <= 0) {
            close(subFd);
            break;
        }
        taosWrite(subFd, tblBuf, readLen);
        close(subFd);
H
Hui Li 已提交
1110 1111
    }

1112 1113
    sprintf(tmpBuf, ".select-tbname.tmp");
    (void)remove(tmpBuf);
1114

1115 1116 1117 1118
    if (fd >= 0) {
        close(fd);
        fd = -1;
    }
1119

1120
    *totalNumOfThread = numOfThread;
1121

1122 1123
    free(tblBuf);
    return 0;
1124 1125
}

1126 1127 1128 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 1176 1177 1178 1179 1180 1181 1182 1183 1184 1185 1186 1187 1188 1189 1190 1191 1192 1193 1194 1195 1196 1197 1198 1199 1200 1201 1202 1203 1204 1205 1206 1207
static int getDbCount()
{
    int count;

    TAOS     *taos = NULL;
    TAOS_RES *result     = NULL;
    char     *command    = NULL;
    TAOS_ROW row;

    command = (char *)malloc(COMMAND_SIZE);
    if (command == NULL) {
        errorPrint("%s() LN%d, failed to allocate command buffer\n", __func__, __LINE__);
        return 0;
    }

    /* Connect to server */
    taos = taos_connect(g_args.host, g_args.user, g_args.password,
            NULL, g_args.port);
    if (NULL == taos) {
        errorPrint("Failed to connect to TDengine server %s\n", g_args.host);
        free(command);
        return 0;
    }

    sprintf(command, "show databases");
    result = taos_query(taos, command);
    int32_t code = taos_errno(result);

    if (0 != code) {
        errorPrint("%s() LN%d, failed to run command: %s, reason: %s\n",
                __func__, __LINE__, command, taos_errstr(result));
        free(command);
        return 0;
    }

    TAOS_FIELD *fields = taos_fetch_fields(result);

    while ((row = taos_fetch_row(result)) != NULL) {
        // sys database name : 'log', but subsequent version changed to 'log'
        if ((strncasecmp(row[TSDB_SHOW_DB_NAME_INDEX], "log",
                        fields[TSDB_SHOW_DB_NAME_INDEX].bytes) == 0)
                && (!g_args.allow_sys)) {
            continue;
        }

        if (g_args.databases) {  // input multi dbs
            for (int i = 0; g_args.arg_list[i]; i++) {
                if (strncasecmp(g_args.arg_list[i],
                            (char *)row[TSDB_SHOW_DB_NAME_INDEX],
                            fields[TSDB_SHOW_DB_NAME_INDEX].bytes) == 0)
                    goto _dump_db_point;
            }
            continue;
        } else if (!g_args.all_databases) {  // only input one db
            if (strncasecmp(g_args.arg_list[0],
                        (char *)row[TSDB_SHOW_DB_NAME_INDEX],
                        fields[TSDB_SHOW_DB_NAME_INDEX].bytes) == 0)
                goto _dump_db_point;
            else
                continue;
        }

_dump_db_point:

        count++;

        if (g_args.databases) {
            if (count > g_args.arg_list_len) break;

        } else if (!g_args.all_databases) {
            if (count >= 1) break;
        }
    }

    if (count == 0) {
        errorPrint("%d databases valid to dump\n", count);
    }

    free(command);
    return count;
}

1208 1209 1210 1211
static int taosDumpOut() {
    TAOS     *taos       = NULL;
    TAOS_RES *result     = NULL;
    char     *command    = NULL;
1212

1213 1214 1215 1216
    TAOS_ROW row;
    FILE *fp = NULL;
    int32_t count = 0;
    STableRecordInfo tableRecordInfo;
1217

1218 1219 1220 1221 1222 1223
    char tmpBuf[4096] = {0};
    if (g_args.outpath[0] != 0) {
        sprintf(tmpBuf, "%s/dbs.sql", g_args.outpath);
    } else {
        sprintf(tmpBuf, "dbs.sql");
    }
1224

1225 1226 1227 1228 1229 1230
    fp = fopen(tmpBuf, "w");
    if (fp == NULL) {
        errorPrint("%s() LN%d, failed to open file %s\n",
                __func__, __LINE__, tmpBuf);
        return -1;
    }
1231

1232 1233 1234 1235 1236 1237 1238 1239
    g_args.dbCount = getDbCount();

    if (0 == g_args.dbCount) {
        errorPrint("%d databases valid to dump\n", g_args.dbCount);
        return -1;
    }

    g_dbInfos = (SDbInfo **)calloc(g_args.dbCount, sizeof(SDbInfo *));
1240 1241 1242 1243 1244
    if (g_dbInfos == NULL) {
        errorPrint("%s() LN%d, failed to allocate memory\n",
                __func__, __LINE__);
        goto _exit_failure;
    }
1245

1246 1247 1248 1249 1250
    command = (char *)malloc(COMMAND_SIZE);
    if (command == NULL) {
        errorPrint("%s() LN%d, failed to allocate memory\n", __func__, __LINE__);
        goto _exit_failure;
    }
1251

1252 1253 1254 1255 1256 1257 1258
    /* Connect to server */
    taos = taos_connect(g_args.host, g_args.user, g_args.password,
            NULL, g_args.port);
    if (taos == NULL) {
        errorPrint("Failed to connect to TDengine server %s\n", g_args.host);
        goto _exit_failure;
    }
1259

1260 1261 1262 1263
    /* --------------------------------- Main Code -------------------------------- */
    /* if (g_args.databases || g_args.all_databases) { // dump part of databases or all databases */
    /*  */
    taosDumpCharset(fp);
1264

1265 1266 1267
    sprintf(command, "show databases");
    result = taos_query(taos, command);
    int32_t code = taos_errno(result);
1268

1269 1270 1271 1272 1273 1274 1275
    if (code != 0) {
        errorPrint("%s() LN%d, failed to run command: %s, reason: %s\n",
                __func__, __LINE__, command, taos_errstr(result));
        goto _exit_failure;
    }

    TAOS_FIELD *fields = taos_fetch_fields(result);
1276

1277 1278 1279
    while ((row = taos_fetch_row(result)) != NULL) {
        // sys database name : 'log', but subsequent version changed to 'log'
        if ((strncasecmp(row[TSDB_SHOW_DB_NAME_INDEX], "log",
1280
                        fields[TSDB_SHOW_DB_NAME_INDEX].bytes) == 0)
1281 1282 1283 1284 1285 1286 1287 1288 1289 1290 1291 1292 1293 1294 1295 1296 1297 1298 1299 1300
                && (!g_args.allow_sys)) {
            continue;
        }

        if (g_args.databases) {  // input multi dbs
            for (int i = 0; g_args.arg_list[i]; i++) {
                if (strncasecmp(g_args.arg_list[i],
                            (char *)row[TSDB_SHOW_DB_NAME_INDEX],
                            fields[TSDB_SHOW_DB_NAME_INDEX].bytes) == 0)
                    goto _dump_db_point;
            }
            continue;
        } else if (!g_args.all_databases) {  // only input one db
            if (strncasecmp(g_args.arg_list[0],
                        (char *)row[TSDB_SHOW_DB_NAME_INDEX],
                        fields[TSDB_SHOW_DB_NAME_INDEX].bytes) == 0)
                goto _dump_db_point;
            else
                continue;
        }
1301

1302
_dump_db_point:
1303

1304 1305 1306 1307 1308 1309
        g_dbInfos[count] = (SDbInfo *)calloc(1, sizeof(SDbInfo));
        if (g_dbInfos[count] == NULL) {
            errorPrint("%s() LN%d, failed to allocate %"PRIu64" memory\n",
                    __func__, __LINE__, (uint64_t)sizeof(SDbInfo));
            goto _exit_failure;
        }
1310

1311
        tstrncpy(g_dbInfos[count]->name, (char *)row[TSDB_SHOW_DB_NAME_INDEX],
1312
                min(TSDB_DB_NAME_LEN, fields[TSDB_SHOW_DB_NAME_INDEX].bytes + 1));
1313 1314 1315 1316 1317 1318 1319
        if (g_args.with_property) {
            g_dbInfos[count]->ntables = *((int32_t *)row[TSDB_SHOW_DB_NTABLES_INDEX]);
            g_dbInfos[count]->vgroups = *((int32_t *)row[TSDB_SHOW_DB_VGROUPS_INDEX]);
            g_dbInfos[count]->replica = *((int16_t *)row[TSDB_SHOW_DB_REPLICA_INDEX]);
            g_dbInfos[count]->quorum = *((int16_t *)row[TSDB_SHOW_DB_QUORUM_INDEX]);
            g_dbInfos[count]->days = *((int16_t *)row[TSDB_SHOW_DB_DAYS_INDEX]);

1320
            tstrncpy(g_dbInfos[count]->keeplist, (char *)row[TSDB_SHOW_DB_KEEP_INDEX],
1321
                    min(32, fields[TSDB_SHOW_DB_KEEP_INDEX].bytes + 1));
1322 1323 1324 1325 1326 1327 1328 1329 1330 1331 1332 1333
            //g_dbInfos[count]->daysToKeep = *((int16_t *)row[TSDB_SHOW_DB_KEEP_INDEX]);
            //g_dbInfos[count]->daysToKeep1;
            //g_dbInfos[count]->daysToKeep2;
            g_dbInfos[count]->cache = *((int32_t *)row[TSDB_SHOW_DB_CACHE_INDEX]);
            g_dbInfos[count]->blocks = *((int32_t *)row[TSDB_SHOW_DB_BLOCKS_INDEX]);
            g_dbInfos[count]->minrows = *((int32_t *)row[TSDB_SHOW_DB_MINROWS_INDEX]);
            g_dbInfos[count]->maxrows = *((int32_t *)row[TSDB_SHOW_DB_MAXROWS_INDEX]);
            g_dbInfos[count]->wallevel = *((int8_t *)row[TSDB_SHOW_DB_WALLEVEL_INDEX]);
            g_dbInfos[count]->fsync = *((int32_t *)row[TSDB_SHOW_DB_FSYNC_INDEX]);
            g_dbInfos[count]->comp = (int8_t)(*((int8_t *)row[TSDB_SHOW_DB_COMP_INDEX]));
            g_dbInfos[count]->cachelast = (int8_t)(*((int8_t *)row[TSDB_SHOW_DB_CACHELAST_INDEX]));

1334 1335 1336
            tstrncpy(g_dbInfos[count]->precision,
                    (char *)row[TSDB_SHOW_DB_PRECISION_INDEX],
                    DB_PRECISION_LEN);
1337 1338 1339
            g_dbInfos[count]->update = *((int8_t *)row[TSDB_SHOW_DB_UPDATE_INDEX]);
        }
        count++;
1340

1341 1342
        if (g_args.databases) {
            if (count > g_args.arg_list_len) break;
1343

1344 1345 1346
        } else if (!g_args.all_databases) {
            if (count >= 1) break;
        }
1347
    }
1348

1349 1350
    if (count == 0) {
        errorPrint("%d databases valid to dump\n", count);
1351
        goto _exit_failure;
1352
    }
1353

1354 1355 1356
    if (g_args.databases || g_args.all_databases) { // case: taosdump --databases dbx dby ...   OR  taosdump --all-databases
        for (int i = 0; i < count; i++) {
            taosDumpDb(g_dbInfos[i], fp, taos);
1357
        }
1358 1359 1360 1361 1362 1363 1364 1365 1366 1367 1368 1369 1370 1371 1372 1373
    } else {
        if (g_args.arg_list_len == 1) {             // case: taosdump <db>
            taosDumpDb(g_dbInfos[0], fp, taos);
        } else {                                        // case: taosdump <db> tablex tabley ...
            taosDumpCreateDbClause(g_dbInfos[0], g_args.with_property, fp);
            fprintf(g_fpOfResult, "\n#### database:                       %s\n",
                    g_dbInfos[0]->name);
            g_resultStatistics.totalDatabasesOfDumpOut++;

            sprintf(command, "use %s", g_dbInfos[0]->name);

            result = taos_query(taos, command);
            code = taos_errno(result);
            if (code != 0) {
                errorPrint("invalid database %s\n", g_dbInfos[0]->name);
                goto _exit_failure;
1374
            }
1375

1376 1377
            fprintf(fp, "USE %s;\n\n", g_dbInfos[0]->name);

1378
            int32_t totalNumOfThread = 1;  // 0: all normal table into .tables.tmp.0
1379 1380 1381 1382 1383 1384
            int  normalTblFd = -1;
            int32_t retCode;
            int superTblCnt = 0 ;
            for (int i = 1; g_args.arg_list[i]; i++) {
                if (taosGetTableRecordInfo(g_args.arg_list[i],
                            &tableRecordInfo, taos) < 0) {
1385
                    errorPrint("input the invalid table %s\n",
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
                            g_args.arg_list[i]);
                    continue;
                }

                if (tableRecordInfo.isMetric) {  // dump all table of this metric
                    int ret = taosDumpStable(
                            tableRecordInfo.tableRecord.metric,
                            fp, taos, g_dbInfos[0]->name);
                    if (0 == ret) {
                        superTblCnt++;
                    }
                    retCode = taosSaveTableOfMetricToTempFile(
                            taos, tableRecordInfo.tableRecord.metric,
                            &totalNumOfThread);
                } else {
                    if (tableRecordInfo.tableRecord.metric[0] != '\0') {  // dump this sub table and it's metric
                        int ret = taosDumpStable(
                                tableRecordInfo.tableRecord.metric,
                                fp, taos, g_dbInfos[0]->name);
                        if (0 == ret) {
                            superTblCnt++;
                        }
                    }
                    retCode = taosSaveAllNormalTableToTempFile(
                            taos, tableRecordInfo.tableRecord.name,
                            tableRecordInfo.tableRecord.metric, &normalTblFd);
                }

                if (retCode < 0) {
                    if (-1 != normalTblFd){
                        taosClose(normalTblFd);
                    }
                    goto _clean_tmp_file;
                }
            }
1421

1422 1423 1424 1425
            // TODO: save dump super table <superTblCnt> into result_output.txt
            fprintf(g_fpOfResult, "# super table counter:               %d\n",
                    superTblCnt);
            g_resultStatistics.totalSuperTblsOfDumpOut += superTblCnt;
1426

1427 1428 1429
            if (-1 != normalTblFd){
                taosClose(normalTblFd);
            }
1430

1431
            // start multi threads to dumpout
1432

1433
            taosStartDumpOutWorkThreads(totalNumOfThread,
1434 1435
                    g_dbInfos[0]->name,
                    getPrecisionByString(g_dbInfos[0]->precision));
1436

1437 1438 1439 1440 1441 1442 1443
            char tmpFileName[MAX_FILE_NAME_LEN];
_clean_tmp_file:
            for (int loopCnt = 0; loopCnt < totalNumOfThread; loopCnt++) {
                sprintf(tmpFileName, ".tables.tmp.%d", loopCnt);
                remove(tmpFileName);
            }
        }
1444 1445
    }

1446 1447 1448 1449 1450 1451 1452 1453
    /* Close the handle and return */
    fclose(fp);
    taos_close(taos);
    taos_free_result(result);
    tfree(command);
    taosFreeDbInfos();
    fprintf(stderr, "dump out rows: %" PRId64 "\n", g_totalDumpOutRows);
    return 0;
1454 1455

_exit_failure:
1456 1457 1458 1459 1460 1461 1462
    fclose(fp);
    taos_close(taos);
    taos_free_result(result);
    tfree(command);
    taosFreeDbInfos();
    errorPrint("dump out rows: %" PRId64 "\n", g_totalDumpOutRows);
    return -1;
1463 1464
}

1465
static int taosGetTableDes(
1466
        char* dbName, char *table,
1467
        STableDef *stableDes, TAOS* taosCon, bool isSuperTable) {
1468 1469 1470
    TAOS_ROW row = NULL;
    TAOS_RES* res = NULL;
    int count = 0;
1471

1472 1473
    char sqlstr[COMMAND_SIZE];
    sprintf(sqlstr, "describe %s.%s;", dbName, table);
1474

1475 1476 1477 1478 1479 1480 1481 1482
    res = taos_query(taosCon, sqlstr);
    int32_t code = taos_errno(res);
    if (code != 0) {
        errorPrint("%s() LN%d, failed to run command <%s>, reason:%s\n",
                __func__, __LINE__, sqlstr, taos_errstr(res));
        taos_free_result(res);
        return -1;
    }
1483

1484 1485
    TAOS_FIELD *fields = taos_fetch_fields(res);

1486
    tstrncpy(stableDes->name, table, TSDB_TABLE_NAME_LEN);
1487
    while ((row = taos_fetch_row(res)) != NULL) {
1488
        tstrncpy(stableDes->cols[count].field,
1489
                (char *)row[TSDB_DESCRIBE_METRIC_FIELD_INDEX],
1490 1491
                min(TSDB_COL_NAME_LEN + 1,
                    fields[TSDB_DESCRIBE_METRIC_FIELD_INDEX].bytes + 1));
1492
        tstrncpy(stableDes->cols[count].type,
1493
                (char *)row[TSDB_DESCRIBE_METRIC_TYPE_INDEX],
1494
                min(16, fields[TSDB_DESCRIBE_METRIC_TYPE_INDEX].bytes + 1));
1495
        stableDes->cols[count].length =
1496
            *((int *)row[TSDB_DESCRIBE_METRIC_LENGTH_INDEX]);
1497
        tstrncpy(stableDes->cols[count].note,
1498
                (char *)row[TSDB_DESCRIBE_METRIC_NOTE_INDEX],
1499 1500
                min(COL_NOTE_LEN,
                    fields[TSDB_DESCRIBE_METRIC_NOTE_INDEX].bytes + 1));
1501 1502 1503

        count++;
    }
1504

1505 1506
    taos_free_result(res);
    res = NULL;
1507

1508 1509 1510
    if (isSuperTable) {
        return count;
    }
1511

1512
    // if child-table have tag, using  select tagName from table to get tagValue
1513
    for (int i = 0 ; i < count; i++) {
1514
        if (strcmp(stableDes->cols[i].note, "TAG") != 0) continue;
1515

1516
        sprintf(sqlstr, "select %s from %s.%s",
1517
                stableDes->cols[i].field, dbName, table);
1518

1519 1520 1521 1522 1523 1524 1525 1526
        res = taos_query(taosCon, sqlstr);
        code = taos_errno(res);
        if (code != 0) {
            errorPrint("%s() LN%d, failed to run command <%s>, reason:%s\n",
                    __func__, __LINE__, sqlstr, taos_errstr(res));
            taos_free_result(res);
            return -1;
        }
1527

1528
        fields = taos_fetch_fields(res);
1529

1530 1531 1532 1533 1534 1535 1536
        row = taos_fetch_row(res);
        if (NULL == row) {
            errorPrint("%s() LN%d, fetch failed to run command <%s>, reason:%s\n",
                    __func__, __LINE__, sqlstr, taos_errstr(res));
            taos_free_result(res);
            return -1;
        }
1537

1538
        if (row[0] == NULL) {
1539
            sprintf(stableDes->cols[i].note, "%s", "NULL");
1540 1541 1542 1543
            taos_free_result(res);
            res = NULL;
            continue;
        }
1544

1545 1546 1547 1548 1549
        int32_t* length = taos_fetch_lengths(res);

        //int32_t* length = taos_fetch_lengths(tmpResult);
        switch (fields[0].type) {
            case TSDB_DATA_TYPE_BOOL:
1550
                sprintf(stableDes->cols[i].note, "%d",
1551 1552 1553
                        ((((int32_t)(*((char *)row[0]))) == 1) ? 1 : 0));
                break;
            case TSDB_DATA_TYPE_TINYINT:
1554
                sprintf(stableDes->cols[i].note, "%d", *((int8_t *)row[0]));
1555 1556
                break;
            case TSDB_DATA_TYPE_SMALLINT:
1557
                sprintf(stableDes->cols[i].note, "%d", *((int16_t *)row[0]));
1558 1559
                break;
            case TSDB_DATA_TYPE_INT:
1560
                sprintf(stableDes->cols[i].note, "%d", *((int32_t *)row[0]));
1561 1562
                break;
            case TSDB_DATA_TYPE_BIGINT:
1563
                sprintf(stableDes->cols[i].note, "%" PRId64 "", *((int64_t *)row[0]));
1564 1565
                break;
            case TSDB_DATA_TYPE_FLOAT:
1566
                sprintf(stableDes->cols[i].note, "%f", GET_FLOAT_VAL(row[0]));
1567 1568
                break;
            case TSDB_DATA_TYPE_DOUBLE:
1569
                sprintf(stableDes->cols[i].note, "%f", GET_DOUBLE_VAL(row[0]));
1570 1571 1572
                break;
            case TSDB_DATA_TYPE_BINARY:
                {
1573 1574
                    memset(stableDes->cols[i].note, 0, sizeof(stableDes->cols[i].note));
                    stableDes->cols[i].note[0] = '\'';
1575 1576
                    char tbuf[COL_NOTE_LEN];
                    converStringToReadable((char *)row[0], length[0], tbuf, COL_NOTE_LEN);
1577
                    char* pstr = stpcpy(&(stableDes->cols[i].note[1]), tbuf);
1578 1579 1580 1581 1582
                    *(pstr++) = '\'';
                    break;
                }
            case TSDB_DATA_TYPE_NCHAR:
                {
1583
                    memset(stableDes->cols[i].note, 0, sizeof(stableDes->cols[i].note));
1584 1585
                    char tbuf[COL_NOTE_LEN-2];    // need reserve 2 bytes for ' '
                    convertNCharToReadable((char *)row[0], length[0], tbuf, COL_NOTE_LEN);
1586
                    sprintf(stableDes->cols[i].note, "\'%s\'", tbuf);
1587 1588 1589
                    break;
                }
            case TSDB_DATA_TYPE_TIMESTAMP:
1590
                sprintf(stableDes->cols[i].note, "%" PRId64 "", *(int64_t *)row[0]);
1591 1592 1593 1594 1595 1596 1597 1598 1599 1600 1601 1602 1603 1604 1605
#if 0
                if (!g_args.mysqlFlag) {
                    sprintf(tableDes->cols[i].note, "%" PRId64 "", *(int64_t *)row[0]);
                } else {
                    char buf[64] = "\0";
                    int64_t ts = *((int64_t *)row[0]);
                    time_t tt = (time_t)(ts / 1000);
                    struct tm *ptm = localtime(&tt);
                    strftime(buf, 64, "%y-%m-%d %H:%M:%S", ptm);
                    sprintf(tableDes->cols[i].note, "\'%s.%03d\'", buf, (int)(ts % 1000));
                }
#endif
                break;
            default:
                break;
1606
        }
1607 1608 1609

        taos_free_result(res);
        res = NULL;
1610
    }
1611

1612 1613
    return count;
}
1614

1615
static int convertSchemaToAvroSchema(STableDef *stableDes, char **avroSchema)
1616 1617 1618 1619
{
    errorPrint("%s() LN%d TODO: covert table schema to avro schema\n",
            __func__, __LINE__);
    return 0;
1620 1621
}

1622
static int32_t taosDumpTable(
1623
        char *tbName, char *metric,
1624
        FILE *fp, TAOS* taosCon, char* dbName, int precision) {
1625
    int count = 0;
1626

1627 1628
    STableDef *tableDes = (STableDef *)calloc(1, sizeof(STableDef)
            + sizeof(SColDes) * TSDB_MAX_COLUMNS);
1629

1630 1631 1632
    if (metric != NULL && metric[0] != '\0') {  // dump table schema which is created by using super table
        /*
           count = taosGetTableDes(metric, tableDes, taosCon);
1633

1634 1635 1636 1637
           if (count < 0) {
           free(tableDes);
           return -1;
           }
1638

1639
           taosDumpCreateTableClause(tableDes, count, fp);
1640

1641 1642
           memset(tableDes, 0, sizeof(STableDef) + sizeof(SColDes) * TSDB_MAX_COLUMNS);
           */
1643

1644
        count = taosGetTableDes(dbName, tbName, tableDes, taosCon, false);
1645

1646 1647 1648 1649
        if (count < 0) {
            free(tableDes);
            return -1;
        }
1650

1651 1652
        // create child-table using super-table
        taosDumpCreateMTableClause(tableDes, metric, count, fp, dbName);
1653

1654
    } else {  // dump table definition
1655
        count = taosGetTableDes(dbName, tbName, tableDes, taosCon, false);
1656

1657 1658 1659 1660 1661 1662 1663
        if (count < 0) {
            free(tableDes);
            return -1;
        }

        // create normal-table or super-table
        taosDumpCreateTableClause(tableDes, count, fp, dbName);
1664 1665
    }

1666 1667 1668 1669 1670 1671
    char *jsonAvroSchema = NULL;
    if (g_args.avro) {
        convertSchemaToAvroSchema(tableDes, &jsonAvroSchema);
    }

    free(tableDes);
1672

1673 1674
    int32_t ret = 0;
    if (!g_args.schemaonly) {
1675
        ret = taosDumpTableData(fp, tbName, taosCon, dbName, precision,
1676 1677
            jsonAvroSchema);
    }
1678

1679
    return ret;
1680 1681
}

1682 1683 1684 1685 1686 1687 1688 1689 1690 1691 1692 1693 1694 1695 1696 1697 1698
static void taosDumpCreateDbClause(
        SDbInfo *dbInfo, bool isDumpProperty, FILE *fp) {
    char sqlstr[TSDB_MAX_SQL_LEN] = {0};

    char *pstr = sqlstr;
    pstr += sprintf(pstr, "CREATE DATABASE IF NOT EXISTS %s ", dbInfo->name);
    if (isDumpProperty) {
        pstr += sprintf(pstr,
                "REPLICA %d QUORUM %d DAYS %d KEEP %s CACHE %d BLOCKS %d MINROWS %d MAXROWS %d FSYNC %d CACHELAST %d COMP %d PRECISION '%s' UPDATE %d",
                dbInfo->replica, dbInfo->quorum, dbInfo->days,
                dbInfo->keeplist,
                dbInfo->cache,
                dbInfo->blocks, dbInfo->minrows, dbInfo->maxrows,
                dbInfo->fsync,
                dbInfo->cachelast,
                dbInfo->comp, dbInfo->precision, dbInfo->update);
    }
1699

1700 1701
    pstr += sprintf(pstr, ";");
    fprintf(fp, "%s\n\n", sqlstr);
1702 1703
}

1704
static void* taosDumpOutWorkThreadFp(void *arg)
1705
{
1706 1707 1708 1709
    SThreadParaObj *pThread = (SThreadParaObj*)arg;
    STableRecord    tableRecord;
    int fd;

1710 1711
    setThreadName("dumpOutWorkThrd");

1712 1713 1714 1715 1716 1717 1718 1719
    char tmpBuf[4096] = {0};
    sprintf(tmpBuf, ".tables.tmp.%d", pThread->threadIndex);
    fd = open(tmpBuf, O_RDWR | O_CREAT, S_IRWXU | S_IRGRP | S_IXGRP | S_IROTH);
    if (fd == -1) {
        errorPrint("%s() LN%d, failed to open temp file: %s\n",
                __func__, __LINE__, tmpBuf);
        return NULL;
    }
1720

1721 1722
    FILE *fp = NULL;
    memset(tmpBuf, 0, 4096);
1723

1724 1725 1726 1727 1728 1729 1730
    if (g_args.outpath[0] != 0) {
        sprintf(tmpBuf, "%s/%s.tables.%d.sql",
                g_args.outpath, pThread->dbName, pThread->threadIndex);
    } else {
        sprintf(tmpBuf, "%s.tables.%d.sql",
                pThread->dbName, pThread->threadIndex);
    }
1731

1732 1733 1734 1735 1736 1737 1738
    fp = fopen(tmpBuf, "w");
    if (fp == NULL) {
        errorPrint("%s() LN%d, failed to open file %s\n",
                __func__, __LINE__, tmpBuf);
        close(fd);
        return NULL;
    }
1739

1740 1741
    memset(tmpBuf, 0, 4096);
    sprintf(tmpBuf, "use %s", pThread->dbName);
H
Hui Li 已提交
1742

1743 1744 1745 1746 1747 1748
    TAOS_RES* tmpResult = taos_query(pThread->taosCon, tmpBuf);
    int32_t code = taos_errno(tmpResult);
    if (code != 0) {
        errorPrint("%s() LN%d, invalid database %s. reason: %s\n",
                __func__, __LINE__, pThread->dbName, taos_errstr(tmpResult));
        taos_free_result(tmpResult);
H
Hui Li 已提交
1749
        fclose(fp);
1750 1751 1752
        close(fd);
        return NULL;
    }
1753

1754 1755 1756 1757 1758 1759 1760 1761 1762 1763 1764 1765
#if 0
    int     fileNameIndex = 1;
    int     tablesInOneFile = 0;
#endif
    int64_t lastRowsPrint = 5000000;
    fprintf(fp, "USE %s;\n\n", pThread->dbName);
    while (1) {
        ssize_t readLen = read(fd, &tableRecord, sizeof(STableRecord));
        if (readLen <= 0) break;

        int ret = taosDumpTable(
                tableRecord.name, tableRecord.metric,
1766 1767
                fp, pThread->taosCon, pThread->dbName,
                pThread->precision);
1768 1769 1770 1771 1772 1773 1774 1775 1776 1777 1778 1779 1780 1781 1782 1783 1784 1785 1786 1787 1788 1789 1790 1791 1792 1793 1794 1795 1796 1797 1798 1799 1800 1801 1802 1803 1804 1805
        if (ret >= 0) {
            // TODO: sum table count and table rows by self
            pThread->tablesOfDumpOut++;
            pThread->rowsOfDumpOut += ret;

            if (pThread->rowsOfDumpOut >= lastRowsPrint) {
                printf(" %"PRId64 " rows already be dumpout from database %s\n",
                        pThread->rowsOfDumpOut, pThread->dbName);
                lastRowsPrint += 5000000;
            }

#if 0
            tablesInOneFile++;
            if (tablesInOneFile >= g_args.table_batch) {
                fclose(fp);
                tablesInOneFile = 0;

                memset(tmpBuf, 0, 4096);
                if (g_args.outpath[0] != 0) {
                    sprintf(tmpBuf, "%s/%s.tables.%d-%d.sql",
                            g_args.outpath, pThread->dbName,
                            pThread->threadIndex, fileNameIndex);
                } else {
                    sprintf(tmpBuf, "%s.tables.%d-%d.sql",
                            pThread->dbName, pThread->threadIndex, fileNameIndex);
                }
                fileNameIndex++;

                fp = fopen(tmpBuf, "w");
                if (fp == NULL) {
                    errorPrint("%s() LN%d, failed to open file %s\n",
                            __func__, __LINE__, tmpBuf);
                    close(fd);
                    taos_free_result(tmpResult);
                    return NULL;
                }
            }
#endif
H
Hui Li 已提交
1806
        }
1807
    }
1808

1809 1810 1811
    taos_free_result(tmpResult);
    close(fd);
    fclose(fp);
1812

1813
    return NULL;
1814 1815
}

1816
static void taosStartDumpOutWorkThreads(int32_t  numOfThread, char *dbName, int precision)
1817
{
1818 1819 1820 1821 1822 1823 1824 1825
    pthread_attr_t thattr;
    SThreadParaObj *threadObj =
        (SThreadParaObj *)calloc(numOfThread, sizeof(SThreadParaObj));

    if (threadObj == NULL) {
        errorPrint("%s() LN%d, memory allocation failed!\n",
                __func__, __LINE__);
        return;
1826 1827
    }

1828 1829 1830 1831 1832 1833 1834
    for (int t = 0; t < numOfThread; ++t) {
        SThreadParaObj *pThread = threadObj + t;
        pThread->rowsOfDumpOut = 0;
        pThread->tablesOfDumpOut = 0;
        pThread->threadIndex = t;
        pThread->totalThreads = numOfThread;
        tstrncpy(pThread->dbName, dbName, TSDB_DB_NAME_LEN);
1835
        pThread->precision = precision;
1836 1837 1838 1839
        pThread->taosCon = taos_connect(g_args.host, g_args.user, g_args.password,
            NULL, g_args.port);
        if (pThread->taosCon == NULL) {
            errorPrint("Failed to connect to TDengine server %s\n", g_args.host);
1840
            free(threadObj);
1841 1842 1843 1844 1845 1846 1847 1848 1849 1850 1851 1852 1853
            return;
        }
        pthread_attr_init(&thattr);
        pthread_attr_setdetachstate(&thattr, PTHREAD_CREATE_JOINABLE);

        if (pthread_create(&(pThread->threadID), &thattr,
                    taosDumpOutWorkThreadFp,
                    (void*)pThread) != 0) {
            errorPrint("%s() LN%d, thread:%d failed to start\n",
                   __func__, __LINE__, pThread->threadIndex);
            exit(-1);
        }
    }
1854

1855 1856 1857
    for (int32_t t = 0; t < numOfThread; ++t) {
        pthread_join(threadObj[t].threadID, NULL);
    }
1858

1859 1860 1861 1862 1863 1864 1865 1866 1867 1868 1869 1870 1871 1872 1873
    // TODO: sum all thread dump table count and rows of per table, then save into result_output.txt
    int64_t   totalRowsOfDumpOut = 0;
    int64_t   totalChildTblsOfDumpOut = 0;
    for (int32_t t = 0; t < numOfThread; ++t) {
        totalChildTblsOfDumpOut += threadObj[t].tablesOfDumpOut;
        totalRowsOfDumpOut      += threadObj[t].rowsOfDumpOut;
    }

    fprintf(g_fpOfResult, "# child table counter:               %"PRId64"\n",
            totalChildTblsOfDumpOut);
    fprintf(g_fpOfResult, "# row counter:                       %"PRId64"\n",
            totalRowsOfDumpOut);
    g_resultStatistics.totalChildTblsOfDumpOut += totalChildTblsOfDumpOut;
    g_resultStatistics.totalRowsOfDumpOut      += totalRowsOfDumpOut;
    free(threadObj);
1874 1875
}

1876 1877
static int32_t taosDumpStable(char *table, FILE *fp,
        TAOS* taosCon, char* dbName) {
1878

1879 1880 1881 1882
    uint64_t sizeOfTableDes =
        (uint64_t)(sizeof(STableDef) + sizeof(SColDes) * TSDB_MAX_COLUMNS);
    STableDef *stableDes = (STableDef *)calloc(1, sizeOfTableDes);
    if (NULL == stableDes) {
1883 1884 1885 1886
        errorPrint("%s() LN%d, failed to allocate %"PRIu64" memory\n",
                __func__, __LINE__, sizeOfTableDes);
        exit(-1);
    }
1887

1888
    int count = taosGetTableDes(dbName, table, stableDes, taosCon, true);
1889

1890
    if (count < 0) {
1891
        free(stableDes);
1892 1893 1894 1895
        errorPrint("%s() LN%d, failed to get stable[%s] schema\n",
               __func__, __LINE__, table);
        exit(-1);
    }
1896

1897
    taosDumpCreateTableClause(stableDes, count, fp, dbName);
1898

1899
    free(stableDes);
1900
    return 0;
1901 1902
}

1903
static int32_t taosDumpCreateSuperTableClause(TAOS* taosCon, char* dbName, FILE *fp)
1904
{
1905 1906 1907 1908
    TAOS_ROW row;
    int fd = -1;
    STableRecord tableRecord;
    char sqlstr[TSDB_MAX_SQL_LEN] = {0};
1909

1910
    sprintf(sqlstr, "show %s.stables", dbName);
1911

1912 1913 1914 1915 1916 1917 1918 1919
    TAOS_RES* res = taos_query(taosCon, sqlstr);
    int32_t  code = taos_errno(res);
    if (code != 0) {
        errorPrint("%s() LN%d, failed to run command <%s>, reason: %s\n",
                __func__, __LINE__, sqlstr, taos_errstr(res));
        taos_free_result(res);
        exit(-1);
    }
1920

1921 1922 1923 1924 1925 1926 1927 1928 1929 1930 1931 1932 1933
    TAOS_FIELD *fields = taos_fetch_fields(res);

    char     tmpFileName[MAX_FILE_NAME_LEN];
    memset(tmpFileName, 0, MAX_FILE_NAME_LEN);
    sprintf(tmpFileName, ".stables.tmp");
    fd = open(tmpFileName, O_RDWR | O_CREAT, S_IRWXU | S_IRGRP | S_IXGRP | S_IROTH);
    if (fd == -1) {
        errorPrint("%s() LN%d, failed to open temp file: %s\n",
                __func__, __LINE__, tmpFileName);
        taos_free_result(res);
        (void)remove(".stables.tmp");
        exit(-1);
    }
1934

1935 1936
    while ((row = taos_fetch_row(res)) != NULL) {
        memset(&tableRecord, 0, sizeof(STableRecord));
1937 1938
        tstrncpy(tableRecord.name, (char *)row[TSDB_SHOW_TABLES_NAME_INDEX],
                min(TSDB_TABLE_NAME_LEN,
1939
                    fields[TSDB_SHOW_TABLES_NAME_INDEX].bytes + 1));
1940 1941
        taosWrite(fd, &tableRecord, sizeof(STableRecord));
    }
1942

1943 1944
    taos_free_result(res);
    (void)lseek(fd, 0, SEEK_SET);
1945

1946 1947 1948 1949
    int superTblCnt = 0;
    while (1) {
        ssize_t readLen = read(fd, &tableRecord, sizeof(STableRecord));
        if (readLen <= 0) break;
1950

1951 1952 1953 1954
        int ret = taosDumpStable(tableRecord.name, fp, taosCon, dbName);
        if (0 == ret) {
            superTblCnt++;
        }
1955
    }
1956

1957 1958 1959
    // TODO: save dump super table <superTblCnt> into result_output.txt
    fprintf(g_fpOfResult, "# super table counter:               %d\n", superTblCnt);
    g_resultStatistics.totalSuperTblsOfDumpOut += superTblCnt;
1960

1961 1962
    close(fd);
    (void)remove(".stables.tmp");
1963

1964
    return 0;
1965 1966 1967
}


1968 1969 1970 1971
static int taosDumpDb(SDbInfo *dbInfo, FILE *fp, TAOS *taosCon) {
    TAOS_ROW row;
    int fd = -1;
    STableRecord tableRecord;
1972

1973
    taosDumpCreateDbClause(dbInfo, g_args.with_property, fp);
1974

1975 1976 1977
    fprintf(g_fpOfResult, "\n#### database:                       %s\n",
            dbInfo->name);
    g_resultStatistics.totalDatabasesOfDumpOut++;
1978

1979
    char sqlstr[TSDB_MAX_SQL_LEN] = {0};
1980

1981
    fprintf(fp, "USE %s;\n\n", dbInfo->name);
1982

1983
    (void)taosDumpCreateSuperTableClause(taosCon, dbInfo->name, fp);
1984

1985
    sprintf(sqlstr, "show %s.tables", dbInfo->name);
1986

1987 1988 1989 1990 1991 1992 1993 1994
    TAOS_RES* res = taos_query(taosCon, sqlstr);
    int code = taos_errno(res);
    if (code != 0) {
        errorPrint("%s() LN%d, failed to run command <%s>, reason:%s\n",
                __func__, __LINE__, sqlstr, taos_errstr(res));
        taos_free_result(res);
        return -1;
    }
1995

1996 1997 1998 1999 2000 2001 2002 2003 2004 2005
    char tmpBuf[MAX_FILE_NAME_LEN];
    memset(tmpBuf, 0, MAX_FILE_NAME_LEN);
    sprintf(tmpBuf, ".show-tables.tmp");
    fd = open(tmpBuf, O_RDWR | O_CREAT | O_TRUNC, S_IRWXU | S_IRGRP | S_IXGRP | S_IROTH);
    if (fd == -1) {
        errorPrint("%s() LN%d, failed to open temp file: %s\n",
                __func__, __LINE__, tmpBuf);
        taos_free_result(res);
        return -1;
    }
2006

2007
    TAOS_FIELD *fields = taos_fetch_fields(res);
2008

2009 2010 2011 2012
    int32_t  numOfTable  = 0;
    while ((row = taos_fetch_row(res)) != NULL) {
        memset(&tableRecord, 0, sizeof(STableRecord));
        tstrncpy(tableRecord.name, (char *)row[TSDB_SHOW_TABLES_NAME_INDEX],
2013
                min(TSDB_TABLE_NAME_LEN,
2014
                    fields[TSDB_SHOW_TABLES_NAME_INDEX].bytes + 1));
2015
        tstrncpy(tableRecord.metric, (char *)row[TSDB_SHOW_TABLES_METRIC_INDEX],
2016
                min(TSDB_TABLE_NAME_LEN,
2017
                    fields[TSDB_SHOW_TABLES_METRIC_INDEX].bytes + 1));
H
Hui Li 已提交
2018

2019
        taosWrite(fd, &tableRecord, sizeof(STableRecord));
2020

2021
        numOfTable++;
2022
    }
2023 2024
    taos_free_result(res);
    lseek(fd, 0, SEEK_SET);
H
Hui Li 已提交
2025

2026 2027 2028 2029 2030 2031 2032 2033 2034 2035
    int maxThreads = g_args.thread_num;
    int tableOfPerFile ;
    if (numOfTable <= g_args.thread_num) {
        tableOfPerFile = 1;
        maxThreads = numOfTable;
    } else {
        tableOfPerFile = numOfTable / g_args.thread_num;
        if (0 != numOfTable % g_args.thread_num) {
            tableOfPerFile += 1;
        }
H
Hui Li 已提交
2036
    }
2037

2038 2039 2040 2041 2042 2043 2044
    char* tblBuf = (char*)calloc(1, tableOfPerFile * sizeof(STableRecord));
    if (NULL == tblBuf){
        errorPrint("failed to calloc %" PRIzu "\n",
                tableOfPerFile * sizeof(STableRecord));
        close(fd);
        return -1;
    }
H
Hui Li 已提交
2045

2046 2047 2048 2049 2050 2051 2052 2053 2054 2055 2056 2057 2058 2059 2060 2061 2062 2063 2064
    int32_t  numOfThread = 0;
    int      subFd = -1;
    for (numOfThread = 0; numOfThread < maxThreads; numOfThread++) {
        memset(tmpBuf, 0, MAX_FILE_NAME_LEN);
        sprintf(tmpBuf, ".tables.tmp.%d", numOfThread);
        subFd = open(tmpBuf, O_RDWR | O_CREAT | O_TRUNC, S_IRWXU | S_IRGRP | S_IXGRP | S_IROTH);
        if (subFd == -1) {
            errorPrint("%s() LN%d, failed to open temp file: %s\n",
                    __func__, __LINE__, tmpBuf);
            for (int32_t loopCnt = 0; loopCnt < numOfThread; loopCnt++) {
                sprintf(tmpBuf, ".tables.tmp.%d", loopCnt);
                (void)remove(tmpBuf);
            }
            sprintf(tmpBuf, ".show-tables.tmp");
            (void)remove(tmpBuf);
            free(tblBuf);
            close(fd);
            return -1;
        }
2065

2066 2067 2068 2069 2070 2071 2072 2073 2074
        // read tableOfPerFile for fd, write to subFd
        ssize_t readLen = read(fd, tblBuf, tableOfPerFile * sizeof(STableRecord));
        if (readLen <= 0) {
            close(subFd);
            break;
        }
        taosWrite(subFd, tblBuf, readLen);
        close(subFd);
    }
2075

2076
    sprintf(tmpBuf, ".show-tables.tmp");
H
Hui Li 已提交
2077
    (void)remove(tmpBuf);
2078

2079 2080 2081 2082 2083 2084
    if (fd >= 0) {
        close(fd);
        fd = -1;
    }

    // start multi threads to dumpout
2085 2086
    taosStartDumpOutWorkThreads(numOfThread, dbInfo->name,
            getPrecisionByString(dbInfo->precision));
2087 2088 2089 2090 2091 2092 2093
    for (int loopCnt = 0; loopCnt < numOfThread; loopCnt++) {
        sprintf(tmpBuf, ".tables.tmp.%d", loopCnt);
        (void)remove(tmpBuf);
    }

    free(tblBuf);
    return 0;
2094 2095
}

2096 2097
static void taosDumpCreateTableClause(STableDef *tableDes, int numOfCols,
        FILE *fp, char* dbName) {
2098 2099 2100
    int counter = 0;
    int count_temp = 0;
    char sqlstr[COMMAND_SIZE];
2101

2102
    char* pstr = sqlstr;
2103

2104 2105
    pstr += sprintf(sqlstr, "CREATE TABLE IF NOT EXISTS %s.%s",
            dbName, tableDes->name);
2106

2107 2108
    for (; counter < numOfCols; counter++) {
        if (tableDes->cols[counter].note[0] != '\0') break;
2109

2110 2111 2112 2113 2114 2115 2116
        if (counter == 0) {
            pstr += sprintf(pstr, " (%s %s",
                    tableDes->cols[counter].field, tableDes->cols[counter].type);
        } else {
            pstr += sprintf(pstr, ", %s %s",
                    tableDes->cols[counter].field, tableDes->cols[counter].type);
        }
2117

2118 2119 2120 2121
        if (strcasecmp(tableDes->cols[counter].type, "binary") == 0 ||
                strcasecmp(tableDes->cols[counter].type, "nchar") == 0) {
            pstr += sprintf(pstr, "(%d)", tableDes->cols[counter].length);
        }
2122 2123
    }

2124
    count_temp = counter;
2125

2126 2127 2128 2129 2130 2131 2132 2133
    for (; counter < numOfCols; counter++) {
        if (counter == count_temp) {
            pstr += sprintf(pstr, ") TAGS (%s %s",
                    tableDes->cols[counter].field, tableDes->cols[counter].type);
        } else {
            pstr += sprintf(pstr, ", %s %s",
                    tableDes->cols[counter].field, tableDes->cols[counter].type);
        }
2134

2135 2136 2137 2138
        if (strcasecmp(tableDes->cols[counter].type, "binary") == 0 ||
                strcasecmp(tableDes->cols[counter].type, "nchar") == 0) {
            pstr += sprintf(pstr, "(%d)", tableDes->cols[counter].length);
        }
2139 2140
    }

2141
    pstr += sprintf(pstr, ");");
2142

2143
    fprintf(fp, "%s\n\n", sqlstr);
2144 2145
}

2146 2147
static void taosDumpCreateMTableClause(STableDef *tableDes, char *metric,
        int numOfCols, FILE *fp, char* dbName) {
2148 2149 2150 2151 2152 2153 2154 2155 2156
    int counter = 0;
    int count_temp = 0;

    char* tmpBuf = (char *)malloc(COMMAND_SIZE);
    if (tmpBuf == NULL) {
        errorPrint("%s() LN%d, failed to allocate %d memory\n",
               __func__, __LINE__, COMMAND_SIZE);
        return;
    }
2157

2158 2159
    char *pstr = NULL;
    pstr = tmpBuf;
2160

2161 2162 2163
    pstr += sprintf(tmpBuf,
            "CREATE TABLE IF NOT EXISTS %s.%s USING %s.%s TAGS (",
            dbName, tableDes->name, dbName, metric);
2164

2165 2166 2167
    for (; counter < numOfCols; counter++) {
        if (tableDes->cols[counter].note[0] != '\0') break;
    }
2168

2169 2170 2171 2172 2173 2174 2175 2176 2177 2178 2179 2180 2181 2182 2183 2184 2185 2186 2187 2188 2189 2190
    assert(counter < numOfCols);
    count_temp = counter;

    for (; counter < numOfCols; counter++) {
        if (counter != count_temp) {
            if (strcasecmp(tableDes->cols[counter].type, "binary") == 0 ||
                    strcasecmp(tableDes->cols[counter].type, "nchar") == 0) {
                //pstr += sprintf(pstr, ", \'%s\'", tableDes->cols[counter].note);
                pstr += sprintf(pstr, ", %s", tableDes->cols[counter].note);
            } else {
                pstr += sprintf(pstr, ", %s", tableDes->cols[counter].note);
            }
        } else {
            if (strcasecmp(tableDes->cols[counter].type, "binary") == 0 ||
                    strcasecmp(tableDes->cols[counter].type, "nchar") == 0) {
                //pstr += sprintf(pstr, "\'%s\'", tableDes->cols[counter].note);
                pstr += sprintf(pstr, "%s", tableDes->cols[counter].note);
            } else {
                pstr += sprintf(pstr, "%s", tableDes->cols[counter].note);
            }
            /* pstr += sprintf(pstr, "%s", tableDes->cols[counter].note); */
        }
2191

2192 2193 2194 2195
        /* if (strcasecmp(tableDes->cols[counter].type, "binary") == 0 || strcasecmp(tableDes->cols[counter].type, "nchar")
         * == 0) { */
        /*     pstr += sprintf(pstr, "(%d)", tableDes->cols[counter].length); */
        /* } */
2196 2197
    }

2198
    pstr += sprintf(pstr, ");");
2199

2200 2201
    fprintf(fp, "%s\n", tmpBuf);
    free(tmpBuf);
2202 2203
}

2204 2205 2206 2207 2208 2209
static int writeSchemaToAvro(char *jsonAvroSchema)
{
    errorPrint("%s() LN%d, TODO: implement write schema to avro",
            __func__, __LINE__);
    return 0;
}
2210

2211 2212 2213
static int64_t writeResultToAvro(TAOS_RES *res)
{
    errorPrint("%s() LN%d, TODO: implementation need\n", __func__, __LINE__);
2214
    return 0;
2215
}
2216

2217 2218 2219
static int64_t writeResultToSql(TAOS_RES *res, FILE *fp, char *dbName, char *tbName)
{
    int64_t    totalRows     = 0;
2220

2221 2222 2223 2224 2225 2226
    int32_t  sql_buf_len = g_args.max_sql_len;
    char* tmpBuffer = (char *)calloc(1, sql_buf_len + 128);
    if (tmpBuffer == NULL) {
        errorPrint("failed to allocate %d memory\n", sql_buf_len + 128);
        return -1;
    }
2227

2228
    char *pstr = tmpBuffer;
2229

2230 2231 2232 2233 2234
    TAOS_ROW row = NULL;
    int numFields = 0;
    int rowFlag = 0;
    int64_t    lastRowsPrint = 5000000;
    int count = 0;
2235

2236 2237 2238
    numFields = taos_field_count(res);
    assert(numFields > 0);
    TAOS_FIELD *fields = taos_fetch_fields(res);
2239

2240 2241
    int32_t  curr_sqlstr_len = 0;
    int32_t  total_sqlstr_len = 0;
2242

2243 2244
    while ((row = taos_fetch_row(res)) != NULL) {
        curr_sqlstr_len = 0;
2245

2246 2247 2248 2249 2250 2251
        int32_t* length = taos_fetch_lengths(res);   // act len

        if (count == 0) {
            total_sqlstr_len = 0;
            curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len,
                    "INSERT INTO %s.%s VALUES (", dbName, tbName);
2252
        } else {
2253 2254 2255 2256 2257 2258 2259 2260 2261 2262
            if (g_args.mysqlFlag) {
                if (0 == rowFlag) {
                    curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len, "(");
                    rowFlag++;
                } else {
                    curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len, ", (");
                }
            } else {
                curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len, "(");
            }
2263 2264
        }

2265 2266
        for (int col = 0; col < numFields; col++) {
            if (col != 0) curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len, ", ");
2267

2268 2269 2270 2271
            if (row[col] == NULL) {
                curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len, "NULL");
                continue;
            }
2272

2273 2274 2275 2276 2277 2278 2279 2280 2281 2282 2283 2284 2285 2286 2287 2288 2289 2290 2291 2292 2293 2294 2295 2296 2297 2298 2299 2300 2301 2302 2303 2304 2305 2306 2307 2308 2309 2310 2311 2312 2313 2314 2315 2316 2317 2318 2319 2320 2321 2322 2323 2324 2325 2326 2327 2328 2329 2330
            switch (fields[col].type) {
                case TSDB_DATA_TYPE_BOOL:
                    curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len, "%d",
                            ((((int32_t)(*((char *)row[col]))) == 1) ? 1 : 0));
                    break;
                case TSDB_DATA_TYPE_TINYINT:
                    curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len, "%d", *((int8_t *)row[col]));
                    break;
                case TSDB_DATA_TYPE_SMALLINT:
                    curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len, "%d", *((int16_t *)row[col]));
                    break;
                case TSDB_DATA_TYPE_INT:
                    curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len, "%d", *((int32_t *)row[col]));
                    break;
                case TSDB_DATA_TYPE_BIGINT:
                    curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len, "%" PRId64 "",
                            *((int64_t *)row[col]));
                    break;
                case TSDB_DATA_TYPE_FLOAT:
                    curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len, "%f", GET_FLOAT_VAL(row[col]));
                    break;
                case TSDB_DATA_TYPE_DOUBLE:
                    curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len, "%f", GET_DOUBLE_VAL(row[col]));
                    break;
                case TSDB_DATA_TYPE_BINARY:
                    {
                        char tbuf[COMMAND_SIZE] = {0};
                        //*(pstr++) = '\'';
                        converStringToReadable((char *)row[col], length[col], tbuf, COMMAND_SIZE);
                        //pstr = stpcpy(pstr, tbuf);
                        //*(pstr++) = '\'';
                        curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len, "\'%s\'", tbuf);
                        break;
                    }
                case TSDB_DATA_TYPE_NCHAR:
                    {
                        char tbuf[COMMAND_SIZE] = {0};
                        convertNCharToReadable((char *)row[col], length[col], tbuf, COMMAND_SIZE);
                        curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len, "\'%s\'", tbuf);
                        break;
                    }
                case TSDB_DATA_TYPE_TIMESTAMP:
                    if (!g_args.mysqlFlag) {
                        curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len, "%" PRId64 "",
                                *(int64_t *)row[col]);
                    } else {
                        char buf[64] = "\0";
                        int64_t ts = *((int64_t *)row[col]);
                        time_t tt = (time_t)(ts / 1000);
                        struct tm *ptm = localtime(&tt);
                        strftime(buf, 64, "%y-%m-%d %H:%M:%S", ptm);
                        curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len, "\'%s.%03d\'",
                                buf, (int)(ts % 1000));
                    }
                    break;
                default:
                    break;
            }
2331
        }
2332 2333 2334 2335 2336 2337 2338 2339 2340 2341 2342 2343 2344 2345 2346 2347 2348 2349 2350

        curr_sqlstr_len += sprintf(pstr + curr_sqlstr_len, ")");

        totalRows++;
        count++;
        fprintf(fp, "%s", tmpBuffer);

        if (totalRows >= lastRowsPrint) {
            printf(" %"PRId64 " rows already be dumpout from %s.%s\n",
                    totalRows, dbName, tbName);
            lastRowsPrint += 5000000;
        }

        total_sqlstr_len += curr_sqlstr_len;

        if ((count >= g_args.data_batch)
                || (sql_buf_len - total_sqlstr_len < TSDB_MAX_BYTES_PER_ROW)) {
            fprintf(fp, ";\n");
            count = 0;
2351
        }
2352 2353
    }

2354
    debugPrint("total_sqlstr_len: %d\n", total_sqlstr_len);
2355

2356 2357 2358
    fprintf(fp, "\n");
    atomic_add_fetch_64(&g_totalDumpOutRows, totalRows);
    free(tmpBuffer);
2359

2360 2361
    return 0;
}
2362

2363
static int taosDumpTableData(FILE *fp, char *tbName,
2364
        TAOS* taosCon, char* dbName, int precision,
2365 2366
        char *jsonAvroSchema) {
    int64_t    totalRows     = 0;
2367

2368
    char sqlstr[1024] = {0};
2369 2370 2371 2372 2373 2374 2375 2376 2377 2378 2379 2380 2381 2382 2383 2384 2385 2386 2387 2388 2389 2390 2391 2392

    int64_t start_time, end_time;
    if (strlen(g_args.humanStartTime)) {
        if (TSDB_CODE_SUCCESS != taosParseTime(
                g_args.humanStartTime, &start_time, strlen(g_args.humanStartTime),
                precision, 0)) {
            errorPrint("Input %s, time format error!\n", g_args.humanStartTime);
            return -1;
        }
    } else {
        start_time = g_args.start_time;
    }

    if (strlen(g_args.humanEndTime)) {
        if (TSDB_CODE_SUCCESS != taosParseTime(
                g_args.humanEndTime, &end_time, strlen(g_args.humanEndTime),
                precision, 0)) {
            errorPrint("Input %s, time format error!\n", g_args.humanEndTime);
            return -1;
        }
    } else {
        end_time = g_args.end_time;
    }

2393 2394
    sprintf(sqlstr,
            "select * from %s.%s where _c0 >= %" PRId64 " and _c0 <= %" PRId64 " order by _c0 asc;",
2395
            dbName, tbName, start_time, end_time);
2396

2397 2398 2399 2400 2401 2402 2403 2404
    TAOS_RES* res = taos_query(taosCon, sqlstr);
    int32_t code = taos_errno(res);
    if (code != 0) {
        errorPrint("failed to run command %s, reason: %s\n",
                sqlstr, taos_errstr(res));
        taos_free_result(res);
        return -1;
    }
2405

2406 2407 2408 2409 2410 2411
    if (g_args.avro) {
        writeSchemaToAvro(jsonAvroSchema);
        totalRows = writeResultToAvro(res);
    } else {
        totalRows = writeResultToSql(res, fp, dbName, tbName);
    }
2412

2413 2414
    taos_free_result(res);
    return totalRows;
2415 2416
}

2417
static int taosCheckParam(struct arguments *arguments) {
2418 2419 2420 2421
    if (g_args.all_databases && g_args.databases) {
        fprintf(stderr, "conflict option --all-databases and --databases\n");
        return -1;
    }
2422

2423 2424 2425 2426
    if (g_args.start_time > g_args.end_time) {
        fprintf(stderr, "start time is larger than end time\n");
        return -1;
    }
2427

2428 2429
    if (g_args.arg_list_len == 0) {
        if ((!g_args.all_databases) && (!g_args.isDumpIn)) {
2430
            errorPrint("%s", "taosdump requires parameters for database and operation\n");
2431 2432 2433 2434 2435 2436 2437 2438 2439 2440 2441 2442
            return -1;
        }
    }
    /*
       if (g_args.isDumpIn && (strcmp(g_args.outpath, DEFAULT_DUMP_FILE) != 0)) {
       fprintf(stderr, "duplicate parameter input and output file path\n");
       return -1;
       }
       */
    if (!g_args.isDumpIn && g_args.encode != NULL) {
        fprintf(stderr, "invalid option in dump out\n");
        return -1;
2443 2444
    }

2445 2446 2447 2448
    if (g_args.table_batch <= 0) {
        fprintf(stderr, "invalid option in dump out\n");
        return -1;
    }
2449

2450
    return 0;
2451 2452
}

2453 2454
/*
static bool isEmptyCommand(char *cmd) {
2455 2456 2457 2458 2459 2460 2461 2462 2463 2464
  char *pchar = cmd;

  while (*pchar != '\0') {
    if (*pchar != ' ') return false;
    pchar++;
  }

  return true;
}

2465 2466
static void taosReplaceCtrlChar(char *str) {
  bool ctrlOn = false;
2467 2468 2469 2470 2471 2472 2473 2474 2475 2476 2477 2478 2479 2480 2481 2482 2483 2484 2485 2486 2487 2488 2489 2490 2491 2492 2493 2494 2495 2496 2497 2498 2499 2500 2501 2502 2503 2504 2505 2506 2507
  char *pstr = NULL;

  for (pstr = str; *str != '\0'; ++str) {
    if (ctrlOn) {
      switch (*str) {
        case 'n':
          *pstr = '\n';
          pstr++;
          break;
        case 'r':
          *pstr = '\r';
          pstr++;
          break;
        case 't':
          *pstr = '\t';
          pstr++;
          break;
        case '\\':
          *pstr = '\\';
          pstr++;
          break;
        case '\'':
          *pstr = '\'';
          pstr++;
          break;
        default:
          break;
      }
      ctrlOn = false;
    } else {
      if (*str == '\\') {
        ctrlOn = true;
      } else {
        *pstr = *str;
        pstr++;
      }
    }
  }

  *pstr = '\0';
}
2508
*/
2509 2510 2511 2512 2513 2514 2515 2516 2517 2518 2519 2520 2521 2522 2523 2524 2525 2526 2527 2528 2529 2530 2531

char *ascii_literal_list[] = {
    "\\x00", "\\x01", "\\x02", "\\x03", "\\x04", "\\x05", "\\x06", "\\x07", "\\x08", "\\t",   "\\n",   "\\x0b", "\\x0c",
    "\\r",   "\\x0e", "\\x0f", "\\x10", "\\x11", "\\x12", "\\x13", "\\x14", "\\x15", "\\x16", "\\x17", "\\x18", "\\x19",
    "\\x1a", "\\x1b", "\\x1c", "\\x1d", "\\x1e", "\\x1f", " ",     "!",     "\\\"",  "#",     "$",     "%",     "&",
    "\\'",   "(",     ")",     "*",     "+",     ",",     "-",     ".",     "/",     "0",     "1",     "2",     "3",
    "4",     "5",     "6",     "7",     "8",     "9",     ":",     ";",     "<",     "=",     ">",     "?",     "@",
    "A",     "B",     "C",     "D",     "E",     "F",     "G",     "H",     "I",     "J",     "K",     "L",     "M",
    "N",     "O",     "P",     "Q",     "R",     "S",     "T",     "U",     "V",     "W",     "X",     "Y",     "Z",
    "[",     "\\\\",  "]",     "^",     "_",     "`",     "a",     "b",     "c",     "d",     "e",     "f",     "g",
    "h",     "i",     "j",     "k",     "l",     "m",     "n",     "o",     "p",     "q",     "r",     "s",     "t",
    "u",     "v",     "w",     "x",     "y",     "z",     "{",     "|",     "}",     "~",     "\\x7f", "\\x80", "\\x81",
    "\\x82", "\\x83", "\\x84", "\\x85", "\\x86", "\\x87", "\\x88", "\\x89", "\\x8a", "\\x8b", "\\x8c", "\\x8d", "\\x8e",
    "\\x8f", "\\x90", "\\x91", "\\x92", "\\x93", "\\x94", "\\x95", "\\x96", "\\x97", "\\x98", "\\x99", "\\x9a", "\\x9b",
    "\\x9c", "\\x9d", "\\x9e", "\\x9f", "\\xa0", "\\xa1", "\\xa2", "\\xa3", "\\xa4", "\\xa5", "\\xa6", "\\xa7", "\\xa8",
    "\\xa9", "\\xaa", "\\xab", "\\xac", "\\xad", "\\xae", "\\xaf", "\\xb0", "\\xb1", "\\xb2", "\\xb3", "\\xb4", "\\xb5",
    "\\xb6", "\\xb7", "\\xb8", "\\xb9", "\\xba", "\\xbb", "\\xbc", "\\xbd", "\\xbe", "\\xbf", "\\xc0", "\\xc1", "\\xc2",
    "\\xc3", "\\xc4", "\\xc5", "\\xc6", "\\xc7", "\\xc8", "\\xc9", "\\xca", "\\xcb", "\\xcc", "\\xcd", "\\xce", "\\xcf",
    "\\xd0", "\\xd1", "\\xd2", "\\xd3", "\\xd4", "\\xd5", "\\xd6", "\\xd7", "\\xd8", "\\xd9", "\\xda", "\\xdb", "\\xdc",
    "\\xdd", "\\xde", "\\xdf", "\\xe0", "\\xe1", "\\xe2", "\\xe3", "\\xe4", "\\xe5", "\\xe6", "\\xe7", "\\xe8", "\\xe9",
    "\\xea", "\\xeb", "\\xec", "\\xed", "\\xee", "\\xef", "\\xf0", "\\xf1", "\\xf2", "\\xf3", "\\xf4", "\\xf5", "\\xf6",
    "\\xf7", "\\xf8", "\\xf9", "\\xfa", "\\xfb", "\\xfc", "\\xfd", "\\xfe", "\\xff"};

2532
static int converStringToReadable(char *str, int size, char *buf, int bufsize) {
2533 2534 2535 2536 2537 2538 2539 2540 2541 2542
    char *pstr = str;
    char *pbuf = buf;
    while (size > 0) {
        if (*pstr == '\0') break;
        pbuf = stpcpy(pbuf, ascii_literal_list[((uint8_t)(*pstr))]);
        pstr++;
        size--;
    }
    *pbuf = '\0';
    return 0;
2543 2544
}

2545
static int convertNCharToReadable(char *str, int size, char *buf, int bufsize) {
2546 2547 2548 2549 2550 2551 2552 2553 2554 2555 2556 2557 2558 2559 2560 2561 2562 2563 2564
    char *pstr = str;
    char *pbuf = buf;
    // TODO
    wchar_t wc;
    while (size > 0) {
        if (*pstr == '\0') break;
        int byte_width = mbtowc(&wc, pstr, MB_CUR_MAX);
        if (byte_width < 0) {
            errorPrint("%s() LN%d, mbtowc() return fail.\n", __func__, __LINE__);
            exit(-1);
        }

        if ((int)wc < 256) {
            pbuf = stpcpy(pbuf, ascii_literal_list[(int)wc]);
        } else {
            memcpy(pbuf, pstr, byte_width);
            pbuf += byte_width;
        }
        pstr += byte_width;
2565 2566
    }

2567
    *pbuf = '\0';
2568

2569
    return 0;
2570 2571
}

2572
static void taosDumpCharset(FILE *fp) {
2573
    char charsetline[256];
2574

2575 2576 2577
    (void)fseek(fp, 0, SEEK_SET);
    sprintf(charsetline, "#!%s\n", tsCharset);
    (void)fwrite(charsetline, strlen(charsetline), 1, fp);
2578 2579
}

2580
static void taosLoadFileCharset(FILE *fp, char *fcharset) {
2581 2582
    char * line = NULL;
    size_t line_size = 0;
2583

2584 2585 2586 2587 2588
    (void)fseek(fp, 0, SEEK_SET);
    ssize_t size = getline(&line, &line_size, fp);
    if (size <= 2) {
        goto _exit_no_charset;
    }
2589

2590 2591 2592 2593 2594 2595 2596 2597
    if (strncmp(line, "#!", 2) != 0) {
        goto _exit_no_charset;
    }
    if (line[size - 1] == '\n') {
        line[size - 1] = '\0';
        size--;
    }
    strcpy(fcharset, line + 2);
2598

2599 2600
    tfree(line);
    return;
2601 2602

_exit_no_charset:
2603 2604 2605 2606
    (void)fseek(fp, 0, SEEK_SET);
    *fcharset = '\0';
    tfree(line);
    return;
2607 2608 2609 2610
}

// ========  dumpIn support multi threads functions ================================//

2611 2612 2613 2614 2615 2616 2617
static char    **g_tsDumpInSqlFiles   = NULL;
static int32_t   g_tsSqlFileNum = 0;
static char      g_tsDbSqlFile[MAX_FILE_NAME_LEN] = {0};
static char      g_tsCharset[64] = {0};

static int taosGetFilesNum(const char *directoryName,
        const char *prefix, const char *prefix2)
2618
{
2619
    char cmd[1024] = { 0 };
2620

2621 2622 2623 2624 2625
    if (prefix2)
        sprintf(cmd, "ls %s/*.%s %s/*.%s | wc -l ",
                directoryName, prefix, directoryName, prefix2);
    else
        sprintf(cmd, "ls %s/*.%s | wc -l ", directoryName, prefix);
2626

2627 2628 2629 2630 2631
    FILE *fp = popen(cmd, "r");
    if (fp == NULL) {
        errorPrint("failed to execute:%s, error:%s\n", cmd, strerror(errno));
        exit(-1);
    }
2632

2633 2634 2635 2636 2637 2638 2639
    int fileNum = 0;
    if (fscanf(fp, "%d", &fileNum) != 1) {
        errorPrint("failed to execute:%s, parse result error\n", cmd);
        exit(-1);
    }

    if (fileNum <= 0) {
2640
        errorPrint("directory:%s is empty\n", directoryName);
2641 2642
        exit(-1);
    }
2643

2644 2645
    pclose(fp);
    return fileNum;
2646 2647
}

2648 2649 2650
static void taosParseDirectory(const char *directoryName,
        const char *prefix, const char *prefix2,
        char **fileArray, int totalFiles)
2651
{
2652
    char cmd[1024] = { 0 };
2653

2654 2655 2656 2657 2658 2659
    if (prefix2) {
        sprintf(cmd, "ls %s/*.%s %s/*.%s | sort",
                directoryName, prefix, directoryName, prefix2);
    } else {
        sprintf(cmd, "ls %s/*.%s | sort", directoryName, prefix);
    }
2660

2661 2662 2663 2664 2665 2666 2667 2668 2669 2670 2671 2672 2673 2674
    FILE *fp = popen(cmd, "r");
    if (fp == NULL) {
        errorPrint("failed to execute:%s, error:%s\n", cmd, strerror(errno));
        exit(-1);
    }

    int fileNum = 0;
    while (fscanf(fp, "%128s", fileArray[fileNum++])) {
        if (strcmp(fileArray[fileNum-1], g_tsDbSqlFile) == 0) {
            fileNum--;
        }
        if (fileNum >= totalFiles) {
            break;
        }
2675
    }
2676 2677 2678 2679 2680

    if (fileNum != totalFiles) {
        errorPrint("directory:%s changed while read\n", directoryName);
        pclose(fp);
        exit(-1);
2681 2682 2683 2684 2685
    }

    pclose(fp);
}

2686
static void taosCheckDatabasesSQLFile(const char *directoryName)
2687
{
2688 2689
    char cmd[1024] = { 0 };
    sprintf(cmd, "ls %s/dbs.sql", directoryName);
2690

2691 2692 2693 2694 2695
    FILE *fp = popen(cmd, "r");
    if (fp == NULL) {
        errorPrint("failed to execute:%s, error:%s\n", cmd, strerror(errno));
        exit(-1);
    }
2696

2697 2698 2699
    while (fscanf(fp, "%128s", g_tsDbSqlFile)) {
        break;
    }
2700

2701
    pclose(fp);
2702 2703
}

2704
static void taosMallocDumpFiles()
2705
{
2706 2707 2708 2709
    g_tsDumpInSqlFiles = (char**)calloc(g_tsSqlFileNum, sizeof(char*));
    for (int i = 0; i < g_tsSqlFileNum; i++) {
        g_tsDumpInSqlFiles[i] = calloc(1, MAX_FILE_NAME_LEN);
    }
2710 2711
}

2712
static void taosFreeDumpFiles()
2713
{
2714 2715 2716 2717
    for (int i = 0; i < g_tsSqlFileNum; i++) {
        tfree(g_tsDumpInSqlFiles[i]);
    }
    tfree(g_tsDumpInSqlFiles);
2718 2719 2720 2721
}

static void taosGetDirectoryFileList(char *inputDir)
{
2722 2723 2724 2725
    struct stat fileStat;
    if (stat(inputDir, &fileStat) < 0) {
        errorPrint("%s not exist\n", inputDir);
        exit(-1);
2726
    }
2727 2728 2729 2730 2731 2732 2733 2734 2735 2736 2737 2738 2739 2740 2741 2742 2743 2744 2745 2746 2747 2748 2749 2750 2751 2752 2753

    if (fileStat.st_mode & S_IFDIR) {
        taosCheckDatabasesSQLFile(inputDir);
        if (g_args.avro)
            g_tsSqlFileNum = taosGetFilesNum(inputDir, "sql", "avro");
        else
            g_tsSqlFileNum += taosGetFilesNum(inputDir, "sql", NULL);

        int tsSqlFileNumOfTbls = g_tsSqlFileNum;
        if (g_tsDbSqlFile[0] != 0) {
            tsSqlFileNumOfTbls--;
        }
        taosMallocDumpFiles();
        if (0 != tsSqlFileNumOfTbls) {
            if (g_args.avro) {
                taosParseDirectory(inputDir, "sql", "avro",
                        g_tsDumpInSqlFiles, tsSqlFileNumOfTbls);
            } else {
                taosParseDirectory(inputDir, "sql", NULL,
                        g_tsDumpInSqlFiles, tsSqlFileNumOfTbls);
            }
        }
        fprintf(stdout, "\nstart to dispose %d files in %s\n",
                g_tsSqlFileNum, inputDir);
    } else {
        errorPrint("%s is not a directory\n", inputDir);
        exit(-1);
2754
    }
2755 2756
}

2757
static FILE* taosOpenDumpInFile(char *fptr) {
2758
    wordexp_t full_path;
2759

2760 2761 2762 2763
    if (wordexp(fptr, &full_path, 0) != 0) {
        errorPrint("illegal file name: %s\n", fptr);
        return NULL;
    }
2764

2765
    char *fname = full_path.we_wordv[0];
2766

2767 2768 2769 2770 2771 2772 2773 2774
    FILE *f = NULL;
    if ((fname) && (strlen(fname) > 0)) {
        f = fopen(fname, "r");
        if (f == NULL) {
            errorPrint("%s() LN%d, failed to open file %s\n",
                    __func__, __LINE__, fname);
        }
    }
2775

2776 2777
    wordfree(&full_path);
    return f;
2778 2779
}

2780 2781
static int taosDumpInOneFile(TAOS* taos, FILE* fp, char* fcharset,
        char* encode, char* fileName) {
2782 2783 2784 2785 2786 2787 2788 2789 2790 2791 2792 2793
    int       read_len = 0;
    char *    cmd      = NULL;
    size_t    cmd_len  = 0;
    char *    line     = NULL;
    size_t    line_len = 0;

    cmd  = (char *)malloc(TSDB_MAX_ALLOWED_SQL_LEN);
    if (cmd == NULL) {
        errorPrint("%s() LN%d, failed to allocate memory\n",
                __func__, __LINE__);
        return -1;
    }
2794

2795 2796 2797 2798 2799 2800
    int lastRowsPrint = 5000000;
    int lineNo = 0;
    while ((read_len = getline(&line, &line_len, fp)) != -1) {
        ++lineNo;
        if (read_len >= TSDB_MAX_ALLOWED_SQL_LEN) continue;
        line[--read_len] = '\0';
2801

2802 2803 2804 2805
        //if (read_len == 0 || isCommentLine(line)) {  // line starts with #
        if (read_len == 0 ) {
            continue;
        }
2806

2807 2808 2809 2810 2811 2812
        if (line[read_len - 1] == '\\') {
            line[read_len - 1] = ' ';
            memcpy(cmd + cmd_len, line, read_len);
            cmd_len += read_len;
            continue;
        }
2813

2814 2815 2816
        memcpy(cmd + cmd_len, line, read_len);
        cmd[read_len + cmd_len]= '\0';
        if (queryDbImpl(taos, cmd)) {
2817
            errorPrint("%s() LN%d, error sql: lineno:%d, file:%s\n",
2818
                    __func__, __LINE__, lineNo, fileName);
2819
            fprintf(g_fpOfResult, "error sql: lineno:%d, file:%s\n", lineNo, fileName);
2820
        }
2821

2822 2823
        memset(cmd, 0, TSDB_MAX_ALLOWED_SQL_LEN);
        cmd_len = 0;
2824

2825 2826 2827 2828
        if (lineNo >= lastRowsPrint) {
            printf(" %d lines already be executed from file %s\n", lineNo, fileName);
            lastRowsPrint += 5000000;
        }
H
Hui Li 已提交
2829
    }
2830

2831 2832 2833 2834
    tfree(cmd);
    tfree(line);
    fclose(fp);
    return 0;
2835 2836
}

2837
static void* taosDumpInWorkThreadFp(void *arg)
2838
{
2839
    SThreadParaObj *pThread = (SThreadParaObj*)arg;
2840 2841
    setThreadName("dumpInWorkThrd");

2842 2843 2844 2845 2846 2847 2848 2849 2850 2851 2852
    for (int32_t f = 0; f < g_tsSqlFileNum; ++f) {
        if (f % pThread->totalThreads == pThread->threadIndex) {
            char *SQLFileName = g_tsDumpInSqlFiles[f];
            FILE* fp = taosOpenDumpInFile(SQLFileName);
            if (NULL == fp) {
                continue;
            }
            fprintf(stderr, ", Success Open input file: %s\n",
                    SQLFileName);
            taosDumpInOneFile(pThread->taosCon, fp, g_tsCharset, g_args.encode, SQLFileName);
        }
2853 2854
    }

2855
    return NULL;
2856 2857
}

2858
static void taosStartDumpInWorkThreads()
2859
{
2860 2861 2862
    pthread_attr_t  thattr;
    SThreadParaObj *pThread;
    int32_t         totalThreads = g_args.thread_num;
2863

2864 2865 2866
    if (totalThreads > g_tsSqlFileNum) {
        totalThreads = g_tsSqlFileNum;
    }
2867

2868 2869
    SThreadParaObj *threadObj = (SThreadParaObj *)calloc(
            totalThreads, sizeof(SThreadParaObj));
2870

2871 2872 2873
    if (NULL == threadObj) {
        errorPrint("%s() LN%d, memory allocation failed\n", __func__, __LINE__);
    }
2874

2875 2876 2877 2878 2879 2880 2881 2882
    for (int32_t t = 0; t < totalThreads; ++t) {
        pThread = threadObj + t;
        pThread->threadIndex = t;
        pThread->totalThreads = totalThreads;
        pThread->taosCon = taos_connect(g_args.host, g_args.user, g_args.password,
            NULL, g_args.port);
        if (pThread->taosCon == NULL) {
            errorPrint("Failed to connect to TDengine server %s\n", g_args.host);
2883
            free(threadObj);
2884 2885 2886 2887 2888 2889 2890 2891 2892 2893 2894
            return;
        }
        pthread_attr_init(&thattr);
        pthread_attr_setdetachstate(&thattr, PTHREAD_CREATE_JOINABLE);

        if (pthread_create(&(pThread->threadID), &thattr,
                    taosDumpInWorkThreadFp, (void*)pThread) != 0) {
            errorPrint("%s() LN%d, thread:%d failed to start\n",
                    __func__, __LINE__, pThread->threadIndex);
            exit(0);
        }
2895 2896
    }

2897 2898 2899
    for (int t = 0; t < totalThreads; ++t) {
        pthread_join(threadObj[t].threadID, NULL);
    }
2900

2901 2902 2903 2904
    for (int t = 0; t < totalThreads; ++t) {
        taos_close(threadObj[t].taosCon);
    }
    free(threadObj);
2905 2906
}

2907 2908
static int taosDumpIn() {
    assert(g_args.isDumpIn);
2909

2910 2911
    TAOS     *taos    = NULL;
    FILE     *fp      = NULL;
2912

2913 2914 2915 2916 2917 2918 2919 2920
    taos = taos_connect(
            g_args.host, g_args.user, g_args.password,
            NULL, g_args.port);
    if (taos == NULL) {
        errorPrint("%s() LN%d, failed to connect to TDengine server\n",
                __func__, __LINE__);
        return -1;
    }
2921

2922
    taosGetDirectoryFileList(g_args.inpath);
2923

2924 2925 2926
    int32_t  tsSqlFileNumOfTbls = g_tsSqlFileNum;
    if (g_tsDbSqlFile[0] != 0) {
        tsSqlFileNumOfTbls--;
2927

2928 2929 2930 2931 2932 2933 2934
        fp = taosOpenDumpInFile(g_tsDbSqlFile);
        if (NULL == fp) {
            errorPrint("%s() LN%d, failed to open input file %s\n",
                    __func__, __LINE__, g_tsDbSqlFile);
            return -1;
        }
        fprintf(stderr, "Success Open input file: %s\n", g_tsDbSqlFile);
2935

2936
        taosLoadFileCharset(fp, g_tsCharset);
2937

2938 2939 2940
        taosDumpInOneFile(taos, fp, g_tsCharset, g_args.encode,
                g_tsDbSqlFile);
    }
2941

2942
    taos_close(taos);
2943

2944 2945 2946 2947 2948 2949
    if (0 != tsSqlFileNumOfTbls) {
        taosStartDumpInWorkThreads();
    }

    taosFreeDumpFiles();
    return 0;
2950 2951
}