Files

235 lines
7.8 KiB
Python

"""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()