# 流式调用-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()' 以运行演示")



更多推荐