feat:用户画像功能

This commit is contained in:
2026-09-14 00:19:31 +08:00
parent cf8fcf63c9
commit a7f9e182a4
17 changed files with 971 additions and 37 deletions
+42
View File
@@ -0,0 +1,42 @@
"""最终用户画像路由:仅员工/管理员可查询三源整合画像。
权限约定:
- 客户(CUSTOMER):**无权访问本接口**,整合画像属于内部经营数据;
- 员工/管理员(EMPLOYEE/ADMIN):必须显式指定 customer_id,可查任意客户。
整合逻辑在 ComposedProfileService(LLM 整合 + Redis 缓存 + 降级),
本路由只做鉴权与编排;画像仅供参考,不用于适当性校验。
"""
from fastapi import APIRouter, Depends, Query
from sqlalchemy.ext.asyncio import AsyncSession
from api.deps import get_current_user
from config.deps import get_db
from model.sys_user import SysUser
from service.memory.composed_profile import ComposedProfileService
from utils.exceptions import ForbiddenError
from utils.response import success
router = APIRouter()
async def require_profile_viewer(user: SysUser = Depends(get_current_user)) -> SysUser:
"""仅员工或管理员可查询最终用户画像,客户一律拒绝。"""
if user.user_type == "ADMIN":
return user
if user.user_type != "EMPLOYEE":
raise ForbiddenError("仅员工可查询用户画像")
return user
@router.get("/profile/composed", summary="查询最终用户画像(画像+记忆+持仓三源整合,仅员工)")
async def get_composed_profile(
customer_id: int = Query(..., description="目标客户 ID"),
user: SysUser = Depends(require_profile_viewer),
db: AsyncSession = Depends(get_db),
):
"""返回整合后的最终用户画像 JSON 及降级 warnings。"""
service = ComposedProfileService()
# 完整响应体先查 Redis(20 分钟过期),重复请求直接命中缓存
data, _warnings = await service.compose_response(db, customer_id=customer_id)
return success(data)
+6 -4
View File
@@ -6,22 +6,24 @@
"enabled": true,
"aliases": ["基金产品", "产品", "基金", "基金名称", "产品名称"],
"tables": ["fin_product"],
"fields": ["product_code", "product_name", "product_type", "risk_level", "status"]
"fields": ["product_code", "product_name", "product_type", "risk_level", "status"],
"value_hint": "product_type 实际枚举值为:货币型/债券型/混合型/股票型/指数型/QDII(注意是'XX型',不是 Schema 注释里的'股票基金');status 实际枚举值为:在售/暂停/到期/清盘,查在售产品用 status = '在售'"
},
{
"term": "股票型基金",
"enabled": true,
"aliases": ["股票型", "股票基金", "权益类"],
"aliases": ["股票型", "股票基金", "股票型基金", "权益类"],
"tables": ["fin_product"],
"fields": ["product_type"],
"value_hint": "product_type 通常使用业务库中的股票型标准值,不能自行创造枚举值"
"value_hint": "fin_product.product_type 实际枚举值为:货币型/债券型/混合型/股票型/指数型/QDII;股票型基金必须用 product_type = '股票型',不要使用 Schema 注释里的'股票基金'"
},
{
"term": "客户",
"enabled": true,
"aliases": ["客户", "投资人", "持有人"],
"tables": ["fin_customer_profile", "fin_holdings", "fin_risk_assessment"],
"fields": ["customer_id", "risk_level", "customer_level"]
"fields": ["customer_id", "risk_level", "customer_level"],
"value_hint": "fin_customer_profile.risk_level / fin_risk_assessment.risk_level 实际枚举值为 R1-R5(R1 最保守,R5 最激进,与产品风险等级同一口径);查'保守型/稳健型客户'对应 R1/R2,'激进型客户'对应 R5"
},
{
"term": "持仓",
+4 -2
View File
@@ -107,10 +107,12 @@ def resolve_semantics(question: str, *, catalog: dict[str, Any] | None = None) -
{
"term": item["term"],
"field": item["fields"][0],
"hint": item.get("metric_hint", item.get("value_hint", "")),
"hint": item.get("metric_hint") or item.get("value_hint", ""),
}
# metric_hint(指标口径)与 value_hint(字段取值口径)都要进入
# SQL 生成上下文;只过滤 metric_hint 会让枚举值提示永远丢失。
for item in matched
if "metric_hint" in item
if "metric_hint" in item or "value_hint" in item
]
return {
"terms": [item["term"] for item in matched],
+6 -4
View File
@@ -50,12 +50,14 @@ INTENT_SYSTEM_PROMPT = (
"- guide_purchase: 用户询问如何购买基金、开户、注册等引导类问题\n"
"- want_advisor: 用户希望获得个性化基金推荐或投资顾问服务\n"
"- knowledge_qa: 用户询问基金相关的知识性问题,如净值、费率、风险、申赎规则等\n"
"- company_info: 用户询问华夏科技公司本身的信息,如公司全称、成立时间、牌照、总部地址、"
"客服电话、服务时间、官网、投诉渠道等\n"
"- company_info: 用户询问华夏科技公司本身的静态档案信息,如公司全称、成立时间、牌照、总部地址、"
"客服电话、服务时间、官网、投诉渠道等;"
"注意:公司的产品清单、产品数量、产品数据不属于公司信息,应归为 nl2sql_request\n"
"- nl2sql_request: 用户要求查询具体数据或账户信息,"
"如“我的持仓有哪些”“我买了多少XX基金”“我的交易记录”"
"“最近一周XX基金的净值数据”“XX基金最新的规模/费率数据”等"
"要求数据本身而非知识解释的问题\n"
"“最近一周XX基金的净值数据”“XX基金最新的规模/费率数据”"
"“你们公司有哪些股票型基金产品”“公司在售的债券基金有哪些”等"
"询问公司产品清单/产品数据或个人账户数据、要求数据本身而非知识解释的问题\n"
"- chitchat: 普通寒暄,如问候、致谢、告别、询问你是谁/你能做什么、在吗等一两句话的闲聊\n"
"- off_topic: 用户要求你实质性地处理与金融、基金、公司业务无关的事情,"
"如写代码、讲笑话、写作文、问天气、聊政治、情感咨询、做数学题等\n"
+2 -4
View File
@@ -7,13 +7,11 @@ Agent 复用同一实现、禁止各自复制。在 Agent 侧共享包发布前
"""
from __future__ import annotations
# 风险等级 → 序号(C_n / R_n 的 n)。兼容两套口径:标准 C1-C5/R1-R5 与历史中文等级,
# 避免问卷侧尚未完成 C1-C5 归一化时发送终审误判。
# 假设:中文等级与 C 级一一对应(保守=C1/稳健=C2/平衡=C3/进取=C4/激进=C5)。
# 风险等级 → 序号(C_n / R_n 的 n)。客户风险能力口径 C1-C5,客户画像与产品口径
# 统一为 R1-R5(历史中文等级已在数据迁移中转换,故不再收录中文键)。
_RISK_RANK = {
"C1": 1, "C2": 2, "C3": 3, "C4": 4, "C5": 5,
"R1": 1, "R2": 2, "R3": 3, "R4": 4, "R5": 5,
"保守": 1, "稳健": 2, "平衡": 3, "进取": 4, "激进": 5,
}
+436
View File
@@ -0,0 +1,436 @@
"""最终用户画像整合服务:fin_customer_profile + memory_unit + 持仓 三源加权。
- 权重语义是 LLM 整合时的注意力权重(配置于 sys_config:profile.weight.*),
不做算术加权——三源是异构数据。
- 只做派生视图:结果缓存在 Redis(成功 30 分钟 / 降级 5 分钟),
不落库、不回写 fin_customer_profile,适当性校验仍以问卷风评原值为准。
- LLM 失败时降级为规则拼接摘要,沿项目 warnings 机制上报,绝不阻塞调用方。
"""
from __future__ import annotations
import json
import logging
import re
from datetime import datetime
from typing import Any
from sqlalchemy.ext.asyncio import AsyncSession
from repositories.sys_config import SysConfigRepo
from tool.confidence_rank import FinalConfidenceRankTool
from .holdings import CustomerHoldingsMemory
from .long_term import LongTermMemoryService
from .profile import CustomerProfileMemory
logger = logging.getLogger(__name__)
class ComposedProfileService:
"""拉取三路数据,调用统一 LLM 整合为最终用户画像。"""
CACHE_TTL = 20 * 60
DEGRADED_CACHE_TTL = 20 * 60
RESPONSE_CACHE_TTL = 20 * 60
RESPONSE_CACHE_PREFIX = "composed_profile_response"
DEFAULT_WEIGHTS = {"profile": 0.5, "memory": 0.2, "holdings": 0.3}
WEIGHT_KEYS = {
"profile": "profile.weight.profile",
"memory": "profile.weight.memory",
"holdings": "profile.weight.holdings",
}
ENABLED_KEY = "profile.composed.enabled"
MEMORY_TOP_K_KEY = "profile.composed.memory_top_k"
DEFAULT_MEMORY_TOP_K = 10
SYSTEM_PROMPT = """
你是用户画像整合器。根据三路客户数据生成最终用户画像 JSON,只输出 JSON,不要 Markdown:
{{"summary": "一段话画像结论", "risk_tendency": "风险偏好结论",
"preference_tags": ["标签"], "behavior_traits": ["行为特征"], "advice_focus": "服务建议"}}
权重指引:结构化画像信息可信度最高(权重 {w_profile});持仓反映真实风险行为(权重 {w_holdings});
对话记忆是口头表达补充(权重 {w_memory})。三源冲突时按权重取舍,且必须在 summary 中显式说明冲突。
status 为 candidate 的记忆只能作参考,不能作为确定性结论。
只能基于输入事实归纳,禁止编造;数据缺失的字段留空数组或空字符串,不要推测。
禁止在结果中出现内部字段名、ID、SQL 或证据引用。
""".strip()
def __init__(
self,
*,
profile: CustomerProfileMemory | None = None,
long_term: LongTermMemoryService | None = None,
holdings: CustomerHoldingsMemory | None = None,
redis=None,
llm_client=None,
rank_tool: FinalConfidenceRankTool | None = None,
config_repo_cls=SysConfigRepo,
):
from config.database.redis import client as redis_client
from tool.llm import llm as default_llm
self.profile = profile or CustomerProfileMemory()
self.long_term = long_term or LongTermMemoryService()
self.holdings = holdings or CustomerHoldingsMemory()
self.redis = redis or redis_client()
self.llm_client = llm_client or default_llm
self.rank_tool = rank_tool or FinalConfidenceRankTool()
self.config_repo_cls = config_repo_cls
@staticmethod
def cache_key(customer_id: int) -> str:
"""生成画像负载缓存 Key。"""
return f"composed_profile:{customer_id}"
@staticmethod
def response_cache_key(customer_id: int) -> str:
"""生成接口完整响应缓存 Key。"""
return f"{ComposedProfileService.RESPONSE_CACHE_PREFIX}:{customer_id}"
async def compose(
self,
db: AsyncSession,
*,
customer_id: int,
query: str | None = None,
profile: dict[str, Any] | None = None,
memories: list[dict[str, Any]] | None = None,
holdings_summary: dict[str, Any] | None = None,
) -> tuple[dict[str, Any] | None, list[str]]:
"""整合最终画像;可传入 facade 已召回的数据避免重复取数。"""
warnings: list[str] = []
enabled = await self._enabled(db, warnings)
if not enabled:
return None, warnings
cached = await self._cache_get(customer_id)
if cached is not None:
return cached, warnings
weights, weight_warnings = await self._load_weights(db)
warnings.extend(weight_warnings)
top_k = await self._memory_top_k(db, warnings)
if profile is None:
profile, profile_warnings = await self.profile.get(db, customer_id)
warnings.extend(profile_warnings)
if memories is None:
memories = await self._recall_memories(db, customer_id, query, top_k, warnings)
if holdings_summary is None:
holdings_summary, holdings_warnings = await self.holdings.summary(db, customer_id)
warnings.extend(holdings_warnings)
slim_memories = self._slim_memories(memories, top_k)
if profile is None and not slim_memories and holdings_summary is None:
return None, warnings
try:
result = await self._integrate(weights, profile, slim_memories, holdings_summary)
cache_ttl = self.CACHE_TTL
except Exception as exc: # noqa: BLE001 降级不阻塞
warnings.append(f"composed_profile_llm_failed:{type(exc).__name__}")
result = self._degraded(profile, slim_memories, holdings_summary)
cache_ttl = self.DEGRADED_CACHE_TTL
payload = {
"profile": result,
"weights": weights,
"generated_at": datetime.now().isoformat(timespec="seconds"),
}
await self._cache_set(customer_id, payload, cache_ttl)
return payload, warnings
async def invalidate(self, customer_id: int) -> list[str]:
"""删除画像负载与接口响应两级缓存,供画像/持仓变更后调用。"""
warnings: list[str] = []
for key in (self.cache_key(customer_id), self.response_cache_key(customer_id)):
try:
await self.redis.delete(key)
except Exception as exc: # noqa: BLE001
warnings.append(f"composed_profile_cache_invalidate_failed:{type(exc).__name__}")
return warnings
async def compose_response(
self, db: AsyncSession, *, customer_id: int
) -> tuple[dict[str, Any] | None, list[str]]:
"""接口层入口:完整响应体先查 Redis(20 分钟过期),未命中才整合。
返回结构即接口 data:
{"customer_id": ..., "profile": <画像负载或 None>, "warnings": [...]}
profile 为 None(功能关闭或无任何输入)时不缓存,便于开关立即生效。
"""
key = self.response_cache_key(customer_id)
try:
raw = await self.redis.get(key)
if raw:
return json.loads(raw), []
except Exception: # noqa: BLE001 缓存读取失败视为未命中
pass
payload, warnings = await self.compose(db, customer_id=customer_id)
data = {
"customer_id": customer_id,
"profile": payload,
"warnings": warnings,
}
if payload is not None:
try:
await self.redis.set(
key,
json.dumps(data, ensure_ascii=False, default=str),
ex=self.RESPONSE_CACHE_TTL,
)
except Exception: # noqa: BLE001 缓存写入失败不影响返回
pass
return data, warnings
# ---- 配置 -----------------------------------------------------------
async def _enabled(self, db: AsyncSession, warnings: list[str]) -> bool:
"""读取功能开关;配置读取失败时默认开启,保证功能可用。"""
try:
value = await self.config_repo_cls(db).get_value(self.ENABLED_KEY, "true")
except Exception as exc: # noqa: BLE001
warnings.append(f"composed_config_failed:{type(exc).__name__}")
return True
return str(value).strip().lower() not in {"false", "0", "off"}
async def _load_weights(self, db: AsyncSession) -> tuple[dict[str, float], list[str]]:
"""读取三源权重;缺失用默认值,非法或权重和不为 1 时整体回退默认。"""
warnings: list[str] = []
try:
repo = self.config_repo_cls(db)
raw = {
name: await repo.get_value(key, str(default))
for name, key, default in (
(name, key, self.DEFAULT_WEIGHTS[name])
for name, key in self.WEIGHT_KEYS.items()
)
}
except Exception as exc: # noqa: BLE001
return dict(self.DEFAULT_WEIGHTS), [f"composed_config_failed:{type(exc).__name__}"]
weights: dict[str, float] = {}
valid = True
for name, value in raw.items():
try:
parsed = float(value)
except (TypeError, ValueError):
valid = False
break
if not 0.0 <= parsed <= 1.0:
valid = False
break
weights[name] = parsed
if not valid or abs(sum(weights.values()) - 1.0) > 0.01:
warnings.append("composed_weight_invalid")
return dict(self.DEFAULT_WEIGHTS), warnings
return weights, warnings
async def _memory_top_k(self, db: AsyncSession, warnings: list[str]) -> int:
"""读取进入整合的记忆条数上限,非法回退默认。"""
try:
raw = await self.config_repo_cls(db).get_value(
self.MEMORY_TOP_K_KEY, str(self.DEFAULT_MEMORY_TOP_K)
)
top_k = int(str(raw))
except Exception: # noqa: BLE001
return self.DEFAULT_MEMORY_TOP_K
if top_k <= 0:
warnings.append("composed_top_k_invalid")
return self.DEFAULT_MEMORY_TOP_K
return top_k
# ---- 数据与整合 ------------------------------------------------------
async def _recall_memories(
self,
db: AsyncSession,
customer_id: int,
query: str | None,
top_k: int,
warnings: list[str],
) -> list[dict[str, Any]]:
"""召回长期记忆并按综合置信分重排;失败降级为空列表。"""
try:
dtos, memory_warnings = await self.long_term.recall(
db, customer_id, limit=max(top_k, 10), query=query
)
warnings.extend(memory_warnings)
ranked = self.rank_tool.rank(
[dto.model_dump(mode="json") for dto in dtos], top_k=top_k
)
return ranked
except Exception as exc: # noqa: BLE001
warnings.append(f"composed_memory_recall_failed:{type(exc).__name__}")
return []
@staticmethod
def _slim_memories(
memories: list[dict[str, Any]] | None, top_k: int
) -> list[dict[str, Any]]:
"""裁剪记忆字段:只保留整合所需的最小集合,不带证据引用。"""
slim = []
for item in (memories or [])[:top_k]:
slim.append(
{
"tag": item.get("tag"),
"memory_type": item.get("memory_type"),
"content": item.get("content"),
"status": item.get("status"),
"confidence": item.get("confidence"),
}
)
return slim
async def _integrate(
self,
weights: dict[str, float],
profile: dict[str, Any] | None,
memories: list[dict[str, Any]],
holdings: dict[str, Any] | None,
) -> dict[str, Any]:
"""调用统一 LLM 整合三源数据,解析失败时重试一次。"""
prompt = [
{
"role": "system",
"content": self.SYSTEM_PROMPT.format(
w_profile=weights["profile"],
w_memory=weights["memory"],
w_holdings=weights["holdings"],
),
},
{
"role": "user",
"content": json.dumps(
{
"weights": weights,
"profile": profile,
"memories": memories,
"holdings": holdings,
},
ensure_ascii=False,
default=str,
),
},
]
last_error: Exception | None = None
for _ in range(2):
try:
# max_tokens 传 None:沿用全局 LLM_MAX_TOKENS。思考型模型(如
# qwen3)会先输出推理段,限太小会截断 JSON 导致解析失败。
text = await self.llm_client.chat(
prompt, temperature=0, max_tokens=None
)
result = self._parse_json(text)
except Exception as exc: # noqa: BLE001 网络或解析失败都重试一次
last_error = exc
logger.warning(
"composed profile integrate attempt failed: %s: %s; raw=%r",
type(exc).__name__,
exc,
(text if "text" in locals() else "")[:300],
)
continue
if isinstance(result.get("summary"), str) and result["summary"].strip():
return result
last_error = ValueError("整合结果缺少 summary 字段")
raise last_error or ValueError("LLM 整合结果非法")
_THINK_PATTERN = re.compile(r"<think>.*?</think>", re.S | re.I)
@classmethod
def _parse_json(cls, text: str) -> dict[str, Any]:
"""解析 LLM JSON 输出:剥离思考标签和围栏,容忍正文夹杂说明文字。"""
payload = text.strip()
payload = cls._THINK_PATTERN.sub("", payload).strip()
fenced = re.search(r"```(?:json)?\s*(.*?)\s*```", payload, re.S | re.I)
if fenced:
payload = fenced.group(1).strip()
try:
data = json.loads(payload)
except json.JSONDecodeError:
# 模型可能在 JSON 前后加了说明文字:提取首个 { 到最后一个 } 的片段重试
start = payload.find("{")
end = payload.rfind("}")
if start < 0 or end <= start:
raise
data = json.loads(payload[start : end + 1])
if not isinstance(data, dict):
raise ValueError("整合结果必须是 JSON 对象")
return data
@staticmethod
def _degraded(
profile: dict[str, Any] | None,
memories: list[dict[str, Any]],
holdings: dict[str, Any] | None,
) -> dict[str, Any]:
"""LLM 失败时的规则拼接摘要,保证接口始终有可用输出。"""
parts: list[str] = []
if profile:
parts.append(
f"画像:风险等级 {profile.get('risk_level') or '未知'}"
f"(评分 {profile.get('risk_score') or '未知'})"
)
if memories:
confirmed = [m["content"] for m in memories if m.get("status") == "confirmed"]
chosen = confirmed or [m["content"] for m in memories]
parts.append("记忆:" + ";".join(chosen[:3]))
if holdings:
parts.append(
f"持仓:{holdings.get('holding_count', 0)} 只在持"
f",总市值 {holdings.get('total_market_value', 0)}"
)
return {
"summary": "。".join(parts) if parts else "暂无可用客户信息",
"risk_tendency": (profile or {}).get("risk_level") or "",
"preference_tags": [m.get("tag") for m in memories if m.get("tag")][:5],
"behavior_traits": [],
"advice_focus": "",
"degraded": True,
"sources": {
"profile": profile is not None,
"memory": bool(memories),
"holdings": holdings is not None,
},
}
# ---- 缓存 -------------------------------------------------------------
async def _cache_get(self, customer_id: int) -> dict[str, Any] | None:
"""读取缓存;读取或解析失败视为未命中。"""
try:
raw = await self.redis.get(self.cache_key(customer_id))
if raw:
return json.loads(raw)
except Exception: # noqa: BLE001
return None
return None
async def _cache_set(
self, customer_id: int, payload: dict[str, Any], ttl: int
) -> None:
"""写入缓存;失败静默(缓存只是加速手段)。"""
try:
await self.redis.set(
self.cache_key(customer_id),
json.dumps(payload, ensure_ascii=False, default=str),
ex=ttl,
)
except Exception: # noqa: BLE001
return
async def invalidate_composed_profile(customer_id: int) -> None:
"""轻量失效入口:画像回写、申购/赎回落账后调用,静默失败不阻塞业务。"""
from config.database.redis import client as redis_client
try:
redis = redis_client()
await redis.delete(ComposedProfileService.cache_key(customer_id))
await redis.delete(ComposedProfileService.response_cache_key(customer_id))
except Exception: # noqa: BLE001 缓存失效失败只影响新鲜度,不阻塞业务事务
return
__all__ = ["ComposedProfileService", "invalidate_composed_profile"]
+2
View File
@@ -15,6 +15,7 @@ def build_customer_memory_context(
long_term_memories: list[MemoryUnitDTO],
customer_relations: list[dict] | None = None,
customer_products: list[dict] | None = None,
composed_profile: dict | None = None,
warnings: list[str] | None = None,
) -> CustomerMemoryContext:
"""构造稳定的客服记忆上下文结构。"""
@@ -23,6 +24,7 @@ def build_customer_memory_context(
session_id=session_id,
short_term_messages=short_term_messages,
customer_profile=customer_profile,
composed_profile=composed_profile,
work_orders=work_orders,
long_term_memories=long_term_memories,
customer_relations=customer_relations or [],
+39
View File
@@ -9,9 +9,11 @@ from config.database.mysql import get_session_factory
from tool.confidence_rank import FinalConfidenceRankTool
from .archive import ConversationArchiver
from .composed_profile import ComposedProfileService
from .customer_relation import CustomerRelationMemory
from .customer_product import CustomerProductMemory
from .context_builder import build_customer_memory_context
from .holdings import CustomerHoldingsMemory
from .long_term import LongTermMemoryService
from .interest_topic import InterestTopicTracker
from .profile import CustomerProfileMemory
@@ -32,6 +34,8 @@ class MemoryService:
work_orders: WorkOrderMemory | None = None,
relations: CustomerRelationMemory | None = None,
products: CustomerProductMemory | None = None,
holdings: CustomerHoldingsMemory | None = None,
composed: ComposedProfileService | None = None,
long_term: LongTermMemoryService | None = None,
archiver: ConversationArchiver | None = None,
rank_tool: FinalConfidenceRankTool | None = None,
@@ -43,6 +47,13 @@ class MemoryService:
self.work_orders = work_orders or WorkOrderMemory()
self.relations = relations or CustomerRelationMemory()
self.products = products or CustomerProductMemory()
self.holdings = holdings or CustomerHoldingsMemory()
self.composed = composed or ComposedProfileService(
profile=self.profile,
long_term=self.long_term,
holdings=self.holdings,
redis=self.short_term.redis,
)
self.long_term = long_term or LongTermMemoryService()
self.archiver = archiver or ConversationArchiver(short_term=self.short_term)
self.rank_tool = rank_tool or FinalConfidenceRankTool()
@@ -109,11 +120,35 @@ class MemoryService:
top_k=limit,
)
ranked_memories = [MemoryUnitDTO.model_validate(item) for item in ranked]
holdings_summary = None
try:
holdings_summary, holdings_warnings = await self.holdings.summary(
db, customer_id
)
warnings.extend(holdings_warnings)
except Exception as exc:
warnings.append(f"holdings_recall_failed:{type(exc).__name__}")
composed_payload: dict[str, Any] | None = None
try:
composed_payload, composed_warnings = await self.composed.compose(
db,
customer_id=customer_id,
query=query,
profile=profile,
memories=ranked,
holdings_summary=holdings_summary,
)
warnings.extend(composed_warnings)
except Exception as exc:
warnings.append(f"composed_profile_failed:{type(exc).__name__}")
return build_customer_memory_context(
customer_id=customer_id,
session_id=session_id,
short_term_messages=short_term_messages,
customer_profile=profile,
composed_profile=composed_payload,
work_orders=work_orders,
customer_relations=customer_relations,
customer_products=customer_products,
@@ -155,6 +190,10 @@ class MemoryService:
"""记录一次重复兴趣主题信号,达到阈值后允许写入长期记忆。"""
return await self.interest_tracker.record(customer_id, tag)
async def invalidate_composed(self, customer_id: int) -> list[str]:
"""失效最终画像缓存,供画像回写/持仓变更后调用。"""
return await self.composed.invalidate(customer_id)
async def close_session(
self,
*,
+112
View File
@@ -0,0 +1,112 @@
"""客户持仓记忆:fin_holdings 聚合摘要,供最终画像整合使用。
只输出聚合结果(总市值、产品类型分布、盈亏概览、前三大持仓),
不向 LLM 或接口消费方暴露原始持仓行,控制 token 并避免明细泄露。
"""
from __future__ import annotations
from decimal import Decimal
from typing import Any
from sqlalchemy.ext.asyncio import AsyncSession
from repositories.fin_holdings import FinHoldingsRepo
class CustomerHoldingsMemory:
"""读取并聚合客户当前持仓,输出画像整合所需的持仓摘要。"""
TOP_HOLDINGS_LIMIT = 3
def __init__(self, *, repository_factory=FinHoldingsRepo):
self.repository_factory = repository_factory
async def summary(
self, db: AsyncSession, customer_id: int
) -> tuple[dict[str, Any] | None, list[str]]:
"""聚合持有中持仓;无持仓返回 None,异常返回降级 warnings。"""
try:
rows = await self.repository_factory(db).list_with_products(
customer_id, include_closed=False
)
except Exception as exc: # noqa: BLE001 沿用记忆模块降级约定
return None, [f"holdings_recall_failed:{type(exc).__name__}"]
if not rows:
return None, []
total_value = Decimal("0")
total_cost = Decimal("0")
total_pnl = Decimal("0")
win_count = 0
lose_count = 0
mix: dict[str, Decimal] = {}
ranked: list[tuple[Decimal, str]] = []
for holding, product in rows:
value = self._to_decimal(holding.current_value)
total_value += value
total_cost += self._to_decimal(holding.cost_amount)
total_pnl += self._to_decimal(holding.profit_loss)
ratio = holding.profit_ratio
if ratio is not None:
if ratio > 0:
win_count += 1
elif ratio < 0:
lose_count += 1
product_type = (product.product_type if product else None) or "未知"
mix[product_type] = mix.get(product_type, Decimal("0")) + value
name = product.product_name if product else f"产品{holding.product_id}"
ranked.append((value, name))
payload: dict[str, Any] = {
"holding_count": len(rows),
"total_market_value": float(total_value),
"total_cost": float(total_cost),
"product_type_mix": self._ratio_map(mix, total_value),
"profit_summary": {
"total_pnl": float(total_pnl),
"total_pnl_ratio": (
round(float(total_pnl / total_cost), 4)
if total_cost > 0
else None
),
"win_count": win_count,
"lose_count": lose_count,
},
"top_holdings": self._top_holdings(ranked, total_value),
}
return payload, []
@staticmethod
def _to_decimal(value: Any) -> Decimal:
"""把 ORM Decimal / float / None 统一转换为 Decimal,防御脏数据类型。"""
if value is None:
return Decimal("0")
try:
return value if isinstance(value, Decimal) else Decimal(str(value))
except Exception: # noqa: BLE001
return Decimal("0")
@staticmethod
def _ratio_map(values: dict[str, Decimal], total: Decimal) -> dict[str, float]:
"""把各维度市值转换为相对总市值的占比,总量为 0 时返回空表。"""
if total <= 0:
return {}
return {
key: round(float(value / total), 4)
for key, value in sorted(values.items(), key=lambda item: -item[1])
}
def _top_holdings(
self, ranked: list[tuple[Decimal, str]], total: Decimal
) -> list[dict[str, Any]]:
"""按市值取前三大持仓,返回名称与占比。"""
if total <= 0:
return []
ranked.sort(key=lambda item: -item[0])
return [
{"name": name, "weight": round(float(value / total), 4)}
for value, name in ranked[: self.TOP_HOLDINGS_LIMIT]
]
__all__ = ["CustomerHoldingsMemory"]
+1
View File
@@ -99,6 +99,7 @@ class CustomerMemoryContext(BaseModel):
session_id: str
short_term_messages: list[ShortTermMessage] = Field(default_factory=list)
customer_profile: dict[str, Any] | None = None
composed_profile: dict[str, Any] | None = None
work_orders: list[dict[str, Any]] = Field(default_factory=list)
long_term_memories: list[MemoryUnitDTO] = Field(default_factory=list)
customer_relations: list[dict[str, Any]] = Field(default_factory=list)
+5 -5
View File
@@ -17,6 +17,7 @@ from repositories.fin_product import FinProductRepo
from repositories.trade_order import TradeOrderRepo
from schemas.holdings import HoldingResp
from schemas.purchase import PurchaseResp
from service.memory.composed_profile import invalidate_composed_profile
from service.risk.engine import RiskEngine, summarize
from service.risk.settle import settle
from utils.exceptions import (
@@ -30,11 +31,8 @@ from utils.order_no import gen_order_no
_MONEY = Decimal("0.01")
_SHARES = Decimal("0.0001")
# 风险等级 → 序号。兼容两套口径:R1~R5 与 保守~激进(同一映射)。
_RISK_RANK = {
"R1": 1, "R2": 2, "R3": 3, "R4": 4, "R5": 5,
"保守": 1, "稳健": 2, "平衡": 3, "进取": 4, "激进": 5,
}
# 风险等级 → 序号。客户画像与产品统一 R1~R5 口径(历史中文等级已在数据迁移中转换)。
_RISK_RANK = {"R1": 1, "R2": 2, "R3": 3, "R4": 4, "R5": 5}
def _risk_rank(level: str | None) -> int | None:
@@ -137,6 +135,8 @@ async def purchase(
order.id, status="已确认", confirm_time=now
)
await db.commit()
# 持仓变更后失效最终画像缓存(失败静默,不影响交易结果)
await invalidate_composed_profile(user.id)
account = await account_repo.get_by_customer_id(user.id)
holding = await holdings_repo.get_by_customer_product(user.id, product_id)
return PurchaseResp(
+10 -6
View File
@@ -8,15 +8,17 @@ from sqlalchemy.ext.asyncio import AsyncSession
from model.fin_risk_assessment import FinRiskAssessment
from repositories.questionnaire import QuestionRepo, QuestionnaireRepo
from repositories.risk_assessment import CustomerProfileRepo, RiskAssessmentRepo
from service.memory.composed_profile import invalidate_composed_profile
from utils.exceptions import NotFoundError, ParamError
# 分数 → 风险等级映射(总分 0-100,20 分一档;后续可迁移到 sys_config 运营化)
# 等级口径统一为 R1-R5:R1 最保守、R5 最激进,与产品 fin_product.risk_level 对齐
_SCORE_LEVELS = (
(80, "激进"),
(60, "进取"),
(40, "平衡"),
(20, "稳健"),
(0, "保守"),
(80, "R5"),
(60, "R4"),
(40, "R3"),
(20, "R2"),
(0, "R1"),
)
@@ -24,7 +26,7 @@ def _score_to_level(score: int) -> str:
for threshold, level in _SCORE_LEVELS:
if score >= threshold:
return level
return "保守"
return "R1"
async def submit_assessment(
@@ -76,6 +78,8 @@ async def submit_assessment(
# 回写画像:无画像初始化,有画像更新风险等级/评分并递增版本号
await CustomerProfileRepo(db).upsert_risk(customer_id, risk_level, total_score)
# 画像变更后失效最终画像缓存(失败静默,不影响问卷提交)
await invalidate_composed_profile(customer_id)
return {
"assessment_id": record.id,
+3
View File
@@ -16,6 +16,7 @@ from repositories.fin_product import FinProductRepo
from repositories.trade_order import TradeOrderRepo
from schemas.holdings import HoldingResp
from schemas.redeem import RedeemResp
from service.memory.composed_profile import invalidate_composed_profile
from service.risk.engine import RiskEngine, summarize
from service.risk.settle import settle
from utils.exceptions import ForbiddenError, NotFoundError, ParamError
@@ -118,6 +119,8 @@ async def redeem(
order.id, status="已确认", confirm_time=now
)
await db.commit()
# 持仓变更后失效最终画像缓存(失败静默,不影响交易结果)
await invalidate_composed_profile(user.id)
account = await account_repo.get_by_customer_id(user.id)
current_holding = await holdings_repo.get_by_customer_product(
user.id, product_id
+2 -5
View File
@@ -31,11 +31,8 @@ from repositories.fin_product import FinProductRepo
from repositories.fin_transaction import FinTransactionRepo
from repositories.risk_rule import RiskRuleRepo
# 风险等级 → 序号(兼容 R1~R5 与 保守~激进 两套口径)。
_RISK_RANK = {
"R1": 1, "R2": 2, "R3": 3, "R4": 4, "R5": 5,
"保守": 1, "稳健": 2, "平衡": 3, "进取": 4, "激进": 5,
}
# 风险等级 → 序号(客户画像与产品统一 R1~R5 口径,历史中文等级已迁移转换)。
_RISK_RANK = {"R1": 1, "R2": 2, "R3": 3, "R4": 4, "R5": 5}
_LEVEL_RANK = {"低": 1, "中": 2, "高": 3}
+38
View File
@@ -0,0 +1,38 @@
-- ---------------------------------------------------------------------
-- 风险等级口径统一迁移:保守/稳健/平衡/进取/激进 → R1-R5
-- 映射:保守=R1 / 稳健=R2 / 平衡=R3 / 进取=R4 / 激进=R5
-- 幂等:UPDATE 仅命中中文值,重复执行无副作用
-- 覆盖表:fin_customer_profile / fin_risk_assessment / portfolio_benchmark
-- ---------------------------------------------------------------------
-- 01 数据迁移(中文值 → R 级)
UPDATE fin_customer_profile SET risk_level = 'R1' WHERE risk_level = '保守';
UPDATE fin_customer_profile SET risk_level = 'R2' WHERE risk_level = '稳健';
UPDATE fin_customer_profile SET risk_level = 'R3' WHERE risk_level = '平衡';
UPDATE fin_customer_profile SET risk_level = 'R4' WHERE risk_level = '进取';
UPDATE fin_customer_profile SET risk_level = 'R5' WHERE risk_level = '激进';
UPDATE fin_risk_assessment SET risk_level = 'R1' WHERE risk_level = '保守';
UPDATE fin_risk_assessment SET risk_level = 'R2' WHERE risk_level = '稳健';
UPDATE fin_risk_assessment SET risk_level = 'R3' WHERE risk_level = '平衡';
UPDATE fin_risk_assessment SET risk_level = 'R4' WHERE risk_level = '进取';
UPDATE fin_risk_assessment SET risk_level = 'R5' WHERE risk_level = '激进';
UPDATE portfolio_benchmark SET risk_level = 'R1' WHERE risk_level = '保守';
UPDATE portfolio_benchmark SET risk_level = 'R2' WHERE risk_level = '稳健';
UPDATE portfolio_benchmark SET risk_level = 'R3' WHERE risk_level = '平衡';
UPDATE portfolio_benchmark SET risk_level = 'R4' WHERE risk_level = '进取';
UPDATE portfolio_benchmark SET risk_level = 'R5' WHERE risk_level = '激进';
-- 02 列注释与数据口径对齐(MODIFY 需带完整列定义,不改动类型/约束)
ALTER TABLE fin_customer_profile
MODIFY COLUMN risk_level VARCHAR(16) NULL
COMMENT '风险等级:R1-R5(R1 最保守,R5 最激进,与产品 risk_level 对齐)';
ALTER TABLE fin_risk_assessment
MODIFY COLUMN risk_level VARCHAR(16) NOT NULL
COMMENT '评定:R1-R5(R1 最保守,R5 最激进)';
ALTER TABLE portfolio_benchmark
MODIFY COLUMN risk_level VARCHAR(16) NOT NULL
COMMENT '客户风险等级:R1-R5(R1 最保守,R5 最激进)';
+262 -6
View File
@@ -1,8 +1,9 @@
-- =====================================================================
-- 智能公募基金系统 · MySQL 建表脚本(24 张,见开发计划 §3.1)
-- 智能公募基金系统 · MySQL 建表脚本(36 张,见开发计划 §3.1,与 model/ 目录 ORM 对齐)
-- 执行:mysql --default-character-set=utf8mb4 -u<user> -p <database> < sql/schema.sql
-- 约定:InnoDB / utf8mb4;采用"逻辑外键"(不建物理 FK,便于 Mock 数据灌入),
-- 注释即文档,NL2SQL 以 COMMENT 注入 Schema。
-- 维护:新增/调整 ORM 模型时同步更新本脚本(01-24 基础业务域,25-36 扩展域)。
-- =====================================================================
SET NAMES utf8mb4;
@@ -36,6 +37,18 @@ DROP TABLE IF EXISTS fund_nav_history;
DROP TABLE IF EXISTS fin_product;
DROP TABLE IF EXISTS fin_customer_profile;
DROP TABLE IF EXISTS sys_user;
DROP TABLE IF EXISTS nl2sql_query_history;
DROP TABLE IF EXISTS nl2sql_sensitive_field;
DROP TABLE IF EXISTS nl2sql_role_column_permission;
DROP TABLE IF EXISTS nl2sql_role_table_permission;
DROP TABLE IF EXISTS nl2sql_query_role;
DROP TABLE IF EXISTS event_log;
DROP TABLE IF EXISTS advisor_visit_record;
DROP TABLE IF EXISTS advisor_report;
DROP TABLE IF EXISTS advisor_draft;
DROP TABLE IF EXISTS advisor_todo;
DROP TABLE IF EXISTS sensitive_word;
DROP TABLE IF EXISTS fin_account;
SET FOREIGN_KEY_CHECKS = 1;
-- ---------------------------------------------------------------------
@@ -64,7 +77,7 @@ CREATE TABLE IF NOT EXISTS sys_user (
-- ---------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS fin_customer_profile (
customer_id BIGINT UNSIGNED NOT NULL COMMENT '客户ID,一对一关联 sys_user.id',
risk_level VARCHAR(16) NULL COMMENT '风险等级:保守/稳健/平衡/进取/激进',
risk_level VARCHAR(16) NULL COMMENT '风险等级:R1-R5(R1 最保守,R5 最激进,与产品 risk_level 对齐)',
risk_score INT NULL COMMENT '风险评分 0-100,越高越激进',
investment_experience VARCHAR(16) NULL COMMENT '投资经验:0-1/1-3/3-5/5-10/10年以上',
annual_income_range VARCHAR(32) NULL COMMENT '年收入区间',
@@ -87,7 +100,7 @@ CREATE TABLE IF NOT EXISTS fin_product (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
product_code VARCHAR(32) NOT NULL COMMENT '基金代码,唯一,如 F000001',
product_name VARCHAR(128) NOT NULL COMMENT '基金名称',
product_type VARCHAR(32) NOT NULL COMMENT '类型:货币/债券/混合/股票基金',
product_type VARCHAR(32) NOT NULL COMMENT '类型:货币型/债券型/混合型/股票型/指数型/QDII',
risk_level VARCHAR(8) NOT NULL COMMENT '风险等级 R1-R5',
expected_return DECIMAL(7,4) NULL COMMENT '预期年化收益率(%)',
nav DECIMAL(12,6) NULL COMMENT '最新单位净值',
@@ -191,7 +204,7 @@ CREATE TABLE IF NOT EXISTS fin_risk_assessment (
assessment_date DATE NOT NULL COMMENT '评估日期',
question_version VARCHAR(16) NULL COMMENT '问卷版本(题目变更后按版本追溯)',
total_score INT NOT NULL COMMENT '总分 0-100',
risk_level VARCHAR(16) NOT NULL COMMENT '评定:保守/稳健/平衡/进取/激进',
risk_level VARCHAR(16) NOT NULL COMMENT '评定:R1-R5(R1 最保守,R5 最激进)',
answers JSON NULL COMMENT '答题详情 [{q:1,a:B,score:10}]',
assessor_type VARCHAR(16) NOT NULL DEFAULT 'AI评估' COMMENT 'AI评估/人工评估',
valid_until DATE NOT NULL COMMENT '有效期至(默认 +1年,监管一年一评)',
@@ -308,7 +321,7 @@ CREATE TABLE IF NOT EXISTS conversation_archive (
KEY idx_agent (agent_type),
UNIQUE KEY uk_session_message (session_id, message_id),
KEY idx_agent_run (agent_run_id)
) COMMENT='会话归档表(审计回溯 + Agent 持续学习素材,归档前脱敏)';
) COMMENT='会话归档表(审计回溯 + Agent 持续学习素材)';
-- ---------------------------------------------------------------------
-- 15 知识元数据表(原文存 MySQL,替代 MinIO)
@@ -492,7 +505,7 @@ CREATE TABLE IF NOT EXISTS fund_performance (
-- ---------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS portfolio_benchmark (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
risk_level VARCHAR(16) NOT NULL COMMENT '客户风险等级:保守/稳健/平衡/进取/激进',
risk_level VARCHAR(16) NOT NULL COMMENT '客户风险等级:R1-R5(R1 最保守,R5 最激进)',
target_allocation JSON NOT NULL COMMENT '目标配置 {货币:50,债券:40,股票:10}',
drift_threshold DECIMAL(5,2) NOT NULL DEFAULT 5.00 COMMENT '偏离阈值(%),超阈值触发调仓建议',
status VARCHAR(8) NOT NULL DEFAULT '启用' COMMENT '启用/停用',
@@ -500,3 +513,246 @@ CREATE TABLE IF NOT EXISTS portfolio_benchmark (
update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_risk_level (risk_level)
) COMMENT='组合基准配置表(投顾Agent 再平衡参照,运营可调整)';
-- =====================================================================
-- 以下为扩展域建表(25-36),与 model/ 目录 ORM 一一对齐
-- =====================================================================
-- ---------------------------------------------------------------------
-- 25 客户资金账户表(现金余额,一人一户)
-- ---------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS fin_account (
customer_id BIGINT UNSIGNED NOT NULL COMMENT '客户ID,一人一户,关联 sys_user.id',
balance DECIMAL(18,2) NOT NULL DEFAULT 0 COMMENT '账户总余额(含冻结部分)',
frozen_amount DECIMAL(18,2) NOT NULL DEFAULT 0 COMMENT '冻结金额,可用余额=balance-frozen_amount(不落库)',
currency VARCHAR(8) NOT NULL DEFAULT 'CNY' COMMENT '币种',
status VARCHAR(16) NOT NULL DEFAULT '正常' COMMENT '账户状态,默认正常',
version INT NOT NULL DEFAULT 0 COMMENT '乐观锁版本号(并发扣款校验)',
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
PRIMARY KEY (customer_id)
) COMMENT='客户资金账户表(现金余额,申购扣款/赎回入账的账务载体)';
-- ---------------------------------------------------------------------
-- 26 敏感词库表(合规拦截唯一来源)
-- ---------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS sensitive_word (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
word VARCHAR(128) NOT NULL COMMENT '敏感词',
category VARCHAR(32) NOT NULL COMMENT '敏感词分类(工作台可配置)',
level VARCHAR(8) NOT NULL DEFAULT '高' COMMENT '级别,默认高',
status VARCHAR(8) NOT NULL DEFAULT '启用' COMMENT '启用/停用',
update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_word (word)
) COMMENT='敏感词库表(工作台与Agent同源,合规拦截唯一来源)';
-- ---------------------------------------------------------------------
-- 27 投顾待办表(事件/定时/计算三类来源统一承载)
-- ---------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS advisor_todo (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
todo_type VARCHAR(32) NOT NULL COMMENT '待办类型(事件/定时/计算来源的细分类型)',
customer_id BIGINT UNSIGNED NULL COMMENT '关联客户(可为空)',
advisor_id BIGINT UNSIGNED NOT NULL COMMENT '所属投顾',
source VARCHAR(32) NOT NULL COMMENT '来源:事件/定时/计算',
biz_id VARCHAR(64) NULL COMMENT '关联业务ID(与 todo_type 组合定位业务)',
priority VARCHAR(8) NOT NULL DEFAULT '普通' COMMENT '普通/紧急/特急',
status VARCHAR(16) NOT NULL DEFAULT '待处理' COMMENT '待办状态,默认待处理',
due_at DATETIME NULL COMMENT '截止时间',
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
handle_time DATETIME NULL COMMENT '处理时间',
UNIQUE KEY uk_todo (todo_type, customer_id, biz_id),
KEY idx_advisor_status (advisor_id, status)
) COMMENT='投顾待办表(事件/定时/计算三类来源统一承载,uk_todo 幂等防重复建待办)';
-- ---------------------------------------------------------------------
-- 28 投顾 Agent 草稿表(不存 sent,不生成交易指令)
-- ---------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS advisor_draft (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
draft_id VARCHAR(64) NOT NULL COMMENT '草稿ID,全局唯一',
customer_id BIGINT UNSIGNED NOT NULL COMMENT '目标客户',
advisor_id BIGINT UNSIGNED NOT NULL COMMENT '所属投顾',
intent VARCHAR(32) NOT NULL COMMENT '报告意图,如 调仓建议',
title VARCHAR(128) NOT NULL COMMENT '草稿标题',
content TEXT NOT NULL COMMENT '报告正文',
structured_data JSON NULL COMMENT '结构化数据(调仓明细等)',
status VARCHAR(16) NOT NULL DEFAULT 'draft' COMMENT '草稿状态,默认draft',
deviation DECIMAL(10,4) NULL COMMENT '组合偏离度(%)',
disclaimer_ok TINYINT(1) NOT NULL DEFAULT 0 COMMENT '免责声明确认',
warning VARCHAR(512) NULL COMMENT '合规提示/警告信息',
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_draft_id (draft_id),
CONSTRAINT ck_advisor_draft_status CHECK (status IN ('draft', 'discarded')),
KEY idx_customer (customer_id),
KEY idx_advisor (advisor_id),
KEY idx_intent (intent),
KEY idx_status (status),
KEY idx_create_time (create_time)
) COMMENT='投顾 Agent 草稿(不存 sent,不生成交易指令)';
-- ---------------------------------------------------------------------
-- 29 建议报告表(工作台本地镜像:sent 状态 + 编辑/发送留痕)
-- ---------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS advisor_report (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
report_id VARCHAR(64) NOT NULL COMMENT '报告ID,全局唯一',
draft_id VARCHAR(64) NOT NULL COMMENT '来源草稿ID(advisor_draft.draft_id)',
customer_id BIGINT UNSIGNED NOT NULL COMMENT '目标客户',
advisor_id BIGINT UNSIGNED NOT NULL COMMENT '发送投顾',
intent VARCHAR(32) NOT NULL COMMENT '报告意图(继承草稿)',
title VARCHAR(128) NULL COMMENT '报告标题',
content TEXT NULL COMMENT '报告正文',
edit_history JSON NULL COMMENT '编辑留痕(JSON数组)',
send_status VARCHAR(16) NOT NULL DEFAULT 'draft' COMMENT 'draft/sent/discarded,sent 为工作台本地状态不回写投顾Agent',
send_time DATETIME NULL COMMENT '发送时间',
send_by BIGINT UNSIGNED NULL COMMENT '发送人',
msg_id BIGINT UNSIGNED NULL COMMENT '发送后落 sys_message 的站内信ID',
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_report_id (report_id),
KEY idx_customer (customer_id),
KEY idx_advisor (advisor_id),
KEY idx_send_status (send_status)
) COMMENT='建议报告表(工作台本地镜像,承载 sent 状态与发送/编辑留痕)';
-- ---------------------------------------------------------------------
-- 30 投顾回访记录表(人工留痕归档)
-- ---------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS advisor_visit_record (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
customer_id BIGINT UNSIGNED NOT NULL COMMENT '回访客户',
advisor_id BIGINT UNSIGNED NOT NULL COMMENT '回访投顾',
visit_type VARCHAR(32) NOT NULL COMMENT '回访类型',
visit_time DATETIME NOT NULL COMMENT '回访时间',
summary TEXT NULL COMMENT '回访纪要',
audio_url VARCHAR(512) NULL COMMENT '录音文件地址',
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
KEY idx_customer (customer_id),
KEY idx_advisor (advisor_id)
) COMMENT='投顾回访记录表(人工留痕归档,不触发Agent记忆抽取)';
-- ---------------------------------------------------------------------
-- 31 事件日志表(事件双写落库,消费以 event_id 幂等)
-- ---------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS event_log (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
event_id VARCHAR(64) NOT NULL COMMENT '事件ID,全局唯一,消费幂等键',
event_name VARCHAR(64) NOT NULL COMMENT '事件名',
payload JSON NULL COMMENT '事件负载',
trace_id VARCHAR(64) NULL COMMENT '请求链路追踪ID',
status VARCHAR(16) NOT NULL DEFAULT 'pending' COMMENT '消费状态,默认pending',
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
consume_time DATETIME NULL COMMENT '消费完成时间',
UNIQUE KEY uk_event_id (event_id),
KEY idx_event_name (event_name),
KEY idx_status (status)
) COMMENT='事件日志表(事件双写落库,Pub/Sub 仅作实时通知,消费幂等)';
-- ---------------------------------------------------------------------
-- 32 NL2SQL 查询角色(按 sys_user.employee_role 映射)
-- ---------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS nl2sql_query_role (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
role_code VARCHAR(64) NOT NULL COMMENT '角色编码,全局唯一',
role_name VARCHAR(128) NOT NULL COMMENT '角色名称',
employee_role VARCHAR(32) NOT NULL COMMENT '对应 sys_user.employee_role,一对一映射',
can_query TINYINT(1) NOT NULL DEFAULT 0 COMMENT '是否允许 NL2SQL 查询:0/1',
max_rows BIGINT NOT NULL DEFAULT 1000 COMMENT '单次查询最大返回行数',
daily_quota BIGINT NOT NULL DEFAULT 0 COMMENT '每日查询配额',
status VARCHAR(16) NOT NULL DEFAULT 'active' COMMENT 'active/inactive',
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_role_code (role_code),
UNIQUE KEY uk_employee_role (employee_role),
CONSTRAINT ck_nl2sql_role_can_query CHECK (can_query IN (0, 1)),
CONSTRAINT ck_nl2sql_role_status CHECK (status IN ('active', 'inactive')),
KEY idx_status (status)
) COMMENT='NL2SQL 查询角色,按 sys_user.employee_role 映射';
-- ---------------------------------------------------------------------
-- 33 NL2SQL 角色的表级 + 行级权限
-- ---------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS nl2sql_role_table_permission (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
role_id BIGINT UNSIGNED NOT NULL COMMENT '角色ID(nl2sql_query_role.id)',
table_name VARCHAR(128) NOT NULL COMMENT '授权表名',
permission VARCHAR(16) NOT NULL DEFAULT 'SELECT' COMMENT '授权动作,仅 SELECT',
row_scope_type VARCHAR(32) NOT NULL DEFAULT 'none' COMMENT '行级范围:none/customer_ids/product_ids',
row_scope_column VARCHAR(128) NULL COMMENT '行级过滤列(如 customer_id)',
status VARCHAR(16) NOT NULL DEFAULT 'active' COMMENT 'active/inactive',
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
CONSTRAINT ck_nl2sql_table_permission CHECK (permission = 'SELECT'),
CONSTRAINT ck_nl2sql_row_scope_type CHECK (row_scope_type IN ('none', 'customer_ids', 'product_ids')),
CONSTRAINT ck_nl2sql_table_status CHECK (status IN ('active', 'inactive')),
KEY idx_role (role_id),
KEY idx_status (status)
) COMMENT='NL2SQL 查询角色的表和行级权限';
-- ---------------------------------------------------------------------
-- 34 NL2SQL 角色的字段权限
-- ---------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS nl2sql_role_column_permission (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
role_id BIGINT UNSIGNED NOT NULL COMMENT '角色ID(nl2sql_query_role.id)',
table_name VARCHAR(128) NOT NULL COMMENT '表名',
column_name VARCHAR(128) NOT NULL COMMENT '列名',
access_mode VARCHAR(16) NOT NULL DEFAULT 'allow' COMMENT 'allow/deny/mask',
mask_type VARCHAR(32) NULL COMMENT '脱敏类型(access_mode=mask 时生效)',
status VARCHAR(16) NOT NULL DEFAULT 'active' COMMENT 'active/inactive',
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
CONSTRAINT ck_nl2sql_column_access_mode CHECK (access_mode IN ('allow', 'deny', 'mask')),
CONSTRAINT ck_nl2sql_column_status CHECK (status IN ('active', 'inactive')),
KEY idx_role (role_id),
KEY idx_status (status)
) COMMENT='NL2SQL 查询角色的字段权限';
-- ---------------------------------------------------------------------
-- 35 NL2SQL 全局敏感字段规则
-- ---------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS nl2sql_sensitive_field (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
table_name VARCHAR(128) NOT NULL COMMENT '表名',
column_name VARCHAR(128) NOT NULL COMMENT '列名',
mask_type VARCHAR(32) NOT NULL DEFAULT 'partial' COMMENT '脱敏方式,如 partial',
status VARCHAR(16) NOT NULL DEFAULT 'active' COMMENT 'active/inactive',
description VARCHAR(255) NULL COMMENT '规则说明',
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
CONSTRAINT ck_nl2sql_sensitive_status CHECK (status IN ('active', 'inactive')),
KEY idx_table_col (table_name, column_name),
KEY idx_status (status)
) COMMENT='NL2SQL 全局敏感字段规则';
-- ---------------------------------------------------------------------
-- 36 NL2SQL 查询归档(不保存完整结果行)
-- ---------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS nl2sql_query_history (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
query_id VARCHAR(64) NOT NULL COMMENT '查询ID,全局唯一',
user_id BIGINT UNSIGNED NOT NULL COMMENT '发起用户',
session_id VARCHAR(64) NULL COMMENT '会话ID',
caller_agent VARCHAR(64) NULL COMMENT '调用方Agent标识',
question TEXT NOT NULL COMMENT '自然语言问题',
generated_sql TEXT NULL COMMENT '生成的SQL(未执行则为空)',
access_tables JSON NULL COMMENT '实际访问的表列表',
status VARCHAR(16) NOT NULL COMMENT 'success/failed/blocked/timeout',
error_code VARCHAR(64) NULL COMMENT '错误码',
error_message VARCHAR(512) NULL COMMENT '错误信息',
row_count BIGINT NOT NULL DEFAULT 0 COMMENT '返回行数',
truncated TINYINT(1) NOT NULL DEFAULT 0 COMMENT '结果是否被截断',
elapsed_ms FLOAT NULL COMMENT '耗时(毫秒)',
trace_id VARCHAR(64) NULL COMMENT '请求链路追踪ID',
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
UNIQUE KEY uk_query_id (query_id),
CONSTRAINT ck_nl2sql_history_status CHECK (status IN ('success', 'failed', 'blocked', 'timeout')),
KEY idx_user (user_id),
KEY idx_session (session_id),
KEY idx_caller_agent (caller_agent),
KEY idx_status (status),
KEY idx_trace (trace_id),
KEY idx_create_time (create_time)
) COMMENT='NL2SQL 查询归档,不保存完整结果行';
+1 -1
View File
@@ -148,7 +148,7 @@ Agent 编排层(统一执行骨架: 记忆召回→意图路由→核心逻辑
| `sys_message` | 站内信/消息中心(user_id, type, is_read) |
| `memory_unit` | **记忆单元**(tag, content, info_type=FACT/OPINION, source, confidence, evidence/conflict/recall_count, status, valid_until)——记忆架构主体,见 §4 |
| `fund_performance` | 基金业绩指标缓存(近1月/3月/6月/1年/成立以来收益、年化波动率、最大回撤、夏普比率、计算日期)——定时任务从净值计算 |
| `portfolio_benchmark` | **组合基准配置**(risk_level → 目标资产类别权重,如 保守=货币50/债券40/股票10,可运营调整)——再平衡的参照系 |
| `portfolio_benchmark` | **组合基准配置**(risk_level → 目标资产类别权重,如 R1=货币50/债券40/股票10,可运营调整)——再平衡的参照系 |
| `sys_config` | 运营参数 KV(**高净值门槛**、偏离度阈值、风控开关等,免发版调整) |
### 3.2 Milvus 集合(3 个业务集 + 1 个记忆集)