本文目录导读:

我来提供一个完整的冷热数据分离存储脚本示例,这里以MySQL数据库为例,展示一个实用的冷热数据分离方案:
基础冷热数据分离脚本
-- 冷热数据分离存储过程
DELIMITER //
CREATE PROCEDURE `cold_data_separation`(
IN p_days_threshold INT, -- 多少天前的数据算冷数据
IN p_batch_size INT, -- 每批处理数量
IN p_table_name VARCHAR(100) -- 源表名
)
BEGIN
DECLARE v_affected_rows INT DEFAULT 0;
DECLARE v_total_rows INT DEFAULT 0;
DECLARE v_cold_table VARCHAR(100);
-- 构建冷数据表名
SET v_cold_table = CONCAT(p_table_name, '_cold');
-- 创建冷数据表(如果不存在)
SET @create_sql = CONCAT(
'CREATE TABLE IF NOT EXISTS ', v_cold_table, ' ',
'LIKE ', p_table_name
);
PREPARE stmt FROM @create_sql;
EXECUTE stmt;
DEALLOCATE PREPARE stmt;
-- 循环迁移冷数据
REPEAT
-- 迁移数据到冷表
SET @move_sql = CONCAT(
'INSERT INTO ', v_cold_table, ' ',
'SELECT * FROM ', p_table_name, ' ',
'WHERE create_time < DATE_SUB(NOW(), INTERVAL ', p_days_threshold, ' DAY) ',
'LIMIT ', p_batch_size
);
PREPARE stmt FROM @move_sql;
EXECUTE stmt;
SET v_affected_rows = ROW_COUNT();
DEALLOCATE PREPARE stmt;
-- 删除已迁移的热数据
IF v_affected_rows > 0 THEN
SET @delete_sql = CONCAT(
'DELETE FROM ', p_table_name, ' ',
'WHERE create_time < DATE_SUB(NOW(), INTERVAL ', p_days_threshold, ' DAY) ',
'LIMIT ', p_batch_size
);
PREPARE stmt FROM @delete_sql;
EXECUTE stmt;
DEALLOCATE PREPARE stmt;
SET v_total_rows = v_total_rows + v_affected_rows;
SELECT CONCAT('Moved ', v_affected_rows, ' rows to cold storage. Total: ', v_total_rows) AS progress;
END IF;
-- 避免长时间锁表
DO SLEEP(0.1);
UNTIL v_affected_rows = 0 END REPEAT;
SELECT CONCAT('Completed. Total rows moved: ', v_total_rows) AS result;
END//
DELIMITER ;
-- 调用示例:迁移30天前的数据,每批1000条
CALL cold_data_separation(30, 1000, 'orders');
完善的Python脚本版本
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
冷热数据分离存储脚本
支持多种数据库(MySQL, PostgreSQL)和多种存储后端
"""
import os
import sys
import time
import logging
import argparse
from datetime import datetime, timedelta
import json
import pymysql
import boto3
from botocore.exceptions import ClientError
# 配置日志
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(levelname)s - %(message)s',
handlers=[
logging.FileHandler('cold_data_migration.log'),
logging.StreamHandler()
]
)
logger = logging.getLogger(__name__)
class ColdDataMigration:
"""冷热数据迁移类"""
def __init__(self, config):
self.config = config
self.source_conn = None
self.target_conn = None
def connect(self):
"""建立数据库连接"""
try:
self.source_conn = pymysql.connect(
host=self.config['source_host'],
port=self.config.get('source_port', 3306),
user=self.config['source_user'],
password=self.config['source_password'],
database=self.config['source_database'],
charset='utf8mb4'
)
if self.config.get('target_type') == 'database':
self.target_conn = pymysql.connect(
host=self.config['target_host'],
port=self.config.get('target_port', 3306),
user=self.config['target_user'],
password=self.config['target_password'],
database=self.config['target_database'],
charset='utf8mb4'
)
logger.info("数据库连接成功")
except Exception as e:
logger.error(f"数据库连接失败: {e}")
raise
def get_cold_data_query(self, table_config):
"""生成冷数据查询SQL"""
days_threshold = table_config.get('days_threshold', 90)
date_column = table_config.get('date_column', 'create_time')
return f"""
SELECT * FROM {table_config['table_name']}
WHERE {date_column} < DATE_SUB(NOW(), INTERVAL {days_threshold} DAY)
AND {date_column} IS NOT NULL
"""
def migrate_to_database(self, table_config):
"""迁移到数据库"""
try:
# 创建目标表(如果不存在)
create_target_sql = f"""
CREATE TABLE IF NOT EXISTS {table_config.get('target_table', table_config['table_name'] + '_cold')}
LIKE {table_config['table_name']}
"""
with self.target_conn.cursor() as cursor:
cursor.execute(create_target_sql)
# 获取冷数据
cold_data_query = self.get_cold_data_query(table_config)
with self.source_conn.cursor() as source_cursor:
source_cursor.execute(cold_data_query)
batch_size = table_config.get('batch_size', 1000)
total_moved = 0
while True:
rows = source_cursor.fetchmany(batch_size)
if not rows:
break
# 插入目标表
columns = [desc[0] for desc in source_cursor.description]
placeholders = ','.join(['%s'] * len(columns))
columns_str = ','.join(columns)
insert_sql = f"""
INSERT INTO {table_config.get('target_table', table_config['table_name'] + '_cold')}
({columns_str})
VALUES ({placeholders})
"""
with self.target_conn.cursor() as target_cursor:
target_cursor.executemany(insert_sql, rows)
self.target_conn.commit()
# 删除源数据
ids_to_delete = [row[0] for row in rows] # 假设第一列是ID
delete_sql = f"""
DELETE FROM {table_config['table_name']}
WHERE id IN ({','.join(['%s'] * len(ids_to_delete))})
"""
with self.source_conn.cursor() as delete_cursor:
delete_cursor.execute(delete_sql, ids_to_delete)
self.source_conn.commit()
total_moved += len(rows)
logger.info(f"已迁移 {total_moved} 条数据到冷存储")
# 避免长时间锁表
time.sleep(0.1)
logger.info(f"表 {table_config['table_name']} 迁移完成,共迁移 {total_moved} 条")
except Exception as e:
logger.error(f"数据库迁移失败: {e}")
self.source_conn.rollback()
if self.target_conn:
self.target_conn.rollback()
raise
def migrate_to_s3(self, table_config):
"""迁移到AWS S3"""
try:
s3_client = boto3.client(
's3',
aws_access_key_id=self.config.get('aws_access_key_id'),
aws_secret_access_key=self.config.get('aws_secret_access_key'),
region_name=self.config.get('aws_region', 'us-east-1')
)
bucket_name = table_config.get('s3_bucket')
prefix = table_config.get('s3_prefix', 'cold_data/')
# 获取冷数据
cold_data_query = self.get_cold_data_query(table_config)
with self.source_conn.cursor() as cursor:
cursor.execute(cold_data_query)
columns = [desc[0] for desc in cursor.description]
batch_size = table_config.get('batch_size', 1000)
total_moved = 0
while True:
rows = cursor.fetchmany(batch_size)
if not rows:
break
# 转换为JSON并上传到S3
for row in rows:
data = dict(zip(columns, row))
file_key = f"{prefix}{table_config['table_name']}/{datetime.now().strftime('%Y/%m/%d')}/{data['id']}.json"
s3_client.put_object(
Bucket=bucket_name,
Key=file_key,
Body=json.dumps(data, default=str)
)
# 删除源数据
ids_to_delete = [row[0] for row in rows]
delete_sql = f"""
DELETE FROM {table_config['table_name']}
WHERE id IN ({','.join(['%s'] * len(ids_to_delete))})
"""
with self.source_conn.cursor() as delete_cursor:
delete_cursor.execute(delete_sql, ids_to_delete)
self.source_conn.commit()
total_moved += len(rows)
logger.info(f"已上传 {total_moved} 条数据到S3")
time.sleep(0.1)
logger.info(f"S3迁移完成,共迁移 {total_moved} 条")
except Exception as e:
logger.error(f"S3迁移失败: {e}")
self.source_conn.rollback()
raise
def run(self):
"""运行迁移"""
try:
self.connect()
for table_config in self.config['tables']:
target_type = self.config.get('target_type', 'database')
if target_type == 'database':
self.migrate_to_database(table_config)
elif target_type == 's3':
self.migrate_to_s3(table_config)
else:
logger.error(f"不支持的目标类型: {target_type}")
except Exception as e:
logger.error(f"迁移过程出错: {e}")
raise
finally:
if self.source_conn:
self.source_conn.close()
if self.target_conn:
self.target_conn.close()
def main():
"""主函数"""
parser = argparse.ArgumentParser(description='冷热数据分离迁移脚本')
parser.add_argument('-c', '--config', required=True, help='配置文件路径')
parser.add_argument('-d', '--dry-run', action='store_true', help='模拟运行')
args = parser.parse_args()
# 加载配置
with open(args.config, 'r') as f:
config = json.load(f)
if args.dry_run:
logger.info("模拟运行模式 - 仅显示将要执行的操作")
for table_config in config['tables']:
logger.info(f"将迁移表: {table_config['table_name']}")
logger.info(f"条件: {table_config.get('days_threshold', 90)} 天前的数据")
logger.info(f"每次批处理: {table_config.get('batch_size', 1000)} 条")
return
# 执行迁移
migration = ColdDataMigration(config)
migration.run()
if __name__ == "__main__":
main()
配置文件示例 (config.json)
{
"source_host": "localhost",
"source_port": 3306,
"source_user": "root",
"source_password": "password",
"source_database": "production_db",
"target_type": "database",
"target_host": "cold-storage-host",
"target_port": 3306,
"target_user": "cold_user",
"target_password": "cold_password",
"target_database": "cold_storage_db",
"aws_access_key_id": "YOUR_AWS_KEY",
"aws_secret_access_key": "YOUR_AWS_SECRET",
"aws_region": "us-east-1",
"tables": [
{
"table_name": "orders",
"date_column": "create_time",
"days_threshold": 90,
"batch_size": 1000,
"target_table": "orders_archive",
"s3_bucket": "my-cold-storage-bucket",
"s3_prefix": "orders/"
},
{
"table_name": "logs",
"date_column": "log_time",
"days_threshold": 30,
"batch_size": 5000,
"target_table": "logs_archive",
"s3_bucket": "my-cold-storage-bucket",
"s3_prefix": "logs/"
}
]
}
使用说明
安装依赖
pip install pymysql boto3
运行脚本
# 实际运行 python cold_data_migration.py -c config.json # 模拟运行 python cold_data_migration.py -c config.json --dry-run
定时执行(crontab)
# 每天凌晨2点执行 0 2 * * * /usr/bin/python3 /path/to/cold_data_migration.py -c /path/to/config.json >> /var/log/cold_migration.log 2>&1
这个脚本提供了完善的冷热数据分离功能,支持:
- 多种数据库源
- 多种存储目标(数据库、S3等)
- 批量处理
- 错误处理和日志
- 配置化
- 模拟运行模式
可以根据实际需求调整参数和存储策略。