数据导入

更新时间:
复制 MD 格式

本文介绍如何使用Direct Load快速装载机制向PolarDB MySQL列存表中高效导入数据,支持从OSS和本地文件导入CSV格式数据。

概述

Direct Load是一种快速装载机制,为X-Engine列存表提供快速数据加载能力。Direct Load支持从OSS导入数据、从本地文件导入数据以及导入格式为CSV(默认格式)这三种导入方式。除Direct Load外,X-Engine全列存还支持以下数据导入方式:

会话参数

Direct Load提供以下会话级参数,用于控制导入行为。

参数名

说明

xengine_dload_parallel

存储层并发度,0表示关闭快速装载。

  • 取值范围:0~1024。

  • 默认值:0

xengine_dload_check_unique_mod

唯一性检查模式。取值范围:

  • OFF(默认值):关闭唯一性检查(主键重复仍报错)。

  • KEEPONE:发现重复时随机保留一条。

  • ERROR:发现重复立即报错。

语法

Direct Load使用LOAD DATA语法,支持从OSS或本地文件导入数据。

LOAD DATA [OSS|LOCAL]
    INFILE 'file_name'
    [REPLACE | IGNORE]
    INTO TABLE tbl_name
    [CHARACTER SET charset_name]
    [{FIELDS | COLUMNS}
        [TERMINATED BY 'string']
        [[OPTIONALLY] ENCLOSED BY 'char']
        [ESCAPED BY 'char']
    ]
    [LINES
        [STARTING BY 'string']
        [TERMINATED BY 'string']
    ]
    [IGNORE number {LINES | ROWS}]
    [(col_name_or_user_var [, ...])]
    [SET col_name={expr | DEFAULT} [, ...]]

限制条件

