快速装载测试方法

更新时间:
复制 MD 格式

本文将基于Airline数据集以及TPC-H数据集,对PolarDB MySQL8.0集群的全列存表快速装载功能进行性能测试。您可以按照本文所述的方法自行进行测试对比,从而快速了解该集群的性能。

注意事项

  • 本文的测试基于Airline公开数据集和TPC-H基准数据集,测试结果仅供参考。

    说明

    本文的TPC-H的实现基于TPC-H的基准测试,并不能与已发布的TPC-H基准测试结果相比较,本文中的测试并不完全符合TPC-H的所有要求。

  • 快速装载功能需要PolarDB MySQL8.0.2.2.34及以上版本。

    说明

    如何查看集群版本号,请参见查询版本号

测试环境

  • ECS实例和PolarDB MySQL集群位于同一地域,同一专有网络VPC内。

  • PolarDB MySQL集群配置如下:

    • 集群个数:1个。

    • 数据库引擎:MySQL 8.0.2

    • 产品版本:企业版。

    • 系列:集群版,子系列为独享规格。

    • 节点规格:polar.mysql.x8.4xlarge(32256 GB)。

    • 节点数量:2个(1个读写节点、1个只读节点)。

    说明

    快速装载测试使用的连接地址为主地址

  • ECS实例配置如下:

    • 实例个数:1台。

    • 实例规格:ecs.c5.4xlarge(16vCPU 32 GiB)。

    • 镜像:CentOS 7.0 64位。

    • 系统盘:ESSD云盘1000 GB。

环境准备

  • OSS准备:确保有一个具有ListObjects以及GetObject权限的OSS Bucket,准备对应的AccessKey ID以及AccessKey Secret。

  • 客户端工具:确认ECS已经安装mysql-clientossutil

Airline数据

Airline数据集源自美国交通部(U.S. Department of Transportation, Bureau of Transportation Statistics)的官方统计,记录了从1987年至今美国国内所有商业航班的详细起降信息。该数据集规模庞大,能够很好地模拟企业数据分析场景。

  • 原始格式:CSV格式。

  • 数据量:1.7亿行,110列(宽表场景)。

  • 数据大小:压缩前(CSV)约77 GB,压缩后(Gzip)约6 GB。

  • 数据类型:包含整型(Year, Month)、浮点型(DepDelay, AirTime)、长字符串(Carrier, Origin)、短字符串(TailNum)及日期/时间戳等混合数据类型。

准备数据集

使用以下脚本下载官方数据集并解压成CSV文件,最后将数据存放于OSS上。假设OSS目录为oss://<your_bucket_name>/airline/

  1. 下载数据集。

    #下载数据集
    for s in `seq 1987 2018`
    do
    for m in `seq 1 12`
    do
    wget https://transtats.bts.gov/PREZIP/On_Time_Reporting_Carrier_On_Time_Performance_1987_present_${s}_${m}.zip
    done
    done
  2. 解压并重命名文件。

    # 解压
    unzip *.zip -d /path/to/csv/
    
    # 重命名为简单格式(可选操作,为了方便)
    cd /path/to/csv/
    i=1
    for f in $(ls *.csv | sort); do
        mv "$f" "airline${i}.csv"
        i=$((i+1))
    done
  3. 将数据集上传至OSS。

    ossutil cp ./ oss://<your_bucket_name>/airline/ -r

创建表

说明

创建表时指定ENGINE=xengine TABLE_FORMAT=column代表创建的是全列存表。如需测试InnoDB引擎的导入性能作为对比,请将建表语句中的引擎改为engine=innodb

CREATE DATABASE tpch1000;
USE tpch1000;
-- 关闭hidden key
SET xengine_enable_ctable_hidden_key=off;
DROP TABLE IF EXISTS ontime;

