Documentation
¶
Overview ¶
Package downloader 实现了 scrapy-go 框架的下载器系统。
概述 ¶
downloader 包负责执行 HTTP 请求并返回响应,是框架中实际发起网络请求的组件。 通过 Slot 机制实现按域名分组的并发控制和请求延迟。 对应 Scrapy Python 版本中 scrapy.core.downloader 模块的功能。
核心类型 ¶
本包提供以下核心类型:
- Downloader:下载器,管理 Slot 和全局并发控制
- Slot:域名级别的并发和延迟控制器
- DownloadHandler:下载处理器接口(协议适配层)
- HTTPDownloadHandler:基于 net/http 的 HTTP/1.1 和 HTTP/2 下载处理器(通过 ALPN 自动协商)
- ProgressHTTPDownloadHandler:支持下载进度回调的处理器
- ConnPoolConfig:连接池精细化配置(14 项参数)
- ConnPoolStats:连接池运行时统计(atomic 无锁)
- ManagedTransport:带统计功能的 http.Transport 包装
- MiddlewareManager:下载器中间件管理器
- MiddlewareEntry:带优先级的中间件条目
架构设计 ¶
Downloader 采用 Slot 机制实现精细化的并发控制:
Engine
│
▼
┌─────────────────────────────────────────┐
│ Downloader │
│ ┌─────────────────────────────────┐ │
│ │ 全局 active 集合 │ │
│ │ (CONCURRENT_REQUESTS 上限) │ │
│ └─────────────────────────────────┘ │
│ │ │ │ │
│ ▼ ▼ ▼ │
│ ┌────────┐ ┌────────┐ ┌────────┐ │
│ │ Slot A │ │ Slot B │ │ Slot C │ │
│ │(域名A) │ │(域名B) │ │(域名C) │ │
│ └────────┘ └────────┘ └────────┘ │
└─────────────────────────────────────────┘
│
▼
DownloadHandler
┌──────────┼──────────┐
▼ ▼ ▼
HTTP/1.1 Progress WithConfig
Handler Handler Handler
(默认) (进度回调) (连接池统计+h2c)
Slot 调度模型 ¶
每个 Slot 对应一个域名/IP,内部采用队列驱动模型:
- 请求入队到 Slot 的 channel
- processQueue goroutine 串行出队
- 通过 lastSeen 时间戳精确控制请求间隔(DOWNLOAD_DELAY)
- 通过 [semaphore.Weighted] 控制并发传输数(CONCURRENT_REQUESTS_PER_DOMAIN)
- 不同 Slot 之间完全并行
中间件管理器 ¶
MiddlewareManager 管理下载器中间件链,支持接口隔离(ISP):
- ProcessRequest:按优先级正序调用
- ProcessResponse:按优先级逆序调用
- ProcessException:按优先级逆序调用
中间件只需实现关心的接口,Manager 通过类型断言自动适配。
下载处理器选择 ¶
框架根据配置自动选择合适的下载处理器:
- HTTP2_ENABLED=true → HTTPDownloadHandler(通过 NewHTTPDownloadHandlerWithConfig 创建,启用 ForceAttemptHTTP2)
- DOWNLOAD_PROGRESS_ENABLED=true → ProgressHTTPDownloadHandler(进度回调)
- 默认 → HTTPDownloadHandler(HTTP/1.1 标准处理器)
HTTPDownloadHandler 通过 Go 标准库 net/http 原生支持 HTTP/2 ALPN 自动协商, 无需独立的 HTTP/2 Handler。当 HTTP2_ENABLED=true 时,通过 ConnPoolConfig.ForceHTTP2 设置 ForceAttemptHTTP2=true,在自定义 DialContext 场景下也能启用 HTTP/2。 当 HTTP2_ALLOW_H2C=true 时,注册 http2.Transport 作为 http:// scheme 的 handler, 支持内网/测试场景的 HTTP/2 over cleartext。
ProgressHTTPDownloadHandler 通过 request.Meta["download_progress_callback"] 设置回调函数, 支持已知/未知大小响应的进度报告,可配置最小报告间隔避免回调过于频繁。
连接池管理 ¶
ConnPoolConfig 提供 14 项连接池参数的精细化配置:
- MaxIdleConns / MaxIdleConnsPerHost / MaxConnsPerHost:连接数控制
- IdleConnTimeout / DialTimeout / TLSHandshakeTimeout:超时控制
- DisableKeepAlives / ForceHTTP2 / AllowH2C:协议行为控制
- WriteBufferSize / ReadBufferSize:缓冲区大小
- TLSInsecureSkipVerify:TLS 验证控制
ManagedTransport 包装 http.Transport,通过 ConnPoolStats 提供 连接创建/复用/关闭/TLS 握手等运行时统计(atomic 无锁,零竞争开销)。
连接池配置通过 CONNPOOL_* 前缀的 Settings 键名读取, 使用 ConnPoolConfigFromSettings 函数从 Settings 中构建配置。
配置项 ¶
Downloader 通过 Settings 读取以下配置:
- CONCURRENT_REQUESTS:全局最大并发数(默认 16)
- CONCURRENT_REQUESTS_PER_DOMAIN:每域名最大并发数(默认 8)
- DOWNLOAD_DELAY:请求间隔(默认 0)
- RANDOMIZE_DOWNLOAD_DELAY:是否随机化延迟(默认 true)
- DOWNLOAD_TIMEOUT:下载超时时间(默认 180s)
- HTTP2_ENABLED:启用 HTTP/2 协议支持(默认 false,通过 ALPN 自动协商)
- HTTP2_ALLOW_H2C:启用 HTTP/2 over cleartext 支持(默认 false,适用于内网/测试场景)
- DOWNLOAD_PROGRESS_ENABLED:启用下载进度回调(默认 false)
- DOWNLOAD_PROGRESS_MIN_INTERVAL:进度报告最小间隔(默认 100ms)
- CONNPOOL_MAX_IDLE_CONNS:最大空闲连接总数(默认 100)
- CONNPOOL_MAX_IDLE_CONNS_PER_HOST:每 host 最大空闲连接数(默认 10)
- CONNPOOL_MAX_CONNS_PER_HOST:每 host 最大连接数(默认 0,不限制)
- CONNPOOL_IDLE_CONN_TIMEOUT:空闲连接超时(默认 90s)
- CONNPOOL_TLS_HANDSHAKE_TIMEOUT:TLS 握手超时(默认 10s)
- CONNPOOL_DIAL_TIMEOUT:TCP 连接超时(默认 30s)
- CONNPOOL_DISABLE_KEEPALIVES:禁用 HTTP keep-alive(默认 false)
信号系统集成 ¶
Downloader 在关键节点发送以下信号:
- RequestReachedDownloader:请求到达下载器
- RequestLeftDownloader:请求离开下载器
- ResponseDownloaded:响应下载完成
与 Scrapy 的差异 ¶
- 使用 [semaphore.Weighted] 替代 Twisted DeferredSemaphore 控制并发
- 使用 goroutine + channel 替代 Twisted callLater 实现延迟控制
- 使用 sync.RWMutex 保护共享状态,替代 Twisted 的协作式并发
- 全局 active 集合的添加/移除由 Engine 在同步路径中完成(对齐 Scrapy 原版)
- Slot GC 机制自动清理空闲超时的 Slot,防止内存泄漏
- 禁用 net/http 自动重定向和自动解压,由中间件统一管理
- 支持通过 request.Meta["proxy"] 动态设置代理
- 内置 HTTP/2 支持(Go 标准库 net/http 原生 ALPN 自动协商,无需独立 Handler)
- 支持 h2c(HTTP/2 over cleartext)配置,适用于内网/测试场景
- 连接池参数可通过 Settings 精细化配置(Scrapy 受限于 Twisted reactor)
- 下载进度回调通过 Meta 传递(Scrapy 通过 Signal 实现,粒度较粗)
并发安全 ¶
Downloader 的所有公共方法均为并发安全。 Slot 内部通过 channel 和 semaphore 实现并发控制,无需外部同步。
Package downloader 实现了 scrapy-go 框架的下载器系统。
下载器负责执行 HTTP 请求并返回响应,通过 Slot 机制控制并发和延迟。 对应 Scrapy Python 版本中 scrapy.core.downloader 模块的功能。
Index ¶
- Constants
- type ConnPoolConfig
- type ConnPoolStats
- type DownloadHandler
- type DownloadProgressCallback
- type Downloader
- func (d *Downloader) ActiveCount() int
- func (d *Downloader) AddActive(request *shttp.Request)
- func (d *Downloader) AdjustDelay(slotKey string, delay time.Duration)
- func (d *Downloader) Close() error
- func (d *Downloader) Download(ctx context.Context, request *shttp.Request) (*shttp.Response, error)
- func (d *Downloader) NeedsBackout() bool
- func (d *Downloader) RemoveActive(request *shttp.Request)
- func (d *Downloader) StartSlotGC()
- func (d *Downloader) TotalConcurrency() int
- type HTTPDownloadHandler
- type ManagedTransport
- type MiddlewareEntry
- type MiddlewareManager
- type ProgressHTTPDownloadHandler
- type Slot
- func (s *Slot) ActiveCount() int
- func (s *Slot) AddActive(request *shttp.Request)
- func (s *Slot) AddTransferring(request *shttp.Request)
- func (s *Slot) Close()
- func (s *Slot) DownloadDelay() time.Duration
- func (s *Slot) Enqueue(ctx context.Context, request *shttp.Request) (*shttp.Response, error)
- func (s *Slot) FreeTransferSlots() int
- func (s *Slot) IsIdle() bool
- func (s *Slot) LastSeen() time.Time
- func (s *Slot) RemoveActive(request *shttp.Request)
- func (s *Slot) RemoveTransferring(request *shttp.Request)
- func (s *Slot) SetDelay(delay time.Duration)
- func (s *Slot) TransferringCount() int
Examples ¶
Constants ¶
const DownloadProgressMetaKey = "download_progress_callback"
DownloadProgressMetaKey 是请求 Meta 中存储进度回调的键名。
const (
// DownloadSlotMetaKey 是请求 Meta 中存储 Slot key 的键名。
DownloadSlotMetaKey = "download_slot"
)
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type ConnPoolConfig ¶ added in v1.0.2
type ConnPoolConfig struct {
// MaxIdleConns 控制所有 host 的最大空闲连接总数。
// 默认值:100。设为 0 表示不限制。
MaxIdleConns int
// MaxIdleConnsPerHost 控制每个 host 的最大空闲连接数。
// 默认值:10。对于高并发单站点爬取,建议设为与 CONCURRENT_REQUESTS_PER_DOMAIN 相同。
MaxIdleConnsPerHost int
// MaxConnsPerHost 控制每个 host 的最大连接数(含活跃和空闲)。
// 默认值:0(不限制)。设置此值可防止对单个站点建立过多连接。
MaxConnsPerHost int
// IdleConnTimeout 控制空闲连接在被关闭前的最大存活时间。
// 默认值:90s。对于长时间运行的爬虫,可适当增大以复用连接。
IdleConnTimeout time.Duration
// TLSHandshakeTimeout 控制 TLS 握手的超时时间。
// 默认值:10s。
TLSHandshakeTimeout time.Duration
// ResponseHeaderTimeout 控制等待响应头的超时时间。
// 默认值:0(不限制,由 DOWNLOAD_TIMEOUT 统一控制)。
ResponseHeaderTimeout time.Duration
// ExpectContinueTimeout 控制发送 Expect: 100-continue 后等待服务器响应的超时。
// 默认值:1s。
ExpectContinueTimeout time.Duration
// DisableKeepAlives 禁用 HTTP keep-alive,每次请求使用新连接。
// 默认值:false。仅在调试或特殊场景下启用。
DisableKeepAlives bool
// ForceHTTP2 强制使用 HTTP/2 协议。
// 默认值:false。启用后将使用 HTTP/2 专用 Transport。
ForceHTTP2 bool
// WriteBufferSize 设置 Transport 的写缓冲区大小。
// 默认值:0(使用标准库默认值 4KB)。
WriteBufferSize int
// ReadBufferSize 设置 Transport 的读缓冲区大小。
// 默认值:0(使用标准库默认值 4KB)。
ReadBufferSize int
// DialTimeout 控制建立 TCP 连接的超时时间。
// 默认值:30s。
DialTimeout time.Duration
// DialKeepAlive 控制 TCP keep-alive 探测间隔。
// 默认值:30s。设为负值禁用 TCP keep-alive。
DialKeepAlive time.Duration
// TLSInsecureSkipVerify 跳过 TLS 证书验证。
// 默认值:false。仅用于测试或信任的内网环境。
TLSInsecureSkipVerify bool
// AllowH2C 启用 HTTP/2 over cleartext(h2c)支持。
// 默认值:false。启用后将注册 http2.Transport 作为 http:// scheme 的 handler,
// 允许在不使用 TLS 的情况下使用 HTTP/2 协议(适用于内网/测试场景)。
AllowH2C bool
}
ConnPoolConfig 定义连接池的精细化配置。 允许用户根据目标站点特性调整连接池参数,以获得最佳性能。
对应 Scrapy 中通过 Twisted reactor 配置的连接池参数, 但 Go 的 net/http.Transport 提供了更丰富的连接池控制能力。
func ConnPoolConfigFromSettings ¶ added in v1.0.2
func ConnPoolConfigFromSettings(getInt func(string, int) int, getDuration func(string, time.Duration) time.Duration, getBool func(string, bool) bool) *ConnPoolConfig
ConnPoolConfigFromSettings 从 Settings 中读取连接池配置。 配置键名以 CONNPOOL_ 为前缀。
func DefaultConnPoolConfig ¶ added in v1.0.2
func DefaultConnPoolConfig() *ConnPoolConfig
DefaultConnPoolConfig 返回默认的连接池配置。
type ConnPoolStats ¶ added in v1.0.2
type ConnPoolStats struct {
// TotalConnsCreated 累计创建的连接数。
TotalConnsCreated atomic.Int64
// TotalConnsReused 累计复用的连接数。
TotalConnsReused atomic.Int64
// TotalConnsClosed 累计关闭的连接数。
TotalConnsClosed atomic.Int64
// TotalTLSHandshakes 累计 TLS 握手次数。
TotalTLSHandshakes atomic.Int64
// ActiveConns 当前活跃连接数。
ActiveConns atomic.Int64
// IdleConns 当前空闲连接数。
IdleConns atomic.Int64
}
ConnPoolStats 记录连接池的运行时统计信息。 通过 atomic 操作保证并发安全,无锁读取。
func (*ConnPoolStats) Snapshot ¶ added in v1.0.2
func (s *ConnPoolStats) Snapshot() map[string]int64
Snapshot 返回连接池统计的快照(用于日志和监控)。
type DownloadHandler ¶
type DownloadHandler interface {
// Download 执行下载请求并返回响应。
Download(ctx context.Context, request *shttp.Request) (*shttp.Response, error)
// Close 关闭处理器,释放资源。
Close() error
}
DownloadHandler 定义下载处理器接口。 不同协议(http、https、ftp 等)可以有不同的处理器实现。
type DownloadProgressCallback ¶ added in v1.0.2
DownloadProgressCallback 定义下载进度回调函数类型。 参数:
- bytesRead: 已读取的字节数
- totalBytes: 总字节数(-1 表示未知,如 chunked 传输)
- request: 关联的请求
type Downloader ¶
type Downloader struct {
// contains filtered or unexported fields
}
Downloader 管理 HTTP 下载,通过 Slot 机制控制并发和延迟。
核心调度策略(对齐 Scrapy 原版):
- 每个域名/IP 对应一个 Slot
- Slot 内部通过队列驱动串行出队,用 lastSeen 时间戳精确控制请求间隔
- 不同 Slot 之间完全并行
- 全局通过 active 集合控制总并发数
对应 Scrapy 的 Downloader 类。
func NewDownloader ¶
func NewDownloader(s *settings.Settings, handler DownloadHandler, signals *signal.Manager, sc stats.Collector, logger *slog.Logger) *Downloader
NewDownloader 创建一个新的下载器。
Example ¶
ExampleNewDownloader 演示创建下载器。
package main
import (
"fmt"
"log/slog"
"time"
"github.com/dplcz/scrapy-go/pkg/downloader"
"github.com/dplcz/scrapy-go/pkg/settings"
"github.com/dplcz/scrapy-go/pkg/signal"
"github.com/dplcz/scrapy-go/pkg/stats"
)
func main() {
// 创建配置
s := settings.New()
s.Set("CONCURRENT_REQUESTS", 16, settings.PriorityDefault)
s.Set("CONCURRENT_REQUESTS_PER_DOMAIN", 8, settings.PriorityDefault)
s.Set("DOWNLOAD_DELAY", "0s", settings.PriorityDefault)
s.Set("DOWNLOAD_TIMEOUT", "30s", settings.PriorityDefault)
// 创建下载器
handler := downloader.NewHTTPDownloadHandler(30 * time.Second)
signals := signal.NewManager(nil)
statsCollector := stats.NewMemoryCollector(false, slog.Default())
dl := downloader.NewDownloader(s, handler, signals, statsCollector, slog.Default())
fmt.Println("Total concurrency:", dl.TotalConcurrency())
fmt.Println("Active count:", dl.ActiveCount())
fmt.Println("Needs backout:", dl.NeedsBackout())
}
Output: Total concurrency: 16 Active count: 0 Needs backout: false
func (*Downloader) ActiveCount ¶
func (d *Downloader) ActiveCount() int
ActiveCount 返回当前活跃请求数。 使用 atomic 计数器实现无锁读取。
func (*Downloader) AddActive ¶
func (d *Downloader) AddActive(request *shttp.Request)
AddActive 将请求添加到全局活跃集合。 由 Engine 在同步调度路径中调用,确保 NeedsBackout() 能立即看到最新计数。 对齐 Scrapy 原版:Downloader.fetch() 入口处同步执行 self.active.add(request)。
先更新 atomic 计数器(NeedsBackout 无锁快速路径立即可见), 再更新 active map(用于精确追踪和调试)。
func (*Downloader) AdjustDelay ¶ added in v1.0.3
func (d *Downloader) AdjustDelay(slotKey string, delay time.Duration)
AdjustDelay 调整指定 Slot 的下载延迟。 由 AutoThrottle 扩展调用,实现 extension.DelayAdjuster 接口。 如果指定的 Slot 不存在,则忽略(Slot 可能已被 GC 回收)。
func (*Downloader) Download ¶
Download 执行下载请求。
处理流程:
- 获取或创建请求对应的 Slot
- 将请求入队到 Slot 的队列中
- Slot 的 processQueue goroutine 负责: a. 计算并等待延迟(基于 lastSeen 时间戳) b. 获取传输信号量(控制并发) c. 执行实际下载
- 等待下载结果返回
注意:全局 active 集合的添加/移除由 Engine 在同步路径中完成(对齐 Scrapy 原版), 以确保 NeedsBackout() 在调度循环中能看到最新的计数。 Slot 级别的 active 仍在此处管理。
func (*Downloader) NeedsBackout ¶
func (d *Downloader) NeedsBackout() bool
NeedsBackout 检查是否需要回退(活跃请求数达到总并发上限)。 使用 atomic 计数器实现无锁快速路径,避免调度循环高频调用时的 RWMutex 竞争。
func (*Downloader) RemoveActive ¶
func (d *Downloader) RemoveActive(request *shttp.Request)
RemoveActive 从全局活跃集合中移除请求。 由 Engine 在请求完成后调用。
先更新 active map,再减少 atomic 计数器, 确保不会出现计数器已减但 map 中仍存在的短暂不一致。
func (*Downloader) StartSlotGC ¶
func (d *Downloader) StartSlotGC()
StartSlotGC 启动 Slot 垃圾回收。 定期清理空闲超过指定时间的 Slot。
func (*Downloader) TotalConcurrency ¶
func (d *Downloader) TotalConcurrency() int
TotalConcurrency 返回全局并发限制值。
type HTTPDownloadHandler ¶
type HTTPDownloadHandler struct {
// contains filtered or unexported fields
}
HTTPDownloadHandler 是基于 net/http 的 HTTP 下载处理器。 支持 HTTP/1.1 和 HTTP/2(通过 ALPN 自动协商)。
当通过 NewHTTPDownloadHandlerWithConfig 创建时,支持连接池精细化配置和运行时统计。 当通过 NewHTTPDownloadHandler 创建时,使用默认配置,行为与重构前完全一致。
func NewHTTPDownloadHandler ¶
func NewHTTPDownloadHandler(timeout time.Duration) *HTTPDownloadHandler
NewHTTPDownloadHandler 创建一个新的 HTTP 下载处理器。 使用默认的连接池配置,不启用连接池统计。
Example ¶
ExampleNewHTTPDownloadHandler 演示创建 HTTP 下载处理器。
package main
import (
"fmt"
"time"
"github.com/dplcz/scrapy-go/pkg/downloader"
)
func main() {
// 创建带 30 秒超时的 HTTP 处理器
handler := downloader.NewHTTPDownloadHandler(30 * time.Second)
// 处理器实现了 DownloadHandler 接口
var _ downloader.DownloadHandler = handler
fmt.Println("Handler created successfully")
// 关闭处理器
handler.Close()
}
Output: Handler created successfully
func NewHTTPDownloadHandlerWithConfig ¶ added in v1.1.6
func NewHTTPDownloadHandlerWithConfig(timeout time.Duration, config *ConnPoolConfig) *HTTPDownloadHandler
NewHTTPDownloadHandlerWithConfig 创建一个带精细化配置的 HTTP 下载处理器。
相比 NewHTTPDownloadHandler,此构造函数支持:
- ConnPoolConfig 配置注入(连接池参数、HTTP/2 控制、h2c 支持等)
- ConnPoolStats 连接池运行时统计
- ForceHTTP2=true 时设置 ForceAttemptHTTP2(用于自定义 DialContext 场景)
- AllowH2C=true 时注册 http2.Transport 作为 http:// scheme 的 handler
参数:
- timeout: 全局下载超时时间
- config: 连接池配置(传 nil 使用默认配置)
func (*HTTPDownloadHandler) Config ¶ added in v1.1.6
func (h *HTTPDownloadHandler) Config() *ConnPoolConfig
Config 返回当前连接池配置。 仅在通过 NewHTTPDownloadHandlerWithConfig 创建时可用,否则返回 nil。
func (*HTTPDownloadHandler) ConnPoolStats ¶ added in v1.1.6
func (h *HTTPDownloadHandler) ConnPoolStats() *ConnPoolStats
ConnPoolStats 返回连接池统计信息。 仅在通过 NewHTTPDownloadHandlerWithConfig 创建时可用,否则返回 nil。
type ManagedTransport ¶ added in v1.0.2
ManagedTransport 是对 http.Transport 的包装,提供连接池统计和精细化配置。
func NewManagedTransport ¶ added in v1.0.2
func NewManagedTransport(config *ConnPoolConfig) *ManagedTransport
NewManagedTransport 根据配置创建一个带统计功能的 Transport。
func (*ManagedTransport) CloseIdleConnections ¶ added in v1.0.2
func (mt *ManagedTransport) CloseIdleConnections()
CloseIdleConnections 关闭所有空闲连接。
func (*ManagedTransport) Stats ¶ added in v1.0.2
func (mt *ManagedTransport) Stats() *ConnPoolStats
Stats 返回连接池统计信息。
type MiddlewareEntry ¶
type MiddlewareEntry struct {
// Middleware 存储中间件实例。
// 至少需要实现 RequestProcessor / ResponseProcessor / ExceptionProcessor 之一。
Middleware any
Name string
Priority int
// contains filtered or unexported fields
}
MiddlewareEntry 表示一个带优先级的中间件条目。
中间件可以是以下任意接口的实现:
- middleware.RequestProcessor:处理请求
- middleware.ResponseProcessor:处理响应
- middleware.ExceptionProcessor:处理异常
- middleware.DownloaderMiddleware:全功能(同时满足以上三者)
Manager 通过类型断言在运行时适配,仅调用中间件实际实现的方法。
type MiddlewareManager ¶
type MiddlewareManager struct {
// contains filtered or unexported fields
}
MiddlewareManager 管理下载器中间件链。 对应 Scrapy 的 DownloaderMiddlewareManager。
支持接口隔离:中间件只需实现关心的细粒度接口, Manager 通过类型断言自动适配,跳过未实现的处理阶段。
func NewMiddlewareManager ¶
func NewMiddlewareManager(logger *slog.Logger) *MiddlewareManager
NewMiddlewareManager 创建一个新的中间件管理器。
Example ¶
ExampleNewMiddlewareManager 演示创建中间件管理器。
package main
import (
"fmt"
"log/slog"
"github.com/dplcz/scrapy-go/pkg/downloader"
)
func main() {
mgr := downloader.NewMiddlewareManager(slog.Default())
fmt.Println("Middleware count:", mgr.Count())
}
Output: Middleware count: 0
func (*MiddlewareManager) AddMiddleware ¶
func (m *MiddlewareManager) AddMiddleware(mw any, name string, priority int)
AddMiddleware 添加一个中间件。 中间件按优先级排序,优先级数值小的先执行 ProcessRequest。
mw 可以是以下类型之一:
- middleware.DownloaderMiddleware(全功能,向后兼容)
- middleware.RequestProcessor(仅处理请求)
- middleware.ResponseProcessor(仅处理响应)
- middleware.ExceptionProcessor(仅处理异常)
- 或以上接口的任意组合
func (*MiddlewareManager) Download ¶
func (m *MiddlewareManager) Download(ctx context.Context, downloadFunc middleware.DownloadFunc, request *shttp.Request) (*shttp.Response, error)
Download 执行完整的中间件链处理流程。 对应 Scrapy 的 DownloaderMiddlewareManager.download_async。
处理流程:
- 正序调用 ProcessRequest(仅对实现了 RequestProcessor 的中间件)
- 调用 downloadFunc 执行实际下载
- 逆序调用 ProcessResponse(仅对实现了 ResponseProcessor 的中间件)
- 异常时逆序调用 ProcessException(仅对实现了 ExceptionProcessor 的中间件)
当中间件返回 NewRequestError 时,该错误会被直接传播给调用方(Engine), 由 Engine 负责将新请求重新调度到 Scheduler。
type ProgressHTTPDownloadHandler ¶ added in v1.0.2
type ProgressHTTPDownloadHandler struct {
// contains filtered or unexported fields
}
ProgressHTTPDownloadHandler 是支持下载进度回调的 HTTP 下载处理器。
在标准 HTTPDownloadHandler 的基础上增加了下载进度通知能力:
- 通过 Request.Meta["download_progress_callback"] 设置进度回调
- 支持已知大小(Content-Length)和未知大小(chunked)的进度报告
- 进度回调在下载 goroutine 中同步调用,不引入额外 goroutine
- 可配置进度报告的最小间隔,避免高频回调影响性能
使用场景:
- 大文件下载进度展示
- 下载速率监控
- 超大响应体的流式处理决策
func NewProgressHTTPDownloadHandler ¶ added in v1.0.2
func NewProgressHTTPDownloadHandler(timeout time.Duration, config *ConnPoolConfig, minReportInterval time.Duration) *ProgressHTTPDownloadHandler
NewProgressHTTPDownloadHandler 创建一个支持进度回调的 HTTP 下载处理器。
参数:
- timeout: 全局下载超时时间
- config: 连接池配置(传 nil 使用默认配置)
- minReportInterval: 进度报告最小间隔(传 0 使用默认值 100ms)
func (*ProgressHTTPDownloadHandler) Close ¶ added in v1.0.2
func (h *ProgressHTTPDownloadHandler) Close() error
Close 关闭处理器,释放所有连接。
func (*ProgressHTTPDownloadHandler) ConnPoolStats ¶ added in v1.0.2
func (h *ProgressHTTPDownloadHandler) ConnPoolStats() *ConnPoolStats
ConnPoolStats 返回连接池统计信息。
func (*ProgressHTTPDownloadHandler) Download ¶ added in v1.0.2
func (h *ProgressHTTPDownloadHandler) Download(ctx context.Context, request *shttp.Request) (*shttp.Response, error)
Download 执行带进度回调的 HTTP 下载。
如果请求 Meta 中设置了 download_progress_callback,则在读取响应体时 定期调用回调函数报告进度。否则行为与标准 HTTPDownloadHandler 完全一致。
下载完成后会将 download_latency(time.Duration)设置到 request.Meta 中, 供 AutoThrottle、Telemetry 等扩展在 RequestLeftDownloader 信号中消费。
type Slot ¶
type Slot struct {
// contains filtered or unexported fields
}
Slot 控制单个域名/IP 的并发和延迟。
Worker Pool 化设计(P4-007i):
- 请求先入队到 queue channel
- 启动固定数量的 worker goroutine(N = concurrency),从 queue 直接消费 task
- worker 复用整个 Slot 生命周期,消除 per-request goroutine 创建/销毁开销
- 通过 gateMu + lastSeen 串行化 delay 等待,保证同一 Slot 内请求间隔语义
- 不同 Slot 之间完全并行(每个 Slot 有自己的 worker pool)
与旧设计(per-request `go func()`)的对比:
- 减少 goroutine 创建/销毁:每请求节省 ~2 us 调度开销
- 稳定的 goroutine 数量,降低 GC 扫描压力
- 替换 transferSem(semaphore.Weighted),由 worker 数量天然限制并发
对应 Scrapy 的 Slot 类。
func NewSlot ¶
func NewSlot( concurrency int, delay time.Duration, randomizeDelay bool, downloadFn func(ctx context.Context, request *shttp.Request) (*shttp.Response, error), ) *Slot
NewSlot 创建一个新的下载 Slot 并启动 worker pool。
worker pool 大小等于 concurrency,每个 worker 是一个长生命周期的 goroutine, 直接从 queue channel 消费 task 并执行下载。
func (*Slot) AddTransferring ¶
AddTransferring 将请求添加到传输中集合。
func (*Slot) DownloadDelay ¶
DownloadDelay 返回配置的下载延迟(公开方法,用于外部查询)。
func (*Slot) Enqueue ¶
Enqueue 将请求入队,阻塞等待结果返回。 这是外部调用的主要接口。 使用 downloadTaskPool 复用 downloadTask 对象和 resultCh channel, 避免每请求分配,减少约 10% 的内存分配开销。
func (*Slot) FreeTransferSlots ¶
FreeTransferSlots 返回可用的传输槽位数。
func (*Slot) RemoveActive ¶
RemoveActive 从活跃集合中移除请求。
func (*Slot) RemoveTransferring ¶
RemoveTransferring 从传输中集合移除请求。