使用Direct Load导入数据时,需要注意以下限制。

  • 表要求

    • Direct Load仅支持ENGINE=XENGINETABLE_FORMAT=COLUMN的表,表必须为空。

    • 表需为普通表(非分区表)。

      说明

      对于分区表,如果需要按分区导入,则可以先把单个分区导入临时表,然后通过EXCHANGE PARTITION语法将临时表转换成目标分区表的分区。EXCHANGE PARTITION语法仅交换元数据速度很快。HASH分区暂不支持。

  • 索引支持

    • 支持带ORDER KEY的表

    • 不支持二级索引、外键或Hidden Key。

      说明

      当参数xengine_enable_ctable_hidden_key=ON时存在Hidden Key限制。

  • Binlog:Direct Load启用时自动关闭Binlog写入(通过临时清除OPTION_BIN_LOG标志),不影响全局配置。

  • 触发器:带触发器的表不支持Direct Load。

  • 并发限制

    • 每张表同时只允许一个Direct Load语句。

    • 执行期间会持有排他MDL锁。

    • 禁止对该表进行DML(数据操作)和DDL(数据定义)操作。

  • 文件通配符

    • OSS模式支持通配符(*),例如/home/t4/data/*.csv

    • LOCAL模式和非Direct Load模式不支持通配符。

使用示例

OSS导入

步骤一:创建OSS Server

首先创建OSS Server,配置OSS连接信息。

权限说明

  • 创建OSS Server时需要SERVERS_ADMIN权限。

  • 可通过以下命令查看当前登录用户是否具有该权限:

    SHOW GRANTS FOR 用户名;
  • 高权限账户可以给低权限账户赋予该权限。

  • 如果没有权限,会报错:

    Access denied; you need (at least one of) the SERVERS_ADMIN OR SUPER privilege(s) for this operation

创建OSS Server示例

DROP SERVER IF EXISTS my_server;
CREATE SERVER my_server FOREIGN DATA WRAPPER oss OPTIONS (
  EXTRA_SERVER_INFO '{"oss_endpoint":"oss-cn-xxxx-internal.aliyuncs.com","oss_bucket":"xxxx","oss_access_key_id":"xxxx","oss_access_key_secret":"xxxx"}'
);
SELECT * FROM mysql.servers;

步骤二:创建目标表

创建X-Engine列存目标表。以下以TPC-H lineitem表为例。

CREATE TABLE lineitem (
    l_orderkey BIGINT NOT NULL,
    l_partkey BIGINT NOT NULL,
    l_suppkey BIGINT NOT NULL,
    l_linenumber INT 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),
    ORDER KEY(l_orderkey)
) ENGINE=XENGINE TABLE_FORMAT=COLUMN;

步骤三:导入数据

设置并行度后,使用LOAD DATA OSSOSS导入数据。支持通配符匹配多个文件。

-- 设置并发度(可选)
SET xengine_dload_parallel = 2;

-- 使用通配符导入(OSS模式支持通配符)
LOAD DATA OSS INFILE 'my_server/lineitem/lineitem_test*.csv'
INTO TABLE lineitem
FIELDS TERMINATED BY '|';

-- 验证导入结果
SELECT * FROM lineitem ORDER BY l_orderkey, l_partkey DESC LIMIT 10;
SELECT COUNT(1) FROM lineitem;

从本地文件导入

说明

从本地文件Direct Load,仅支持单个文件导入,多个文件请使用从OSS导入。

使用LOAD DATA LOCAL从客户端本地文件导入数据。

CREATE TABLE lineitem (
    l_orderkey       BIGINT NOT NULL,
    l_partkey        BIGINT NOT NULL,
    l_suppkey        BIGINT NOT NULL,
    l_linenumber     BIGINT 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),
    ORDER KEY(l_orderkey)
) ENGINE=XENGINE TABLE_FORMAT=COLUMN;

SET xengine_dload_parallel = 1;

-- 从本地文件导入(不支持通配符)
LOAD DATA LOCAL INFILE '/home/data/lineitem/lineitem_test.csv'
INTO TABLE lineitem
FIELDS TERMINATED BY '|';

-- 验证导入结果
SELECT * FROM lineitem ORDER BY l_orderkey, l_partkey DESC LIMIT 10;
SELECT COUNT(1) FROM lineitem;

TPC-H 100G导入示例

以下示例展示如何使用Direct LoadOSS导入TPC-H 100 GB规模的数据。

步骤一:创建目标表

lineitem表为例,创建X-Engine列存表。

SET xengine_enable_ctable_hidden_key = OFF;

DROP DATABASE IF EXISTS tpch100;
CREATE DATABASE tpch100;
USE tpch100;

-- 示例:创建lineitem表
CREATE TABLE lineitem (
    L_ORDERKEY    INTEGER 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

DROP SERVER IF EXISTS my_server;

CREATE SERVER my_server
FOREIGN DATA WRAPPER oss OPTIONS
(
  EXTRA_SERVER_INFO '{"oss_endpoint": "oss-cn-xxxx.aliyuncs.com","oss_bucket": "xxxx","oss_access_key_id": "xxxx","oss_access_key_secret": "xxxx"}'
);

批量导入脚本

使用Bash脚本批量导入所有TPC-H表的数据。

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

SOCK=/tmp/mysql29.sock
OUT=output.txt
PREFIX="my_server/imci/tpch100g"
DB=tpch100
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_parallel=${PARA};
         LOAD DATA OSS INFILE '${PREFIX}/${t}.*'
         INTO TABLE ${DB}.${t}
         FIELDS TERMINATED BY '|';"
    { time mysql -uroot -S "$SOCK" -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"

错误处理

使用Direct Load导入数据时,可能遇到以下常见错误。

错误码

说明

解决方案

ER_FEATURE_UNSUPPORTED

不支持的特性。

检查表是否满足Direct Load要求:

  • ENGINE=XENGINE

  • TABLE_FORMAT=COLUMN

  • 确认不存在二级索引、外键或触发器。

ER_DUP_KEY

主键冲突。

检查导入数据中是否存在重复主键。可通过设置xengine_dload_check_unique_mod=KEEPONE自动去重。

最佳实践

  • Binlog自动关闭:Direct Load自动关闭Binlog写入,无需手动操作。导入完成后自动恢复,不影响全局配置。

  • 设置合适的并发度:OSS导入场景通常建议设置xengine_dload_parallel32 ~ 64,本地文件导入建议设置较低的并发度。示例:

    SET xengine_dload_parallel = 64;
  • 确保表为空:Direct Load要求目标表为空,导入前请确认表中无数据。

  • 导入前移除二级索引:Direct Load不支持二级索引,请在导入前删除二级索引,导入完成后再重新创建。

  • OSS模式使用通配符:利用通配符(*)一次性匹配多个文件,简化大规模数据导入流程。示例:

    LOAD DATA OSS INFILE 'my_server/lineitem/lineitem_test*.csv'
    INTO TABLE lineitem
    FIELDS TERMINATED BY '|';