retriever.py 3.3 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495
  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 .embeddings import get_embeddings
  7. from app.config import settings
  8. from app.tools.java_client import JavaClient
  9. from typing import Optional
  10. import logging
  11. logger = logging.getLogger(__name__)
  12. class RagRetriever:
  13. """混合检索器: ChromaDB 向量 + BM25 关键词 + 可选 LLM 压缩"""
  14. def __init__(self, collection_name: str = "cfc_knowledge"):
  15. embeddings = get_embeddings()
  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. self._bm25_retriever = None
  23. self._bm25_texts = []
  24. async def initialize(self):
  25. """从 Java 侧拉取知识库数据, 构建 BM25 索引"""
  26. try:
  27. articles = await self.java_client.get_published_articles()
  28. self._bm25_texts = [
  29. f"{a['title']} {a['summary']} {a.get('tags', '')}"
  30. for a in articles
  31. ]
  32. if self._bm25_texts:
  33. self._bm25_retriever = BM25Retriever.from_texts(
  34. self._bm25_texts, metadatas=articles
  35. )
  36. logger.info("BM25 索引就绪: %d 条", len(self._bm25_texts))
  37. except Exception as e:
  38. logger.warning("BM25 初始化失败(不影响向量检索): %s", e)
  39. async def retrieve(
  40. self,
  41. query: str,
  42. filters: Optional[dict] = None,
  43. k: int = 5,
  44. use_compression: bool = False,
  45. ) -> list[dict]:
  46. """混合检索, 返回 [{content, metadata, score}]"""
  47. retrievers = []
  48. # 向量检索
  49. vector_retriever = self.vectorstore.as_retriever(
  50. search_kwargs={"k": k, "filter": filters}
  51. )
  52. retrievers.append(vector_retriever)
  53. # BM25 检索
  54. if self._bm25_retriever:
  55. retrievers.append(self._bm25_retriever)
  56. if len(retrievers) == 1:
  57. docs = await retrievers[0].ainvoke(query)
  58. else:
  59. ensemble = EnsembleRetriever(
  60. retrievers=retrievers, weights=[0.6, 0.4]
  61. )
  62. docs = await ensemble.ainvoke(query)
  63. # 可选: LLM 压缩去噪
  64. if use_compression and docs:
  65. from langchain_openai import ChatOpenAI
  66. llm = ChatOpenAI(
  67. model=settings.llm_model,
  68. api_key=settings.llm_api_key,
  69. base_url=settings.llm_base_url,
  70. )
  71. compressor = LLMChainExtractor.from_llm(llm)
  72. compression_retriever = ContextualCompressionRetriever(
  73. base_compressor=compressor,
  74. base_retriever=self.vectorstore.as_retriever(),
  75. )
  76. docs = await compression_retriever.ainvoke(query)
  77. results = []
  78. for doc in docs:
  79. results.append({
  80. "content": doc.page_content,
  81. "metadata": doc.metadata,
  82. "score": getattr(doc, "metadata", {}).get("score", 0),
  83. })
  84. return results[:k]