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