batch

package
v0.5.12 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 13, 2026 License: Apache-2.0 Imports: 9 Imported by: 0

Documentation

Overview

Package batch 提供 LLM 请求批处理功能

本包实现了 LLM 请求的批量处理和优化:

  • 请求合并:将多个相似请求合并处理

  • 请求队列:管理并发请求队列

  • 速率限制:控制 API 调用频率

  • 自动重试:处理临时错误

  • OpenAI Batch API: 批量处理

  • gRPC: 请求流和批处理

使用示例(Batcher 实例方式,需手动 Start/Stop):

batcher := batch.NewBatcher(provider, batch.DefaultConfig())
batcher.Start()
defer batcher.Stop()
results := batcher.BatchSubmit(ctx, requests)

使用示例(包级辅助函数,自动管理生命周期):

results := batch.BatchComplete(ctx, provider, requests)

Index

Constants

This section is empty.

Variables

View Source
var (
	// ErrBatchFailed 批处理失败
	ErrBatchFailed = errors.New("batch processing failed")

	// ErrQueueFull 队列已满
	ErrQueueFull = errors.New("request queue is full")

	// ErrRateLimited 被限流
	ErrRateLimited = errors.New("rate limited")

	// ErrTimeout 超时
	ErrTimeout = errors.New("timeout")

	// ErrProviderNotSet 未设置 Provider
	ErrProviderNotSet = errors.New("provider not set")
)

Functions

This section is empty.

Types

type Batcher

type Batcher struct {
	// contains filtered or unexported fields
}

Batcher 批处理器

func NewBatcher

func NewBatcher(provider Provider, config ...*Config) *Batcher

NewBatcher 创建批处理器

func (*Batcher) BatchSubmit

func (b *Batcher) BatchSubmit(ctx context.Context, requests []*Request) []*Response

BatchSubmit 批量提交请求

func (*Batcher) GetStats

func (b *Batcher) GetStats() Stats

Stats 获取统计信息

func (*Batcher) Start

func (b *Batcher) Start()

Start 启动批处理器

func (*Batcher) Stop

func (b *Batcher) Stop()

Stop 停止批处理器

停止流程:先关闭 stopChan 通知 collector 退出,collector 退出前会把 残留在 batch 缓冲里的请求做一次最终 flush(阻塞发送,确保不丢弃), 随后 close(batchChan) 让所有 worker 自然退出,最后 wg.Wait 等待收尾。

func (*Batcher) Submit

func (b *Batcher) Submit(ctx context.Context, req *Request) (*Response, error)

Submit 提交请求

type Config

type Config struct {
	// MaxBatchSize 最大批量大小
	MaxBatchSize int

	// MaxConcurrent 最大并发数
	MaxConcurrent int

	// QueueSize 队列大小
	QueueSize int

	// FlushInterval 刷新间隔
	FlushInterval time.Duration

	// Timeout 请求超时
	Timeout time.Duration

	// MaxRetries 最大重试次数
	MaxRetries int

	// RetryDelay 重试延迟
	RetryDelay time.Duration

	// RateLimit 速率限制(每秒请求数)
	RateLimit int

	// OnBatchStart 批次开始回调
	OnBatchStart func(batchID string, count int)

	// OnBatchComplete 批次完成回调
	OnBatchComplete func(batchID string, count int, duration time.Duration)

	// OnRequestComplete 请求完成回调
	OnRequestComplete func(req *Request, resp *Response)

	// OnError 错误回调
	OnError func(req *Request, err error)
}

Config 批处理配置

func DefaultConfig

func DefaultConfig() *Config

DefaultConfig 默认配置

type Message

type Message struct {
	Role    string `json:"role"`
	Content string `json:"content"`
}

Message 消息

type Provider

type Provider interface {
	// Complete 执行补全
	Complete(ctx context.Context, req *Request) (*Response, error)
}

Provider LLM 提供者接口(简化版)

type Request

type Request struct {
	// ID 请求 ID
	ID string

	// Messages 消息列表
	Messages []Message

	// Model 模型名称(可选,使用默认)
	Model string

	// MaxTokens 最大 token 数
	MaxTokens int

	// Temperature 温度
	Temperature float64

	// Metadata 元数据
	Metadata map[string]any
}

Request 批处理请求

type RequestBuilder

type RequestBuilder struct {
	// contains filtered or unexported fields
}

RequestBuilder 请求构建器

func NewRequestBuilder

func NewRequestBuilder() *RequestBuilder

NewRequestBuilder 创建请求构建器

func (*RequestBuilder) AddAssistantMessage

func (rb *RequestBuilder) AddAssistantMessage(content string) *RequestBuilder

AddAssistantMessage 添加助手消息

func (*RequestBuilder) AddSystemMessage

func (rb *RequestBuilder) AddSystemMessage(content string) *RequestBuilder

AddSystemMessage 添加系统消息

func (*RequestBuilder) AddUserMessage

func (rb *RequestBuilder) AddUserMessage(content string) *RequestBuilder

AddUserMessage 添加用户消息

func (*RequestBuilder) Build

func (rb *RequestBuilder) Build() *Request

Build 构建请求

func (*RequestBuilder) WithID

func (rb *RequestBuilder) WithID(id string) *RequestBuilder

WithID 设置 ID

func (*RequestBuilder) WithMaxTokens

func (rb *RequestBuilder) WithMaxTokens(max int) *RequestBuilder

WithMaxTokens 设置最大 token 数

func (*RequestBuilder) WithMetadata

func (rb *RequestBuilder) WithMetadata(key string, value any) *RequestBuilder

WithMetadata 设置元数据

func (*RequestBuilder) WithModel

func (rb *RequestBuilder) WithModel(model string) *RequestBuilder

WithModel 设置模型

func (*RequestBuilder) WithTemperature

func (rb *RequestBuilder) WithTemperature(temp float64) *RequestBuilder

WithTemperature 设置温度

type Response

type Response struct {
	// ID 请求 ID
	ID string

	// Content 响应内容
	Content string

	// TokensUsed 使用的 token 数
	TokensUsed int

	// FinishReason 完成原因
	FinishReason string

	// Error 错误(如果有)
	Error error

	// Latency 延迟(毫秒)
	Latency int64

	// Metadata 元数据
	Metadata map[string]any
}

Response 批处理响应

func BatchComplete

func BatchComplete(ctx context.Context, provider Provider, requests []*Request, config ...*Config) []*Response

BatchComplete 批量补全(简化接口)

func ParallelComplete

func ParallelComplete(ctx context.Context, provider Provider, requests []*Request, maxConcurrent int) []*Response

ParallelComplete 并行补全(不使用批处理器)

type Stats

type Stats struct {
	TotalRequests   int64         `json:"total_requests"`
	SuccessRequests int64         `json:"success_requests"`
	FailedRequests  int64         `json:"failed_requests"`
	TotalBatches    int64         `json:"total_batches"`
	TotalTokens     int64         `json:"total_tokens"`
	TotalLatency    time.Duration `json:"total_latency"`
	// contains filtered or unexported fields
}

Stats 统计信息

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL