大模型调用的方法流式输出stream()
·
# 流式调用-stream()
from langchain.chat_models import init_chat_model
import os
from dotenv import load_dotenv
from langsmith.evaluation import _name_generation
from openrouter.components import responseinputfile
# 加载配置文件
load_dotenv(override=True)
DEEPSEEK_API_KEY = os.getenv("DEEPSEEK_API_KEY")
DEEPSEEK_BASE_URL = os.getenv("DEEPSEEK_BASE_URL")
# 获取大模型
model = init_chat_model(
model="deepseek-v4-flash",
model_provider="deepseek",
temperature=1.5,
api_key=DEEPSEEK_API_KEY,
base_url=DEEPSEEK_BASE_URL,
)
for chunk in model.stream("帮我写一个程序员的歌词"):
print(chunk.text,end="",flush=True)
# 批量调用-batch()
from langchain.chat_models import init_chat_model
import os
from dotenv import load_dotenv
# 加载配置文件
load_dotenv(override=True)
DEEPSEEK_API_KEY = os.getenv("DEEPSEEK_API_KEY")
DEEPSEEK_BASE_URL = os.getenv("DEEPSEEK_BASE_URL")
# 获取大模型
model = init_chat_model(
model="deepseek-v4-flash",
model_provider="deepseek",
temperature=1.5,
api_key=DEEPSEEK_API_KEY,
base_url=DEEPSEEK_BASE_URL,
)
messages =[
"你好,你是谁?",
"中国的首都在哪里?"
]
responseS =model.batch(messages)
for response in responseS:
print(response)
# 按完成顺序接收响应-batch_as_completed()
from langchain.chat_models import init_chat_model
import os
from dotenv import load_dotenv
# 加载配置文件
load_dotenv(override=True)
DEEPSEEK_API_KEY = os.getenv("DEEPSEEK_API_KEY")
DEEPSEEK_BASE_URL = os.getenv("DEEPSEEK_BASE_URL")
# 获取大模型
model = init_chat_model(
model="deepseek-v4-flash",
model_provider="deepseek",
temperature=1.5,
api_key=DEEPSEEK_API_KEY,
base_url=DEEPSEEK_BASE_URL,
)
messages =[
"你好,你是谁?",
"中国的首都在哪里?"
]
responseS =model.batch_as_completed(messages)
for response in responseS:
print(response)
#%%
#异步调用
#5.4异步调用
#复习:同步vs异步
#同步(sync)
#表现:当前执行流会被『阻塞』
#概念:发起一个任务之后,需要等待该任务完成后,才能继续执行后续任务。
##异步(async)
#概念:发起一个任务之后,不必等该任务完成,就可以继续执行其他任务。
#备注:虽然不必等待任务完成,但任务完成后,仍然可以通过特定方式获取结果。
#表现:当前执行流不会被『阻塞』
#举例:
# 在LangChain框架中,异步方法(ainvoke、astream、abatch)与它的同步版本(invoke,stream、batch)相比,具备如下特点:
# 避免阻塞主线程:同步调用会阻塞程序执行,而异步方法让应用程序在等待API响应时保持响应生。
# 优化资源利用:异步操作可以更高效地利用系统资源,减少空闲等待时间
# 例一: astream()
from langchain.chat_models import init_chat_model
import asyncio
import os
from dotenv import load_dotenv
import time
# 加载配置文件
load_dotenv(override=True)
DEEPSEEK_API_KEY = os.getenv("DEEPSEEK_API_KEY")
DEEPSEEK_BASE_URL = os.getenv("DEEPSEEK_BASE_URL")
# 获取大模型(若模型名不存在,请改为 "deepseek-chat" 等实际名称)
model = init_chat_model(
model="deepseek-v4-flash", # 请确认模型名称是否正确
model_provider="deepseek",
temperature=1.5,
api_key=DEEPSEEK_API_KEY,
base_url=DEEPSEEK_BASE_URL,
)
async def demo_async_stream():
"""演示异步流式调用的非阻塞特性"""
print("=== 演示:astream 的异步(非阻塞)效果 ===")
start_time = time.perf_counter()
print("程序开始...")
# 1. 发起异步流式请求(立即返回异步生成器,不阻塞)
print(">>> 发起异步流式调用 (astream)...")
stream_resp = model.astream("请用一句话解释机器学习的基本概念。")
print(">>> 流式请求已发送,程序无需等待,继续执行其他异步任务...")
# 2. 模拟并发任务(利用等待时间处理其他事务)
for i in range(3):
await asyncio.sleep(1) # 使用 asyncio.sleep 让出控制权,允许网络 IO 后台进行
print(f">>> 正在执行第 {i+1} 个并发任务... (已耗时 {time.perf_counter() - start_time:.2f}s)")
# 3. 准备读取流式结果
print(">>> 模拟任务已完成,开始读取缓冲区中的流式结果...")
print(">>> 流式输出: ", end="", flush=True)
# 4. 逐块读取并打印(此时网络数据已部分或全部缓冲)
async for chunk in stream_resp:
content = chunk.content if hasattr(chunk, 'content') else str(chunk)
print(content, end="", flush=True)
# 5. 读取完成后记录结束时间
end_time = time.perf_counter()
print("\n>>> 流式输出结束")
print(f"=== 总运行耗时(含流式读取): {end_time - start_time:.2f}s ===")
async def main():
await demo_async_stream()
# if __name__ == "__main__":
# asyncio.run(main())
await main()
# abatch()
import asyncio
import os
import time
from langchain.chat_models import init_chat_model
from dotenv import load_dotenv
load_dotenv(override=True)
DEEPSEEK_API_KEY = os.getenv("DEEPSEEK_API_KEY")
DEEPSEEK_BASE_URL = os.getenv("DEEPSEEK_BASE_URL")
if not DEEPSEEK_API_KEY:
raise ValueError("请在 .env 中设置 DEEPSEEK_API_KEY")
try:
model = init_chat_model(
model="deepseek-chat", # 注意此处改为有效名称
model_provider="deepseek",
temperature=1.5,
api_key=DEEPSEEK_API_KEY,
base_url=DEEPSEEK_BASE_URL,
)
print("✅ 模型初始化成功")
except Exception as e:
print(f"❌ 模型初始化失败: {e}")
raise
async def demo_async_batch():
print("=== 演示:abatch 的异步(非阻塞)效果 ===")
start_time = time.perf_counter()
print("程序开始...")
questions = [
"用一句话说明深度学习与传统机器学习的区别",
"中国首都在哪里?"
]
print(">>> 发起异步批量调用 (abatch)...")
batch_task = asyncio.create_task(model.abatch(questions))
print(">>> 批量任务已在后台运行,主程序继续执行...")
for i in range(3):
await asyncio.sleep(1)
elapsed = time.perf_counter() - start_time
print(f">>> 正在执行第 {i+1} 个并发任务... (已耗时 {elapsed:.2f}s)")
print(">>> 其他任务已完成,现在获取后台批量任务的结果...")
responses = await batch_task
for idx, response in enumerate(responses):
content = response.content if hasattr(response, 'content') else str(response)
print(f">>> 问题 {idx+1} 的响应: {content}")
end_time = time.perf_counter()
print(f"=== 总运行耗时: {end_time - start_time:.2f}s ===")
async def main():
await demo_async_batch()
# 在 Jupyter 中直接执行此 Cell 后,请在下一个 Cell 执行:await main()
print("✅ 代码已加载,请执行 'await main()' 以运行演示")
更多推荐
所有评论(0)