适合Gemini!Open WebUI 函数:优化模型流式输出效果,大型响应块转逐字符输出

在部分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”路径(如有)。

感觉功能还不是很完善,欢迎佬友们提出修改建议! :tieba_087:

9 个赞

@TRG :tieba_087:

1 个赞

来了来了,支持一下

1 个赞

@Throttle 来看看吧 :tieba_087:

1 个赞

学到了 :hugs:

1 个赞

学到了,支持一下,谢谢

1 个赞

Claude也适合用,不过Claude的响应块比较小,观感影响不大

请问为什么要优化某个模型,而不是直接优化全部模型?我有点搞不懂

好注意欸,相当于把Fluidly stream large external response chunks优化又加回来了,不过有一点比较麻烦,如何优化使得不同模型,不同语言,不同速度生成的内容都尽可能流畅(就是如何根据缓存控制生成速度)

因为实际使用时并非所有模型都需要优化,且每个模型输出速度不同,建议自己调试一下。

或者说有没有中转api管理软件可以做到优化 :bili_102:
like:

  • 把脚本写到中转的 Cloudflare Worker 里。 用别人的中转也可以。
  • 直接魔改 NextChat 等api中转客户端。

workers要是实现这个功能得用到kv吧?那kv的存取写入量不得爆了

1 个赞

感谢大佬!

1 个赞

我已经用Cloudflare Workers实现了,而且能根据响应块大小和间隔时间自动调整输出延迟,且无需KV,我再测试一会会~

2 个赞

在函式中可以為Gemini 2.0 Flash Thinking添加可折疊思考過程的方法嗎?

目前gemini thinking在api中已经无法显示思考过程了!所以没办法做到

感觉还可以

这个项目不知道为什么又死灰复燃了:joy:已经做了cf workers版本

好的,因為我只有對Gemini這個流式有需求,謝謝師兄喇!

1 个赞

此话题已在最后回复的 30 天后被自动关闭。不再允许新回复。