Compare commits

..

1 Commits

Author SHA1 Message Date
jqb
732b8f32ca 111 2026-06-15 10:50:54 +08:00
7 changed files with 14 additions and 88 deletions

1
.gitignore vendored
View File

@ -2,4 +2,3 @@
.trae .trae
.venv .venv
docker/ docker/
.idea

View File

@ -1 +1 @@
【日期】 2026年6月11日 \n【沟通对象】 客户Janice \n** 核心诉求 (User Intent)** \n咨询线上角色建模相关课程销售进一步了解其年级与专业背景。 \n** 处理结果 (Resolution)** \n初步确认客户对角色建模有兴趣正在收集基本信息。 \n** 待办事项 (Action Items)** \n \n** 关键标签**#角色建模#初步咨询 【日期】 2026年6月11日 \n【沟通对象】 客户Janice \n** 核心诉求 (User Intent)** \n咨询线上角色建模相关课程销售进一步了解其年级与专业背景。 \n** 处理结果 (Resolution)** \n初步确认客户对角色建模有兴趣正在收集基本信息。 \n** 待办事项 (Action Items)** \n \n** 关键标签**#角色建模#初步咨询 11

View File

View File

@ -1,4 +0,0 @@
from contextvars import ContextVar
request_id_ctx_var = ContextVar("request_id", default="1")

View File

@ -1,68 +0,0 @@
import asyncio
import sys
from pathlib import Path
from loguru import logger
from app.conf.app_config import app_config
from app.core.context import request_id_ctx_var
log_format = (
"<green>{time:YYYY-MM-DD HH:mm:ss.SSS}</green> | "
"<level>{level: <8}</level> | "
"<magenta>request_id - {extra[request_id]}</magenta> | "
"<cyan>{name}</cyan>:<cyan>{function}</cyan>:<cyan>{line}</cyan> - "
"<level>{message}</level>"
)
def inject_request_id(record):
request_id = request_id_ctx_var.get()
record["extra"]["request_id"] = request_id
logger.remove()
logger = logger.patch(inject_request_id)
if app_config.logging.console.enable:
logger.add(sink=sys.stdout, level=app_config.logging.console.level, format=log_format)
if app_config.logging.file.enable:
path = Path(app_config.logging.file.path)
path.mkdir(parents=True, exist_ok=True)
logger.add(
sink=path / "app.log",
level=app_config.logging.file.level,
format=log_format,
rotation=app_config.logging.file.rotation,
retention=app_config.logging.file.retention,
encoding="utf-8"
)
if __name__ == '__main__':
async def graph(request: str):
# 打印日志
logger.info(request)
async def test1():
# 接收到请求
request_id_ctx_var.set("request-1")
# 模拟处理
await asyncio.sleep(1)
await graph("request-1")
async def test2():
# 接收到请求
request_id_ctx_var.set("request-2")
# 模拟处理
await asyncio.sleep(1)
await graph("request-2")
async def main():
await asyncio.gather(test1(), test2())
asyncio.run(main())

View File

