在部分API中调用模型时(如Gemini、Github Models等)流式输出并非逐字符输出,而是多个大型响应块,这样很影响输出时的观感。于是用Sonnet 3.7捏了一个pipe函数优化了这个问题。
该函数默认适用于OpenAI兼容API中的Gemini模型,可自行修改以适配添加其他模型。
函数如下:
import os
import json
import time
from pydantic import BaseModel, Field
import requests
from typing import List, Union, Iterator, Optional
# 设置 DEBUG 为 True 以启用详细日志记录
DEBUG = os.getenv("OPENAI_DEBUG", "False").lower() in ["true", "1", "yes"]
class Pipe:
class Valves(BaseModel):
# API 凭据
OPENAI_API_KEY: str = Field(
default="", description="API 密钥(必填)"
)
# 自定义端点配置
OPENAI_API_BASE: str = Field(
default="",
description="自定义 API 端点 URL,含/v1(留空为默认端点)",
)
# 流式传输设置,实现流畅输出
STREAMING_CHUNK_SIZE: int = Field(
default=1, description="每个流式传输块中发送的字符数(1 为逐字符)", ge=1
)
STREAMING_DELAY: float = Field(
default=0.005,
description="未指定模型的字间延迟(秒)(0 为无延迟)",
ge=0.0,
)
def __init__(self):
self.id = "better_stream"
self.type = "manifold"
self.name = "Google: "
# 可用模型,逗号分隔
self.available_models = "gemini-2.0-flash,gemini-2.0-flash-lite,gemini-2.0-pro-exp,gemini-2.0-flash-thinking-exp"
# 定义模型特定的延迟
self.model_delays = {
"gemini-2.0-flash": 0.002,
"gemini-2.0-flash-lite": 0.001,
"gemini-2.0-pro-exp": 0.005,
"gemini-2.0-flash-thinking-exp": 0.0005,
}
self.valves = self.Valves(
**{
"OPENAI_API_KEY": os.getenv("OPENAI_API_KEY", ""),
"OPENAI_API_BASE": os.getenv("OPENAI_API_BASE", ""),
"STREAMING_CHUNK_SIZE": int(
os.getenv("GEMINI_STREAMING_CHUNK_SIZE", "1")
),
"STREAMING_DELAY": float(os.getenv("GEMINI_STREAMING_DELAY", "0.01")),
}
)
if DEBUG and self.valves.OPENAI_API_BASE:
print(f"设置自定义 API 端点: {self.valves.OPENAI_API_BASE}")
def pipes(self) -> List[dict]:
"""返回可用模型列表"""
if not self.valves.OPENAI_API_KEY:
return [
{
"id": "error",
"name": "OPENAI_API_KEY 未设置。请在变量中更新 API 密钥。",
}
]
models = []
for model_id in self.available_models.split(","):
model_id = model_id.strip()
if model_id == "gemini-2.0-flash":
name = "Gemini 2.0 Flash"
elif model_id == "gemini-2.0-flash-lite":
name = "Gemini 2.0 Flash Lite"
elif model_id == "gemini-2.0-pro-exp":
name = "Gemini 2.0 Pro Experimental"
elif model_id == "gemini-2.0-flash-thinking-exp":
name = "Gemini 2.0 Flash Thinking Experimantal"
else:
name = model_id
models.append({"id": model_id, "name": name})
return models
def pipe(self, body: dict) -> Union[str, Iterator[str]]:
"""处理输入并使用 OpenAI 兼容 API 生成响应
参数:
body: 包含请求参数的字典
- model: 要使用的模型 ID
- messages: 包含角色和内容的message对象数组
- stream: 指示是否需要流式传输的布尔值
- temperature, top_p, top_k, max_tokens: 生成参数
- stop: 停止序列数组
返回:
如果未启用流式传输,则返回字符串响应;如果启用了流式传输,则返回字符串块的迭代器
"""
if not self.valves.OPENAI_API_KEY:
return "错误: OPENAI_API_KEY 未设置"
try:
model_id = body["model"]
# 处理不同的模型 ID 格式
if model_id.startswith("better_stream."):
model_id = model_id[14:]
model_id = model_id.lstrip(".")
messages = body["messages"]
stream = body.get("stream", False)
if DEBUG:
print("传入的 body:", json.dumps(body, indent=2))
# 构建 OpenAI API 请求数据
request_data = {
"model": model_id,
"messages": messages,
"temperature": body.get("temperature", 0.7),
"top_p": body.get("top_p", 0.9),
"max_tokens": body.get("max_tokens", 8192),
"stream": stream,
}
# 添加 stop 参数(如果存在)
if "stop" in body and body["stop"]:
request_data["stop"] = body["stop"]
# 设置 API URL
api_base = self.valves.OPENAI_API_BASE or "https://api.openai.com/v1"
api_url = f"{api_base}/chat/completions"
headers = {
"Content-Type": "application/json",
"Authorization": f"Bearer {self.valves.OPENAI_API_KEY}",
}
if DEBUG:
print("OpenAI 兼容 API 请求:")
print(" URL:", api_url)
print(" 模型:", model_id)
print(" 内容:", json.dumps(request_data, indent=2))
print(" 流式传输:", stream)
# 处理流式传输与非流式传输
if stream:
def stream_generator():
"""用于处理 OpenAI 流式响应的生成器"""
try:
response = requests.post(
api_url,
headers=headers,
json=request_data,
stream=True,
)
response.raise_for_status()
# 处理 OpenAI 流式响应格式
buffer = ""
model_delay = self.model_delays.get(
model_id, self.valves.STREAMING_DELAY
)
for line in response.iter_lines():
if not line:
continue
line = line.decode("utf-8")
# 处理 SSE 格式
if line.startswith("data: "):
line = line[6:] # 删除 'data: ' 前缀
if line == "[DONE]":
break
try:
data = json.loads(line)
if "choices" in data and len(data["choices"]) > 0:
delta = data["choices"][0].get("delta", {})
if "content" in delta and delta["content"]:
text = delta["content"]
buffer += text
# 按照设定的块大小发送文本
chunk_size = max(
1, self.valves.STREAMING_CHUNK_SIZE
)
while len(buffer) >= chunk_size:
sub_chunk = buffer[:chunk_size]
buffer = buffer[chunk_size:]
yield sub_chunk
# 使用模型特定的延迟
time.sleep(model_delay)
except json.JSONDecodeError as e:
if DEBUG:
print(f"JSON 解析错误: {e}, 行: {line}")
# 发送缓冲区中的剩余文本
if buffer:
yield buffer
except Exception as e:
error_msg = f"流式传输错误: {str(e)}"
if DEBUG:
print(error_msg)
yield error_msg
return stream_generator()
else:
# 非流式传输响应
response = requests.post(
api_url,
headers=headers,
json=request_data,
)
response.raise_for_status()
response_data = response.json()
if "choices" in response_data and len(response_data["choices"]) > 0:
return response_data["choices"][0]["message"]["content"]
else:
return "错误: 无效的 API 响应"
except Exception as e:
error_msg = f"错误: {str(e)}"
if DEBUG:
print(f"pipe 方法中发生错误: {e}")
return error_msg
- 在添加函数时,函数id必需设置为
better_stream。 - 请在函数代码
self.available_models中修改使用的模型(默认为Gemini 2.0系列)。 - 可自定义每个模型的输出延迟(修改
self.model_delays部分)。 - 变量API端点需包含“/v1”路径(如有)。
感觉功能还不是很完善,欢迎佬友们提出修改建议! ![]()