"""Step 4: ERP数据同步执行器 直接连接ERP SQL Server,执行KPI-SQL并写入 kpi_values 运行: python3 scripts/erp_data_sync.py [period] [--dry-run] """ import sys, os, json, logging, re, argparse from datetime import datetime sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) from app.database import get_session_local from sqlalchemy import text, create_engine logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s") logger = logging.getLogger("erp_data_sync") # ── ERP SQL Server 直连配置 ── ERP_DB_HOST = "211.149.143.215" ERP_DB_PORT = 1433 ERP_DB_NAME = "SUBzxbtest" ERP_DB_USER = "zxbtest" ERP_DB_PASS = os.environ.get("ERP_DB_PASS") # ── KPI-SQL 映射(从 generate_kpi_sql.py 复制核心映射) ── KPI_SQL_MAP = { "F_REVENUE_001": { "sql": """SELECT COALESCE(SUM(SumMoney), 0) as value FROM MasterBill WHERE BillType=1 AND BillState>=3 AND Period=:period""" }, "F_PROFIT_001": { "sql": """SELECT CASE WHEN SUM(SumMoney) > 0 THEN ROUND((SUM(SumMoney) - COALESCE(SUM(SumCostMoney),0)) / SUM(SumMoney) * 100, 2) ELSE 0 END as value FROM MasterBill WHERE BillType=1 AND BillState>=3 AND Period=:period""" }, "C_CUST_001": { "sql": """SELECT COUNT(DISTINCT Unit_ID) as value FROM MasterBill WHERE BillType=1 AND BillState>=3 AND Period=:period""" }, "C_CUST_002": { "sql": """SELECT CASE WHEN total_sales > 0 THEN ROUND(top5_sales / total_sales * 100, 2) ELSE 0 END as value FROM ( SELECT SUM(CASE WHEN rn <= 5 THEN SumMoney ELSE 0 END) as top5_sales, SUM(SumMoney) as total_sales FROM ( SELECT SumMoney, ROW_NUMBER() OVER (ORDER BY SumMoney DESC) as rn FROM ( SELECT SUM(SumMoney) as SumMoney FROM MasterBill WHERE BillType=1 AND BillState>=3 AND Period=:period GROUP BY Unit_ID ) t ) t2 ) t3""" }, "F_AR_002": { "sql": """SELECT CASE WHEN total_receivable > 0 THEN ROUND(overdue_receivable / total_receivable * 100, 2) ELSE 0 END as value FROM ( SELECT SUM(CASE WHEN BillType=1 THEN SumMoney ELSE 0 END) as total_receivable, SUM(CASE WHEN BillType=1 AND DATEDIFF(day, BillDate, GETDATE()) > 30 THEN SumMoney ELSE 0 END) as overdue_receivable FROM MasterBill WHERE Period<=:period AND BillState>=3 ) t""", }, # ── 杜邦分析 ── "F_ASSET_TOTAL": { "sql": """SELECT COALESCE( (SELECT SUM(CAST(Act_Tot AS FLOAT)) FROM BalanceInfo WHERE Act_ID=4 AND Period=:period) + (SELECT SUM(CAST(Act_Tot AS FLOAT)) FROM BalanceInfo WHERE Act_ID=5 AND Period=:period) , 0) as value""" }, "F_EQUITY_TOTAL": { "sql": """SELECT COALESCE( (SELECT SUM(CAST(Act_Tot AS FLOAT)) FROM BalanceInfo WHERE Act_ID=3 AND Period=:period) - (SELECT SUM(CAST(Act_Tot AS FLOAT)) FROM BalanceInfo WHERE Act_ID=2 AND Period=:period) , 0) as value""" }, } def get_erp_engine(): """创建ERP直连引擎""" conn_str = f"mssql+pymssql://{ERP_DB_USER}:{ERP_DB_PASS}@{ERP_DB_HOST}:{ERP_DB_PORT}/{ERP_DB_NAME}" return create_engine(conn_str, pool_size=2, max_overflow=5, pool_pre_ping=True) def get_all_kpis(db) -> list: """获取所有标记了erp的KPI定义""" rows = db.execute(text(""" SELECT id, kpi_code, kpi_name, formula FROM kpi_definitions WHERE data_source_type = 'erp' ORDER BY id """)).fetchall() return [dict(r._mapping) for r in rows] def get_periods_to_sync(db) -> list: """确定需要同步的期间(最近12个月)""" rows = db.execute(text(""" SELECT DISTINCT period FROM kpi_values WHERE source_type = 'erp' ORDER BY period DESC """)).fetchall() existing = set(r[0] for r in rows) # 生成最近12个月的期间 periods = [] now = datetime.now() for i in range(12): m = now.month - i y = now.year if m <= 0: m += 12 y -= 1 period = f"{y}-{m:02d}" periods.append(period) # 只同步已有期间中没有数据或需要更新的 # ERP中 Period 列: 1=1月, 2=2月 ... 12=12月(年度期间) return periods def erp_period_to_int(period: str) -> str: """将 2026-05 转为ERP的期间数字""" return period.split("-")[1] # "05" def execute_erp_sql(erp_engine, sql: str, period: str) -> float: """在ERP SQL Server上执行SQL""" period_month = period.split("-")[1] period_year = period.split("-")[0] # 替换参数 exec_sql = sql.replace(":period", period_month) exec_sql = re.sub(r"GETDATE\(\)", f"'{datetime.now().strftime('%Y-%m-%d')}'", exec_sql) with erp_engine.connect() as conn: result = conn.execute(text(exec_sql)) row = result.fetchone() return float(row[0]) if row and row[0] is not None else 0.0 def main(): parser = argparse.ArgumentParser(description="ERP数据同步") parser.add_argument("period", nargs="?", default=None, help="期间,如 2026-05") parser.add_argument("--dry-run", action="store_true", help="试运行,不写入数据库") args = parser.parse_args() db = get_session_local()() # 获取KPI kpis = get_all_kpis(db) logger.info(f"待同步KPI: {len(kpis)} 个") for k in kpis: has_sql = "✅" if k["kpi_code"] in KPI_SQL_MAP else "❌" logger.info(f" {has_sql} [{k['kpi_code']}] {k['kpi_name']}") # 获取期间 if args.period: periods = [args.period] else: periods = get_periods_to_sync(db) periods = periods[:3] # 先只同步最近3个月 logger.info(f"期间: {periods}") # 连接ERP logger.info("连接ERP SQL Server...") try: erp_engine = get_erp_engine() with erp_engine.connect() as conn: conn.execute(text("SELECT 1")) logger.info("✅ ERP连接成功") except Exception as e: logger.error(f"❌ ERP连接失败: {e}") db.close() return # 逐KPI逐期间执行 total_written = 0 for kpi in kpis: code = kpi["kpi_code"] mapping = KPI_SQL_MAP.get(code) if not mapping: logger.info(f" [{code}] 跳过(无SQL映射)") continue for period in periods: try: value = execute_erp_sql(erp_engine, mapping["sql"], period) logger.info(f" [{code}] {period} = {value}") if not args.dry_run: # 写入 kpi_values db.execute(text(""" INSERT INTO kpi_values (kpi_id, period, actual_value, source_type, source_batch, data_status, calculated_at) VALUES (:kpi_id, :period, :value, 'erp', :batch, 'verified', NOW()) ON DUPLICATE KEY UPDATE actual_value=:value2, source_batch=:batch2, data_status='verified', calculated_at=NOW() """), { "kpi_id": kpi["id"], "period": period, "value": value, "batch": f"erp_sync_{datetime.now().strftime('%Y%m%d_%H%M')}", "value2": value, "batch2": f"erp_sync_{datetime.now().strftime('%Y%m%d_%H%M')}", }) db.commit() total_written += 1 except Exception as e: logger.error(f" ❌ [{code}] {period} 失败: {e}") db.close() logger.info(f"\n同步完成! 写入 {total_written} 条, 期间: {periods}") if args.dry_run: logger.info("(试运行模式,未写入数据库)") if __name__ == "__main__": main()