CREATE TABLE `ontime` (
  `Year` year(4) DEFAULT NULL,
  `Quarter` tinyint(4) DEFAULT NULL,
  `Month` tinyint(4) DEFAULT NULL,
  `DayofMonth` tinyint(4) DEFAULT NULL,
  `DayOfWeek` tinyint(4) DEFAULT NULL,
  `FlightDate` date DEFAULT NULL,
  `UniqueCarrier` char(7) DEFAULT NULL,
  `AirlineID` int(11) DEFAULT NULL,
  `Carrier` char(2) DEFAULT NULL,
  `TailNum` varchar(50) DEFAULT NULL,
  `FlightNum` varchar(10) DEFAULT NULL,
  `Origin` char(5) DEFAULT NULL,
  `OriginCityName` varchar(100) DEFAULT NULL,
  `OriginState` char(2) DEFAULT NULL,
  `OriginStateFips` varchar(10) DEFAULT NULL,
  `OriginStateName` varchar(100) DEFAULT NULL,
  `OriginWac` int(11) DEFAULT NULL,
  `Dest` char(5) DEFAULT NULL,
  `DestCityName` varchar(100) DEFAULT NULL,
  `DestState` char(2) DEFAULT NULL,
  `DestStateFips` varchar(10) DEFAULT NULL,
  `DestStateName` varchar(100) DEFAULT NULL,
  `DestWac` int(11) DEFAULT NULL,
  `CRSDepTime` int(11) DEFAULT NULL,
  `DepTime` int(11) DEFAULT NULL,
  `DepDelay` int(11) DEFAULT NULL,
  `DepDelayMinutes` int(11) DEFAULT NULL,
  `DepDel15` int(11) DEFAULT NULL,
  `DepartureDelayGroups` int(11) DEFAULT NULL,
  `DepTimeBlk` varchar(20) DEFAULT NULL,
  `TaxiOut` int(11) DEFAULT NULL,
  `WheelsOff` int(11) DEFAULT NULL,
  `WheelsOn` int(11) DEFAULT NULL,
  `TaxiIn` int(11) DEFAULT NULL,
  `CRSArrTime` int(11) DEFAULT NULL,
  `ArrTime` int(11) DEFAULT NULL,
  `ArrDelay` int(11) DEFAULT NULL,
  `ArrDelayMinutes` int(11) DEFAULT NULL,
  `ArrDel15` int(11) DEFAULT NULL,
  `ArrivalDelayGroups` int(11) DEFAULT NULL,
  `ArrTimeBlk` varchar(20) DEFAULT NULL,
  `Cancelled` tinyint(4) DEFAULT NULL,
  `CancellationCode` char(1) DEFAULT NULL,
  `Diverted` tinyint(4) DEFAULT NULL,
  `CRSElapsedTime` INT(11) DEFAULT NULL,
  `ActualElapsedTime` INT(11) DEFAULT NULL,
  `AirTime` INT(11) DEFAULT NULL,
  `Flights` INT(11) DEFAULT NULL,
  `Distance` INT(11) DEFAULT NULL,
  `DistanceGroup` TINYINT(4) DEFAULT NULL,
  `CarrierDelay` INT(11) DEFAULT NULL,
  `WeatherDelay` INT(11) DEFAULT NULL,
  `NASDelay` INT(11) DEFAULT NULL,
  `SecurityDelay` INT(11) DEFAULT NULL,
  `LateAircraftDelay` INT(11) DEFAULT NULL,
  `FirstDepTime` varchar(10) DEFAULT NULL,
  `TotalAddGTime` varchar(10) DEFAULT NULL,
  `LongestAddGTime` varchar(10) DEFAULT NULL,
  `DivAirportLandings` varchar(10) DEFAULT NULL,
  `DivReachedDest` varchar(10) DEFAULT NULL,
  `DivActualElapsedTime` varchar(10) DEFAULT NULL,
  `DivArrDelay` varchar(10) DEFAULT NULL,
  `DivDistance` varchar(10) DEFAULT NULL,
  `Div1Airport` varchar(10) DEFAULT NULL,
  `Div1WheelsOn` varchar(10) DEFAULT NULL,
  `Div1TotalGTime` varchar(10) DEFAULT NULL,
  `Div1LongestGTime` varchar(10) DEFAULT NULL,
  `Div1WheelsOff` varchar(10) DEFAULT NULL,
  `Div1TailNum` varchar(10) DEFAULT NULL,
  `Div2Airport` varchar(10) DEFAULT NULL,
  `Div2WheelsOn` varchar(10) DEFAULT NULL,
  `Div2TotalGTime` varchar(10) DEFAULT NULL,
  `Div2LongestGTime` varchar(10) DEFAULT NULL,
  `Div2WheelsOff` varchar(10) DEFAULT NULL,
  `Div2TailNum` varchar(10) DEFAULT NULL,
  `Div3Airport` varchar(10) DEFAULT NULL,
  `Div3WheelsOn` varchar(10) DEFAULT NULL,
  `Div3TotalGTime` varchar(10) DEFAULT NULL,
  `Div3LongestGTime` varchar(10) DEFAULT NULL,
  `Div3WheelsOff` varchar(10) DEFAULT NULL,
  `Div3TailNum` varchar(10) DEFAULT NULL,
  `Div4Airport` varchar(10) DEFAULT NULL,
  `Div4WheelsOn` varchar(10) DEFAULT NULL,
  `Div4TotalGTime` varchar(10) DEFAULT NULL,
  `Div4LongestGTime` varchar(10) DEFAULT NULL,
  `Div4WheelsOff` varchar(10) DEFAULT NULL,
  `Div4TailNum` varchar(10) DEFAULT NULL,
  `Div5Airport` varchar(10) DEFAULT NULL,
  `Div5WheelsOn` varchar(10) DEFAULT NULL,
  `Div5TotalGTime` varchar(10) DEFAULT NULL,
  `Div5LongestGTime` varchar(10) DEFAULT NULL,
  `Div5WheelsOff` varchar(10) DEFAULT NULL,
  `Div5TailNum` varchar(10) DEFAULT NULL,
  `Dummy` varchar(10) DEFAULT NULL
) ENGINE=xengine TABLE_FORMAT=column;

