""" ERP数据源配置首批 — 任务8 配置data_source_config + 同步端点 + 定时任务 """ import sys, os sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) from dotenv import load_dotenv load_dotenv() from app.database import get_engine from sqlalchemy import text from datetime import datetime engine = get_engine() # 首批8个ERP-KPI的数据源配置 ERP_SOURCES = [ { "name": "ERP总账-营业收入", "kpi_codes": ["F_REVENUE"], "source_type": "erp", "api_endpoint": "http://127.0.0.1:8300/api/v1/stats/monthly", "query_sql": "SELECT COALESCE(SUM(SumMoney), 0) as value FROM MasterBill WHERE BillType=1 AND BillState>=3 AND Period=:period", "sync_type": "daily", }, { "name": "ERP总账-毛利率", "kpi_codes": ["F_GROSS_MARGIN", "F_PROFIT_RATE"], "source_type": "erp", "api_endpoint": "http://127.0.0.1:8300/api/v1/stats/gross-profit", "query_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", "sync_type": "daily", }, { "name": "ERP总账-净利润", "kpi_codes": ["F_NET_PROFIT", "F_NET_PROFIT_RATE"], "source_type": "erp", "api_endpoint": "http://127.0.0.1:8300/api/v1/stats/monthly", "query_sql": "SELECT COALESCE(SUM(SumMoney),0)-COALESCE(SUM(SumCostMoney),0)-COALESCE(SUM(SumExpense),0) as value FROM MasterBill WHERE BillType IN (1,2,3) AND BillState>=3 AND Period=:period", "sync_type": "daily", }, { "name": "ERP总账-费用率", "kpi_codes": ["F_COST_RATIO"], "source_type": "erp", "api_endpoint": "http://127.0.0.1:8300/api/v1/stats/monthly", "query_sql": "SELECT CASE WHEN SUM(CASE WHEN BillType=1 THEN SumMoney ELSE 0 END) > 0 THEN ROUND((COALESCE(SUM(CASE WHEN BillType=2 THEN SumMoney ELSE 0 END),0)+COALESCE(SUM(CASE WHEN BillType=3 THEN SumMoney ELSE 0 END),0))/SUM(CASE WHEN BillType=1 THEN SumMoney ELSE 0 END)*100,2) ELSE 0 END as value FROM MasterBill WHERE BillState>=3 AND Period=:period", "sync_type": "daily", }, { "name": "ERP总账-经营性现金流", "kpi_codes": ["F_OP_CFLOW", "F_CASH_FLOW"], "source_type": "erp", "api_endpoint": "http://127.0.0.1:8300/api/v1/cashflow", "query_sql": "SELECT COALESCE(SUM(CashIn),0)-COALESCE(SUM(CashOut),0) as value FROM CashFlow WHERE Period=:period", "sync_type": "daily", }, { "name": "ERP应收-周转天数", "kpi_codes": ["F_AR_DAYS", "F_AR_TURNOVER"], "source_type": "erp", "api_endpoint": "http://127.0.0.1:8300/api/v1/ar/aging", "query_sql": "SELECT CASE WHEN total_sales > 0 THEN ROUND(AVG_receivable/total_sales*365,0) ELSE 0 END as value FROM (SELECT SUM(CASE WHEN BillType=1 THEN SumMoney ELSE 0 END) as total_sales, AVG(CASE WHEN BillType=1 THEN SumMoney ELSE 0 END) as AVG_receivable FROM MasterBill WHERE BillState>=3 AND Period<=:period) t", "sync_type": "daily", }, { "name": "ERP交付-及时率", "kpi_codes": ["P_DELIVERY", "P_DELIVERY_ON_TIME"], "source_type": "erp", "api_endpoint": "http://127.0.0.1:8300/api/v1/delivery/rate", "query_sql": "SELECT CASE WHEN total_orders > 0 THEN ROUND(on_time_orders/total_orders*100,2) ELSE 0 END as value FROM (SELECT COUNT(*) as total_orders, SUM(CASE WHEN ActualDelivery<=PlanDelivery THEN 1 ELSE 0 END) as on_time_orders FROM OrderDelivery WHERE Period=:period) t", "sync_type": "daily", }, { "name": "ERP质量-缺陷率", "kpi_codes": ["P_BUG_RATE", "P_DEFECT_RATE"], "source_type": "erp", "api_endpoint": "http://127.0.0.1:8300/api/v1/quality/defect", "query_sql": "SELECT CASE WHEN total_qty > 0 THEN ROUND(defect_qty/total_qty*100,2) ELSE 0 END as value FROM (SELECT COUNT(*) as total_qty, SUM(CASE WHEN QualityStatus='NG' THEN 1 ELSE 0 END) as defect_qty FROM QualityInspection WHERE Period=:period) t", "sync_type": "daily", }, ] def log(msg): print(f"[{datetime.now():%H:%M:%S}] {msg}") with engine.connect() as conn: print("=" * 70) log("配置ERP数据源 (data_source_config)") print("=" * 70) existing = conn.execute(text("SELECT COUNT(*) FROM data_source_config")).scalar() log(f"当前已有 {existing} 条数据源配置") created = 0 for src in ERP_SOURCES: exists = conn.execute( text("SELECT id FROM data_source_config WHERE name=:name"), {"name": src["name"]} ).fetchone() if exists: log(f" ⏭️ 已存在: {src['name']} (id={exists[0]})") continue conn.execute(text(""" INSERT INTO data_source_config (name, source_type, api_endpoint, query_sql, sync_type, status, created_at) VALUES (:name, :source_type, :api_endpoint, :query_sql, :sync_type, 'active', NOW()) """), { "name": src["name"], "source_type": src["source_type"], "api_endpoint": src["api_endpoint"], "query_sql": src["query_sql"], "sync_type": src["sync_type"], }) created += 1 log(f" ✅ 新增: {src['name']} — 关联KPI: {', '.join(src['kpi_codes'])}") conn.commit() log(f"\n✅ 数据源配置完成: 新增 {created} 条, 当前共 {existing + created} 条") # 打印所有数据源 rows = conn.execute(text("SELECT id, name, source_type, sync_type, status FROM data_source_config ORDER BY id")).fetchall() print("\n数据源清单:") for r in rows: print(f" [{r[0]}] {r[1]:30s} type={r[2]:10s} sync={r[3]:10s} status={r[4]}")