## 来源
组员通过腾讯会议发来的 `seed_risk_demo_data.py`,按约定放入 `tools/`。
**脚本内容是组员写的**,我只做了两件事:放进仓库、把 docstring 里的用法示例
统一成实际文件名(组员原文写的是 `seed_risk_scan_demo_data.py`,与文件名不符,
已改成 `tools/seed_risk_demo_data.py`),**未改动任何逻辑**。
## 它做什么
固定号段:客户 **12001–12005**、产品 **159991–159995**(独立号段,不覆盖人工数据
与其他演示数据)。写入客户 / 画像 / 产品 / 账户 / 持仓 / 交易 / 资金流水 / 登录 /
工单,**刻意不写 `fin_risk_alert` 与 `fin_risk_notification`** —— 预警必须由风控
扫描器按规则生成,这样演示才能真实覆盖规则、证据、通知三条链路。
五个场景各自瞄准一条规则:
| 客户 | 场景 | 预期规则 |
|---|---|---|
| 12001 | 大额快进快出(3 天内入金 80 万后赎回 75 万,比例 93.75%) | RW-003 |
| 12002 | 高风险适当性错配(缺风险揭示 + 二次确认) | RW-007 |
| 12003 | 中风险适当性错配(等级差 1 且缺风险揭示) | RW-007 |
| 12004 | 72 岁客户赎回 75 万,达历史均值 3 倍且用非常用设备 | RW-012 |
| 12005 | 凌晨 2:15 小额自动定投(关联有效工单,预期合并预警) | RW-015、RW-018 |
模式:`--check`(只读核对)、默认 `--dry-run`(只打印计划、**不连数据库**)、
`--apply`(写入)、`--dry-run-scan`(跑规则引擎后回滚)。
## 我跑了什么、核到什么
按约定的命令 `--apply --dry-run-scan`:
[命中] 客户 12001:RW-003 (高)
[命中] 客户 12002:RW-007 (高)
[命中] 客户 12003:RW-007 (中)
[命中] 客户 12004:RW-012 (高)
[命中] 客户 12005:RW-015、RW-018 (低)
规则干跑完成,共命中 5 条演示预警,已回滚,不会写库。
**"已回滚"这句话我没有只信它的打印**,而是直接查了库:
- 上游数据确实写入:客户 5 / 产品 5 / 账户 5 / 持仓 5 / 委托 6 / 成交 6 /
资金流水 1 / 登录 1 / 工单 3,演示产品行情 5 行(`source=custom_seed`);
- **`fin_risk_alert` = 0、`fin_risk_notification` = 0** —— 回滚确实生效。
跑之前还先确认了它承诺回滚是**成立**的:`RiskRuleEngine.refresh_alerts()`
(`risk_scan_service.py:101`)到 565 行之间**没有 `commit()`**,全文件仅有的
`commit/rollback` 在 `:566-568`,属于另一个方法(正式扫描 `scan()`)。
否则 `dry_run_scan` 的 `session.rollback()` 根本回滚不掉,那就是"假干跑"。
## 复用与兼容
脚本 import 了 `tools/seed_custom_holdings.py` 的 `ensure_user` / `ensure_account` /
`ensure_holding` / `ensure_product` / `ensure_nav_history` / `ensure_market_price`,
我逐个核对了签名,**全部匹配**(组员看过那个脚本)。
它自己的 `_next_id` 用 `SELECT MAX(id)+1` 发号:单线程脚本里没问题,且
`sys_user` 与 `fin_product` 本来就没有 AUTO_INCREMENT(被外键引用,迁移
`20260914_baseline_auto_increment` 已说明为何排除)。**这不是缺陷**,不改。
## 实测
- `ruff check tools/seed_risk_demo_data.py` -> All checks passed
- `--apply --dry-run-scan` -> 5/5 场景命中预期规则,退出码 0
- `pytest tests/unit tests/contract` -> **1428 passed, 2 skipped, 1 failed**
(那个失败是投顾工作台页面被替换所致,与本次无关)
863 lines
28 KiB
Python
863 lines
28 KiB
Python
"""生成风控演示所需的上游业务数据。
|
||
|
||
本脚本只写入客户、画像、产品、账户、持仓、交易、资金流水、登录和工单,
|
||
不直接写 `fin_risk_alert` 或 `fin_risk_notification`。预警必须由风控扫描器
|
||
按规则生成,确保演示内容能真实覆盖规则、证据和通知链路。
|
||
|
||
固定演示客户为 12001 至 12005,产品编码为 159991 至 159995。
|
||
重复执行按业务编号更新,不重复新增同一批演示记录。
|
||
|
||
用法:
|
||
python tools/seed_risk_demo_data.py
|
||
python tools/seed_risk_demo_data.py --check
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import argparse
|
||
import asyncio
|
||
import sys
|
||
from dataclasses import dataclass
|
||
from datetime import UTC, date, datetime, time, timedelta
|
||
from decimal import Decimal
|
||
from pathlib import Path
|
||
|
||
from sqlalchemy import func, select, update
|
||
from sqlalchemy.ext.asyncio import AsyncSession
|
||
|
||
PROJECT_ROOT = Path(__file__).resolve().parents[1]
|
||
if str(PROJECT_ROOT) not in sys.path:
|
||
sys.path.insert(0, str(PROJECT_ROOT))
|
||
|
||
from app.core.timeutil import local_zone, to_utc_naive # noqa: E402
|
||
from app.infrastructure.db import SessionFactory # noqa: E402
|
||
from app.model.fund import ( # noqa: E402
|
||
FundCapitalFlow,
|
||
FundCustomerProfile,
|
||
FundProduct,
|
||
FundRiskAssessment,
|
||
FundSimAccount,
|
||
FundSimOrder,
|
||
FundTransaction,
|
||
)
|
||
from app.model.risk import RiskLoginRecord, RiskUser, RiskWorkOrder # noqa: E402
|
||
from app.service.risk_scan_service import RiskRuleEngine # noqa: E402
|
||
from tools.seed_custom_holdings import ( # noqa: E402
|
||
ensure_account,
|
||
ensure_holding,
|
||
ensure_market_price,
|
||
ensure_nav_history,
|
||
ensure_product,
|
||
ensure_user,
|
||
)
|
||
|
||
if hasattr(sys.stdout, "reconfigure"):
|
||
sys.stdout.reconfigure(errors="replace") # type: ignore[union-attr]
|
||
|
||
|
||
DEMO_PASSWORD = "risk12345"
|
||
DEMO_SOURCE = "risk_demo_seed"
|
||
DEMO_QUESTIONNAIRE_VERSION = "risk-demo-v1"
|
||
|
||
#: 客户编号和产品编号都固定在独立号段,避免覆盖人工数据和其他演示数据。
|
||
DEMO_CUSTOMERS = (12001, 12002, 12003, 12004, 12005)
|
||
DEMO_PRODUCTS = ("159991", "159992", "159993", "159994", "159995")
|
||
|
||
|
||
@dataclass(frozen=True)
|
||
class CustomerSpec:
|
||
customer_id: int
|
||
name: str
|
||
age: int
|
||
investor_type: str
|
||
total_asset: Decimal
|
||
behavior_score: int = 20
|
||
|
||
|
||
@dataclass(frozen=True)
|
||
class ProductSpec:
|
||
code: str
|
||
name: str
|
||
risk_level: str
|
||
nav: Decimal
|
||
risk_disclosure_required: int = 0
|
||
second_confirmation_required: int = 0
|
||
recording_required: int = 0
|
||
|
||
|
||
@dataclass(frozen=True)
|
||
class Scenario:
|
||
customer_id: int
|
||
label: str
|
||
expected_rules: tuple[str, ...]
|
||
|
||
|
||
CUSTOMER_SPECS: dict[int, CustomerSpec] = {
|
||
12001: CustomerSpec(12001, "演示客户一", 42, "C4", Decimal("2000000.00")),
|
||
12002: CustomerSpec(12002, "演示客户二", 38, "C1", Decimal("1500000.00")),
|
||
12003: CustomerSpec(12003, "演示客户三", 46, "C3", Decimal("1800000.00")),
|
||
12004: CustomerSpec(12004, "演示客户四", 72, "C3", Decimal("5000000.00")),
|
||
12005: CustomerSpec(12005, "演示客户五", 51, "C4", Decimal("1200000.00")),
|
||
}
|
||
|
||
PRODUCT_SPECS: dict[str, ProductSpec] = {
|
||
"159991": ProductSpec("159991", "风控演示稳健ETF", "R3", Decimal("1.000000")),
|
||
"159992": ProductSpec(
|
||
"159992",
|
||
"风控演示高风险ETF",
|
||
"R4",
|
||
Decimal("1.000000"),
|
||
risk_disclosure_required=1,
|
||
second_confirmation_required=1,
|
||
),
|
||
"159993": ProductSpec(
|
||
"159993",
|
||
"风控演示进取ETF",
|
||
"R4",
|
||
Decimal("1.000000"),
|
||
risk_disclosure_required=1,
|
||
),
|
||
"159994": ProductSpec("159994", "风控演示低风险ETF", "R2", Decimal("1.000000")),
|
||
"159995": ProductSpec("159995", "风控演示平衡ETF", "R3", Decimal("1.000000")),
|
||
}
|
||
|
||
SCENARIOS: tuple[Scenario, ...] = (
|
||
Scenario(12001, "大额快进快出", ("RW-003",)),
|
||
Scenario(12002, "高风险适当性错配", ("RW-007",)),
|
||
Scenario(12003, "中风险适当性错配", ("RW-007",)),
|
||
Scenario(12004, "高龄客户异常赎回", ("RW-012",)),
|
||
Scenario(12005, "凌晨小额自动定投", ("RW-015", "RW-018")),
|
||
)
|
||
|
||
|
||
def _now() -> datetime:
|
||
return datetime.now(UTC).replace(tzinfo=None)
|
||
|
||
|
||
def _birth_date(age: int) -> date:
|
||
today = date.today()
|
||
return date(today.year - age, 1, 1)
|
||
|
||
|
||
def _local_timestamp(hour: int, minute: int = 0) -> datetime:
|
||
zone = local_zone()
|
||
now = datetime.now(zone)
|
||
target = datetime.combine(now.date(), time(hour, minute), tzinfo=zone)
|
||
if target > now:
|
||
target -= timedelta(days=1)
|
||
return to_utc_naive(target)
|
||
|
||
|
||
async def _next_id(session: AsyncSession, model: type) -> int:
|
||
pk = model.__table__.primary_key.columns[0] # type: ignore[attr-defined]
|
||
value = await session.scalar(select(func.coalesce(func.max(pk), 0)))
|
||
return int(value or 0) + 1
|
||
|
||
|
||
async def _ensure_assessment(
|
||
session: AsyncSession,
|
||
customer_id: int,
|
||
investor_type: str,
|
||
now: datetime,
|
||
) -> None:
|
||
existing = await session.scalar(
|
||
select(FundRiskAssessment.id).where(
|
||
FundRiskAssessment.customer_id == customer_id,
|
||
FundRiskAssessment.questionnaire_version == DEMO_QUESTIONNAIRE_VERSION,
|
||
)
|
||
)
|
||
answers = {f"q{index}": 1 for index in range(1, 14)}
|
||
values = {
|
||
"customer_id": customer_id,
|
||
"questionnaire_version": DEMO_QUESTIONNAIRE_VERSION,
|
||
"answers": answers,
|
||
"total_score": 30,
|
||
"investor_type": investor_type,
|
||
"assessed_at": now,
|
||
"valid_until": now + timedelta(days=365),
|
||
"created_at": now,
|
||
}
|
||
if existing is None:
|
||
session.add(FundRiskAssessment(id=await _next_id(session, FundRiskAssessment), **values))
|
||
else:
|
||
await session.execute(
|
||
update(FundRiskAssessment).where(FundRiskAssessment.id == int(existing)).values(**values)
|
||
)
|
||
|
||
|
||
async def _ensure_customer(
|
||
session: AsyncSession,
|
||
spec: CustomerSpec,
|
||
*,
|
||
now: datetime,
|
||
) -> tuple[str, int]:
|
||
await ensure_user(
|
||
session,
|
||
spec.customer_id,
|
||
password=DEMO_PASSWORD,
|
||
password_changed=True,
|
||
)
|
||
account_no, _, _ = await ensure_account(session, spec.customer_id)
|
||
account_id = await session.scalar(
|
||
select(FundSimAccount.id).where(FundSimAccount.customer_id == spec.customer_id)
|
||
)
|
||
if account_id is None:
|
||
raise RuntimeError(f"客户 {spec.customer_id} 的模拟账户未创建")
|
||
|
||
await session.execute(
|
||
update(RiskUser)
|
||
.where(RiskUser.id == spec.customer_id)
|
||
.values(
|
||
investor_type=spec.investor_type,
|
||
investor_type_assessed_at=now,
|
||
updated_at=now,
|
||
)
|
||
)
|
||
|
||
profile = await session.get(FundCustomerProfile, spec.customer_id)
|
||
profile_values = {
|
||
"trade_account": account_no,
|
||
"real_name": spec.name,
|
||
"birth_date": _birth_date(spec.age),
|
||
"occupation": "演示职业",
|
||
"mobile_masked": "138****0000",
|
||
"investor_type": spec.investor_type,
|
||
"investment_horizon": "3至5年",
|
||
"preferred_asset_class": ["固定收益类", "权益类"],
|
||
"trading_frequency": "中",
|
||
"last_active_at": now,
|
||
"total_asset": spec.total_asset,
|
||
"behavior_score": spec.behavior_score,
|
||
"risk_tags": [],
|
||
"opened_at": now - timedelta(days=365),
|
||
"updated_at": now,
|
||
}
|
||
if profile is None:
|
||
session.add(FundCustomerProfile(customer_id=spec.customer_id, **profile_values))
|
||
else:
|
||
await session.execute(
|
||
update(FundCustomerProfile)
|
||
.where(FundCustomerProfile.customer_id == spec.customer_id)
|
||
.values(**profile_values)
|
||
)
|
||
|
||
await _ensure_assessment(session, spec.customer_id, spec.investor_type, now)
|
||
return account_no, int(account_id)
|
||
|
||
|
||
async def _ensure_product(
|
||
session: AsyncSession,
|
||
spec: ProductSpec,
|
||
*,
|
||
now: datetime,
|
||
) -> int:
|
||
product_id, _ = await ensure_product(
|
||
session,
|
||
spec.code,
|
||
name=spec.name,
|
||
exchange="SSE",
|
||
category="ETF",
|
||
risk_level=spec.risk_level,
|
||
nav=spec.nav,
|
||
)
|
||
await session.execute(
|
||
update(FundProduct)
|
||
.where(FundProduct.id == product_id)
|
||
.values(
|
||
risk_level=spec.risk_level,
|
||
risk_disclosure_required=spec.risk_disclosure_required,
|
||
second_confirmation_required=spec.second_confirmation_required,
|
||
recording_required=spec.recording_required,
|
||
status="上市",
|
||
updated_at=now,
|
||
)
|
||
)
|
||
await ensure_nav_history(session, product_id, spec.nav, 120)
|
||
await ensure_market_price(session, product_id, spec.nav)
|
||
return int(product_id)
|
||
|
||
|
||
async def _ensure_work_order(
|
||
session: AsyncSession,
|
||
*,
|
||
work_order_no: str,
|
||
customer_id: int,
|
||
product_id: int,
|
||
amount: Decimal,
|
||
channel: str,
|
||
risk_disclosure_ack_at: datetime | None,
|
||
second_confirmation_at: datetime | None,
|
||
now: datetime,
|
||
) -> int:
|
||
values = {
|
||
"customer_id": customer_id,
|
||
"order_type": "风险演示",
|
||
"product_id": product_id,
|
||
"amount": amount,
|
||
"channel": channel,
|
||
"risk_rule_hits": None,
|
||
"risk_disclosure_ack_at": risk_disclosure_ack_at,
|
||
"second_confirmation_at": second_confirmation_at,
|
||
"recording_reference": None,
|
||
"submitted_at": now,
|
||
"status": "已完成",
|
||
"work_order_type": "演示工单",
|
||
"submitter_id": customer_id,
|
||
"priority": "普通",
|
||
"request_detail": {"source": DEMO_SOURCE},
|
||
"created_at": now,
|
||
"updated_at": now,
|
||
}
|
||
existing = await session.scalar(
|
||
select(RiskWorkOrder.id).where(RiskWorkOrder.work_order_no == work_order_no)
|
||
)
|
||
if existing is None:
|
||
work_order_id = await _next_id(session, RiskWorkOrder)
|
||
session.add(RiskWorkOrder(id=work_order_id, work_order_no=work_order_no, **values))
|
||
return int(work_order_id)
|
||
await session.execute(
|
||
update(RiskWorkOrder).where(RiskWorkOrder.id == int(existing)).values(**values)
|
||
)
|
||
return int(existing)
|
||
|
||
|
||
async def _ensure_order(
|
||
session: AsyncSession,
|
||
*,
|
||
order_no: str,
|
||
customer_id: int,
|
||
account_id: int,
|
||
product_id: int,
|
||
transaction_type: str,
|
||
amount: Decimal,
|
||
nav: Decimal,
|
||
confirmed_at: datetime,
|
||
) -> int:
|
||
side = "buy" if transaction_type == "申购" else "sell"
|
||
quantity = (amount / nav).quantize(Decimal("0.0001"))
|
||
values = {
|
||
"customer_id": customer_id,
|
||
"account_id": account_id,
|
||
"product_id": product_id,
|
||
"order_side": side,
|
||
"price_type": "market",
|
||
"quantity": quantity,
|
||
"limit_price": None,
|
||
"quote_price": nav,
|
||
"quote_at": confirmed_at,
|
||
"quote_source": DEMO_SOURCE,
|
||
"channel": "risk_demo",
|
||
"advisor_id": None,
|
||
"filled_quantity": quantity,
|
||
"average_executed_price": nav,
|
||
"status": "已成交",
|
||
"risk_rule_hits": None,
|
||
"risk_disclosure_ack_at": None,
|
||
"second_confirmation_at": None,
|
||
"recording_reference": None,
|
||
"submitted_at": confirmed_at,
|
||
"cancelled_at": None,
|
||
"created_at": confirmed_at,
|
||
"updated_at": confirmed_at,
|
||
}
|
||
existing = await session.scalar(
|
||
select(FundSimOrder.id).where(FundSimOrder.order_no == order_no)
|
||
)
|
||
if existing is None:
|
||
order_id = await _next_id(session, FundSimOrder)
|
||
session.add(FundSimOrder(id=order_id, order_no=order_no, **values))
|
||
return int(order_id)
|
||
await session.execute(
|
||
update(FundSimOrder).where(FundSimOrder.id == int(existing)).values(**values)
|
||
)
|
||
return int(existing)
|
||
|
||
|
||
async def _ensure_transaction(
|
||
session: AsyncSession,
|
||
*,
|
||
transaction_no: str,
|
||
order_id: int,
|
||
work_order_id: int | None,
|
||
customer_id: int,
|
||
account_id: int,
|
||
product_id: int,
|
||
transaction_type: str,
|
||
amount: Decimal,
|
||
nav: Decimal,
|
||
confirmed_at: datetime,
|
||
) -> int:
|
||
side = "buy" if transaction_type == "申购" else "sell"
|
||
quantity = (amount / nav).quantize(Decimal("0.0001"))
|
||
values = {
|
||
"order_id": order_id,
|
||
"work_order_id": work_order_id,
|
||
"customer_id": customer_id,
|
||
"account_id": account_id,
|
||
"product_id": product_id,
|
||
"order_side": side,
|
||
"transaction_type": transaction_type,
|
||
"executed_price": nav,
|
||
"nav": nav,
|
||
"executed_quantity": quantity,
|
||
"shares": quantity,
|
||
"gross_amount": amount,
|
||
"amount": amount,
|
||
"fee_rule_id": None,
|
||
"fee_rate_snapshot": Decimal("0.000000"),
|
||
"fee_amount": Decimal("0.00"),
|
||
"fee": Decimal("0.00"),
|
||
"net_amount": amount,
|
||
"quote_at": confirmed_at,
|
||
"quote_source": DEMO_SOURCE,
|
||
"executed_at": confirmed_at,
|
||
"confirmed_at": confirmed_at,
|
||
"confirmed_by": None,
|
||
"auto_confirmed": 1,
|
||
"created_at": confirmed_at,
|
||
}
|
||
existing = await session.scalar(
|
||
select(FundTransaction.id).where(FundTransaction.transaction_no == transaction_no)
|
||
)
|
||
if existing is None:
|
||
transaction_id = await _next_id(session, FundTransaction)
|
||
session.add(FundTransaction(id=transaction_id, transaction_no=transaction_no, **values))
|
||
return int(transaction_id)
|
||
await session.execute(
|
||
update(FundTransaction).where(FundTransaction.id == int(existing)).values(**values)
|
||
)
|
||
return int(existing)
|
||
|
||
|
||
async def _ensure_capital_flow(
|
||
session: AsyncSession,
|
||
*,
|
||
flow_no: str,
|
||
customer_id: int,
|
||
account_id: int,
|
||
amount: Decimal,
|
||
settled_at: datetime,
|
||
now: datetime,
|
||
) -> None:
|
||
values = {
|
||
"customer_id": customer_id,
|
||
"account_id": account_id,
|
||
"transaction_id": None,
|
||
"flow_type": "入金",
|
||
"amount": amount,
|
||
"balance_after": amount,
|
||
"status": "成功",
|
||
"settled_at": settled_at,
|
||
"occurred_at": settled_at,
|
||
"payer_name": "演示资金方",
|
||
"source_type": "BANK",
|
||
"related_work_order_id": None,
|
||
"match_status": "已匹配",
|
||
"created_at": now,
|
||
"updated_at": now,
|
||
}
|
||
existing = await session.scalar(
|
||
select(FundCapitalFlow.id).where(FundCapitalFlow.flow_no == flow_no)
|
||
)
|
||
if existing is None:
|
||
session.add(
|
||
FundCapitalFlow(
|
||
id=await _next_id(session, FundCapitalFlow),
|
||
flow_no=flow_no,
|
||
**values,
|
||
)
|
||
)
|
||
else:
|
||
await session.execute(
|
||
update(FundCapitalFlow).where(FundCapitalFlow.id == int(existing)).values(**values)
|
||
)
|
||
|
||
|
||
async def _ensure_login(
|
||
session: AsyncSession,
|
||
*,
|
||
customer_id: int,
|
||
device_id: str,
|
||
login_at: datetime,
|
||
is_common_device: bool,
|
||
) -> None:
|
||
values = {
|
||
"user_id": customer_id,
|
||
"login_at": login_at,
|
||
"login_result": "成功",
|
||
"ip_region": "上海",
|
||
"device_id": device_id,
|
||
"is_common_device": int(is_common_device),
|
||
"failure_reason": None,
|
||
"created_at": login_at,
|
||
}
|
||
existing = await session.scalar(
|
||
select(RiskLoginRecord.id).where(
|
||
RiskLoginRecord.user_id == customer_id,
|
||
RiskLoginRecord.device_id == device_id,
|
||
)
|
||
)
|
||
if existing is None:
|
||
session.add(
|
||
RiskLoginRecord(
|
||
id=await _next_id(session, RiskLoginRecord),
|
||
**values,
|
||
)
|
||
)
|
||
else:
|
||
await session.execute(
|
||
update(RiskLoginRecord).where(RiskLoginRecord.id == int(existing)).values(**values)
|
||
)
|
||
|
||
|
||
async def seed(session: AsyncSession) -> None:
|
||
now = _now()
|
||
customer_data: dict[int, tuple[str, int]] = {}
|
||
product_ids: dict[str, int] = {}
|
||
|
||
for spec in CUSTOMER_SPECS.values():
|
||
customer_data[spec.customer_id] = await _ensure_customer(session, spec, now=now)
|
||
for spec in PRODUCT_SPECS.values():
|
||
product_ids[spec.code] = await _ensure_product(session, spec, now=now)
|
||
|
||
# 每个客户都有一笔持仓,保证客户证据和持仓证据不是空表。
|
||
for customer_id, (account_no, _) in customer_data.items():
|
||
product_code = list(PRODUCT_SPECS)[(customer_id - 12001) % len(PRODUCT_SPECS)]
|
||
spec = PRODUCT_SPECS[product_code]
|
||
await ensure_holding(
|
||
session,
|
||
customer_id,
|
||
account_no,
|
||
product_ids[product_code],
|
||
quantity=Decimal("10000"),
|
||
nav=spec.nav,
|
||
)
|
||
|
||
# RW-003:3 天内入金 80 万,随后赎回 75 万,赎回比例 93.75%。
|
||
_, account_id = customer_data[12001]
|
||
product_id = product_ids["159991"]
|
||
await _ensure_capital_flow(
|
||
session,
|
||
flow_no="RISKDEMO-12001-FLOW-001",
|
||
customer_id=12001,
|
||
account_id=account_id,
|
||
amount=Decimal("800000.00"),
|
||
settled_at=now - timedelta(hours=24),
|
||
now=now,
|
||
)
|
||
await _ensure_transaction(
|
||
session,
|
||
transaction_no="RISKDEMO-12001-TXN-001",
|
||
order_id=await _ensure_order(
|
||
session,
|
||
order_no="RISKDEMO-12001-ORDER-001",
|
||
customer_id=12001,
|
||
account_id=account_id,
|
||
product_id=product_id,
|
||
transaction_type="赎回",
|
||
amount=Decimal("750000.00"),
|
||
nav=PRODUCT_SPECS["159991"].nav,
|
||
confirmed_at=now,
|
||
),
|
||
work_order_id=None,
|
||
customer_id=12001,
|
||
account_id=account_id,
|
||
product_id=product_id,
|
||
transaction_type="赎回",
|
||
amount=Decimal("750000.00"),
|
||
nav=PRODUCT_SPECS["159991"].nav,
|
||
confirmed_at=now,
|
||
)
|
||
|
||
# RW-007:高风险适当性错配,缺失风险揭示和二次确认。
|
||
_, account_id = customer_data[12002]
|
||
work_order_id = await _ensure_work_order(
|
||
session,
|
||
work_order_no="RISKDEMO-12002-WO-001",
|
||
customer_id=12002,
|
||
product_id=product_ids["159992"],
|
||
amount=Decimal("300000.00"),
|
||
channel="APP",
|
||
risk_disclosure_ack_at=None,
|
||
second_confirmation_at=None,
|
||
now=now,
|
||
)
|
||
await _ensure_transaction(
|
||
session,
|
||
transaction_no="RISKDEMO-12002-TXN-001",
|
||
order_id=await _ensure_order(
|
||
session,
|
||
order_no="RISKDEMO-12002-ORDER-001",
|
||
customer_id=12002,
|
||
account_id=account_id,
|
||
product_id=product_ids["159992"],
|
||
transaction_type="申购",
|
||
amount=Decimal("300000.00"),
|
||
nav=PRODUCT_SPECS["159992"].nav,
|
||
confirmed_at=now - timedelta(hours=2),
|
||
),
|
||
work_order_id=work_order_id,
|
||
customer_id=12002,
|
||
account_id=account_id,
|
||
product_id=product_ids["159992"],
|
||
transaction_type="申购",
|
||
amount=Decimal("300000.00"),
|
||
nav=PRODUCT_SPECS["159992"].nav,
|
||
confirmed_at=now - timedelta(hours=2),
|
||
)
|
||
|
||
# RW-007:中风险适当性错配,等级差 1 且缺失风险揭示。
|
||
_, account_id = customer_data[12003]
|
||
work_order_id = await _ensure_work_order(
|
||
session,
|
||
work_order_no="RISKDEMO-12003-WO-001",
|
||
customer_id=12003,
|
||
product_id=product_ids["159993"],
|
||
amount=Decimal("200000.00"),
|
||
channel="APP",
|
||
risk_disclosure_ack_at=None,
|
||
second_confirmation_at=None,
|
||
now=now,
|
||
)
|
||
await _ensure_transaction(
|
||
session,
|
||
transaction_no="RISKDEMO-12003-TXN-001",
|
||
order_id=await _ensure_order(
|
||
session,
|
||
order_no="RISKDEMO-12003-ORDER-001",
|
||
customer_id=12003,
|
||
account_id=account_id,
|
||
product_id=product_ids["159993"],
|
||
transaction_type="申购",
|
||
amount=Decimal("200000.00"),
|
||
nav=PRODUCT_SPECS["159993"].nav,
|
||
confirmed_at=now - timedelta(hours=3),
|
||
),
|
||
work_order_id=work_order_id,
|
||
customer_id=12003,
|
||
account_id=account_id,
|
||
product_id=product_ids["159993"],
|
||
transaction_type="申购",
|
||
amount=Decimal("200000.00"),
|
||
nav=PRODUCT_SPECS["159993"].nav,
|
||
confirmed_at=now - timedelta(hours=3),
|
||
)
|
||
|
||
# RW-012:72 岁客户赎回 75 万,金额达到历史均值 3 倍且使用非常用设备。
|
||
_, account_id = customer_data[12004]
|
||
history_confirmed_at = now - timedelta(days=60)
|
||
await _ensure_transaction(
|
||
session,
|
||
transaction_no="RISKDEMO-12004-HIST-001",
|
||
order_id=await _ensure_order(
|
||
session,
|
||
order_no="RISKDEMO-12004-HIST-ORDER-001",
|
||
customer_id=12004,
|
||
account_id=account_id,
|
||
product_id=product_ids["159995"],
|
||
transaction_type="申购",
|
||
amount=Decimal("100000.00"),
|
||
nav=PRODUCT_SPECS["159995"].nav,
|
||
confirmed_at=history_confirmed_at,
|
||
),
|
||
work_order_id=None,
|
||
customer_id=12004,
|
||
account_id=account_id,
|
||
product_id=product_ids["159995"],
|
||
transaction_type="申购",
|
||
amount=Decimal("100000.00"),
|
||
nav=PRODUCT_SPECS["159995"].nav,
|
||
confirmed_at=history_confirmed_at,
|
||
)
|
||
await _ensure_login(
|
||
session,
|
||
customer_id=12004,
|
||
device_id="RISKDEMO-DEVICE-12004",
|
||
login_at=now - timedelta(hours=2),
|
||
is_common_device=False,
|
||
)
|
||
await _ensure_transaction(
|
||
session,
|
||
transaction_no="RISKDEMO-12004-TXN-001",
|
||
order_id=await _ensure_order(
|
||
session,
|
||
order_no="RISKDEMO-12004-ORDER-001",
|
||
customer_id=12004,
|
||
account_id=account_id,
|
||
product_id=product_ids["159995"],
|
||
transaction_type="赎回",
|
||
amount=Decimal("750000.00"),
|
||
nav=PRODUCT_SPECS["159995"].nav,
|
||
confirmed_at=now,
|
||
),
|
||
work_order_id=None,
|
||
customer_id=12004,
|
||
account_id=account_id,
|
||
product_id=product_ids["159995"],
|
||
transaction_type="赎回",
|
||
amount=Decimal("750000.00"),
|
||
nav=PRODUCT_SPECS["159995"].nav,
|
||
confirmed_at=now,
|
||
)
|
||
|
||
# RW-015 + RW-018:凌晨小额交易,关联有效自动定投工单,预期合并预警。
|
||
_, account_id = customer_data[12005]
|
||
night_at = _local_timestamp(2, 15)
|
||
work_order_id = await _ensure_work_order(
|
||
session,
|
||
work_order_no="RISKDEMO-12005-WO-001",
|
||
customer_id=12005,
|
||
product_id=product_ids["159994"],
|
||
amount=Decimal("5000.00"),
|
||
channel="自动定投",
|
||
risk_disclosure_ack_at=None,
|
||
second_confirmation_at=None,
|
||
now=now,
|
||
)
|
||
await _ensure_transaction(
|
||
session,
|
||
transaction_no="RISKDEMO-12005-TXN-001",
|
||
order_id=await _ensure_order(
|
||
session,
|
||
order_no="RISKDEMO-12005-ORDER-001",
|
||
customer_id=12005,
|
||
account_id=account_id,
|
||
product_id=product_ids["159994"],
|
||
transaction_type="申购",
|
||
amount=Decimal("5000.00"),
|
||
nav=PRODUCT_SPECS["159994"].nav,
|
||
confirmed_at=night_at,
|
||
),
|
||
work_order_id=work_order_id,
|
||
customer_id=12005,
|
||
account_id=account_id,
|
||
product_id=product_ids["159994"],
|
||
transaction_type="申购",
|
||
amount=Decimal("5000.00"),
|
||
nav=PRODUCT_SPECS["159994"].nav,
|
||
confirmed_at=night_at,
|
||
)
|
||
|
||
|
||
async def check(session: AsyncSession) -> int:
|
||
"""只读检查演示数据是否完整,返回非零表示存在缺口。"""
|
||
missing: list[str] = []
|
||
for spec in CUSTOMER_SPECS.values():
|
||
user = await session.get(RiskUser, spec.customer_id)
|
||
profile = await session.get(FundCustomerProfile, spec.customer_id)
|
||
account = await session.scalar(
|
||
select(FundSimAccount.id).where(FundSimAccount.customer_id == spec.customer_id)
|
||
)
|
||
if user is None or profile is None or account is None:
|
||
missing.append(f"客户 {spec.customer_id} 的账号、画像或账户缺失")
|
||
continue
|
||
if user.investor_type != spec.investor_type or profile.investor_type != spec.investor_type:
|
||
missing.append(f"客户 {spec.customer_id} 的风险等级与 sys_user / 画像不一致")
|
||
|
||
for scenario in SCENARIOS:
|
||
count = await session.scalar(
|
||
select(func.count(FundTransaction.id)).where(
|
||
FundTransaction.customer_id == scenario.customer_id,
|
||
FundTransaction.transaction_no.like(f"RISKDEMO-{scenario.customer_id}-%"),
|
||
)
|
||
)
|
||
if not count:
|
||
missing.append(
|
||
f"场景 {scenario.label} 缺少交易数据,预期规则 {'、'.join(scenario.expected_rules)}"
|
||
)
|
||
|
||
if missing:
|
||
for item in missing:
|
||
print(f"[缺失] {item}")
|
||
return 1
|
||
|
||
print("风控演示数据检查通过。")
|
||
for scenario in SCENARIOS:
|
||
print(
|
||
f" 客户 {scenario.customer_id}:{scenario.label},"
|
||
f"预期规则 {'、'.join(scenario.expected_rules)}"
|
||
)
|
||
return 0
|
||
|
||
|
||
async def dry_run_scan(session: AsyncSession) -> int:
|
||
"""在事务内运行规则引擎,打印结果后回滚,用于校验场景是否真正命中。"""
|
||
alerts = await RiskRuleEngine(session).refresh_alerts()
|
||
demo_alerts = [alert for alert in alerts if alert.customer_id in DEMO_CUSTOMERS]
|
||
for alert in demo_alerts:
|
||
print(
|
||
f" [命中] 客户 {alert.customer_id}:{'、'.join(alert.trigger_rule_codes)} "
|
||
f"({alert.alert_level})"
|
||
)
|
||
await session.rollback()
|
||
if not demo_alerts:
|
||
print("没有命中新的演示预警。")
|
||
return 1
|
||
print(f"规则干跑完成,共命中 {len(demo_alerts)} 条演示预警,已回滚,不会写库。")
|
||
return 0
|
||
|
||
|
||
def print_plan() -> None:
|
||
"""打印默认演示计划,默认模式不连接数据库。"""
|
||
print("风控演示数据脚本,当前为 dry-run 模式,不会写入数据库。")
|
||
print("固定演示号段:客户 12001-12005,产品 159991-159995。")
|
||
for scenario in SCENARIOS:
|
||
print(
|
||
f" 客户 {scenario.customer_id}:{scenario.label},"
|
||
f"预期规则 {'、'.join(scenario.expected_rules)}"
|
||
)
|
||
print("确认无误后使用 --apply 正式写入。")
|
||
|
||
|
||
async def run(args: argparse.Namespace) -> int:
|
||
if args.check:
|
||
async with SessionFactory() as session:
|
||
return await check(session)
|
||
|
||
if args.apply:
|
||
async with SessionFactory() as session:
|
||
async with session.begin():
|
||
await seed(session)
|
||
status = await check(session)
|
||
if status != 0:
|
||
return status
|
||
if args.dry_run_scan:
|
||
return await dry_run_scan(session)
|
||
return 0
|
||
|
||
if args.dry_run_scan:
|
||
async with SessionFactory() as session:
|
||
status = await check(session)
|
||
if status != 0:
|
||
return status
|
||
return await dry_run_scan(session)
|
||
|
||
print_plan()
|
||
return 0
|
||
|
||
|
||
def main() -> int:
|
||
parser = argparse.ArgumentParser(description="生成风控扫描演示上游数据")
|
||
parser.add_argument("--check", action="store_true", help="只检查现有数据,不写入")
|
||
parser.add_argument(
|
||
"--dry-run",
|
||
action="store_true",
|
||
help="只打印生成计划,不连接数据库(默认行为)",
|
||
)
|
||
parser.add_argument(
|
||
"--apply",
|
||
action="store_true",
|
||
help="正式写入上游演示数据",
|
||
)
|
||
parser.add_argument(
|
||
"--dry-run-scan",
|
||
action="store_true",
|
||
help="运行规则干跑并在结束后回滚,不生成正式预警",
|
||
)
|
||
args = parser.parse_args()
|
||
return asyncio.run(run(args))
|
||
|
||
|
||
if __name__ == "__main__":
|
||
sys.exit(main())
|