创建OSS Server

OSS Server 对应前面存放数据的OSS Bucket,这里命名为my_server

DROP SERVER my_server;

CREATE SERVER my_server
FOREIGN DATA WRAPPER oss OPTIONS
(
  EXTRA_SERVER_INFO '{"oss_endpoint": "xxx","oss_bucket": "xxx","oss_access_key_id": "xxxx","oss_access_key_secret": "xxxx"}'
);

-- 验证server创建成功
SELECT * FROM mysql.servers;

导入数据

LOAD DATA语句中OSS关键词表示从OSS中导入,路径支持*通配符一次性导入多个文件。

-- 开启快速装载相关优化参数
-- 开启异步索引构建,导入完成即可查,索引在后台静默构建
SET xengine_dload_async_build_nci = ON;
-- 设置解析与写入的并行度,建议设为节点核数的 1-2 倍
SET xengine_dload_parallel = 64;
LOAD DATA OSS INFILE 'my_server/airline/*.csv' INTO TABLE ontime fields terminated by ',' OPTIONALLY ENCLOSED BY '"'
LINES TERMINATED BY '\n'
IGNORE 1 LINES;
说明

上述参数仅对X-Engine全列存表生效,对其它引擎无效。

InnoDB不支持*通配符,需要通过脚本逐个文件导入。

单击展开查看InnoDB导入参考脚本

#!/usr/bin/env python3
"""
并发导入 airline CSV 到 PolarDB

airline1.csv ~ airline375.csv
OSS 路径: my_server/airlineX.csv
"""

import sys
import time
import subprocess
import logging
import threading
from concurrent.futures import ThreadPoolExecutor, as_completed

# ==================== 配置区 ====================
CONCURRENCY    = 32
MYSQL_HOST     = "127.0.0.1"
MYSQL_PORT     = 3306
MYSQL_USER     = "root"
MYSQL_PASSWORD = "your_password"
MYSQL_CLI      = "mysql"
TIMEOUT        = 7200

