重构项目结构
This commit is contained in:
@@ -0,0 +1,44 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
数据库迁移脚本:将 actual_delivery_date 和 actual_payment_date 字段从 DATE 改为 TIMESTAMP
|
||||
"""
|
||||
import asyncio
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
project_root = Path(__file__).parent
|
||||
sys.path.insert(0, str(project_root))
|
||||
sys.path.insert(0, str(project_root / "src"))
|
||||
|
||||
from sqlalchemy.ext.asyncio import create_async_engine
|
||||
from sqlalchemy import text
|
||||
from config.settings import settings
|
||||
|
||||
|
||||
async def migrate():
|
||||
"""执行数据库迁移"""
|
||||
if not settings.DATABASE_URL:
|
||||
print("错误:未配置数据库连接")
|
||||
return
|
||||
|
||||
engine = create_async_engine(
|
||||
settings.DATABASE_URL,
|
||||
echo=True
|
||||
)
|
||||
|
||||
async with engine.begin() as conn:
|
||||
# 修改日期字段类型
|
||||
alter_sql = """
|
||||
ALTER TABLE sales_orders
|
||||
ALTER COLUMN actual_delivery_date TYPE TIMESTAMP,
|
||||
ALTER COLUMN actual_payment_date TYPE TIMESTAMP
|
||||
"""
|
||||
await conn.execute(text(alter_sql))
|
||||
print("成功将实际交付和实际收款日期字段类型改为 TIMESTAMP")
|
||||
|
||||
await engine.dispose()
|
||||
print("迁移完成")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(migrate())
|
||||
@@ -0,0 +1,46 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
数据库迁移脚本:将销售订单中的日期字段从 TIMESTAMP 改为 DATE 类型
|
||||
"""
|
||||
import asyncio
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
project_root = Path(__file__).parent
|
||||
sys.path.insert(0, str(project_root))
|
||||
sys.path.insert(0, str(project_root / "src"))
|
||||
|
||||
from sqlalchemy.ext.asyncio import create_async_engine
|
||||
from sqlalchemy import text
|
||||
from config.settings import settings
|
||||
|
||||
|
||||
async def migrate():
|
||||
"""执行数据库迁移"""
|
||||
if not settings.DATABASE_URL:
|
||||
print("错误:未配置数据库连接")
|
||||
return
|
||||
|
||||
engine = create_async_engine(
|
||||
settings.DATABASE_URL,
|
||||
echo=True
|
||||
)
|
||||
|
||||
async with engine.begin() as conn:
|
||||
# 修改日期字段类型
|
||||
alter_sql = """
|
||||
ALTER TABLE sales_orders
|
||||
ALTER COLUMN delivery_date TYPE DATE,
|
||||
ALTER COLUMN manufacturing_date TYPE DATE,
|
||||
ALTER COLUMN actual_delivery_date TYPE DATE,
|
||||
ALTER COLUMN actual_payment_date TYPE DATE
|
||||
"""
|
||||
await conn.execute(text(alter_sql))
|
||||
print("成功将销售订单的日期字段类型改为 DATE")
|
||||
|
||||
await engine.dispose()
|
||||
print("迁移完成")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(migrate())
|
||||
@@ -0,0 +1,43 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
数据库迁移脚本:将采购订单的预计到货日期字段从 TIMESTAMP 改为 DATE 类型
|
||||
"""
|
||||
import asyncio
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
project_root = Path(__file__).parent
|
||||
sys.path.insert(0, str(project_root))
|
||||
sys.path.insert(0, str(project_root / "src"))
|
||||
|
||||
from sqlalchemy.ext.asyncio import create_async_engine
|
||||
from sqlalchemy import text
|
||||
from config.settings import settings
|
||||
|
||||
|
||||
async def migrate():
|
||||
"""执行数据库迁移"""
|
||||
if not settings.DATABASE_URL:
|
||||
print("错误:未配置数据库连接")
|
||||
return
|
||||
|
||||
engine = create_async_engine(
|
||||
settings.DATABASE_URL,
|
||||
echo=True
|
||||
)
|
||||
|
||||
async with engine.begin() as conn:
|
||||
# 修改日期字段类型
|
||||
alter_sql = """
|
||||
ALTER TABLE purchase_orders
|
||||
ALTER COLUMN expected_date TYPE DATE
|
||||
"""
|
||||
await conn.execute(text(alter_sql))
|
||||
print("成功将采购订单的预计到货日期字段类型改为 DATE")
|
||||
|
||||
await engine.dispose()
|
||||
print("迁移完成")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(migrate())
|
||||
@@ -0,0 +1,75 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
数据库迁移脚本:为采购订单表添加状态变更时间字段
|
||||
- received_date: 实际到货时间
|
||||
- paid_date: 实际付款时间
|
||||
"""
|
||||
import asyncio
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
project_root = Path(__file__).parent
|
||||
sys.path.insert(0, str(project_root))
|
||||
sys.path.insert(0, str(project_root / "src"))
|
||||
|
||||
from sqlalchemy.ext.asyncio import create_async_engine
|
||||
from sqlalchemy import text
|
||||
from config.settings import settings
|
||||
|
||||
|
||||
async def migrate():
|
||||
"""执行数据库迁移"""
|
||||
if not settings.DATABASE_URL:
|
||||
print("错误:未配置数据库连接")
|
||||
return
|
||||
|
||||
engine = create_async_engine(
|
||||
settings.DATABASE_URL,
|
||||
echo=True
|
||||
)
|
||||
|
||||
async with engine.begin() as conn:
|
||||
# 检查字段是否已存在
|
||||
check_sql = """
|
||||
SELECT column_name
|
||||
FROM information_schema.columns
|
||||
WHERE table_name = 'purchase_orders'
|
||||
AND column_name = 'received_date'
|
||||
"""
|
||||
result = await conn.execute(text(check_sql))
|
||||
if result.fetchone():
|
||||
print("字段 received_date 已存在,跳过")
|
||||
else:
|
||||
# 添加实际到货时间字段
|
||||
alter_sql = """
|
||||
ALTER TABLE purchase_orders
|
||||
ADD COLUMN received_date TIMESTAMP WITHOUT TIME ZONE
|
||||
"""
|
||||
await conn.execute(text(alter_sql))
|
||||
print("成功添加字段 received_date")
|
||||
|
||||
# 检查 paid_date 字段是否已存在
|
||||
check_sql2 = """
|
||||
SELECT column_name
|
||||
FROM information_schema.columns
|
||||
WHERE table_name = 'purchase_orders'
|
||||
AND column_name = 'paid_date'
|
||||
"""
|
||||
result2 = await conn.execute(text(check_sql2))
|
||||
if result2.fetchone():
|
||||
print("字段 paid_date 已存在,跳过")
|
||||
else:
|
||||
# 添加实际付款时间字段
|
||||
alter_sql2 = """
|
||||
ALTER TABLE purchase_orders
|
||||
ADD COLUMN paid_date TIMESTAMP WITHOUT TIME ZONE
|
||||
"""
|
||||
await conn.execute(text(alter_sql2))
|
||||
print("成功添加字段 paid_date")
|
||||
|
||||
await engine.dispose()
|
||||
print("迁移完成")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(migrate())
|
||||
@@ -0,0 +1,378 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
RustFS 存储结构迁移脚本
|
||||
|
||||
将旧的多桶结构迁移到新的单桶结构:
|
||||
旧结构: moldinsight-geometry, moldinsight-stp-files, moldinsight-mold-cavities, moldinsight-html-files, moldinsight-user-files
|
||||
新结构: moldinsight-storage/
|
||||
├── moldinsight/stp-files/{uuid}.stp
|
||||
├── moldinsight/geometry/{file_hash}.json
|
||||
├── moldinsight/mold-cavities/{file_hash}.json
|
||||
├── moldinsight/html/{file_hash}.json
|
||||
└── moldinsight/user-files/{uuid}.{ext}
|
||||
|
||||
注意:这是rustFS,而不是minio,只是用了minio的通用S3接口
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import sys
|
||||
import os
|
||||
from pathlib import Path
|
||||
from typing import Dict, List, Optional
|
||||
import json
|
||||
from datetime import datetime
|
||||
|
||||
# 添加项目根目录和 src 目录到 Python 路径
|
||||
project_root = Path(__file__).parent
|
||||
src_root = project_root / "src"
|
||||
sys.path.insert(0, str(project_root))
|
||||
sys.path.insert(0, str(src_root))
|
||||
|
||||
from storage.rustfs_storage import RustFSManager
|
||||
from config.settings import settings
|
||||
from utils.logger import get_logger
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
|
||||
class RustFSMigration:
|
||||
"""RustFS 存储结构迁移器"""
|
||||
|
||||
def __init__(self):
|
||||
self.rustfs = RustFSManager()
|
||||
self.is_connected = False
|
||||
|
||||
# 旧桶名称映射
|
||||
self.old_buckets = {
|
||||
'stp_files': 'moldinsight-stp-files',
|
||||
'geometry_data': 'moldinsight-geometry',
|
||||
'mold_cavities': 'moldinsight-mold-cavities',
|
||||
'html_files': 'moldinsight-html-files',
|
||||
'user_files': 'moldinsight-user-files'
|
||||
}
|
||||
|
||||
# 新桶名称
|
||||
self.new_bucket = 'moldinsight'
|
||||
|
||||
# 文件类型前缀映射
|
||||
self.file_type_mapping = {
|
||||
'stp_files': 'stp-files',
|
||||
'geometry_data': 'geometry',
|
||||
'mold_cavities': 'mold-cavities',
|
||||
'html_files': 'html',
|
||||
'user_files': 'user-files'
|
||||
}
|
||||
|
||||
async def connect(self):
|
||||
"""连接到 RustFS"""
|
||||
try:
|
||||
await self.rustfs.connect(
|
||||
endpoint=settings.RUSTFS_ENDPOINT,
|
||||
access_key=settings.RUSTFS_ACCESS_KEY,
|
||||
secret_key=settings.RUSTFS_SECRET_KEY,
|
||||
timeout=settings.RUSTFS_TIMEOUT
|
||||
)
|
||||
self.is_connected = True
|
||||
logger.info("RustFS 连接成功")
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error(f"RustFS 连接失败: {e}")
|
||||
return False
|
||||
|
||||
async def close(self):
|
||||
"""关闭连接"""
|
||||
await self.rustfs.close()
|
||||
self.is_connected = False
|
||||
logger.info("RustFS 连接已关闭")
|
||||
|
||||
async def list_old_buckets(self) -> Dict[str, List[Dict]]:
|
||||
"""列出所有旧桶及其文件"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
buckets_info = {}
|
||||
|
||||
for file_type, bucket_name in self.old_buckets.items():
|
||||
try:
|
||||
# 检查桶是否存在
|
||||
if not self.rustfs.client.bucket_exists(bucket_name):
|
||||
logger.info(f"桶不存在: {bucket_name}")
|
||||
buckets_info[bucket_name] = []
|
||||
continue
|
||||
|
||||
# 列出桶中所有文件
|
||||
objects = self.rustfs.client.list_objects(bucket_name, recursive=True)
|
||||
files = []
|
||||
|
||||
for obj in objects:
|
||||
files.append({
|
||||
'object_key': obj.object_name,
|
||||
'size': obj.size,
|
||||
'last_modified': obj.last_modified,
|
||||
'etag': obj.etag
|
||||
})
|
||||
|
||||
buckets_info[bucket_name] = files
|
||||
logger.info(f"桶 {bucket_name} 包含 {len(files)} 个文件")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"列出桶 {bucket_name} 失败: {e}")
|
||||
buckets_info[bucket_name] = []
|
||||
|
||||
return buckets_info
|
||||
|
||||
async def ensure_new_bucket(self):
|
||||
"""确保新桶存在"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
try:
|
||||
if not self.rustfs.client.bucket_exists(self.new_bucket):
|
||||
self.rustfs.client.make_bucket(self.new_bucket)
|
||||
logger.info(f"创建新桶: {self.new_bucket}")
|
||||
else:
|
||||
logger.info(f"新桶已存在: {self.new_bucket}")
|
||||
except Exception as e:
|
||||
logger.error(f"确保新桶存在失败: {e}")
|
||||
raise
|
||||
|
||||
async def migrate_file(self, old_bucket: str, old_object_key: str, file_type: str) -> bool:
|
||||
"""迁移单个文件到新结构"""
|
||||
try:
|
||||
# 下载旧文件
|
||||
response = self.rustfs.client.get_object(old_bucket, old_object_key)
|
||||
file_data = response.read()
|
||||
response.close()
|
||||
response.release_conn()
|
||||
|
||||
# 生成新对象键
|
||||
if file_type in ['stp_files', 'user_files']:
|
||||
# STP文件和用户文件:使用UUID格式
|
||||
import uuid
|
||||
unique_id = str(uuid.uuid4())
|
||||
ext = Path(old_object_key).suffix or ('.stp' if file_type == 'stp_files' else '')
|
||||
new_object_key = f"moldinsight/{self.file_type_mapping[file_type]}/{unique_id}{ext}"
|
||||
else:
|
||||
# JSON数据文件:使用文件哈希格式
|
||||
import hashlib
|
||||
file_hash = hashlib.sha256(file_data).hexdigest()
|
||||
new_object_key = f"moldinsight/{self.file_type_mapping[file_type]}/{file_hash}.json"
|
||||
|
||||
# 上传到新桶
|
||||
self.rustfs.client.put_object(
|
||||
self.new_bucket,
|
||||
new_object_key,
|
||||
data=file_data,
|
||||
length=len(file_data),
|
||||
content_type='application/octet-stream'
|
||||
)
|
||||
|
||||
logger.info(f"文件迁移成功: {old_bucket}/{old_object_key} -> {self.new_bucket}/{new_object_key}")
|
||||
return True
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"文件迁移失败 {old_bucket}/{old_object_key}: {e}")
|
||||
return False
|
||||
|
||||
async def migrate_bucket(self, old_bucket: str, file_type: str) -> Dict[str, any]:
|
||||
"""迁移整个桶"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
migration_result = {
|
||||
'total_files': 0,
|
||||
'successful': 0,
|
||||
'failed': 0,
|
||||
'failed_files': []
|
||||
}
|
||||
|
||||
try:
|
||||
# 检查桶是否存在
|
||||
if not self.rustfs.client.bucket_exists(old_bucket):
|
||||
logger.info(f"桶不存在,跳过迁移: {old_bucket}")
|
||||
return migration_result
|
||||
|
||||
# 列出桶中所有文件
|
||||
objects = self.rustfs.client.list_objects(old_bucket, recursive=True)
|
||||
files = list(objects)
|
||||
migration_result['total_files'] = len(files)
|
||||
|
||||
logger.info(f"开始迁移桶 {old_bucket}, 包含 {len(files)} 个文件")
|
||||
|
||||
for obj in files:
|
||||
success = await self.migrate_file(old_bucket, obj.object_name, file_type)
|
||||
if success:
|
||||
migration_result['successful'] += 1
|
||||
else:
|
||||
migration_result['failed'] += 1
|
||||
migration_result['failed_files'].append(obj.object_name)
|
||||
|
||||
logger.info(f"桶 {old_bucket} 迁移完成: 成功 {migration_result['successful']}, 失败 {migration_result['failed']}")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"桶 {old_bucket} 迁移失败: {e}")
|
||||
|
||||
return migration_result
|
||||
|
||||
async def delete_old_buckets(self) -> Dict[str, bool]:
|
||||
"""删除所有旧桶(可选操作)"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
deletion_results = {}
|
||||
|
||||
for file_type, bucket_name in self.old_buckets.items():
|
||||
try:
|
||||
# 检查桶是否存在
|
||||
if not self.rustfs.client.bucket_exists(bucket_name):
|
||||
logger.info(f"桶不存在,跳过删除: {bucket_name}")
|
||||
deletion_results[bucket_name] = True
|
||||
continue
|
||||
|
||||
# 删除桶中所有文件
|
||||
objects = self.rustfs.client.list_objects(bucket_name, recursive=True)
|
||||
for obj in objects:
|
||||
self.rustfs.client.remove_object(bucket_name, obj.object_name)
|
||||
|
||||
# 删除空桶
|
||||
self.rustfs.client.remove_bucket(bucket_name)
|
||||
|
||||
deletion_results[bucket_name] = True
|
||||
logger.info(f"桶删除成功: {bucket_name}")
|
||||
|
||||
except Exception as e:
|
||||
deletion_results[bucket_name] = False
|
||||
logger.error(f"桶删除失败 {bucket_name}: {e}")
|
||||
|
||||
return deletion_results
|
||||
|
||||
async def run_migration(self, delete_old_buckets: bool = False) -> Dict[str, any]:
|
||||
"""运行完整迁移流程"""
|
||||
migration_summary = {
|
||||
'start_time': datetime.now().isoformat(),
|
||||
'connection_status': False,
|
||||
'old_buckets_info': {},
|
||||
'migration_results': {},
|
||||
'deletion_results': {},
|
||||
'end_time': None,
|
||||
'status': 'failed'
|
||||
}
|
||||
|
||||
try:
|
||||
# 1. 连接
|
||||
logger.info("=== RustFS 存储结构迁移开始 ===")
|
||||
connection_result = await self.connect()
|
||||
if not connection_result:
|
||||
raise RuntimeError("无法连接到 RustFS")
|
||||
|
||||
migration_summary['connection_status'] = True
|
||||
|
||||
# 2. 列出旧桶信息
|
||||
logger.info("1. 检查旧桶结构...")
|
||||
old_buckets_info = await self.list_old_buckets()
|
||||
migration_summary['old_buckets_info'] = old_buckets_info
|
||||
|
||||
# 3. 确保新桶存在
|
||||
logger.info("2. 确保新桶存在...")
|
||||
await self.ensure_new_bucket()
|
||||
|
||||
# 4. 执行迁移
|
||||
logger.info("3. 开始迁移文件...")
|
||||
migration_results = {}
|
||||
|
||||
for file_type, bucket_name in self.old_buckets.items():
|
||||
if old_buckets_info.get(bucket_name):
|
||||
logger.info(f"迁移桶: {bucket_name}")
|
||||
result = await self.migrate_bucket(bucket_name, file_type)
|
||||
migration_results[bucket_name] = result
|
||||
else:
|
||||
logger.info(f"跳过空桶: {bucket_name}")
|
||||
|
||||
migration_summary['migration_results'] = migration_results
|
||||
|
||||
# 5. 可选:删除旧桶
|
||||
if delete_old_buckets:
|
||||
logger.info("4. 删除旧桶...")
|
||||
deletion_results = await self.delete_old_buckets()
|
||||
migration_summary['deletion_results'] = deletion_results
|
||||
else:
|
||||
logger.info("4. 保留旧桶(跳过删除)")
|
||||
|
||||
# 6. 完成
|
||||
migration_summary['status'] = 'completed'
|
||||
migration_summary['end_time'] = datetime.now().isoformat()
|
||||
|
||||
logger.info("=== RustFS 存储结构迁移完成 ===")
|
||||
|
||||
return migration_summary
|
||||
|
||||
except Exception as e:
|
||||
migration_summary['error'] = str(e)
|
||||
migration_summary['end_time'] = datetime.now().isoformat()
|
||||
logger.error(f"迁移失败: {e}")
|
||||
return migration_summary
|
||||
|
||||
finally:
|
||||
await self.close()
|
||||
|
||||
|
||||
async def main():
|
||||
"""主函数"""
|
||||
migration = RustFSMigration()
|
||||
|
||||
print("=== RustFS 存储结构迁移工具 ===")
|
||||
print("注意:这是rustFS,而不是minio,只是用了minio的通用S3接口")
|
||||
print()
|
||||
|
||||
# 询问是否删除旧桶
|
||||
delete_old = input("是否在迁移完成后删除旧桶?(y/N): ").strip().lower() == 'y'
|
||||
|
||||
print("\n开始迁移...")
|
||||
|
||||
# 运行迁移
|
||||
result = await migration.run_migration(delete_old_buckets=delete_old)
|
||||
|
||||
# 输出结果摘要
|
||||
print("\n=== 迁移结果摘要 ===")
|
||||
print(f"状态: {result['status']}")
|
||||
print(f"开始时间: {result['start_time']}")
|
||||
print(f"结束时间: {result['end_time']}")
|
||||
|
||||
if 'error' in result:
|
||||
print(f"错误: {result['error']}")
|
||||
|
||||
# 旧桶信息
|
||||
print("\n--- 旧桶信息 ---")
|
||||
for bucket_name, files in result['old_buckets_info'].items():
|
||||
print(f"{bucket_name}: {len(files)} 个文件")
|
||||
|
||||
# 迁移结果
|
||||
print("\n--- 迁移结果 ---")
|
||||
total_files = 0
|
||||
total_success = 0
|
||||
total_failed = 0
|
||||
|
||||
for bucket_name, migration_result in result['migration_results'].items():
|
||||
print(f"{bucket_name}:")
|
||||
print(f" 总文件数: {migration_result['total_files']}")
|
||||
print(f" 成功: {migration_result['successful']}")
|
||||
print(f" 失败: {migration_result['failed']}")
|
||||
|
||||
total_files += migration_result['total_files']
|
||||
total_success += migration_result['successful']
|
||||
total_failed += migration_result['failed']
|
||||
|
||||
print(f"\n总计: {total_files} 个文件, 成功 {total_success}, 失败 {total_failed}")
|
||||
|
||||
# 删除结果(如果执行了删除)
|
||||
if result['deletion_results']:
|
||||
print("\n--- 旧桶删除结果 ---")
|
||||
for bucket_name, success in result['deletion_results'].items():
|
||||
status = "成功" if success else "失败"
|
||||
print(f"{bucket_name}: {status}")
|
||||
|
||||
print("\n=== 迁移完成 ===")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(main())
|
||||
@@ -0,0 +1,56 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
数据库迁移脚本:为 sales_orders 表添加时间字段
|
||||
"""
|
||||
import asyncio
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
project_root = Path(__file__).parent
|
||||
sys.path.insert(0, str(project_root))
|
||||
sys.path.insert(0, str(project_root / "src"))
|
||||
|
||||
from sqlalchemy.ext.asyncio import create_async_engine
|
||||
from sqlalchemy import text
|
||||
from config.settings import settings
|
||||
|
||||
|
||||
async def migrate():
|
||||
"""执行数据库迁移"""
|
||||
if not settings.DATABASE_URL:
|
||||
print("错误:未配置数据库连接")
|
||||
return
|
||||
|
||||
engine = create_async_engine(
|
||||
settings.DATABASE_URL,
|
||||
echo=True
|
||||
)
|
||||
|
||||
async with engine.begin() as conn:
|
||||
# 检查列是否存在
|
||||
check_sql = """
|
||||
SELECT column_name
|
||||
FROM information_schema.columns
|
||||
WHERE table_name = 'sales_orders'
|
||||
AND column_name = 'manufacturing_date'
|
||||
"""
|
||||
result = await conn.execute(text(check_sql))
|
||||
if result.scalar():
|
||||
print("列 manufacturing_date 已存在,跳过迁移")
|
||||
else:
|
||||
# 添加新列
|
||||
alter_sql = """
|
||||
ALTER TABLE sales_orders
|
||||
ADD COLUMN manufacturing_date TIMESTAMP,
|
||||
ADD COLUMN actual_delivery_date TIMESTAMP,
|
||||
ADD COLUMN actual_payment_date TIMESTAMP
|
||||
"""
|
||||
await conn.execute(text(alter_sql))
|
||||
print("成功添加时间字段到 sales_orders 表")
|
||||
|
||||
await engine.dispose()
|
||||
print("迁移完成")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(migrate())
|
||||
@@ -0,0 +1,457 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
RustFS 存储结构迁移脚本 - 完整版
|
||||
|
||||
从旧结构迁移到新结构:
|
||||
旧桶: moldinsight-storage
|
||||
新桶: moldinsight/
|
||||
├── stp-files/{uuid}.stp
|
||||
├── geometry/{file_hash}.json
|
||||
├── mold-cavities/{file_hash}.json
|
||||
├── html/{file_hash}.json
|
||||
└── user-files/{uuid}.{ext}
|
||||
|
||||
注意:这是rustFS,而不是minio,只是用了minio的通用S3接口
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import sys
|
||||
import os
|
||||
from pathlib import Path
|
||||
from typing import Dict, List, Optional
|
||||
import json
|
||||
from datetime import datetime
|
||||
import hashlib
|
||||
|
||||
# 添加项目根目录和 src 目录到 Python 路径
|
||||
project_root = Path(__file__).parent
|
||||
src_root = project_root / "src"
|
||||
sys.path.insert(0, str(project_root))
|
||||
sys.path.insert(0, str(src_root))
|
||||
|
||||
from storage.rustfs_storage import RustFSManager
|
||||
from config.settings import settings
|
||||
from utils.logger import get_logger
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
|
||||
class RustFSMigration:
|
||||
"""RustFS 存储结构迁移器"""
|
||||
|
||||
def __init__(self):
|
||||
self.rustfs = RustFSManager()
|
||||
self.is_connected = False
|
||||
|
||||
# 旧桶名称
|
||||
self.old_bucket = 'moldinsight-storage'
|
||||
|
||||
# 新桶名称
|
||||
self.new_bucket = 'moldinsight'
|
||||
|
||||
# 文件类型前缀映射
|
||||
self.file_type_mapping = {
|
||||
'stp-files': 'stp_files',
|
||||
'geometry': 'geometry_data',
|
||||
'mold-cavities': 'mold_cavities',
|
||||
'html': 'html_files',
|
||||
'user-files': 'user_files'
|
||||
}
|
||||
|
||||
async def connect(self):
|
||||
"""连接到 RustFS"""
|
||||
try:
|
||||
await self.rustfs.connect(
|
||||
endpoint=settings.RUSTFS_ENDPOINT,
|
||||
access_key=settings.RUSTFS_ACCESS_KEY,
|
||||
secret_key=settings.RUSTFS_SECRET_KEY,
|
||||
timeout=settings.RUSTFS_TIMEOUT
|
||||
)
|
||||
self.is_connected = True
|
||||
logger.info("RustFS 连接成功")
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error(f"RustFS 连接失败: {e}")
|
||||
return False
|
||||
|
||||
async def close(self):
|
||||
"""关闭连接"""
|
||||
await self.rustfs.close()
|
||||
self.is_connected = False
|
||||
logger.info("RustFS 连接已关闭")
|
||||
|
||||
async def list_all_buckets(self) -> List[str]:
|
||||
"""列出所有桶"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
try:
|
||||
buckets = self.rustfs.client.list_buckets()
|
||||
bucket_names = [bucket.name for bucket in buckets]
|
||||
logger.info(f"当前存在的桶: {bucket_names}")
|
||||
return bucket_names
|
||||
except Exception as e:
|
||||
logger.error(f"列出桶失败: {e}")
|
||||
return []
|
||||
|
||||
async def check_old_bucket_files(self) -> List[Dict]:
|
||||
"""检查旧桶中的所有文件"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
try:
|
||||
# 检查旧桶是否存在
|
||||
if not self.rustfs.client.bucket_exists(self.old_bucket):
|
||||
logger.info(f"旧桶不存在: {self.old_bucket}")
|
||||
return []
|
||||
|
||||
# 列出所有文件
|
||||
objects = self.rustfs.client.list_objects(self.old_bucket, recursive=True)
|
||||
files = []
|
||||
|
||||
for obj in objects:
|
||||
files.append({
|
||||
'object_key': obj.object_name,
|
||||
'size': obj.size,
|
||||
'last_modified': obj.last_modified,
|
||||
'etag': obj.etag
|
||||
})
|
||||
|
||||
logger.info(f"旧桶 {self.old_bucket} 包含 {len(files)} 个文件")
|
||||
return files
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"检查旧桶文件失败: {e}")
|
||||
return []
|
||||
|
||||
async def ensure_new_bucket(self):
|
||||
"""确保新桶存在"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
try:
|
||||
if not self.rustfs.client.bucket_exists(self.new_bucket):
|
||||
self.rustfs.client.make_bucket(self.new_bucket)
|
||||
logger.info(f"创建新桶: {self.new_bucket}")
|
||||
else:
|
||||
logger.info(f"新桶已存在: {self.new_bucket}")
|
||||
except Exception as e:
|
||||
logger.error(f"确保新桶存在失败: {e}")
|
||||
raise
|
||||
|
||||
async def migrate_file(self, old_object_key: str) -> bool:
|
||||
"""迁移单个文件到新结构"""
|
||||
try:
|
||||
# 下载旧文件
|
||||
response = self.rustfs.client.get_object(self.old_bucket, old_object_key)
|
||||
file_data = response.read()
|
||||
response.close()
|
||||
response.release_conn()
|
||||
|
||||
# 解析旧对象键,确定文件类型
|
||||
# 旧格式: moldinsight/{文件类型}/{文件名} 或 {文件类型}/{文件名}
|
||||
parts = old_object_key.split('/')
|
||||
|
||||
# 确定文件类型和新对象键
|
||||
if len(parts) >= 2:
|
||||
# 可能是 moldinsight/{类型}/{文件} 或 {类型}/{文件}
|
||||
if parts[0] == 'moldinsight' and len(parts) >= 3:
|
||||
# moldinsight/{类型}/{文件}
|
||||
old_type = parts[1]
|
||||
filename = parts[2]
|
||||
elif parts[0] in self.file_type_mapping:
|
||||
# {类型}/{文件}
|
||||
old_type = parts[0]
|
||||
filename = parts[1]
|
||||
else:
|
||||
# 无法识别的格式,使用默认
|
||||
old_type = 'misc'
|
||||
filename = parts[-1]
|
||||
else:
|
||||
old_type = 'misc'
|
||||
filename = parts[-1]
|
||||
|
||||
# 根据旧类型确定新类型
|
||||
new_type = old_type # 默认保持不变
|
||||
|
||||
# 生成新对象键
|
||||
if old_type in self.file_type_mapping:
|
||||
# 直接使用类型名作为目录
|
||||
new_object_key = f"{old_type}/{filename}"
|
||||
else:
|
||||
# 其他文件放到misc目录
|
||||
new_object_key = f"misc/{filename}"
|
||||
|
||||
# 上传到新桶
|
||||
self.rustfs.client.put_object(
|
||||
self.new_bucket,
|
||||
new_object_key,
|
||||
data=file_data,
|
||||
length=len(file_data),
|
||||
content_type='application/octet-stream'
|
||||
)
|
||||
|
||||
logger.info(f"文件迁移成功: {self.old_bucket}/{old_object_key} -> {self.new_bucket}/{new_object_key}")
|
||||
return True
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"文件迁移失败 {old_object_key}: {e}")
|
||||
return False
|
||||
|
||||
async def migrate_all_files(self) -> Dict[str, any]:
|
||||
"""迁移所有文件"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
migration_result = {
|
||||
'total_files': 0,
|
||||
'successful': 0,
|
||||
'failed': 0,
|
||||
'failed_files': []
|
||||
}
|
||||
|
||||
try:
|
||||
# 获取所有文件
|
||||
files = await self.check_old_bucket_files()
|
||||
migration_result['total_files'] = len(files)
|
||||
|
||||
if len(files) == 0:
|
||||
logger.info("旧桶中没有文件需要迁移")
|
||||
return migration_result
|
||||
|
||||
logger.info(f"开始迁移 {len(files)} 个文件...")
|
||||
|
||||
for file_info in files:
|
||||
old_object_key = file_info['object_key']
|
||||
success = await self.migrate_file(old_object_key)
|
||||
|
||||
if success:
|
||||
migration_result['successful'] += 1
|
||||
else:
|
||||
migration_result['failed'] += 1
|
||||
migration_result['failed_files'].append(old_object_key)
|
||||
|
||||
logger.info(f"迁移完成: 成功 {migration_result['successful']}, 失败 {migration_result['failed']}")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"文件迁移失败: {e}")
|
||||
|
||||
return migration_result
|
||||
|
||||
async def delete_old_bucket(self) -> bool:
|
||||
"""删除旧桶及其所有文件"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
try:
|
||||
# 检查桶是否存在
|
||||
if not self.rustfs.client.bucket_exists(self.old_bucket):
|
||||
logger.info(f"旧桶不存在,无需删除: {self.old_bucket}")
|
||||
return True
|
||||
|
||||
# 列出所有文件
|
||||
objects = list(self.rustfs.client.list_objects(self.old_bucket, recursive=True))
|
||||
|
||||
if len(objects) > 0:
|
||||
logger.info(f"删除旧桶中的 {len(objects)} 个文件...")
|
||||
|
||||
# 删除所有文件
|
||||
for obj in objects:
|
||||
self.rustfs.client.remove_object(self.old_bucket, obj.object_name)
|
||||
logger.debug(f"删除文件: {obj.object_name}")
|
||||
|
||||
# 删除空桶
|
||||
self.rustfs.client.remove_bucket(self.old_bucket)
|
||||
|
||||
logger.info(f"旧桶删除成功: {self.old_bucket}")
|
||||
return True
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"删除旧桶失败 {self.old_bucket}: {e}")
|
||||
return False
|
||||
|
||||
async def list_new_bucket_structure(self) -> Dict[str, List[str]]:
|
||||
"""列出新桶的文件结构"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
try:
|
||||
if not self.rustfs.client.bucket_exists(self.new_bucket):
|
||||
return {}
|
||||
|
||||
objects = self.rustfs.client.list_objects(self.new_bucket, recursive=True)
|
||||
structure = {
|
||||
'stp-files': [],
|
||||
'geometry': [],
|
||||
'mold-cavities': [],
|
||||
'html': [],
|
||||
'user-files': [],
|
||||
'misc': []
|
||||
}
|
||||
|
||||
for obj in objects:
|
||||
parts = obj.object_name.split('/')
|
||||
if len(parts) >= 2:
|
||||
file_type = parts[0]
|
||||
if file_type in structure:
|
||||
structure[file_type].append(obj.object_name)
|
||||
else:
|
||||
structure['misc'].append(obj.object_name)
|
||||
|
||||
return structure
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"列出新桶结构失败: {e}")
|
||||
return {}
|
||||
|
||||
async def run_migration(self, delete_old_bucket: bool = False) -> Dict[str, any]:
|
||||
"""运行完整迁移流程"""
|
||||
migration_summary = {
|
||||
'start_time': datetime.now().isoformat(),
|
||||
'connection_status': False,
|
||||
'all_buckets': [],
|
||||
'old_bucket_files': [],
|
||||
'migration_results': {},
|
||||
'new_bucket_structure': {},
|
||||
'deletion_result': False,
|
||||
'end_time': None,
|
||||
'status': 'failed'
|
||||
}
|
||||
|
||||
try:
|
||||
# 1. 连接
|
||||
logger.info("=== RustFS 存储结构迁移开始 ===")
|
||||
connection_result = await self.connect()
|
||||
if not connection_result:
|
||||
raise RuntimeError("无法连接到 RustFS")
|
||||
|
||||
migration_summary['connection_status'] = True
|
||||
|
||||
# 2. 列出所有桶
|
||||
logger.info("1. 检查桶状态...")
|
||||
all_buckets = await self.list_all_buckets()
|
||||
migration_summary['all_buckets'] = all_buckets
|
||||
|
||||
# 3. 检查旧桶文件
|
||||
logger.info("2. 检查旧桶文件...")
|
||||
old_bucket_files = await self.check_old_bucket_files()
|
||||
migration_summary['old_bucket_files'] = old_bucket_files
|
||||
|
||||
if not old_bucket_files:
|
||||
logger.info("旧桶中没有文件,跳过迁移")
|
||||
migration_summary['status'] = 'completed'
|
||||
migration_summary['end_time'] = datetime.now().isoformat()
|
||||
return migration_summary
|
||||
|
||||
# 4. 确保新桶存在
|
||||
logger.info("3. 确保新桶存在...")
|
||||
await self.ensure_new_bucket()
|
||||
|
||||
# 5. 执行迁移
|
||||
logger.info("4. 开始迁移文件...")
|
||||
migration_results = await self.migrate_all_files()
|
||||
migration_summary['migration_results'] = migration_results
|
||||
|
||||
# 6. 检查新桶结构
|
||||
logger.info("5. 检查新桶结构...")
|
||||
new_bucket_structure = await self.list_new_bucket_structure()
|
||||
migration_summary['new_bucket_structure'] = new_bucket_structure
|
||||
|
||||
# 7. 可选:删除旧桶
|
||||
if delete_old_bucket:
|
||||
logger.info("6. 删除旧桶...")
|
||||
deletion_result = await self.delete_old_bucket()
|
||||
migration_summary['deletion_result'] = deletion_result
|
||||
else:
|
||||
logger.info("6. 保留旧桶(跳过删除)")
|
||||
|
||||
# 8. 完成
|
||||
migration_summary['status'] = 'completed'
|
||||
migration_summary['end_time'] = datetime.now().isoformat()
|
||||
|
||||
logger.info("=== RustFS 存储结构迁移完成 ===")
|
||||
|
||||
return migration_summary
|
||||
|
||||
except Exception as e:
|
||||
migration_summary['error'] = str(e)
|
||||
migration_summary['end_time'] = datetime.now().isoformat()
|
||||
logger.error(f"迁移失败: {e}")
|
||||
return migration_summary
|
||||
|
||||
finally:
|
||||
await self.close()
|
||||
|
||||
|
||||
async def main():
|
||||
"""主函数"""
|
||||
migration = RustFSMigration()
|
||||
|
||||
print("=== RustFS 存储结构迁移工具 ===")
|
||||
print("注意:这是rustFS,而不是minio,只是用了minio的通用S3接口")
|
||||
print()
|
||||
print("从旧结构迁移到新结构:")
|
||||
print(" 旧桶: moldinsight-storage")
|
||||
print(" 新桶: moldinsight/")
|
||||
print(" ├── stp-files/")
|
||||
print(" ├── geometry/")
|
||||
print(" ├── mold-cavities/")
|
||||
print(" ├── html/")
|
||||
print(" └── user-files/")
|
||||
print()
|
||||
|
||||
# 询问是否删除旧桶
|
||||
delete_old = input("是否在迁移完成后删除旧桶 moldinsight-storage?(y/N): ").strip().lower() == 'y'
|
||||
|
||||
print("\n开始迁移...")
|
||||
|
||||
# 运行迁移
|
||||
result = await migration.run_migration(delete_old_bucket=delete_old)
|
||||
|
||||
# 输出结果摘要
|
||||
print("\n=== 迁移结果摘要 ===")
|
||||
print(f"状态: {result['status']}")
|
||||
print(f"开始时间: {result['start_time']}")
|
||||
print(f"结束时间: {result['end_time']}")
|
||||
|
||||
if 'error' in result:
|
||||
print(f"错误: {result['error']}")
|
||||
|
||||
# 桶列表
|
||||
print("\n--- 当前桶列表 ---")
|
||||
for bucket in result['all_buckets']:
|
||||
print(f" - {bucket}")
|
||||
|
||||
# 旧桶文件统计
|
||||
print(f"\n--- 旧桶文件统计 ---")
|
||||
print(f"旧桶 {migration.old_bucket} 包含 {len(result['old_bucket_files'])} 个文件")
|
||||
|
||||
# 迁移结果
|
||||
print("\n--- 迁移结果 ---")
|
||||
migration_result = result['migration_results']
|
||||
print(f"总文件数: {migration_result['total_files']}")
|
||||
print(f"成功: {migration_result['successful']}")
|
||||
print(f"失败: {migration_result['failed']}")
|
||||
|
||||
if migration_result['failed'] > 0:
|
||||
print("\n失败的文件:")
|
||||
for failed_file in migration_result['failed_files']:
|
||||
print(f" - {failed_file}")
|
||||
|
||||
# 新桶结构
|
||||
print("\n--- 新桶文件结构 ---")
|
||||
for file_type, files in result['new_bucket_structure'].items():
|
||||
if files:
|
||||
print(f"{file_type}: {len(files)} 个文件")
|
||||
|
||||
# 删除结果
|
||||
if 'deletion_result' in result:
|
||||
status = "成功" if result['deletion_result'] else "失败"
|
||||
print(f"\n--- 旧桶删除结果 ---")
|
||||
print(f"旧桶 {migration.old_bucket} 删除: {status}")
|
||||
|
||||
print("\n=== 迁移完成 ===")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(main())
|
||||
@@ -0,0 +1,50 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
检查 FastAPI 返回的时间格式
|
||||
"""
|
||||
import asyncio
|
||||
import sys
|
||||
from pathlib import Path
|
||||
import json
|
||||
|
||||
project_root = Path(__file__).parent
|
||||
sys.path.insert(0, str(project_root))
|
||||
sys.path.insert(0, str(project_root / "src"))
|
||||
|
||||
from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession
|
||||
from sqlalchemy import text, select
|
||||
from config.settings import settings
|
||||
from models.database import SalesOrder, Customer
|
||||
|
||||
|
||||
async def check():
|
||||
"""检查 API 返回的时间格式"""
|
||||
if not settings.DATABASE_URL:
|
||||
print("错误:未配置数据库连接")
|
||||
return
|
||||
|
||||
engine = create_async_engine(settings.DATABASE_URL, echo=False)
|
||||
|
||||
async with AsyncSession(engine) as session:
|
||||
# 查询销售订单
|
||||
result = await session.execute(
|
||||
select(SalesOrder, Customer)
|
||||
.join(Customer, SalesOrder.customer_id == Customer.id)
|
||||
.order_by(SalesOrder.id.desc())
|
||||
.limit(1)
|
||||
)
|
||||
row = result.first()
|
||||
if row:
|
||||
order, customer = row
|
||||
print(f"订单号: {order.order_no}")
|
||||
print(f"created_at (Python): {order.created_at}")
|
||||
print(f"created_at (ISO格式): {order.created_at.isoformat() if order.created_at else 'N/A'}")
|
||||
print(f"created_at (带时区): {order.created_at.isoformat() if order.created_at else 'N/A'}")
|
||||
print(f"order_date (Python): {order.order_date}")
|
||||
print(f"order_date (ISO格式): {order.order_date.isoformat() if order.order_date else 'N/A'}")
|
||||
|
||||
await engine.dispose()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(check())
|
||||
@@ -0,0 +1,41 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
检查销售订单的 created_at 字段
|
||||
"""
|
||||
import asyncio
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
project_root = Path(__file__).parent
|
||||
sys.path.insert(0, str(project_root))
|
||||
sys.path.insert(0, str(project_root / "src"))
|
||||
|
||||
from sqlalchemy.ext.asyncio import create_async_engine
|
||||
from sqlalchemy import text
|
||||
from config.settings import settings
|
||||
|
||||
|
||||
async def check():
|
||||
"""检查数据库中的 created_at 字段"""
|
||||
if not settings.DATABASE_URL:
|
||||
print("错误:未配置数据库连接")
|
||||
return
|
||||
|
||||
engine = create_async_engine(
|
||||
settings.DATABASE_URL,
|
||||
echo=True
|
||||
)
|
||||
|
||||
async with engine.begin() as conn:
|
||||
# 查询销售订单的 created_at 字段
|
||||
result = await conn.execute(text("SELECT id, order_no, created_at, order_date FROM sales_orders ORDER BY id DESC LIMIT 5"))
|
||||
rows = result.fetchall()
|
||||
print("\n销售订单数据:")
|
||||
for row in rows:
|
||||
print(f"ID: {row[0]}, 订单号: {row[1]}, created_at: {row[2]}, order_date: {row[3]}")
|
||||
|
||||
await engine.dispose()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(check())
|
||||
@@ -0,0 +1,51 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
检查数据库时区设置
|
||||
"""
|
||||
import asyncio
|
||||
import sys
|
||||
from pathlib import Path
|
||||
from datetime import datetime
|
||||
|
||||
project_root = Path(__file__).parent
|
||||
sys.path.insert(0, str(project_root))
|
||||
sys.path.insert(0, str(project_root / "src"))
|
||||
|
||||
from sqlalchemy.ext.asyncio import create_async_engine
|
||||
from sqlalchemy import text
|
||||
from config.settings import settings
|
||||
|
||||
|
||||
async def check():
|
||||
"""检查数据库时区设置"""
|
||||
if not settings.DATABASE_URL:
|
||||
print("错误:未配置数据库连接")
|
||||
return
|
||||
|
||||
engine = create_async_engine(settings.DATABASE_URL, echo=False)
|
||||
|
||||
async with engine.begin() as conn:
|
||||
# 检查数据库时区设置
|
||||
result = await conn.execute(text("SHOW timezone"))
|
||||
timezone = result.scalar()
|
||||
print(f"数据库时区: {timezone}")
|
||||
|
||||
# 检查数据库当前时间
|
||||
result = await conn.execute(text("SELECT now()"))
|
||||
db_now = result.scalar()
|
||||
print(f"数据库当前时间: {db_now}")
|
||||
|
||||
# 检查 Python 当前时间
|
||||
python_now = datetime.now()
|
||||
print(f"Python 当前时间: {python_now}")
|
||||
|
||||
# 检查数据库当前时间(转换为本地时区)
|
||||
result = await conn.execute(text("SELECT now() AT TIME ZONE 'Asia/Shanghai'"))
|
||||
db_now_local = result.scalar()
|
||||
print(f"数据库当前时间 (Asia/Shanghai): {db_now_local}")
|
||||
|
||||
await engine.dispose()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(check())
|
||||
@@ -0,0 +1,336 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
RustFS 旧桶清理脚本
|
||||
|
||||
删除所有遗留的旧桶结构,包括:
|
||||
- moldinsight-geometry
|
||||
- moldinsight-stp-files
|
||||
- moldinsight-mold-cavities
|
||||
- moldinsight-html
|
||||
- moldinsight-user-files
|
||||
|
||||
注意:这是rustFS,而不是minio,只是用了minio的通用S3接口
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import sys
|
||||
import os
|
||||
from pathlib import Path
|
||||
from typing import List, Dict
|
||||
from datetime import datetime
|
||||
|
||||
# 添加项目根目录和 src 目录到 Python 路径
|
||||
project_root = Path(__file__).parent
|
||||
src_root = project_root / "src"
|
||||
sys.path.insert(0, str(project_root))
|
||||
sys.path.insert(0, str(src_root))
|
||||
|
||||
from storage.rustfs_storage import RustFSManager
|
||||
from config.settings import settings
|
||||
from utils.logger import get_logger
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
|
||||
class RustFSCleanup:
|
||||
"""RustFS 旧桶清理器"""
|
||||
|
||||
def __init__(self):
|
||||
self.rustfs = RustFSManager()
|
||||
self.is_connected = False
|
||||
|
||||
# 所有需要清理的旧桶
|
||||
self.old_buckets_to_clean = [
|
||||
'moldinsight-geometry',
|
||||
'moldinsight-stp-files',
|
||||
'moldinsight-mold-cavities',
|
||||
'moldinsight-html',
|
||||
'moldinsight-user-files'
|
||||
]
|
||||
|
||||
async def connect(self):
|
||||
"""连接到 RustFS"""
|
||||
try:
|
||||
await self.rustfs.connect(
|
||||
endpoint=settings.RUSTFS_ENDPOINT,
|
||||
access_key=settings.RUSTFS_ACCESS_KEY,
|
||||
secret_key=settings.RUSTFS_SECRET_KEY,
|
||||
timeout=settings.RUSTFS_TIMEOUT
|
||||
)
|
||||
self.is_connected = True
|
||||
logger.info("RustFS 连接成功")
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error(f"RustFS 连接失败: {e}")
|
||||
return False
|
||||
|
||||
async def close(self):
|
||||
"""关闭连接"""
|
||||
await self.rustfs.close()
|
||||
self.is_connected = False
|
||||
logger.info("RustFS 连接已关闭")
|
||||
|
||||
async def list_all_buckets(self) -> List[str]:
|
||||
"""列出所有桶"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
try:
|
||||
buckets = self.rustfs.client.list_buckets()
|
||||
bucket_names = [bucket.name for bucket in buckets]
|
||||
logger.info(f"当前存在的桶: {bucket_names}")
|
||||
return bucket_names
|
||||
except Exception as e:
|
||||
logger.error(f"列出桶失败: {e}")
|
||||
return []
|
||||
|
||||
async def check_bucket_status(self) -> Dict[str, Dict]:
|
||||
"""检查桶状态和文件数量"""
|
||||
if not self.is_connected:
|
||||
raise RuntimeError("RustFS 未连接")
|
||||
|
||||
bucket_status = {}
|
||||
|
||||
for bucket_name in self.old_buckets_to_clean:
|
||||
try:
|
||||
# 检查桶是否存在
|
||||
exists = self.rustfs.client.bucket_exists(bucket_name)
|
||||
|
||||
if exists:
|
||||
# 统计文件数量
|
||||
objects = list(self.rustfs.client.list_objects(bucket_name, recursive=True))
|
||||
file_count = len(objects)
|
||||
|
||||
# 计算总大小
|
||||
total_size = sum(obj.size for obj in objects)
|
||||
|
||||
bucket_status[bucket_name] = {
|
||||
'exists': True,
|
||||
'file_count': file_count,
|
||||
'total_size': total_size,
|
||||
'files': [obj.object_name for obj in objects]
|
||||
}
|
||||
|
||||
logger.info(f"桶 {bucket_name}: 存在, {file_count} 个文件, {total_size} 字节")
|
||||
else:
|
||||
bucket_status[bucket_name] = {
|
||||
'exists': False,
|
||||
'file_count': 0,
|
||||
'total_size': 0,
|
||||
'files': []
|
||||
}
|
||||
logger.info(f"桶 {bucket_name}: 不存在")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"检查桶 {bucket_name} 状态失败: {e}")
|
||||
bucket_status[bucket_name] = {
|
||||
'exists': False,
|
||||
'file_count': 0,
|
||||
'total_size': 0,
|
||||
'files': [],
|
||||
'error': str(e)
|
||||
}
|
||||
|
||||
return bucket_status
|
||||
|
||||
async def delete_bucket(self, bucket_name: str) -> Dict[str, any]:
|
||||
"""删除单个桶及其所有文件"""
|
||||
result = {
|
||||
'bucket_name': bucket_name,
|
||||
'success': False,
|
||||
'files_deleted': 0,
|
||||
'error': None
|
||||
}
|
||||
|
||||
try:
|
||||
# 检查桶是否存在
|
||||
if not self.rustfs.client.bucket_exists(bucket_name):
|
||||
result['success'] = True
|
||||
result['message'] = "桶不存在,无需删除"
|
||||
logger.info(f"桶 {bucket_name} 不存在,跳过删除")
|
||||
return result
|
||||
|
||||
# 列出所有文件
|
||||
objects = list(self.rustfs.client.list_objects(bucket_name, recursive=True))
|
||||
file_count = len(objects)
|
||||
|
||||
if file_count > 0:
|
||||
logger.info(f"开始删除桶 {bucket_name} 中的 {file_count} 个文件")
|
||||
|
||||
# 删除所有文件
|
||||
for obj in objects:
|
||||
try:
|
||||
self.rustfs.client.remove_object(bucket_name, obj.object_name)
|
||||
result['files_deleted'] += 1
|
||||
logger.debug(f"删除文件: {bucket_name}/{obj.object_name}")
|
||||
except Exception as e:
|
||||
logger.error(f"删除文件失败 {bucket_name}/{obj.object_name}: {e}")
|
||||
result['error'] = f"删除文件失败: {e}"
|
||||
return result
|
||||
|
||||
# 删除空桶
|
||||
self.rustfs.client.remove_bucket(bucket_name)
|
||||
|
||||
result['success'] = True
|
||||
logger.info(f"桶 {bucket_name} 删除成功,共删除 {file_count} 个文件")
|
||||
|
||||
except Exception as e:
|
||||
result['success'] = False
|
||||
result['error'] = str(e)
|
||||
logger.error(f"删除桶 {bucket_name} 失败: {e}")
|
||||
|
||||
return result
|
||||
|
||||
async def cleanup_all_old_buckets(self) -> Dict[str, Dict]:
|
||||
"""清理所有旧桶"""
|
||||
cleanup_results = {}
|
||||
|
||||
try:
|
||||
# 检查所有桶状态
|
||||
logger.info("=== 检查旧桶状态 ===")
|
||||
bucket_status = await self.check_bucket_status()
|
||||
|
||||
# 统计需要清理的桶
|
||||
buckets_to_clean = [
|
||||
bucket_name for bucket_name, status in bucket_status.items()
|
||||
if status['exists'] and status['file_count'] > 0
|
||||
]
|
||||
|
||||
if not buckets_to_clean:
|
||||
logger.info("没有需要清理的桶")
|
||||
return {}
|
||||
|
||||
print(f"\n发现 {len(buckets_to_clean)} 个需要清理的桶:")
|
||||
for bucket_name in buckets_to_clean:
|
||||
status = bucket_status[bucket_name]
|
||||
print(f" - {bucket_name}: {status['file_count']} 个文件, {status['total_size']} 字节")
|
||||
|
||||
# 确认清理
|
||||
print("\n警告:此操作将永久删除这些桶及其所有文件!")
|
||||
confirm = input("确认清理?(输入 'DELETE' 确认): ").strip()
|
||||
|
||||
if confirm != 'DELETE':
|
||||
logger.info("用户取消清理操作")
|
||||
return {}
|
||||
|
||||
# 执行清理
|
||||
logger.info("=== 开始清理旧桶 ===")
|
||||
|
||||
for bucket_name in buckets_to_clean:
|
||||
print(f"\n清理桶: {bucket_name}")
|
||||
result = await self.delete_bucket(bucket_name)
|
||||
cleanup_results[bucket_name] = result
|
||||
|
||||
if result['success']:
|
||||
print(f" ✓ 清理成功,删除 {result['files_deleted']} 个文件")
|
||||
else:
|
||||
print(f" ✗ 清理失败: {result['error']}")
|
||||
|
||||
logger.info("=== 旧桶清理完成 ===")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"清理过程失败: {e}")
|
||||
|
||||
return cleanup_results
|
||||
|
||||
async def run_cleanup(self) -> Dict[str, any]:
|
||||
"""运行完整清理流程"""
|
||||
cleanup_summary = {
|
||||
'start_time': datetime.now().isoformat(),
|
||||
'connection_status': False,
|
||||
'all_buckets': [],
|
||||
'bucket_status': {},
|
||||
'cleanup_results': {},
|
||||
'end_time': None,
|
||||
'status': 'failed'
|
||||
}
|
||||
|
||||
try:
|
||||
# 连接
|
||||
logger.info("=== RustFS 旧桶清理开始 ===")
|
||||
connection_result = await self.connect()
|
||||
if not connection_result:
|
||||
raise RuntimeError("无法连接到 RustFS")
|
||||
|
||||
cleanup_summary['connection_status'] = True
|
||||
|
||||
# 列出所有桶
|
||||
all_buckets = await self.list_all_buckets()
|
||||
cleanup_summary['all_buckets'] = all_buckets
|
||||
|
||||
# 检查状态
|
||||
bucket_status = await self.check_bucket_status()
|
||||
cleanup_summary['bucket_status'] = bucket_status
|
||||
|
||||
# 执行清理
|
||||
cleanup_results = await self.cleanup_all_old_buckets()
|
||||
cleanup_summary['cleanup_results'] = cleanup_results
|
||||
|
||||
# 完成
|
||||
cleanup_summary['status'] = 'completed'
|
||||
cleanup_summary['end_time'] = datetime.now().isoformat()
|
||||
|
||||
logger.info("=== RustFS 旧桶清理完成 ===")
|
||||
|
||||
return cleanup_summary
|
||||
|
||||
except Exception as e:
|
||||
cleanup_summary['error'] = str(e)
|
||||
cleanup_summary['end_time'] = datetime.now().isoformat()
|
||||
logger.error(f"清理失败: {e}")
|
||||
return cleanup_summary
|
||||
|
||||
finally:
|
||||
await self.close()
|
||||
|
||||
|
||||
async def main():
|
||||
"""主函数"""
|
||||
cleanup = RustFSCleanup()
|
||||
|
||||
print("=== RustFS 旧桶清理工具 ===")
|
||||
print("注意:这是rustFS,而不是minio,只是用了minio的通用S3接口")
|
||||
print("\n此工具将删除以下遗留桶:")
|
||||
print(" - moldinsight-geometry")
|
||||
print(" - moldinsight-stp-files")
|
||||
print(" - moldinsight-mold-cavities")
|
||||
print(" - moldinsight-html")
|
||||
print(" - moldinsight-user-files")
|
||||
print("\n请确保所有重要数据已备份!")
|
||||
|
||||
# 运行清理
|
||||
result = await cleanup.run_cleanup()
|
||||
|
||||
# 输出结果
|
||||
print("\n=== 清理结果摘要 ===")
|
||||
print(f"状态: {result['status']}")
|
||||
print(f"开始时间: {result['start_time']}")
|
||||
print(f"结束时间: {result['end_time']}")
|
||||
|
||||
if 'error' in result:
|
||||
print(f"错误: {result['error']}")
|
||||
|
||||
# 桶状态
|
||||
print("\n--- 桶状态检查 ---")
|
||||
for bucket_name, status in result['bucket_status'].items():
|
||||
if status['exists']:
|
||||
print(f"{bucket_name}: 存在, {status['file_count']} 个文件")
|
||||
else:
|
||||
print(f"{bucket_name}: 不存在")
|
||||
|
||||
# 清理结果
|
||||
print("\n--- 清理结果 ---")
|
||||
if result['cleanup_results']:
|
||||
for bucket_name, cleanup_result in result['cleanup_results'].items():
|
||||
if cleanup_result['success']:
|
||||
print(f"{bucket_name}: 成功,删除 {cleanup_result['files_deleted']} 个文件")
|
||||
else:
|
||||
print(f"{bucket_name}: 失败 - {cleanup_result.get('error', '未知错误')}")
|
||||
else:
|
||||
print("未执行清理操作")
|
||||
|
||||
print("\n=== 清理完成 ===")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(main())
|
||||
Reference in New Issue
Block a user