Compare commits
1 Commits
586b529933
...
732b8f32ca
| Author | SHA1 | Date | |
|---|---|---|---|
| 732b8f32ca |
3
.gitignore
vendored
3
.gitignore
vendored
@ -1,5 +1,4 @@
|
|||||||
# 忽略Maven编译目录
|
# 忽略Maven编译目录
|
||||||
.trae
|
.trae
|
||||||
.venv
|
.venv
|
||||||
docker/
|
docker/
|
||||||
.idea
|
|
||||||
@ -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
|
||||||
@ -1,4 +0,0 @@
|
|||||||
from contextvars import ContextVar
|
|
||||||
|
|
||||||
request_id_ctx_var = ContextVar("request_id", default="1")
|
|
||||||
|
|
||||||
@ -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())
|
|
||||||
@ -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
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user