FILE_START     = 1
FILE_END       = 375
# ================================================

logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s [%(threadName)s] %(message)s",
    datefmt="%H:%M:%S"
)
log = logging.getLogger()

lock = threading.Lock()
stats = {"ok": 0, "fail": 0, "durations": []}


def load_file(i):
    """导入第 i 个文件到 airline_i.ontime"""
    db_name  = "airline_{i}".format(i=i)
    oss_path = "my_server/airline/airline{i}.csv".format(i=i)
    label    = "airline{i}.csv -> {db}.ontime".format(i=i, db=db_name)

    sql = (
        "LOAD DATA OSS INFILE '{oss}' "
        "INTO TABLE ontime "
        "fields terminated by ',' "
        "OPTIONALLY ENCLOSED BY '\"' "
        "LINES TERMINATED BY '\n' "
        "IGNORE 1 LINES;"
    ).format(oss=oss_path)

    cmd = [
        MYSQL_CLI,
        "-h{host}".format(host=MYSQL_HOST),
        "-P{port}".format(port=MYSQL_PORT),
        "-u{user}".format(user=MYSQL_USER),
        "-p{pwd}".format(pwd=MYSQL_PASSWORD),
        "-D{db}".format(db=db_name),
        "--batch",
        "--quick",
    ]

    t0 = time.time()
    try:
        r = subprocess.run(
            cmd,
            input=sql,
            stdout=subprocess.PIPE,
            stderr=subprocess.PIPE,
            universal_newlines=True,
            timeout=TIMEOUT,
        )
        dur = time.time() - t0
        ok = r.returncode == 0
        if ok:
            log.info("OK  {label}  {dur:.1f}s".format(label=label, dur=dur))
        else:
            log.error("FAIL  {label}  {dur:.1f}s  err={err}".format(
                label=label, dur=dur, err=r.stderr.strip()[:200]
            ))
    except Exception as e:
        dur = time.time() - t0
        ok = False
        log.error("FAIL  {label}  {dur:.1f}s  exception={exc}".format(
            label=label, dur=dur, exc=e
        ))

    with lock:
        if ok:
            stats["ok"] += 1
        else:
            stats["fail"] += 1
        stats["durations"].append(dur)
    return ok


def main():
    total = FILE_END - FILE_START + 1
    log.info("共 {n} 个文件 (airline{s}.csv ~ airline{e}.csv), 并发={c}".format(
        n=total, s=FILE_START, e=FILE_END, c=CONCURRENCY
    ))
    log.info("目标: {h}:{p} / airline_1 ~ airline_{e}".format(
        h=MYSQL_HOST, p=MYSQL_PORT, e=FILE_END
    ))

    t_start = time.time()

    with ThreadPoolExecutor(max_workers=CONCURRENCY, thread_name_prefix="load") as pool:
        futs = {}
        for i in range(FILE_START, FILE_END + 1):
            fut = pool.submit(load_file, i)
            futs[fut] = i
        for fut in as_completed(futs):
            fut.result()

    total_dur = time.time() - t_start

    print("")
    print("=" * 60)
    print("总耗时: {d:.1f}s ({m:.1f}min)".format(d=total_dur, m=total_dur / 60.0))
    print("成功: {ok}  失败: {fail}".format(ok=stats["ok"], fail=stats["fail"]))
    if stats["durations"]:
        print("最快: {mn:.1f}s  最慢: {mx:.1f}s  平均: {avg:.1f}s".format(
            mn=min(stats["durations"]),
            mx=max(stats["durations"]),
            avg=sum(stats["durations"]) / len(stats["durations"]),
        ))
    print("=" * 60)

    if stats["fail"] > 0:
        sys.exit(1)


if __name__ == "__main__":
    main()

TPC-H

