128 lines
5.8 KiB
Python
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]}")
|