retriever.py 2.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687
  1. from langchain_chroma import Chroma
  2. from langchain_openai import ChatOpenAI, OpenAIEmbeddings
  3. from app.config import settings
  4. from app.tools.java_client import JavaClient
  5. from typing import Optional
  6. import logging
  7. logger = logging.getLogger(__name__)
  8. class RagRetriever:
  9. """简单的向量检索器 - 基于 ChromaDB"""
  10. def __init__(self, collection_name: str = "cfc_knowledge"):
  11. embeddings = OpenAIEmbeddings(
  12. model=settings.embedding_model,
  13. api_key=settings.effective_embedding_api_key,
  14. base_url=settings.effective_embedding_base_url,
  15. )
  16. self.vectorstore = Chroma(
  17. collection_name=collection_name,
  18. embedding_function=embeddings,
  19. persist_directory=settings.chroma_db_path,
  20. )
  21. self.java_client = JavaClient()
  22. async def initialize(self):
  23. """从 Java 侧拉取知识库,更新到向量库"""
  24. try:
  25. articles = await self.java_client.get_published_articles()
  26. if articles:
  27. # 构建文本和内容元数据
  28. texts = [f"{a['title']} {a.get('summary', '')} {' '.join(a.get('tags', []))}"
  29. for a in articles]
  30. metadatas = [
  31. {
  32. "id": a["id"],
  33. "title": a["title"],
  34. "summary": a.get("summary", ""),
  35. "tags": a.get("tags", []),
  36. "type": "article",
  37. }
  38. for a in articles
  39. ]
  40. # 添加或更新向量数据库
  41. await self.vectorstore.aadd_texts(texts=texts, metadatas=metadatas)
  42. logger.info("知识库向量化完成:%d 篇文章", len(articles))
  43. except Exception as e:
  44. logger.warning("知识库初始化失败:%s", e)
  45. async def retrieve(
  46. self,
  47. query: str,
  48. filters: Optional[dict] = None,
  49. k: int = 5,
  50. ) -> list[dict]:
  51. """向量相似度检索"""
  52. results = []
  53. # 向量检索
  54. filter_query = {"user_id": filters["user_id"]} if filters and "user_id" in filters else None
  55. docs = self.vectorstore.similarity_search(
  56. query,
  57. k=k * 2, # 多取一些
  58. filter=filter_query,
  59. )
  60. # 格式化结果并去重
  61. seen = set()
  62. for doc in docs:
  63. content_hash = hash(doc.page_content[:100])
  64. if content_hash in seen:
  65. continue
  66. seen.add(content_hash)
  67. results.append({
  68. "content": doc.page_content,
  69. "metadata": doc.metadata,
  70. "score": doc.metadata.get("score", 0) if hasattr(doc, "metadata") else 0,
  71. })
  72. if len(results) >= k:
  73. break
  74. return results[:k]