TPC-H是由事务处理性能委员会(Transaction Processing Performance Council)制定的一套决策支持基准测试。它模拟了一个全球范围内的零配件零售商业务场景,包括零件采购、销售订单、库存管理和客户信息。TPC-H包含8张表,这种多表结构可以测试数据库处理不同大小、不同关联关系的文件的能力。

准备数据集

  1. 下载安装dbgen,可参考并行查询测试方法中的方法进行安装下载。

  2. 安装之后,使用以下脚本创建1 TB数据集。

    #!/bin/bash
    
    SCALE=1000
    TOTAL_CHUNKS=512
    MAX_PARALLEL=16
    OUTPUT_DIR="tpch1000g"
    DBGEN_BIN="./dbgen"
    
    mkdir -p $OUTPUT_DIR
    
    # 定义一个处理单个分片的函数,并导出
    run_dbgen() {
        local i=$1
        local SCALE=$2
        local TOTAL_CHUNKS=$3
        local OUTPUT_DIR=$4
        local DBGEN_BIN=$5
    
        $DBGEN_BIN -s $SCALE -C $TOTAL_CHUNKS -S $i -f > /dev/null 2>&1
        mv *.tbl.$i "$OUTPUT_DIR/" 2>/dev/null
        echo "Chunk $i 完成"
    }
    export -f run_dbgen
    
    echo "开始通过 xargs 并行生成..."
    
    # 使用 xargs 控制并行
    seq 1 $TOTAL_CHUNKS | xargs -I {} -P $MAX_PARALLEL bash -c "run_dbgen {} $SCALE $TOTAL_CHUNKS $OUTPUT_DIR $DBGEN_BIN"
    
    # 最后处理 nation/region
    mv nation.tbl region.tbl $OUTPUT_DIR/ 2>/dev/null
    
    echo "所有任务已完成!"
  3. 将数据集存放于OSS上,例如 oss://<your_bucket_name>/tpch1000g/

创建表

说明

创建表时指定ENGINE=xengine TABLE_FORMAT=column代表创建的是全列存表。

SET xengine_enable_ctable_hidden_key = OFF;
DROP TABLE IF EXISTS nation;

CREATE TABLE nation  ( N_NATIONKEY  INTEGER NOT NULL,
                            N_NAME       CHAR(25) NOT NULL,
                            N_REGIONKEY  INTEGER NOT NULL,
                            N_COMMENT    VARCHAR(152),
PRIMARY KEY (N_NATIONKEY)
) ENGINE=xengine TABLE_FORMAT=column;

DROP TABLE IF EXISTS region;
CREATE TABLE region  ( R_REGIONKEY  INTEGER NOT NULL,
                            R_NAME       CHAR(25) NOT NULL,
                            R_COMMENT    VARCHAR(152),
  PRIMARY KEY (R_REGIONKEY)
) ENGINE=xengine TABLE_FORMAT=column;

DROP TABLE IF EXISTS part;
CREATE TABLE part  ( P_PARTKEY     BIGINT NOT NULL,
                          P_NAME        VARCHAR(55) NOT NULL,
                          P_MFGR        CHAR(25) NOT NULL,
                          P_BRAND       CHAR(10) NOT NULL,
                          P_TYPE        VARCHAR(25) NOT NULL,
                          P_SIZE        INTEGER NOT NULL,
                          P_CONTAINER   CHAR(10) NOT NULL,
                          P_RETAILPRICE DECIMAL(15,2) NOT NULL,
                          P_COMMENT     VARCHAR(23) NOT NULL ,
  PRIMARY KEY (P_PARTKEY)
) ENGINE=xengine TABLE_FORMAT=column;

DROP TABLE IF EXISTS supplier;
CREATE TABLE supplier  ( S_SUPPKEY     BIGINT NOT NULL,
                             S_NAME        CHAR(25) NOT NULL,
                             S_ADDRESS     VARCHAR(40) NOT NULL,
                             S_NATIONKEY   INTEGER NOT NULL,
                             S_PHONE       CHAR(15) NOT NULL,
                             S_ACCTBAL     DECIMAL(15,2) NOT NULL,
                             S_COMMENT     VARCHAR(101) NOT NULL,
  PRIMARY KEY (S_SUPPKEY)
) ENGINE=xengine TABLE_FORMAT=column;

