downloader

package
v1.3.0 Latest Latest
Warning

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

Go to latest
Published: Jun 8, 2026 License: MIT Imports: 23 Imported by: 0

Documentation

Overview

Package downloader 实现了 scrapy-go 框架的下载器系统。

概述

downloader 包负责执行 HTTP 请求并返回响应,是框架中实际发起网络请求的组件。 通过 Slot 机制实现按域名分组的并发控制和请求延迟。 对应 Scrapy Python 版本中 scrapy.core.downloader 模块的功能。

核心类型

本包提供以下核心类型:

架构设计

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,内部采用队列驱动模型:

  1. 请求入队到 Slot 的 channel
  2. processQueue goroutine 串行出队
  3. 通过 lastSeen 时间戳精确控制请求间隔(DOWNLOAD_DELAY)
  4. 通过 [semaphore.Weighted] 控制并发传输数(CONCURRENT_REQUESTS_PER_DOMAIN)
  5. 不同 Slot 之间完全并行

中间件管理器

MiddlewareManager 管理下载器中间件链,支持接口隔离(ISP):

  • ProcessRequest:按优先级正序调用
  • ProcessResponse:按优先级逆序调用
  • ProcessException:按优先级逆序调用

中间件只需实现关心的接口,Manager 通过类型断言自动适配。

下载处理器选择

框架根据配置自动选择合适的下载处理器:

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

Examples

Constants

View Source
const DownloadProgressMetaKey = "download_progress_callback"

DownloadProgressMetaKey 是请求 Meta 中存储进度回调的键名。

View Source
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

type DownloadProgressCallback func(bytesRead int64, totalBytes int64, request *shttp.Request)

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) Close

func (d *Downloader) Close() error

Close 关闭下载器,释放所有资源。

func (*Downloader) Download

func (d *Downloader) Download(ctx context.Context, request *shttp.Request) (*shttp.Response, error)

Download 执行下载请求。

处理流程:

  1. 获取或创建请求对应的 Slot
  2. 将请求入队到 Slot 的队列中
  3. Slot 的 processQueue goroutine 负责: a. 计算并等待延迟(基于 lastSeen 时间戳) b. 获取传输信号量(控制并发) c. 执行实际下载
  4. 等待下载结果返回

注意:全局 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) Close

func (h *HTTPDownloadHandler) Close() error

Close 关闭 HTTP 处理器。

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。

func (*HTTPDownloadHandler) Download

func (h *HTTPDownloadHandler) Download(ctx context.Context, request *shttp.Request) (*shttp.Response, error)

Download 执行 HTTP 下载。 如果 request.Meta["proxy"] 设置了代理 URL,则通过该代理发送请求。

type ManagedTransport added in v1.0.2

type ManagedTransport struct {
	*http.Transport
	// contains filtered or unexported fields
}

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 表示一个带优先级的中间件条目。

中间件可以是以下任意接口的实现:

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) Count

func (m *MiddlewareManager) Count() int

Count 返回中间件数量。

func (*MiddlewareManager) Download

func (m *MiddlewareManager) Download(ctx context.Context, downloadFunc middleware.DownloadFunc, request *shttp.Request) (*shttp.Response, error)

Download 执行完整的中间件链处理流程。 对应 Scrapy 的 DownloaderMiddlewareManager.download_async。

处理流程:

  1. 正序调用 ProcessRequest(仅对实现了 RequestProcessor 的中间件)
  2. 调用 downloadFunc 执行实际下载
  3. 逆序调用 ProcessResponse(仅对实现了 ResponseProcessor 的中间件)
  4. 异常时逆序调用 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

Close 关闭处理器,释放所有连接。

func (*ProgressHTTPDownloadHandler) ConnPoolStats added in v1.0.2

func (h *ProgressHTTPDownloadHandler) ConnPoolStats() *ConnPoolStats

ConnPoolStats 返回连接池统计信息。

func (*ProgressHTTPDownloadHandler) Download added in v1.0.2

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) ActiveCount

func (s *Slot) ActiveCount() int

ActiveCount 返回活跃请求数。

func (*Slot) AddActive

func (s *Slot) AddActive(request *shttp.Request)

AddActive 将请求添加到活跃集合。

func (*Slot) AddTransferring

func (s *Slot) AddTransferring(request *shttp.Request)

AddTransferring 将请求添加到传输中集合。

func (*Slot) Close

func (s *Slot) Close()

Close 关闭 Slot,停止 worker pool 并等待所有 worker 退出。

func (*Slot) DownloadDelay

func (s *Slot) DownloadDelay() time.Duration

DownloadDelay 返回配置的下载延迟(公开方法,用于外部查询)。

func (*Slot) Enqueue

func (s *Slot) Enqueue(ctx context.Context, request *shttp.Request) (*shttp.Response, error)

Enqueue 将请求入队,阻塞等待结果返回。 这是外部调用的主要接口。 使用 downloadTaskPool 复用 downloadTask 对象和 resultCh channel, 避免每请求分配,减少约 10% 的内存分配开销。

func (*Slot) FreeTransferSlots

func (s *Slot) FreeTransferSlots() int

FreeTransferSlots 返回可用的传输槽位数。

func (*Slot) IsIdle

func (s *Slot) IsIdle() bool

IsIdle 检查 Slot 是否空闲(无活跃请求)。

func (*Slot) LastSeen

func (s *Slot) LastSeen() time.Time

LastSeen 返回最后一次活动时间。

func (*Slot) RemoveActive

func (s *Slot) RemoveActive(request *shttp.Request)

RemoveActive 从活跃集合中移除请求。

func (*Slot) RemoveTransferring

func (s *Slot) RemoveTransferring(request *shttp.Request)

RemoveTransferring 从传输中集合移除请求。

func (*Slot) SetDelay added in v1.0.3

func (s *Slot) SetDelay(delay time.Duration)

SetDelay 动态调整下载延迟。 由 AutoThrottle 扩展通过 Downloader.AdjustDelay 间接调用。 线程安全:通过 mu 保护 delay 字段的写入。

func (*Slot) TransferringCount

func (s *Slot) TransferringCount() int

TransferringCount 返回传输中请求数。

Directories

Path Synopsis
Package middleware 定义了下载器中间件的接口和内置实现。
Package middleware 定义了下载器中间件的接口和内置实现。
httpcache
Package httpcache 实现了 HTTP 缓存中间件,对应 Scrapy 的 HttpCacheMiddleware。
Package httpcache 实现了 HTTP 缓存中间件,对应 Scrapy 的 HttpCacheMiddleware。

Jump to

Keyboard shortcuts

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