From b99b6b7dba330832bde8f7712857d86fc8696702 Mon Sep 17 00:00:00 2001 From: "qinyong@9artedu.com" Date: Mon, 15 Jun 2026 12:21:44 +0800 Subject: [PATCH] =?UTF-8?q?=E6=8F=90=E4=BA=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app/core/__init__.py | 0 app/core/context.py | 4 +++ app/core/log.py | 68 ++++++++++++++++++++++++++++++++++++++++++ app/summary/service.py | 24 +++++++-------- 4 files changed, 84 insertions(+), 12 deletions(-) create mode 100644 app/core/__init__.py create mode 100644 app/core/context.py create mode 100644 app/core/log.py diff --git a/app/core/__init__.py b/app/core/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/app/core/context.py b/app/core/context.py new file mode 100644 index 0000000..f98aecb --- /dev/null +++ b/app/core/context.py @@ -0,0 +1,4 @@ +from contextvars import ContextVar + +request_id_ctx_var = ContextVar("request_id", default="1") + diff --git a/app/core/log.py b/app/core/log.py new file mode 100644 index 0000000..1d49f6a --- /dev/null +++ b/app/core/log.py @@ -0,0 +1,68 @@ +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 = ( + "{time:YYYY-MM-DD HH:mm:ss.SSS} | " + "{level: <8} | " + "request_id - {extra[request_id]} | " + "{name}:{function}:{line} - " + "{message}" +) + + +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()) diff --git a/app/summary/service.py b/app/summary/service.py index e0a030b..c222598 100644 --- a/app/summary/service.py +++ b/app/summary/service.py @@ -12,7 +12,7 @@ from app.llm import llm from app.models.mysql import ArchiveMessages from app.repository.archive_messages_repository import ArchiveMessagesRepository from app.repository.milvus.summary_repository import SummaryRepository - +from app.core.log import logger class SummaryService: @@ -28,14 +28,14 @@ class SummaryService: if date: day_ArchiveMessages = await self.archive_messages_repository.get_message_statistics_by_date(start_date=date, end_date=date) - print(f"{date}共产生了{len(day_ArchiveMessages)}条会话") + logger.info(f"{date}共产生了{len(day_ArchiveMessages)}条会话") else: day_ArchiveMessages = await self.archive_messages_repository.get_message_statistics_by_date() # 改成map date_user_message_a_set = set() for idx, day_message in enumerate(day_ArchiveMessages, 1): #打印计数 - print(f"处理第{idx}会话") + logger.info(f"处理第{idx}会话") day_date = str(day_message.created_at) 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: continue - print(f"查询{date_user_a}和{date_user_b}对话开始") + logger.info(f"查询{date_user_a}和{date_user_b}对话开始") list1 = await self.archive_messages_repository.get_message_statistics_by_date_user(day_date, date_user_a, date_user_b) @@ -69,16 +69,16 @@ class SummaryService: date_user_message_a_set.add(date_user_b_clean + date_user_a_append) coro_list = list1 + list2 coro_list.sort(key=lambda x: x.created_at, reverse=False) - print(f"查询{date_user_a}和{date_user_b}对话共:{len(coro_list)}") + logger.info(f"查询{date_user_a}和{date_user_b}对话共:{len(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_res, 'content') else str(filter_gossip_message_res) - print(f"查询{date_user_a}和{date_user_b}对话进行保留语义骨架,去除闲聊填充") + logger.info(f"查询{date_user_a}和{date_user_b}对话进行保留语义骨架,去除闲聊填充") # 摘要 summary_res = await self.date_message_summary(coro_list) summary_text = str(summary_res.content) if hasattr(summary_res, 'content') else str(summary_res) - print(f"查询{date_user_a}和{date_user_b}对话进行摘要") + logger.info(f"查询{date_user_a}和{date_user_b}对话进行摘要") # 摘要向量 batch_embeddings = await self.embedding_client.aembed_documents( [summary_text, filter_gossip_message_text]) @@ -96,10 +96,10 @@ class SummaryService: message_dense_vector=[batch_embeddings[1]], summary_dense_vector=[batch_embeddings[0]] ) - print(f"查询{date_user_a}和{date_user_b}对话成功入 库") + logger.info(f"查询{date_user_a}和{date_user_b}对话成功入库") except Exception as e: - print(e) - print(f"运行结束:{(time.time()-stime)}") + logger.error(f"插入Milvus失败: {e}") + logger.info(f"运行结束:{(time.time()-stime)}") async def filter_gossip_message(self, coro_list: List[ArchiveMessages]) -> str: prompt = """ @@ -126,7 +126,7 @@ class SummaryService: prompt_template = PromptTemplate(template=prompt, input_variables=["message_str"]) chain = prompt_template | llm result = await chain.ainvoke({"message_str": message_str}) - print(result) + logger.info(f"filter_gossip_message result: {result}") return result 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"]) chain = prompt_template | llm result = await chain.ainvoke({"message_str": message_str}) - print(result) + logger.info(f"date_message_summary result: {result}") return result