DROP TABLE IF EXISTS partsupp;
CREATE TABLE partsupp  ( PS_PARTKEY     BIGINT NOT NULL,
                             PS_SUPPKEY     INTEGER NOT NULL,
                             PS_AVAILQTY    INTEGER NOT NULL,
                             PS_SUPPLYCOST  DECIMAL(15,2)  NOT NULL,
                             PS_COMMENT     VARCHAR(199) NOT NULL ,
  PRIMARY KEY (`PS_PARTKEY`,`PS_SUPPKEY`)
) ENGINE=xengine TABLE_FORMAT=column;

DROP TABLE IF EXISTS customer;
CREATE TABLE customer  ( C_CUSTKEY     BIGINT NOT NULL,
                             C_NAME        VARCHAR(25) NOT NULL,
                             C_ADDRESS     VARCHAR(40) NOT NULL,
                             C_NATIONKEY   INTEGER NOT NULL,
                             C_PHONE       CHAR(15) NOT NULL,
                             C_ACCTBAL     DECIMAL(15,2)   NOT NULL,
                             C_MKTSEGMENT  CHAR(10) NOT NULL,
                             C_COMMENT     VARCHAR(117) NOT NULL,
  PRIMARY KEY (`C_CUSTKEY`)
) ENGINE=xengine TABLE_FORMAT=column;

DROP TABLE IF EXISTS orders;
CREATE TABLE orders  ( O_ORDERKEY       BIGINT NOT NULL,
                           O_CUSTKEY        INTEGER NOT NULL,
                           O_ORDERSTATUS    CHAR(1) NOT NULL,
                           O_TOTALPRICE     DECIMAL(15,2) NOT NULL,
                           O_ORDERDATE      DATE NOT NULL,
                           O_ORDERPRIORITY  CHAR(15) NOT NULL,
                           O_CLERK          CHAR(15) NOT NULL,
                           O_SHIPPRIORITY   INTEGER NOT NULL,
                           O_COMMENT        VARCHAR(79) NOT NULL,
  PRIMARY KEY (`O_ORDERKEY`,`O_ORDERDATE`)
) ENGINE=xengine TABLE_FORMAT=column;

DROP TABLE IF EXISTS lineitem;
CREATE TABLE lineitem ( L_ORDERKEY    BIGINT NOT NULL,
                             L_PARTKEY     INTEGER NOT NULL,
                             L_SUPPKEY     INTEGER NOT NULL,
                             L_LINENUMBER  INTEGER NOT NULL,
                             L_QUANTITY    DECIMAL(15,2) NOT NULL,
                             L_EXTENDEDPRICE  DECIMAL(15,2) NOT NULL,
                             L_DISCOUNT    DECIMAL(15,2) NOT NULL,
                             L_TAX         DECIMAL(15,2) NOT NULL,
                             L_RETURNFLAG  CHAR(1) NOT NULL,
                             L_LINESTATUS  CHAR(1) NOT NULL,
                             L_SHIPDATE    DATE NOT NULL,
                             L_COMMITDATE  DATE NOT NULL,
                             L_RECEIPTDATE DATE NOT NULL,
                             L_SHIPINSTRUCT CHAR(25) NOT NULL,
                             L_SHIPMODE     CHAR(10) NOT NULL,
                             L_COMMENT      VARCHAR(44) NOT NULL,
  PRIMARY KEY (`L_ORDERKEY`,`L_LINENUMBER`,`L_SHIPDATE`)
) ENGINE=xengine TABLE_FORMAT=column;

创建OSS Server

OSS Server 对应前面存放数据的OSS Bucket,如果已经存在则无需重复创建。

DROP SERVER my_server;

