ClickHouse(ck)数据集下载实战指南:从零搭建高效数据管道

1次阅读
没有评论

共计 2562 个字符,预计需要花费 7 分钟才能阅读完成。

image.webp

背景痛点:为什么下载 CK 数据集总出问题?

刚开始用 ClickHouse 时,我经常遇到数据集下载的坑。比如:

ClickHouse(ck)数据集下载实战指南:从零搭建高效数据管道

  • 连接中断:下载大表时网络波动直接导致前功尽弃
  • 内存爆炸 :无脑SELECT * 把 32G 内存的机器搞崩
  • 格式混乱:CSV 文件里的特殊字符破坏数据完整性
  • 性能低下:单线程下载 20GB 数据等到天荒地老

这些问题的本质,是没搞清楚 CK 的数据传输特性。下面分享我的实战解决方案。

协议对比:HTTP 接口 vs Native 协议

特性 HTTP 接口 Native 协议
吞吐量 中等(依赖 HTTP 服务器) 高(二进制传输)
压缩率 依赖 Content-Encoding 自带 LZ4/ZSTD 压缩
兼容性 通用(任何 HTTP 客户端) 需专用驱动
开发难度 简单(curl 即可) 中等(需 SDK)
适用场景 小数据量临时导出 生产环境大数据传输

结论:超过 1GB 的数据集优先用 Native 协议。

核心实现:Python 分块下载代码

import os
from clickhouse_driver import Client
from tqdm import tqdm

# 常量定义
CH_HOST = 'localhost'
CH_USER = 'default'
CH_PASSWORD = ''TABLE_NAME ='dataset_large'
BATCH_SIZE = 500000  # 每批 50 万行
OUTPUT_DIR = './downloads'

# 创建下载目录
os.makedirs(OUTPUT_DIR, exist_ok=True)

client = Client(host=CH_HOST, user=CH_USER, password=CH_PASSWORD,
               settings={'max_threads': 4})  # 控制并发线程数

# 获取总行数
total_rows = client.execute(f"SELECT count() FROM {TABLE_NAME}")[0][0]

# 分块下载
with tqdm(total=total_rows, desc='Downloading') as pbar:
    offset = 0
    part_num = 1

    while offset < total_rows:
        try:
            # 使用 WITH TIES 避免 LIMIT 截断数据
            query = f"""
            SELECT * FROM {TABLE_NAME} 
            ORDER BY id 
            LIMIT {BATCH_SIZE} 
            OFFSET {offset} WITH TIES
            """

            data = client.execute(query)

            # 写入分块文件(实际项目建议用 Parquet 格式)part_file = os.path.join(OUTPUT_DIR, f'part_{part_num}.csv')
            with open(part_file, 'w') as f:
                for row in data:
                    f.write(','.join(map(str, row)) + '\n')

            # 更新进度
            offset += len(data)
            part_num += 1
            pbar.update(len(data))

        except Exception as e:
            print(f"Error: {e}, retrying...")
            continue

关键点注释

  1. max_threads=4:控制服务端处理查询的并发度,避免把 CK 服务器拖垮
  2. WITH TIES:确保在排序字段值相同时不会丢失数据
  3. tqdm进度条:直观展示下载进度和预估剩余时间
  4. 异常捕获:网络波动时自动重试当前分块

性能优化技巧

1. 并行下载加速

from concurrent.futures import ThreadPoolExecutor

def download_chunk(args):
    # 实现分块下载逻辑...

# 启动 4 个线程并行下载
with ThreadPoolExecutor(max_workers=4) as executor:
    chunks = [(i, BATCH_SIZE) for i in range(0, total_rows, BATCH_SIZE)]
    list(tqdm(executor.map(download_chunk, chunks), total=len(chunks)))

2. 压缩传输

在 Client 连接参数中添加:

client = Client(
    compression='lz4',  # 启用压缩
    # 其他参数...
)

避坑指南

1. 内存管理

  • 致命错误:直接运行client.execute("SELECT * FROM huge_table")
  • 正确做法
  • 始终添加 LIMIT 测试数据规模
  • system.tables 查看表大小:
    SELECT name, formatReadableSize(total_bytes) 
    FROM system.tables 
    WHERE database = 'default'

2. 分区表下载

对于分区表,建议按分区目录结构存储:

downloads/
   ├── pdate=20230101/
   │   ├── part_1.csv
   │   └── part_2.csv
   └── pdate=20230102/
       ├── part_1.csv
       └── part_2.csv

对应的查询优化:

SELECT * FROM partitioned_table 
WHERE pdate = '2023-01-01'
ORDER BY timestamp

延伸思考

1. 断点续传方案

结合 S3/MinIO 实现:
1. 每个分块上传后记录元数据到 Redis
2. 重启时检查已完成的 chunk
3. 使用 ALTER TABLE ... DETACH PARTITION 冻结正在下载的分区

2. 数据完整性校验

推荐方法:

# 服务端生成校验和
server_checksum = client.execute("SELECT sum(cityHash64(*)) FROM table"
)[0][0]

# 客户端验证
local_checksum = 0
for row in local_data:
    local_checksum += cityhash64(row)
assert server_checksum == local_checksum

写在最后

这套方案在我们生产环境稳定运行了半年,单日下载 TB 级数据从未翻车。核心经验就三点:

  1. 分而治之:大文件一定要拆分成小块
  2. 留好后路:做好重试和校验机制
  3. 量力而行:根据机器配置调整并发度

下次遇到 CK 数据导出问题时,希望这篇指南能帮你少走弯路。如果有更好的优化方案,欢迎交流!

正文完
 0
评论(没有评论)