@ -12,7 +12,7 @@ from app.llm import llm
from app.models.mysql import ArchiveMessages from app.models.mysql import ArchiveMessages
from app.repository.archive_messages_repository import ArchiveMessagesRepository from app.repository.archive_messages_repository import ArchiveMessagesRepository
from app.repository.milvus.summary_repository import SummaryRepository from app.repository.milvus.summary_repository import SummaryRepository
from app.core.log import logger
class SummaryService: class SummaryService:
@ -28,14 +28,14 @@ class SummaryService:
if date: if date:
day_ArchiveMessages = await self.archive_messages_repository.get_message_statistics_by_date(start_date=date, day_ArchiveMessages = await self.archive_messages_repository.get_message_statistics_by_date(start_date=date,
end_date=date) end_date=date)
logger.info(f"{date}共产生了{len(day_ArchiveMessages)}条会话") print(f"{date}共产生了{len(day_ArchiveMessages)}条会话")
else: else:
day_ArchiveMessages = await self.archive_messages_repository.get_message_statistics_by_date() day_ArchiveMessages = await self.archive_messages_repository.get_message_statistics_by_date()
# 改成map # 改成map
date_user_message_a_set = set() date_user_message_a_set = set()
for idx, day_message in enumerate(day_ArchiveMessages, 1): #打印计数 for idx, day_message in enumerate(day_ArchiveMessages, 1): #打印计数
logger.info(f"处理第{idx}会话") print(f"处理第{idx}会话")
day_date = str(day_message.created_at) day_date = str(day_message.created_at)
user_a = day_message.from_user user_a = day_message.from_user
# 找出发送人相关的接收人信息 # 找出发送人相关的接收人信息
@ -56,7 +56,7 @@ class SummaryService:
# 判断只要存在就跳过 # 判断只要存在就跳过
if date_user_a + date_user_b in date_user_message_a_set or date_user_b_clean + date_user_a_append in date_user_message_a_set: if date_user_a + date_user_b in date_user_message_a_set or date_user_b_clean + date_user_a_append in date_user_message_a_set:
continue continue
logger.info(f"查询{date_user_a}{date_user_b}对话开始") print(f"查询{date_user_a}{date_user_b}对话开始")
list1 = await self.archive_messages_repository.get_message_statistics_by_date_user(day_date, list1 = await self.archive_messages_repository.get_message_statistics_by_date_user(day_date,
date_user_a, date_user_a,
date_user_b) date_user_b)
@ -69,16 +69,16 @@ class SummaryService:
date_user_message_a_set.add(date_user_b_clean + date_user_a_append) date_user_message_a_set.add(date_user_b_clean + date_user_a_append)
coro_list = list1 + list2 coro_list = list1 + list2
coro_list.sort(key=lambda x: x.created_at, reverse=False) coro_list.sort(key=lambda x: x.created_at, reverse=False)
logger.info(f"查询{date_user_a}{date_user_b}对话共:{len(coro_list)}") print(f"查询{date_user_a}{date_user_b}对话共:{len(coro_list)}")
# 保留语义骨架,去除闲聊填充 # 保留语义骨架,去除闲聊填充
filter_gossip_message_res = await self.filter_gossip_message(coro_list) filter_gossip_message_res = await self.filter_gossip_message(coro_list)
filter_gossip_message_text = str(filter_gossip_message_res.content) if hasattr( filter_gossip_message_text = str(filter_gossip_message_res.content) if hasattr(
filter_gossip_message_res, 'content') else str(filter_gossip_message_res) filter_gossip_message_res, 'content') else str(filter_gossip_message_res)
logger.info(f"查询{date_user_a}{date_user_b}对话进行保留语义骨架,去除闲聊填充") print(f"查询{date_user_a}{date_user_b}对话进行保留语义骨架,去除闲聊填充")
# 摘要 # 摘要
summary_res = await self.date_message_summary(coro_list) summary_res = await self.date_message_summary(coro_list)
summary_text = str(summary_res.content) if hasattr(summary_res, 'content') else str(summary_res) summary_text = str(summary_res.content) if hasattr(summary_res, 'content') else str(summary_res)
logger.info(f"查询{date_user_a}{date_user_b}对话进行摘要") print(f"查询{date_user_a}{date_user_b}对话进行摘要")
# 摘要向量 # 摘要向量
batch_embeddings = await self.embedding_client.aembed_documents( batch_embeddings = await self.embedding_client.aembed_documents(
[summary_text, filter_gossip_message_text]) [summary_text, filter_gossip_message_text])
@ -96,10 +96,10 @@ class SummaryService:
message_dense_vector=[batch_embeddings[1]], message_dense_vector=[batch_embeddings[1]],
summary_dense_vector=[batch_embeddings[0]] summary_dense_vector=[batch_embeddings[0]]
) )
logger.info(f"查询{date_user_a}{date_user_b}对话成功入") print(f"查询{date_user_a}{date_user_b}对话成功入 ")
except Exception as e: except Exception as e:
logger.error(f"插入Milvus失败: {e}") print(e)
logger.info(f"运行结束:{(time.time()-stime)}") print(f"运行结束:{(time.time()-stime)}")
async def filter_gossip_message(self, coro_list: List[ArchiveMessages]) -> str: async def filter_gossip_message(self, coro_list: List[ArchiveMessages]) -> str:
prompt = """ prompt = """
@ -126,7 +126,7 @@ class SummaryService:
prompt_template = PromptTemplate(template=prompt, input_variables=["message_str"]) prompt_template = PromptTemplate(template=prompt, input_variables=["message_str"])
chain = prompt_template | llm chain = prompt_template | llm
result = await chain.ainvoke({"message_str": message_str}) result = await chain.ainvoke({"message_str": message_str})
logger.info(f"filter_gossip_message result: {result}") print(result)
return result return result
async def date_message_summary(self, coro_list: List[ArchiveMessages]) -> str: async def date_message_summary(self, coro_list: List[ArchiveMessages]) -> str:
@ -155,7 +155,7 @@ class SummaryService:
prompt_template = PromptTemplate(template=prompt, input_variables=["message_str"]) prompt_template = PromptTemplate(template=prompt, input_variables=["message_str"])
chain = prompt_template | llm chain = prompt_template | llm
result = await chain.ainvoke({"message_str": message_str}) result = await chain.ainvoke({"message_str": message_str})
logger.info(f"date_message_summary result: {result}") print(result)
return result return result

View File

@ -26,4 +26,3 @@ async def root():
if __name__ == "__main__": if __name__ == "__main__":
import uvicorn import uvicorn
uvicorn.run(app, host="0.0.0.0", port=8000) uvicorn.run(app, host="0.0.0.0", port=8000)