retriever.py 3.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103
  1. from langchain_chroma import Chroma
  2. from langchain_openai import ChatOpenAI
  3. from app.config import settings
  4. from app.rag.embeddings import get_embeddings
  5. from app.tools.java_client import JavaClient
  6. from typing import Optional
  7. import logging
  8. logger = logging.getLogger(__name__)
  9. class RagRetriever:
  10. """简单的向量检索器 - 基于 ChromaDB"""
  11. def __init__(self, collection_name: str = "cfc_knowledge"):
  12. embeddings = get_embeddings()
  13. self.vectorstore = Chroma(
  14. collection_name=collection_name,
  15. embedding_function=embeddings,
  16. persist_directory=settings.chroma_db_path,
  17. )
  18. self.java_client = JavaClient()
  19. async def initialize(self):
  20. """从 Java 侧拉取知识库,更新到向量库"""
  21. try:
  22. articles = await self.java_client.get_published_articles()
  23. if not articles:
  24. logger.info("知识库初始化:无已发布文章")
  25. return
  26. # 防御式处理:容忍缺字段 / tags 为字符串或列表
  27. texts = []
  28. metadatas = []
  29. for a in articles:
  30. if not isinstance(a, dict):
  31. continue
  32. title = a.get("title", "")
  33. if not title:
  34. continue
  35. tags = a.get("tags", []) or []
  36. tags_str = tags if isinstance(tags, str) else " ".join(str(t) for t in tags)
  37. texts.append(f"{title} {a.get('summary', '')} {tags_str}")
  38. metadatas.append({
  39. "id": str(a["id"]),
  40. "title": title,
  41. "summary": a.get("summary", ""),
  42. "tags": tags_str,
  43. "type": "article",
  44. })
  45. if not texts:
  46. logger.info("知识库初始化:无有效文章内容")
  47. return
  48. # upsert: 先删除旧的粗粒度文章索引(type=article), 再写入最新数据,
  49. # 避免每次 worker 启动重复追加导致 embedding 无限累积
  50. try:
  51. self.vectorstore._collection.delete(where={"type": "article"})
  52. except Exception as e:
  53. logger.warning("知识库初始化: 删除旧文章索引失败: %s", e)
  54. # 添加最新文章索引
  55. await self.vectorstore.aadd_texts(texts=texts, metadatas=metadatas)
  56. logger.info("知识库向量化完成:%d 篇文章", len(metadatas))
  57. except Exception as e:
  58. logger.warning("知识库初始化失败:%s", e)
  59. async def retrieve(
  60. self,
  61. query: str,
  62. filters: Optional[dict] = None,
  63. k: int = 5,
  64. ) -> list[dict]:
  65. """向量相似度检索"""
  66. results = []
  67. # 向量检索
  68. filter_query = {"user_id": filters["user_id"]} if filters and "user_id" in filters else None
  69. docs = self.vectorstore.similarity_search(
  70. query,
  71. k=k * 2, # 多取一些
  72. filter=filter_query,
  73. )
  74. # 格式化结果并去重
  75. seen = set()
  76. for doc in docs:
  77. content_hash = hash(doc.page_content[:100])
  78. if content_hash in seen:
  79. continue
  80. seen.add(content_hash)
  81. results.append({
  82. "content": doc.page_content,
  83. "metadata": doc.metadata,
  84. "score": doc.metadata.get("score", 0) if hasattr(doc, "metadata") else 0,
  85. })
  86. if len(results) >= k:
  87. break
  88. return results[:k]