clickhouse如何从postgres导入
·
从 PostgreSQL 导入数据到 ClickHouse 的完整方法,涵盖多种场景:一次性迁移、实时同步、联邦查询等。
一、ClickHouse 内置 PostgreSQL 引擎(推荐)
ClickHouse 内置 PostgreSQL 引擎,可直接查询 PostgreSQL 数据,无需中间件。
1. 创建 PostgreSQL 数据库引擎
sql
复制
-- 创建虚拟数据库,自动映射所有 PostgreSQL 表
CREATE DATABASE pg_db
ENGINE = PostgreSQL(
'postgres-host:5432', -- PostgreSQL 地址:端口
'source_database', -- PostgreSQL 数据库名
'username', -- 用户名
'password' -- 密码
);
2. 直接查询 PostgreSQL 数据
sql
复制
-- 查询 PostgreSQL 表(无需导入)
SELECT * FROM pg_db.users;
SELECT * FROM pg_db.orders WHERE created_at >= '2024-01-01';
-- 关联 ClickHouse 和 PostgreSQL 数据
SELECT
ch.order_id,
pg.customer_name,
pg.customer_email
FROM clickhouse_orders ch
INNER JOIN pg_db.customers pg ON ch.customer_id = pg.id;
3. 将数据导入 ClickHouse
sql
复制
-- 创建本地表
CREATE TABLE users_local (
id UInt64,
name String,
email String,
created_at DateTime
) ENGINE = MergeTree()
ORDER BY id;
-- 从 PostgreSQL 导入数据
INSERT INTO users_local
SELECT *
FROM pg_db.users
WHERE created_at >= '2024-01-01';
二、PostgreSQL 表函数(单次查询)
无需创建数据库,直接查询:
sql
复制
-- 单次查询 PostgreSQL
SELECT * FROM postgresql(
'postgres-host:5432',
'database',
'table_name',
'user',
'password'
);
-- 导入数据
INSERT INTO clickhouse_table
SELECT * FROM postgresql(
'postgres:5432',
'mydb',
'events',
'admin',
'secret'
) WHERE event_date >= today() - INTERVAL 7 DAY;
三、批量导出/导入(CSV 方式)
适用于大规模一次性迁移:
步骤 1:从 PostgreSQL 导出 CSV
bash
复制
# 使用 psql 导出
psql -h localhost -U postgres -d mydb -c \
"\COPY (SELECT * FROM events WHERE created_at >= '2024-01-01') TO '/tmp/events.csv' WITH CSV HEADER"
# 或使用 COPY 命令(在 psql 内执行)
COPY events TO '/tmp/events.csv' WITH (FORMAT CSV, HEADER);
步骤 2:导入 ClickHouse
sql
复制
-- 创建目标表
CREATE TABLE events (
id UInt64,
user_id UInt32,
event_type String,
created_at DateTime
) ENGINE = MergeTree()
ORDER BY (user_id, created_at);
-- 从 CSV 导入
INSERT INTO events
SELECT *
FROM file('/tmp/events.csv', 'CSVWithNames');
-- 或使用 URL/S3
INSERT INTO events
SELECT * FROM url('http://example.com/events.csv', 'CSVWithNames');
四、实时同步方案(CDC)
方案 1:MaterializedPostgreSQL 引擎(自动同步)
ClickHouse 实验性功能,自动同步 PostgreSQL 数据:
sql
复制
-- 创建物化 PostgreSQL 数据库(自动 CDC)
CREATE DATABASE pg_sync
ENGINE = MaterializedPostgreSQL(
'postgres-host:5432',
'source_db',
'user',
'password'
) SETTINGS
materialized_postgresql_tables_list = 'users,orders,products', -- 指定同步表
materialized_postgresql_schema = 'public';
-- 查看同步状态
SELECT
database, table, is_replicated, replication_lag_seconds
FROM system.materialized_postgresql_tables;
⚠️ 注意:需开启 PostgreSQL 逻辑复制(
wal_level = logical)
方案 2:Debezium + Kafka(生产级方案)
架构流程 :
plain
复制
PostgreSQL → WAL → Debezium → Kafka → ClickHouse Kafka Engine
PostgreSQL 配置:
sql
复制
-- 开启逻辑复制
ALTER SYSTEM SET wal_level = logical;
ALTER SYSTEM SET max_replication_slots = 10;
ALTER SYSTEM SET max_wal_senders = 10;
-- 创建复制槽
SELECT pg_create_logical_replication_slot('ch_slot', 'pgoutput');
ClickHouse Kafka 引擎:
sql
复制
CREATE TABLE events_queue (
id UInt64,
data String
) ENGINE = Kafka(
'kafka-host:9092',
'pg-events',
'ch-group',
'JSONEachRow'
);
CREATE MATERIALIZED VIEW events_mv
TO events AS
SELECT * FROM events_queue;
五、云服务商方案
阿里云 DTS(数据传输服务)
适用于阿里云 RDS PostgreSQL 到 ClickHouse:
表格
| 功能 | 说明 |
|---|---|
| 结构迁移 | 自动转换表结构,添加 _sign、_version 系统字段 |
| 全量迁移 | 批量导出导入历史数据 |
| 增量同步 | 基于逻辑复制槽实时同步 DML/DDL |
| 数据校验 | 自动校验数据一致性 |
关键配置:
-
PostgreSQL 需设置
wal_level = logical -
分区键不支持可空字段,仅支持 BIGINT/INT/TIMESTAMP/DATE
-
目标表自动添加
_sign(Int8)、_version(UInt64)字段
六、NineData 数据复制(第三方工具)
支持可视化配置,全量+增量同步:
特点:
-
结构自动迁移(自动添加
_jz_data_sign、_jz_data_time字段) -
行级并发复制
-
断点续传
-
双向同步支持
sql
复制
-- NineData 自动添加的系统列
_jz_data_sign Int8 DEFAULT 1, -- 1=INSERT, -1=DELETE
_jz_data_time String DEFAULT now() -- 更新时间
七、数据类型映射对照
表格
| PostgreSQL | ClickHouse | 说明 |
|---|---|---|
BIGINT | Int64 / UInt64 | 无符号类型映射为 Uint |
INTEGER | Int32 | |
VARCHAR/TEXT | String | |
TIMESTAMP | DateTime64 | 注意时区处理 |
DATE | Date | |
BOOLEAN | UInt8 | 0/1 存储 |
JSON/JSONB | String | 需手动解析 |
NUMERIC | Decimal(P, S) | 需指定精度 |
八、推荐方案选择
表格
| 场景 | 推荐方案 | 延迟 |
|---|---|---|
| 临时查询/联邦分析 | PostgreSQL 引擎 | 实时 |
| 一次性全量迁移 | CSV 导出导入 | 分钟级 |
| 实时同步(简单) | MaterializedPostgreSQL | 秒级 |
| 生产级实时同步 | Debezium + Kafka | 秒级 |
| 云环境迁移 | 阿里云 DTS / NineData | 秒级 |
需要针对特定场景(如大数据量迁移、实时同步配置等)的详细步骤吗?
更多推荐
所有评论(0)