CREATE SERVER my_server
FOREIGN DATA WRAPPER oss OPTIONS
(
  EXTRA_SERVER_INFO '{"oss_endpoint": "xxx","oss_bucket": "xxx","oss_access_key_id": "xxxx","oss_access_key_secret": "xxxx"}'
);

SELECT * FROM mysql.servers;

导入数据

LOAD DATA语句中OSS关键词表示从OSS中导入,路径支持* 通配符一次性导入多个文件。

快速装载参数说明:

参数名

参数值

含义

xengine_dload_async_build_nci

ON

异步构建非聚簇索引

xengine_dload_parallel

64

导入并发度

#!/usr/bin/env bash
set -u

OUT=output_a.txt
PREFIX="my_server/tpch1000g"
DB=tpch1000
USER=xxxxx
PASSWD=xxxx
HOST=xxxxx
PARA=64

: > "$OUT"

tables=(lineitem nation region part supplier partsupp customer orders)

run() {
  local t="$1"
  (
    echo "===== [$t] START $(date '+%F %T') ====="
    sql="set xengine_dload_async_build_nci=ON; set xengine_dload_parallel=${PARA};
         LOAD DATA OSS INFILE '${PREFIX}/${t}.*'
         INTO TABLE ${DB}.${t}
         FIELDS TERMINATED BY '|';"
    { time mysql -h"$HOST" -u"$USER" -p"$PASSWD" -e "$sql"; } 2>&1
    rc=${PIPESTATUS[0]:-0}
    echo "===== [$t] END   $(date '+%F %T') rc=$rc ====="
    echo
    exit "$rc"
  ) >>"$OUT" 2>&1 &
  echo "[$t] pid=$!"
}

for t in "${tables[@]}"; do
  run "$t"
done

wait
echo "ALL DONE $(date '+%F %T')" >>"$OUT"
说明

上述参数仅对X-Engine全列存表生效,对其它引擎无效。

InnoDB不支持 * 通配符,需要通过脚本逐个文件导入。

单击展开查看InnoDB导入参考脚本

#!/usr/bin/env python3
"""
TPC-H 1000 并发导入 MySQL 脚本

每张表独立使用 CONCURRENCY 个并发导入,所有表同时开始。
8张表 x CONCURRENCY 并发 = 最大并发连接数。

修改下方配置区参数后直接运行: python3 load_tpch.py
"""

import os, sys, time, subprocess, logging
from collections import defaultdict
from concurrent.futures import ThreadPoolExecutor, as_completed
import threading

# ==================== 配置区 ====================
CONCURRENCY    = 64           # 每张表的并发数
MYSQL_HOST     = ""
MYSQL_PORT     = 3306
MYSQL_USER     = ""
MYSQL_PASSWORD = "@"
MYSQL_DATABASE = "tpch1000"
DATA_DIR       = "my_server/tpch1000g"
MYSQL_CLI      = "mysql"
USE_LOCAL      = True         # True=LOAD DATA OSS INFILE
TIMEOUT        = 7200         # 单文件超时(秒)
# ================================================

CHUNKED = ["customer", "lineitem", "orders", "part", "partsupp", "supplier"]
CHUNKS  = 512

logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(threadName)s] %(message)s", datefmt="%H:%M:%S")
log = logging.getLogger()

# 统计
lock = threading.Lock()
table_stats = defaultdict(lambda: {"start": None, "end": None, "durations": [], "errors": []})

