如何写冷热数据分离存储脚本

wen 实用脚本 28

本文目录导读:

如何写冷热数据分离存储脚本

  1. 基础冷热数据分离脚本
  2. 完善的Python脚本版本
  3. 配置文件示例 (config.json)
  4. 使用说明

我来提供一个完整的冷热数据分离存储脚本示例,这里以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等)
  • 批量处理
  • 错误处理和日志
  • 配置化
  • 模拟运行模式

可以根据实际需求调整参数和存储策略。

抱歉,评论功能暂时关闭!