retriever.py 3.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109
  1. from langchain_chroma import Chroma
  2. from langchain.retrievers import EnsembleRetriever
  3. from langchain_community.retrievers import BM25Retriever
  4. from langchain.retrievers.document_compressors import LLMChainExtractor
  5. from langchain.retrievers import ContextualCompressionRetriever
  6. from langchain_openai import ChatOpenAI
  7. from .embeddings import get_embeddings
  8. from app.config import settings
  9. from app.tools.java_client import JavaClient
  10. from typing import Optional
  11. import logging
  12. logger = logging.getLogger(__name__)
  13. class RagRetriever:
  14. """升级版混合检索器: 向量 + BM25 + LLM 压缩重排序"""
  15. def __init__(self, collection_name: str = "cfc_knowledge"):
  16. embeddings = get_embeddings()
  17. self.vectorstore = Chroma(
  18. collection_name=collection_name,
  19. embedding_function=embeddings,
  20. persist_directory=settings.chroma_db_path,
  21. )
  22. self.java_client = JavaClient()
  23. self._bm25_retriever: Optional[BM25Retriever] = None
  24. self._bm25_texts: list[str] = []
  25. async def initialize(self):
  26. """从 Java 侧拉取知识库, 构建 BM25 索引"""
  27. try:
  28. articles = await self.java_client.get_published_articles()
  29. self._bm25_texts = [
  30. f"{a['title']} {a['summary']} {a.get('tags', '')}"
  31. for a in articles
  32. ]
  33. if self._bm25_texts:
  34. self._bm25_retriever = BM25Retriever.from_texts(
  35. self._bm25_texts,
  36. metadatas=articles,
  37. )
  38. logger.info("BM25 索引就绪: %d 条", len(self._bm25_texts))
  39. except Exception as e:
  40. logger.warning("BM25 初始化失败: %s", e)
  41. async def retrieve(
  42. self,
  43. query: str,
  44. filters: Optional[dict] = None,
  45. k: int = 5,
  46. use_compression: bool = True,
  47. ) -> list[dict]:
  48. """混合检索 + 可选 LLM 压缩重排序"""
  49. retrievers = []
  50. # 1. 向量检索 (多取一些方便后续 ensemble 排序)
  51. vector_retriever = self.vectorstore.as_retriever(
  52. search_kwargs={"k": k * 2, "filter": filters},
  53. )
  54. retrievers.append(vector_retriever)
  55. # 2. BM25 关键词检索
  56. if self._bm25_retriever:
  57. bm25_k = self._bm25_retriever.k
  58. self._bm25_retriever.k = k * 2
  59. retrievers.append(self._bm25_retriever)
  60. self._bm25_retriever.k = bm25_k
  61. if len(retrievers) == 1:
  62. docs = await retrievers[0].ainvoke(query)
  63. ensemble = retrievers[0]
  64. else:
  65. ensemble = EnsembleRetriever(
  66. retrievers=retrievers,
  67. weights=[0.6, 0.4],
  68. )
  69. docs = await ensemble.ainvoke(query)
  70. # 3. LLM 压缩 (剔除不相关内容)
  71. if use_compression and docs:
  72. llm = ChatOpenAI(
  73. model=settings.llm_model,
  74. api_key=settings.llm_api_key,
  75. base_url=settings.llm_base_url,
  76. temperature=0,
  77. )
  78. compressor = LLMChainExtractor.from_llm(llm)
  79. compression_retriever = ContextualCompressionRetriever(
  80. base_compressor=compressor,
  81. base_retriever=ensemble if len(retrievers) > 1 else retrievers[0],
  82. )
  83. docs = await compression_retriever.ainvoke(query)
  84. # 4. 格式化为统一输出 + 去重
  85. results = []
  86. seen = set()
  87. for doc in docs:
  88. content_hash = hash(doc.page_content[:100])
  89. if content_hash in seen:
  90. continue
  91. seen.add(content_hash)
  92. results.append({
  93. "content": doc.page_content,
  94. "metadata": doc.metadata,
  95. "score": doc.metadata.get("score", 0) if hasattr(doc, "metadata") else 0,
  96. })
  97. return results[:k]