Files
cma-management/backend/scripts/erp_data_source_config.py

128 lines
5.8 KiB
Python

"""
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]}")