181 lines
6.5 KiB
Python
181 lines
6.5 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
MySQL → PostgreSQL Data Migration Script
|
|
Reads from MySQL, transforms data, writes to PostgreSQL.
|
|
Handles ENUM→VARCHAR, TINYINT→BOOLEAN, JSON→JSONB, AUTO_INCREMENT→SERIAL.
|
|
"""
|
|
|
|
import mysql.connector
|
|
import asyncpg
|
|
import asyncio
|
|
import json
|
|
import os
|
|
from dotenv import load_dotenv
|
|
|
|
load_dotenv()
|
|
|
|
MYSQL_CONFIG = {
|
|
"host": os.getenv("MYSQL_HOST", "localhost"),
|
|
"port": int(os.getenv("MYSQL_PORT", "3306")),
|
|
"user": os.getenv("MYSQL_USER", "root"),
|
|
"password": os.getenv("MYSQL_PASSWORD", ""),
|
|
"database": os.getenv("MYSQL_DATABASE", "aegisone"),
|
|
}
|
|
|
|
PG_DSN = os.getenv("PG_DSN", "postgresql://aegisone:aegisone_pass@localhost:5432/aegisone")
|
|
|
|
TABLES = [
|
|
"users", "customers", "objects", "object_assignments",
|
|
"tasks", "task_comments", "task_photos", "reports",
|
|
"questionnaire_sessions", "questionnaire_answers", "questionnaire_items",
|
|
"object_passports", "sla_contracts", "incidents",
|
|
"engineer_kpi", "shs_records", "blog_posts",
|
|
"formula_coefficients", "cases", "login_attempts", "audit_log",
|
|
]
|
|
|
|
ENUM_MAP = {
|
|
"users.role": {"owner": "owner", "engineer": "engineer", "technician": "technician"},
|
|
"objects.status": {"active": "active", "inactive": "inactive", "audit": "audit", "prospective": "prospective"},
|
|
"tasks.priority": {"P1": "P1", "P2": "P2", "P3": "P3", "P4": "P4"},
|
|
"tasks.status": {"open": "open", "in_progress": "in_progress", "completed": "completed", "cancelled": "cancelled"},
|
|
"reports.report_type": {"inspection": "inspection", "service": "service", "emergency": "emergency", "audit": "audit", "monthly": "monthly"},
|
|
"reports.result_status": {"fixed": "fixed", "partial": "partial", "revisit": "revisit", "ok": "ok"},
|
|
"reports.status": {"draft": "draft", "final": "final", "cancelled": "cancelled"},
|
|
"questionnaire_sessions.status": {"draft": "draft", "completed": "completed", "cancelled": "cancelled"},
|
|
"questionnaire_items.type": {"text": "text", "number": "number", "select": "select", "radio": "radio", "checkbox": "checkbox", "textarea": "textarea"},
|
|
"sla_contracts.service_level": {"start": "start", "business": "business", "enterprise": "enterprise"},
|
|
"sla_contracts.status": {"active": "active", "expired": "expired", "cancelled": "cancelled", "negotiation": "negotiation"},
|
|
"incidents.severity": {"P1": "P1", "P2": "P2", "P3": "P3"},
|
|
"incidents.status": {"open": "open", "in_progress": "in_progress", "resolved": "resolved", "closed": "closed"},
|
|
"blog_posts.category": {"audit": "audit", "sla": "sla", "incident": "incident", "supervision": "supervision", "documentation": "documentation", "risk": "risk", "cases": "cases"},
|
|
"blog_posts.status": {"draft": "draft", "published": "published", "archived": "archived"},
|
|
"object_assignments.role": {"engineer": "engineer", "technician": "technician"},
|
|
"customers.status": {"active": "active", "inactive": "inactive", "prospect": "prospect"},
|
|
"task_photos.type": {"before": "before", "after": "after", "evidence": "evidence", "other": "other"},
|
|
}
|
|
|
|
BOOL_COLUMNS = {
|
|
"users": ["is_active"],
|
|
"tasks": ["sla_compliant"],
|
|
"reports": ["client_confirmed"],
|
|
"questionnaire_items": ["required", "is_active"],
|
|
"object_passports": [],
|
|
"sla_contracts": [],
|
|
"incidents": ["sla_breached"],
|
|
"engineer_kpi": [],
|
|
"shs_records": [],
|
|
"blog_posts": [],
|
|
"formula_coefficients": ["is_active"],
|
|
"cases": ["is_active"],
|
|
"customers": [],
|
|
"objects": [],
|
|
"object_assignments": [],
|
|
"task_comments": [],
|
|
"task_photos": [],
|
|
"questionnaire_sessions": [],
|
|
"questionnaire_answers": [],
|
|
"login_attempts": [],
|
|
"audit_log": [],
|
|
}
|
|
|
|
SKIP_COLUMNS = {
|
|
"users": [],
|
|
"customers": [],
|
|
"objects": [],
|
|
"object_assignments": [],
|
|
"tasks": [],
|
|
"task_comments": [],
|
|
"task_photos": [],
|
|
"reports": [],
|
|
"questionnaire_sessions": [],
|
|
"questionnaire_answers": [],
|
|
"questionnaire_items": [],
|
|
"object_passports": [],
|
|
"sla_contracts": [],
|
|
"incidents": [],
|
|
"engineer_kpi": [],
|
|
"shs_records": [],
|
|
"blog_posts": [],
|
|
"formula_coefficients": [],
|
|
"cases": [],
|
|
"login_attempts": [],
|
|
"audit_log": [],
|
|
}
|
|
|
|
|
|
def transform_row(table: str, row: dict) -> dict:
|
|
result = {}
|
|
for key, value in row.items():
|
|
if key in SKIP_COLUMNS.get(table, []):
|
|
continue
|
|
if value is None:
|
|
result[key] = None
|
|
continue
|
|
enum_key = f"{table}.{key}"
|
|
if enum_key in ENUM_MAP:
|
|
result[key] = ENUM_MAP[enum_key].get(value, value)
|
|
elif table in BOOL_COLUMNS and key in BOOL_COLUMNS[table]:
|
|
result[key] = bool(value)
|
|
elif isinstance(value, (bytes, bytearray)):
|
|
result[key] = value.decode("utf-8")
|
|
elif isinstance(value, dict) or isinstance(value, list):
|
|
result[key] = json.dumps(value, ensure_ascii=False)
|
|
else:
|
|
result[key] = value
|
|
return result
|
|
|
|
|
|
async def migrate_table(mysql_conn, pg_conn, table: str):
|
|
cursor = mysql_conn.cursor(dictionary=True)
|
|
cursor.execute(f"SELECT * FROM {table}")
|
|
rows = cursor.fetchall()
|
|
|
|
if not rows:
|
|
print(f" {table}: нет данных")
|
|
return
|
|
|
|
columns = list(rows[0].keys())
|
|
placeholders = ", ".join([f"${i+1}" for i in range(len(columns))])
|
|
col_names = ", ".join(columns)
|
|
|
|
inserted = 0
|
|
errors = 0
|
|
for row in rows:
|
|
try:
|
|
transformed = transform_row(table, row)
|
|
values = [transformed.get(c) for c in columns]
|
|
await pg_conn.execute(
|
|
f"INSERT INTO {table} ({col_names}) VALUES ({placeholders}) ON CONFLICT DO NOTHING",
|
|
*values,
|
|
)
|
|
inserted += 1
|
|
except Exception as e:
|
|
errors += 1
|
|
print(f" Ошибка в {table}: {e}")
|
|
|
|
print(f" {table}: {inserted} записей, {errors} ошибок")
|
|
|
|
|
|
async def main():
|
|
print("Подключение к MySQL...")
|
|
mysql_conn = mysql.connector.connect(**MYSQL_CONFIG)
|
|
|
|
print("Подключение к PostgreSQL...")
|
|
pg_conn = await asyncpg.connect(PG_DSN)
|
|
|
|
total_tables = len(TABLES)
|
|
for i, table in enumerate(TABLES, 1):
|
|
print(f"[{i}/{total_tables}] Миграция {table}...")
|
|
try:
|
|
await migrate_table(mysql_conn, pg_conn, table)
|
|
except Exception as e:
|
|
print(f" Ошибка миграции {table}: {e}")
|
|
|
|
await pg_conn.close()
|
|
mysql_conn.close()
|
|
print("\nМиграция завершена.")
|
|
|
|
|
|
if __name__ == "__main__":
|
|
asyncio.run(main())
|