LangChain智能文档助手【6】-智能对话系统
·
这是一个基于 LangChain 框架和通义千问(Qwen)大语言模型构建的检索增强生成(RAG)智能问答系统的教程。该系统是基于最基础的功能点,如大语言模型调用,文档处理,嵌入模型向量化,向量存储和检索,智能对话构建而成,最后用streamlit生成一个web界面,将功能可视化。
第六章 智能对话
本章构建了一个基于RAG的智能对话系统,结合了:
- 通义千问大模型(Qwen)作为问答引擎
- FAISS向量数据库作为知识存储
- 检索增强生成(RAG)技术提升回答质量
- 对话记忆系统实现多轮对话
接下来,我们开启构建的过程,按如下层级创建文件conversation_chain.py。

conversation_chain.py文件的完整代码如下:
# 文件: conversation_chain.py
"""
第6章:构建基于RAG的对话链
使用通义千问 + FAISS 向量数据库实现智能问答
"""
import os
import sys
import pickle
from typing import List, Dict, Any, Optional
from datetime import datetime
from dotenv import load_dotenv
# 加载环境变量
load_dotenv()
# 添加项目根目录到Python路径
project_root = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
if project_root not in sys.path:
sys.path.insert(0, project_root)
# 导入必要的模块
try:
# 使用可用的模块
from langchain_community.vectorstores import FAISS
from langchain_core.prompts import PromptTemplate
from langchain_core.retrievers import BaseRetriever
# 检查memory模块是否可用
try:
# 尝试多种导入方式
try:
from langchain.memory import ConversationBufferMemory
except ImportError:
try:
from langchain_community.memory import ConversationBufferMemory
except ImportError:
try:
from langchain.memory.buffer import ConversationBufferMemory
except ImportError:
# 如果都不行,使用简化版本
raise ImportError("无法导入ConversationBufferMemory")
MEMORY_AVAILABLE = True
except ImportError:
MEMORY_AVAILABLE = False
print("⚠️ ConversationBufferMemory 不可用,将使用简化版本")
print(" 原因: LangChain版本兼容性问题或Python 3.14+不兼容")
print(" 解决方案: 使用简化版对话记忆系统")
CONVERSATION_AVAILABLE = True
except ImportError as e:
CONVERSATION_AVAILABLE = False
print(f"❌ LangChain导入错误: {e}")
print(" 请检查LangChain版本兼容性")
# 导入项目自定义模块
try:
from qwen_client import QwenClient
from vector_store import QwenVectorStore
from retriever import QwenRetriever
QWEN_AVAILABLE = True
except ImportError as e:
QWEN_AVAILABLE = False
print(f"❌ 自定义模块导入错误: {e}")
try:
from langchain_community.vectorstores import FAISS
FAISS_AVAILABLE = True
except ImportError:
FAISS_AVAILABLE = False
print("❌ FAISS 不可用")
from langchain_core.documents import Document
class QwenRAGConversation:
"""
基于通义千问和RAG的对话系统
结合FAISS向量检索和对话记忆
"""
def __init__(self, index_dir: str = "faiss_db_qwen"):
"""
初始化对话系统
Args:
index_dir: FAISS索引目录
"""
self.index_dir = index_dir
self.vectorstore = None
self.embeddings = None
self.retriever = None
self.memory = None
self.conversation_chain = None
self.qwen_client = None
print("🔧 初始化Qwen RAG对话系统...")
print(f" 索引目录: {index_dir}")
# 初始化组件
self._initialize_components()
def _initialize_components(self):
"""初始化所有组件"""
print("🔧 使用现有的QwenRetriever类进行初始化...")
# 使用现有的QwenRetriever类
if QWEN_AVAILABLE:
try:
# 复用现有的检索器类
self.retriever_manager = QwenRetriever(self.index_dir)
# 获取检索器
self.retriever = self.retriever_manager.create_basic_retriever(
search_type="similarity", k=4
)
# 获取Qwen客户端
self.qwen_client = self.retriever_manager.qwen_client
# 获取向量存储
self.vectorstore = self.retriever_manager.vectorstore
if self.qwen_client and hasattr(self.qwen_client, 'test_connection'):
if self.qwen_client.test_connection():
print("✅ Qwen客户端连接成功")
else:
print("⚠️ Qwen客户端连接测试失败")
print("✅ 复用现有组件成功")
except Exception as e:
print(f"❌ 复用现有组件失败: {e}")
# 降级为原始初始化方式
self._fallback_initialize()
else:
print("⚠️ 使用原始初始化方式")
self._fallback_initialize()
# 初始化对话记忆
self._initialize_memory()
# 创建对话链
self._create_conversation_chain()
def _fallback_initialize(self):
"""回退初始化方式(当复用失败时使用)"""
print("🔄 使用回退初始化方式...")
# 1. 初始化Qwen客户端
try:
self.qwen_client = QwenClient()
if hasattr(self.qwen_client, 'test_connection'):
if self.qwen_client.test_connection():
print("✅ Qwen客户端连接成功")
else:
print("⚠️ Qwen客户端连接测试失败")
except Exception as e:
print(f"❌ Qwen客户端初始化失败: {e}")
# 2. 加载向量存储
try:
from vector_store import QwenVectorStore
vector_manager = QwenVectorStore(self.index_dir)
self.vectorstore = vector_manager.load_existing_store()
if self.vectorstore:
self.retriever = self.vectorstore.as_retriever(
search_type="similarity",
search_kwargs={
"k": 4, # 每次检索4个文档
"score_threshold": 0.3 # 相似度阈值
}
)
print("✅ 向量存储加载成功")
else:
print("❌ 向量存储加载失败")
except Exception as e:
print(f"❌ 向量存储初始化失败: {e}")
def _initialize_embeddings(self):
"""初始化嵌入模型(已弃用,使用复用方式)"""
print("⚠️ 该方法已弃用,使用复用方式")
return False
def _load_vector_store(self):
"""加载向量存储(已弃用,使用复用方式)"""
print("⚠️ 该方法已弃用,使用复用方式")
return False
def _initialize_memory(self):
"""初始化对话记忆"""
try:
if MEMORY_AVAILABLE:
from langchain.memory import ConversationBufferMemory
self.memory = ConversationBufferMemory(
memory_key="chat_history",
output_key="answer",
return_messages=True,
max_token_limit=2000 # 限制记忆长度
)
print("✅ 对话记忆初始化成功")
print(f" 记忆容量: 2000 tokens")
else:
# 简化版本:使用字典存储对话历史
self.memory = {
"chat_history": [],
"max_tokens": 2000
}
print("✅ 简化版对话记忆初始化成功")
except Exception as e:
print(f"❌ 初始化对话记忆失败: {e}")
self.memory = None
def _create_conversation_chain(self):
"""创建对话链"""
if not self.retriever:
print("❌ 检索器未初始化")
return False
if not self.qwen_client:
print("❌ Qwen客户端未初始化")
return False
try:
# 自定义提示模板
prompt_template = """
你是一个基于知识库的智能助手。请根据以下对话历史和知识库内容,用中文回答用户的问题。
对话历史:
{chat_history}
相关上下文(来自知识库):
{context}
用户问题:{question}
请按照以下要求回答:
1. 基于知识库内容回答,如果知识库中没有相关信息,请说明不清楚
2. 回答要详细、准确、有帮助
3. 使用自然、友好的语气
4. 如果适用,可以给出相关建议或进一步的信息
回答:
"""
PROMPT = PromptTemplate(
template=prompt_template,
input_variables=["chat_history", "context", "question"]
)
# 创建自定义的Qwen LLM包装器
class QwenLLMWrapper:
def __init__(self, qwen_client):
self.qwen_client = qwen_client
def __call__(self, prompt: str, **kwargs) -> str:
try:
response = self.qwen_client.chat_completion([
{"role": "system", "content": "你是一个有帮助的AI助手。"},
{"role": "user", "content": prompt}
])
return response
except Exception as e:
return f"调用Qwen失败: {e}"
# 创建LLM实例
llm = QwenLLMWrapper(self.qwen_client)
# 简化版本:直接使用检索器进行问答
# 在LangChain 1.2.0中,ConversationalRetrievalChain已被重构
# 这里简化实现,直接使用检索和问答
self.conversation_chain = {
"llm": llm,
"retriever": self.retriever,
"memory": self.memory,
"prompt": PROMPT
}
print("✅ 对话链创建成功(简化版本)")
print(" 系统提示: 基于知识库的智能助手")
return True
except Exception as e:
print(f"❌ 创建对话链失败: {e}")
import traceback
traceback.print_exc()
return False
def ask_question(self, question: str) -> Dict[str, Any]:
"""
提问并获取回答
Args:
question: 用户问题
Returns:
包含回答和元数据的字典
"""
if not self.conversation_chain:
print("❌ 对话链未初始化")
return {
"answer": "系统未正确初始化,请检查配置。",
"sources": [],
"error": "对话链未初始化"
}
print(f"\n🤔 用户提问: {question}")
try:
# 简化版本:直接检索和问答
chain = self.conversation_chain
# 1. 检索相关文档
if chain["retriever"]:
try:
# 兼容不同版本的LangChain
source_docs = chain["retriever"].invoke(question)
except AttributeError:
try:
source_docs = chain["retriever"].get_relevant_documents(question)
except AttributeError:
try:
source_docs = chain["retriever"]._get_relevant_documents(question)
except:
source_docs = []
else:
source_docs = []
# 2. 构建上下文
context = "\n".join([doc.page_content for doc in source_docs[:3]])
# 3. 获取对话历史
if isinstance(chain["memory"], dict): # 简化版本
history_text = "\n".join([f"user: {h['user']}\nassistant: {h['assistant']}"
for h in chain["memory"]["chat_history"][-2:]])
elif hasattr(chain["memory"], 'chat_memory'): # 完整版本
chat_history = chain["memory"].chat_memory.messages
history_text = "\n".join([f"{msg.type}: {msg.content}" for msg in chat_history[-4:]])
else:
history_text = ""
# 4. 构建提示词
prompt = chain["prompt"].format(
chat_history=history_text,
context=context,
question=question
)
# 5. 调用LLM
answer = chain["llm"](prompt)
# 6. 更新记忆
if isinstance(chain["memory"], dict): # 简化版本
chain["memory"]["chat_history"].append({"user": question, "assistant": answer})
# 限制记忆长度
if len(chain["memory"]["chat_history"]) > 10:
chain["memory"]["chat_history"] = chain["memory"]["chat_history"][-10:]
elif hasattr(chain["memory"], 'save_context'): # 完整版本
chain["memory"].save_context({"input": question}, {"output": answer})
# 处理回答
if not answer or answer.strip() == "":
answer = "抱歉,我无法回答这个问题。"
# 提取来源信息
sources = []
for i, doc in enumerate(source_docs[:3]): # 最多显示3个来源
source = doc.metadata.get('source', '未知来源')
content_preview = doc.page_content[:100] + "..." if len(doc.page_content) > 100 else doc.page_content
sources.append({
"id": i + 1,
"source": source,
"content_preview": content_preview,
"relevance": "高" if i < 2 else "中"
})
# 生成回答摘要
if len(answer) > 500:
answer_summary = answer[:500] + "..."
else:
answer_summary = answer
response_data = {
"answer": answer,
"answer_summary": answer_summary,
"sources": sources,
"timestamp": datetime.now().isoformat(),
"question": question,
"has_answer": True
}
print(f"💡 AI回答: {answer_summary}")
print(f"📚 参考来源: {len(sources)} 个")
return response_data
except Exception as e:
print(f"❌ 提问失败: {e}")
import traceback
traceback.print_exc()
return {
"answer": f"抱歉,处理问题时出错: {str(e)}",
"sources": [],
"error": str(e),
"has_answer": False
}
def clear_memory(self):
"""清空对话记忆"""
if not self.memory:
print("⚠️ 对话记忆未初始化")
return
try:
# 处理简化版本的记忆系统
if isinstance(self.memory, dict):
self.memory["chat_history"] = []
print("🧹 对话记忆已清空")
# 处理完整版本的LangChain记忆系统
elif hasattr(self.memory, 'clear'):
self.memory.clear()
print("🧹 对话记忆已清空")
else:
print("⚠️ 无法清空对话记忆,系统不支持")
except Exception as e:
print(f"❌ 清空对话记忆失败: {e}")
def get_chat_history(self) -> List[Dict[str, str]]:
"""获取对话历史"""
if not self.memory:
return []
try:
# 处理简化版本的记忆系统
if isinstance(self.memory, dict):
chat_history = self.memory.get("chat_history", [])
history = []
for i, item in enumerate(chat_history):
history.append({
"user": item.get("user", ""),
"assistant": item.get("assistant", ""),
"turn": i + 1
})
return history
# 处理完整版本的LangChain记忆系统
elif hasattr(self.memory, 'chat_memory'):
messages = self.memory.chat_memory.messages
history = []
for i in range(0, len(messages), 2):
if i + 1 < len(messages):
history.append({
"user": messages[i].content,
"assistant": messages[i + 1].content,
"turn": i // 2 + 1
})
return history
else:
return []
except:
return []
def similarity_search(self, query: str, k: int = 3, score_threshold: float = 0.3) -> List[tuple]:
"""
相似性搜索 - 复用现有的检索器测试功能
Args:
query: 查询文本
k: 返回的文档数量
score_threshold: 相似度阈值
Returns:
包含文档和相似度分数的列表
"""
if hasattr(self, 'retriever_manager'):
# 复用现有的检索器测试功能
try:
results = self.retriever_manager.test_retrieval(
query, "basic", k=k, search_type="similarity"
)
# 转换格式为带分数的元组列表
filtered_results = []
if hasattr(self.vectorstore, 'similarity_search_with_score'):
docs_with_scores = self.vectorstore.similarity_search_with_score(query, k=k)
for doc, score in docs_with_scores:
if score >= score_threshold:
filtered_results.append((doc, score))
print(f"🔍 相似性搜索: 查询='{query}', 找到 {len(filtered_results)} 个相关文档")
return filtered_results
except Exception as e:
print(f"❌ 复用检索器搜索失败: {e}")
# 降级为原始方式
pass
# 降级为原始实现
if not self.vectorstore:
print("❌ 向量存储未初始化")
return []
try:
# 使用FAISS的相似性搜索
docs_with_scores = self.vectorstore.similarity_search_with_score(query, k=k)
# 过滤低于阈值的文档
filtered_results = []
for doc, score in docs_with_scores:
if score >= score_threshold:
filtered_results.append((doc, score))
print(f"🔍 相似性搜索: 查询='{query}', 找到 {len(filtered_results)} 个相关文档")
return filtered_results
except Exception as e:
print(f"❌ 相似性搜索失败: {e}")
return []
def analyze_question(self, question: str) -> Dict[str, Any]:
"""
分析问题
Args:
question: 用户问题
Returns:
分析结果
"""
analysis = {
"question": question,
"length": len(question),
"word_count": len(question.split()),
"contains_question_words": any(word in question for word in ["什么", "怎么", "为什么", "如何", "是否", "多少"]),
"is_complex": len(question.split()) > 10,
"suggested_retrieval_k": 4, # 默认检索数量
"timestamp": datetime.now().isoformat()
}
# 根据问题复杂度调整检索数量
if analysis["is_complex"]:
analysis["suggested_retrieval_k"] = 6
return analysis
def interactive_conversation():
"""交互式对话"""
print("\n" + "=" * 70)
print("💬 第6章:Qwen RAG 交互式对话系统")
print("=" * 70)
# 初始化对话系统
print("\n🔄 初始化系统...")
conversation_system = QwenRAGConversation("faiss_db_qwen")
if not conversation_system.conversation_chain:
print("❌ 系统初始化失败")
return
print("\n✅ 系统准备就绪!")
print(" 输入 '退出' 或 'quit' 结束对话")
print(" 输入 '清空' 或 'clear' 清空对话历史")
print(" 输入 '历史' 或 'history' 查看对话历史")
print(" 输入 '帮助' 或 'help' 显示帮助")
print("-" * 70)
conversation_count = 0
while True:
try:
# 获取用户输入
user_input = input(f"\n[{conversation_count + 1}] 你: ").strip()
if not user_input:
continue
# 处理特殊命令
if user_input.lower() in ['退出', 'quit', 'exit', 'q']:
print("👋 再见!")
break
elif user_input.lower() in ['清空', 'clear']:
conversation_system.clear_memory()
conversation_count = 0
continue
elif user_input.lower() in ['历史', 'history']:
history = conversation_system.get_chat_history()
if history:
print("\n📜 对话历史:")
for item in history[-5:]: # 显示最后5轮
print(f" 轮次 {item['turn']}:")
print(f" 你: {item['user'][:50]}...")
print(f" AI: {item['assistant'][:50]}...")
else:
print(" 暂无对话历史")
continue
elif user_input.lower() in ['帮助', 'help']:
print("\n📋 帮助信息:")
print(" 1. 正常提问 - 系统会基于知识库回答")
print(" 2. '退出' - 结束对话")
print(" 3. '清空' - 清空对话历史")
print(" 4. '历史' - 查看对话历史")
print(" 5. '帮助' - 显示此帮助")
continue
# 分析问题
analysis = conversation_system.analyze_question(user_input)
print(f" 🔍 问题分析: {analysis['word_count']} 词, {'复杂' if analysis['is_complex'] else '简单'}问题")
# 获取回答
response = conversation_system.ask_question(user_input)
# 显示回答
if response.get("has_answer", False):
print(f"\n🤖 AI: {response['answer']}")
# 显示来源
if response.get("sources"):
print(f"\n📚 参考来源:")
for source in response["sources"]:
print(f" [{source['id']}] {source['source']} - {source['content_preview']}")
else:
print(f"\n⚠️ AI: {response.get('answer', '无法回答')}")
conversation_count += 1
except KeyboardInterrupt:
print("\n\n👋 用户中断,再见!")
break
except Exception as e:
print(f"\n❌ 发生错误: {e}")
continue
def test_conversation():
"""测试对话系统"""
print("\n🧪 测试对话系统...")
# 初始化
conversation_system = QwenRAGConversation("faiss_db_qwen")
if not conversation_system.conversation_chain:
print("❌ 测试失败:系统未初始化")
return False
# 测试问题
test_questions = [
"通义千问是什么?",
"FAISS 是什么库?",
"LangChain 有什么用?",
"什么是检索增强生成?",
]
print(f"\n📋 测试 {len(test_questions)} 个问题...")
results = []
for i, question in enumerate(test_questions):
print(f"\n[{i+1}] 测试: {question}")
# 分析问题
analysis = conversation_system.analyze_question(question)
print(f" 分析: {analysis['word_count']} 词, 检索 {analysis['suggested_retrieval_k']} 个文档")
# 获取回答
response = conversation_system.ask_question(question)
if response.get("has_answer", False):
print(f" ✅ 回答成功")
print(f" 回答长度: {len(response['answer'])} 字符")
print(f" 参考来源: {len(response['sources'])} 个")
results.append(True)
else:
print(f" ❌ 回答失败")
print(f" 错误: {response.get('error', '未知')}")
results.append(False)
# 短暂暂停
import time
time.sleep(1)
# 统计结果
success_count = sum(results)
print(f"\n📊 测试结果: {success_count}/{len(test_questions)} 成功")
return success_count == len(test_questions)
def main():
"""主函数"""
print("\n" + "=" * 70)
print("🚀 第5章:基于RAG的智能对话系统")
print("=" * 70)
# 检查依赖
if not all([CONVERSATION_AVAILABLE, QWEN_AVAILABLE, FAISS_AVAILABLE]):
print("❌ 缺少必要的依赖")
print(" 请安装: pip install langchain dashscope faiss-cpu")
return
# 检查API密钥
api_key = os.getenv("ALIYUN_API_KEY") or os.getenv("DASHSCOPE_API_KEY")
if not api_key:
print("❌ 未设置API密钥")
print(" 请在 .env 文件中添加: ALIYUN_API_KEY=你的密钥")
return
# 检查索引文件
index_file = "faiss_db_qwen/index.faiss"
if not os.path.exists(index_file):
print(f"❌ 索引文件不存在: {index_file}")
print(" 请先运行第3章代码创建向量索引")
return
# 显示菜单
print("\n📋 选择模式:")
print(" 1. 交互式对话")
print(" 2. 系统测试")
print(" 3. 退出")
choice = input("\n请输入选择 (1-3): ").strip()
if choice == "1":
interactive_conversation()
elif choice == "2":
test_conversation()
elif choice == "3":
print("👋 再见!")
else:
print("❌ 无效选择")
if __name__ == "__main__":
main()
代码运行
python src/conversation_chain.py
代码解析:
1.本章代码结构概览
# conversation_chain.py 文件结构
├── 类定义
│ └── QwenRAGConversation # 核心类(支持API调用)
├── 函数定义
│ ├── interactive_conversation() # 交互式模式函数
│ ├── test_conversation() # 测试模式函数
│ └── main() # 主入口函数
└── 主程序入口
└── if __name__ == "__main__":
main()
核心类:QwenRAGConversation 结构
class QwenRAGConversation:
"""基于通义千问和RAG的对话系统"""
def __init__(self, index_dir: str = "faiss_db_qwen"):
"""初始化 - 可以被其他程序调用"""
self.index_dir = index_dir
self.vectorstore = None
# ... 初始化逻辑
def ask_question(self, question: str) -> Dict[str, Any]:
"""API核心方法 - 可以被其他程序调用"""
# 返回结构化数据
return {
"answer": "回答内容",
"sources": [...],
"timestamp": "..."
}
def clear_memory(self):
"""API方法 - 清空记忆"""
pass
def get_chat_history(self) -> List[Dict]:
"""API方法 - 获取历史"""
pass
def similarity_search(self, query: str, k: int = 3):
"""API方法 - 纯检索"""
pass
def analyze_question(self, question: str):
"""API方法 - 分析问题"""
pass
2.对话流程详解
(1)单轮对话
用户提问 → 向量检索 → 构建提示词 → Qwen生成 → 返回答案
↓ ↓ ↓ ↓ ↓
"什么是AI?" 找4个最相关 结合历史+检索 生成答案 格式化输出
(2)多轮对话
第1轮: 用户问 → 系统答 → 保存到记忆
第2轮: 用户新问题 → 检索+历史记忆 → 系统答
第3轮: 持续对话...
(3)对话中提示词模板:
prompt_template = """
对话历史:{chat_history}
相关上下文(来自知识库):
{context}
用户问题:{question}
要求:
1. 基于知识库回答
2. 详细、准确、有帮助
3. 自然友好的语气
"""
3.本章提供了四种使用模式
(1)交互式对话模式
# 启动对话
$ python conversation_chain.py
# 支持的命令:
你: 什么是机器学习?
你: 清空 # 清空历史
你: 历史 # 查看历史
你: 帮助 # 显示帮助
你: 退出 # 结束对话
(2)系统测试模式
# 自动测试4个标准问题:
test_questions = [
"通义千问是什么?",
"FAISS 是什么库?",
"LangChain 有什么用?",
"什么是检索增强生成?"
]
# 输出:测试结果统计
(3)API调用模式
# 程序化调用
from conversation_chain import QwenRAGConversation
system = QwenRAGConversation("faiss_db_qwen")
response = system.ask_question("你的问题")
print(response["answer"])
(4)检索分析模式
# 纯检索,不带生成
docs = system.similarity_search("查询词", k=5)
# 返回:[文档, 相似度分数]
代码运行结果参考
# 终端中运行
$ python conversation_chain.py
# 输出:
🚀 第5章:基于RAG的智能对话系统
===============================
📋 选择模式:
1. 交互式对话
2. 系统测试
3. 退出
请输入选择 (1-3): 1
💬 第5章:Qwen RAG 交互式对话系统
===============================
输入 '退出' 或 'quit' 结束对话
输入 '清空' 或 'clear' 清空对话历史
------------------------------
[1] 你: 什么是机器学习?
🤖 AI: 机器学习是人工智能的一个分支...
📚 参考来源:
[1] 机器学习教材.pdf
完成以上内容,我们就成功构建了一个基于rag的智能对话系统,最后一章,我们会使用streamlit快速搭建一个web页面,将系统可视化。
更多推荐



所有评论(0)