共计 2562 个字符,预计需要花费 7 分钟才能阅读完成。
背景痛点:为什么下载 CK 数据集总出问题?
刚开始用 ClickHouse 时,我经常遇到数据集下载的坑。比如:

- 连接中断:下载大表时网络波动直接导致前功尽弃
- 内存爆炸 :无脑
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
关键点注释:
max_threads=4:控制服务端处理查询的并发度,避免把 CK 服务器拖垮WITH TIES:确保在排序字段值相同时不会丢失数据tqdm进度条:直观展示下载进度和预估剩余时间- 异常捕获:网络波动时自动重试当前分块
性能优化技巧
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 级数据从未翻车。核心经验就三点:
- 分而治之:大文件一定要拆分成小块
- 留好后路:做好重试和校验机制
- 量力而行:根据机器配置调整并发度
下次遇到 CK 数据导出问题时,希望这篇指南能帮你少走弯路。如果有更好的优化方案,欢迎交流!
正文完
发表至: 技术分享
近一天内