def load_file(table, fpath, fname):
    with lock:
        s = table_stats[table]
        if s["start"] is None:
            s["start"] = time.time()

    local = "OSS " if USE_LOCAL else ""
    sql = f"LOAD DATA {local}INFILE '{fpath}' INTO TABLE `{table}` FIELDS TERMINATED BY '|' LINES TERMINATED BY '|\n';"

    cmd = [MYSQL_CLI, f"-h{MYSQL_HOST}", f"-P{MYSQL_PORT}", f"-u{MYSQL_USER}",
           f"-p{MYSQL_PASSWORD}", f"-D{MYSQL_DATABASE}", "--batch", "--quick"]

    t0 = time.time()
    try:
        r = subprocess.run(cmd, input=sql, capture_output=True, text=True, timeout=TIMEOUT)
        dur = time.time() - t0
        ok = r.returncode == 0
        if not ok:
            log.error(f"FAIL {fname} {dur:.1f}s err={r.stderr.strip()[:150]}")
        else:
            log.info(f"OK {fname} {dur:.1f}s")
    except Exception as e:
        dur = time.time() - t0
        ok = False
        log.error(f"FAIL {fname} {dur:.1f}s exception={e}")

    with lock:
        s = table_stats[table]
        s["end"] = time.time()
        s["durations"].append(dur)
        if not ok:
            s["errors"].append(fname)
    return ok

def load_table(table, tasks):
    """为单张表启动独立线程池,CONCURRENCY 个并发"""
    concurrency = min(CONCURRENCY, len(tasks))
    log.info(f"表 {table} 开始导入: {len(tasks)} 个文件, 并发={concurrency}")

    ok_n = 0
    fail_n = 0
    with ThreadPoolExecutor(max_workers=concurrency, thread_name_prefix=table) as pool:
        futs = {pool.submit(load_file, t, fp, fn): fn for t, fp, fn in tasks}
        for fut in as_completed(futs):
            if fut.result():
                ok_n += 1
            else:
                fail_n += 1

    log.info(f"表 {table} 导入完成: 成功={ok_n} 失败={fail_n}")
    return ok_n, fail_n


def main():
    table_tasks = build_tasks_by_table()
    total_files = sum(len(v) for v in table_tasks.values())
    total_tables = len(table_tasks)
    max_conn = total_tables * CONCURRENCY

    log.info(f"共 {total_files} 个文件, {total_tables} 张表, "
             f"每表并发={CONCURRENCY}, 最大连接数={max_conn}, "
             f"目标={MYSQL_HOST}:{MYSQL_PORT}/{MYSQL_DATABASE}")

    t_start = time.time()

    # 每张表一个线程,在线程内部再开独立线程池
    results = {}
    with ThreadPoolExecutor(max_workers=total_tables, thread_name_prefix="dispatcher") as dispatcher:
        futs = {
            dispatcher.submit(load_table, tbl, tasks): tbl
            for tbl, tasks in table_tasks.items()
        }
        for fut in as_completed(futs):
            tbl = futs[fut]
            results[tbl] = fut.result()  # (ok_n, fail_n)

    total_dur = time.time() - t_start

    # 汇总
    total_ok = sum(r[0] for r in results.values())
    total_fail = sum(r[1] for r in results.values())
    print(f"\n{'='*88}")
    print(f"{'表名':<12} {'并发':>4} {'文件数':>6} {'失败':>4} "
          f"{'Wall耗时(s)':>12} {'累计耗时(s)':>12} {'最慢(s)':>8} {'最快(s)':>8}")
    print(f"{'-'*88}")
    for tbl in CHUNKED + ["nation", "region"]:
        s = table_stats[tbl]
        if not s["durations"]:
            continue
        wall = (s["end"] - s["start"]) if s["start"] and s["end"] else 0
        conc = min(CONCURRENCY, len(s["durations"]))
        print(f"{tbl:<12} {conc:>4} {len(s['durations']):>6} {len(s['errors']):>4} "
              f"{wall:>12.1f} {sum(s['durations']):>12.1f} "
              f"{max(s['durations']):>8.1f} {min(s['durations']):>8.1f}")
    print(f"{'='*88}")
    print(f"总耗时: {total_dur:.1f}s ({total_dur/60:.1f}min)  成功: {total_ok}  失败: {total_fail}")
    print(f"最大并发连接数: {max_conn} ({total_tables}表 x {CONCURRENCY}并发/表)")


if __name__ == "__main__":
    main()

测试结果

快速装载性能测试结果请参见快速装载性能