| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317 |
- #!/usr/bin/env python3
- """
- 京东云对象存储存量迁移脚本。
- 遍历本地 upload.base-dir(默认 /data/cfc-uploads)下的所有文件,
- 上传到京东云 OSS bucket cfc,保持原有目录结构。
- 用法:
- # 查看会迁移哪些文件(不实际上传)
- python3 scripts/migrate-to-jdcloud.py --dry-run
- # 执行迁移(自动跳过云端已存在的文件)
- python3 scripts/migrate-to-jdcloud.py
- # 强制覆盖云端已存在的文件
- python3 scripts/migrate-to-jdcloud.py --overwrite
- # 指定本地目录
- python3 scripts/migrate-to-jdcloud.py --local-dir /data/cfc-uploads
- # 迁移后回填 file_record 表(用于去重)
- python3 scripts/migrate-to-jdcloud.py --backfill
- """
- import os
- import sys
- import boto3
- import hashlib
- import pymysql
- from botocore.config import Config
- from pathlib import Path
- import argparse
- import logging
- from datetime import datetime
- logging.basicConfig(
- level=logging.INFO,
- format="%(asctime)s [%(levelname)s] %(message)s",
- datefmt="%H:%M:%S",
- )
- log = logging.getLogger("migrate")
- # 京东云 OSS 配置
- ENDPOINT = "s3.cn-north-1.jdcloud-oss.com"
- ACCESS_KEY = "JDC_F351EC3CB1F8593204ABCDA2CBF2"
- SECRET_KEY = "970C0C6985A6BF29D0C79AE3693B9899"
- BUCKET = "cfc"
- REGION = "cn-north-1"
- # MySQL 配置(用于回填 file_record 表)
- DB_HOST = "mysql-internet-cn-north-1-23feae22680e4dfa.rds.jdcloud.com"
- DB_PORT = 3306
- DB_USER = "cfc"
- DB_PASS = "cfc@1314"
- DB_NAME = "zxyj"
- # 公网URL前缀
- PUBLIC_URL_PREFIX = f"https://{BUCKET}.{ENDPOINT}/"
- def get_client():
- return boto3.client(
- "s3",
- endpoint_url=f"https://{ENDPOINT}",
- aws_access_key_id=ACCESS_KEY,
- aws_secret_access_key=SECRET_KEY,
- config=Config(s3={"addressing_style": "path"}, signature_version="s3v4"),
- region_name=REGION,
- )
- def get_db():
- return pymysql.connect(
- host=DB_HOST, port=DB_PORT, user=DB_USER, password=DB_PASS,
- database=DB_NAME, charset="utf8mb4", cursorclass=pymysql.cursors.DictCursor,
- )
- def list_cloud_keys(client):
- """列出云端已存在的所有 key"""
- keys = set()
- marker = None
- while True:
- if marker:
- resp = client.list_objects_v2(Bucket=BUCKET, MaxKeys=1000, StartAfter=marker)
- else:
- resp = client.list_objects_v2(Bucket=BUCKET, MaxKeys=1000)
- if "Contents" not in resp:
- break
- for obj in resp["Contents"]:
- keys.add(obj["Key"])
- marker = obj["Key"]
- if not resp.get("IsTruncated"):
- break
- return keys
- def upload_file(client, local_path, key, dry_run=False, overwrite=False, cloud_keys=None):
- """上传单个文件,返回 (成功与否, target_url)"""
- target_url = PUBLIC_URL_PREFIX + key
- if cloud_keys and key in cloud_keys and not overwrite:
- log.info(" [SKIP] 已在云端存在: %s", target_url)
- return False, target_url
- if dry_run:
- log.info(" [DRY-RUN] 将上传: %s -> %s", local_path, target_url)
- return False, target_url
- try:
- content_type = guess_content_type(local_path)
- extra_args = {"ContentType": content_type} if content_type else {}
- client.upload_file(str(local_path), BUCKET, key, ExtraArgs=extra_args)
- log.info(" [OK] %s", target_url)
- return True, target_url
- except Exception as e:
- log.error(" [FAIL] %s: %s", local_path, e)
- return False, target_url
- def guess_content_type(path):
- """根据扩展名猜测 Content-Type"""
- ext = Path(path).suffix.lower()
- mapping = {
- ".jpg": "image/jpeg",
- ".jpeg": "image/jpeg",
- ".png": "image/png",
- ".gif": "image/gif",
- ".webp": "image/webp",
- ".bmp": "image/bmp",
- ".svg": "image/svg+xml",
- ".pdf": "application/pdf",
- ".doc": "application/msword",
- ".docx": "application/vnd.openxmlformats-officedocument.wordprocessingml.document",
- ".xls": "application/vnd.ms-excel",
- ".xlsx": "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
- ".mp4": "video/mp4",
- ".mp3": "audio/mpeg",
- ".wav": "audio/wav",
- ".txt": "text/plain",
- ".json": "application/json",
- ".zip": "application/zip",
- ".html": "text/html",
- ".css": "text/css",
- ".js": "application/javascript",
- }
- return mapping.get(ext, "application/octet-stream")
- def guess_category(local_path, local_dir):
- """根据目录结构猜测文件分类"""
- relative = local_path.relative_to(local_dir)
- parts = relative.parts
- if len(parts) >= 2:
- return parts[0]
- return "general"
- def backfill_file_record(db, local_path, url, local_dir):
- """回填 file_record 表"""
- category = guess_category(local_path, local_dir)
- file_size = local_path.stat().st_size
- content_type = guess_content_type(local_path)
- original_name = local_path.name
- # 计算 SHA-256
- sha256_hash = hashlib.sha256()
- with open(local_path, "rb") as f:
- sha256_hash.update(f.read())
- file_hash = sha256_hash.hexdigest()
- now = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
- with db.cursor() as cursor:
- # 检查是否已存在
- cursor.execute("SELECT id FROM file_record WHERE hash = %s", (file_hash,))
- if cursor.fetchone():
- log.debug(" [SKIP-BACKFILL] hash 已存在: %s", file_hash)
- return False
- sql = (
- "INSERT INTO file_record (hash, url, file_size, content_type, "
- "original_filename, category, created_at, updated_at) "
- "VALUES (%s, %s, %s, %s, %s, %s, %s, %s)"
- )
- cursor.execute(sql, (
- file_hash, url, file_size, content_type,
- original_name, category, now, now,
- ))
- db.commit()
- log.info(" [BACKFILL] %s -> %s", file_hash, url)
- return True
- def main():
- parser = argparse.ArgumentParser(description="迁移本地文件到京东云对象存储")
- parser.add_argument("--local-dir", default="/data/cfc-uploads", help="本地文件目录")
- parser.add_argument("--dry-run", action="store_true", help="仅预览,不实际上传")
- parser.add_argument("--overwrite", action="store_true", help="覆盖云端已存在的文件")
- parser.add_argument("--backfill", action="store_true", help="上传后回填 file_record 表(用于去重)")
- args = parser.parse_args()
- local_dir = Path(args.local_dir)
- if not local_dir.is_dir():
- log.error("本地目录不存在: %s", local_dir)
- sys.exit(1)
- # 收集所有本地文件
- all_files = sorted(local_dir.rglob("*"))
- files = [f for f in all_files if f.is_file()]
- if not files:
- log.info("本地目录为空,无需迁移")
- return
- log.info("本地目录: %s", local_dir)
- log.info("目标 Bucket: %s", BUCKET)
- log.info("目标 Endpoint: %s", ENDPOINT)
- log.info("共发现 %d 个文件", len(files))
- if args.dry_run:
- log.info("模式: DRY-RUN(仅预览不上传)")
- elif args.overwrite:
- log.info("模式: 覆盖已存在文件")
- else:
- log.info("模式: 跳过已存在文件")
- if args.backfill:
- log.info("回填: 启用(写入 file_record 表)")
- # 连接京东云
- if not args.dry_run:
- log.info("正在连接京东云 OSS...")
- client = get_client()
- try:
- client.head_bucket(Bucket=BUCKET)
- log.info("Bucket %s 连接成功", BUCKET)
- except Exception as e:
- log.error("Bucket 连接失败: %s", e)
- sys.exit(1)
- cloud_keys = list_cloud_keys(client) if not args.overwrite else set()
- log.info("云端已有 %d 个文件", len(cloud_keys))
- else:
- client = None
- cloud_keys = set()
- # 连接 MySQL(回填模式)
- db = None
- if args.backfill and not args.dry_run:
- try:
- db = get_db()
- log.info("MySQL 连接成功")
- except Exception as e:
- log.error("MySQL 连接失败: %s", e)
- sys.exit(1)
- # 上传文件
- success = 0
- skipped = 0
- failed = 0
- backfill_ok = 0
- backfill_skip = 0
- backfill_fail = 0
- for file_path in files:
- relative = file_path.relative_to(local_dir)
- key = "uploads/" + relative.as_posix()
- if client and key in cloud_keys and not args.overwrite:
- skipped += 1
- # 即使跳过上传,也回填 file_record(如果启用)
- if db:
- try:
- if backfill_file_record(db, file_path, PUBLIC_URL_PREFIX + key, local_dir):
- backfill_ok += 1
- else:
- backfill_skip += 1
- except Exception as e:
- log.error(" [BACKFILL-FAIL] %s: %s", file_path, e)
- backfill_fail += 1
- continue
- ok, url = upload_file(client, file_path, key, args.dry_run, args.overwrite, cloud_keys)
- if ok:
- success += 1
- if db:
- try:
- if backfill_file_record(db, file_path, url, local_dir):
- backfill_ok += 1
- else:
- backfill_skip += 1
- except Exception as e:
- log.error(" [BACKFILL-FAIL] %s: %s", file_path, e)
- backfill_fail += 1
- else:
- if not args.dry_run:
- failed += 1
- # 总结
- if not args.dry_run:
- log.info("=" * 50)
- parts = [f"成功={success}", f"跳过={skipped}", f"失败={failed}"]
- if db:
- parts.append(f"回填={backfill_ok}")
- parts.append(f"回填跳过={backfill_skip}")
- if backfill_fail:
- parts.append(f"回填失败={backfill_fail}")
- log.info("迁移完成: " + ", ".join(parts))
- if failed > 0:
- log.warning("有 %d 个文件上传失败,请检查日志", failed)
- if db:
- db.close()
- else:
- log.info("=" * 50)
- log.info("DRY-RUN 完成: 将上传 %d 个文件,跳过 %d 个(已存在)", success, skipped)
- if __name__ == "__main__":
- main()
|