diff --git a/cmd/root.go b/cmd/root.go index c85f9f9..553b37f 100644 --- a/cmd/root.go +++ b/cmd/root.go @@ -4,20 +4,17 @@ import ( "fmt" "log" "os" - "time" "github.com/spf13/cobra" - "gorm.io/driver/mysql" - "gorm.io/gorm" - "gorm.io/gorm/logger" "my-spiders/internal/config" + "my-spiders/internal/ingest" ) var ( - DB *gorm.DB - cfgFile string - debug bool // 命令行 debug 标志 + IngestClient *ingest.Client // 入库管线客户端(enabled 时非 nil) + cfgFile string + debug bool // 命令行 debug 标志 rootCmd = &cobra.Command{ Use: "spider-cli", @@ -39,10 +36,10 @@ func init() { // 新增全局 -d / --debug 参数 rootCmd.PersistentFlags().BoolVarP(&debug, "debug", "d", false, "调试模式 (仅输出不入库)") - cobra.OnInitialize(initConfigAndDB) + cobra.OnInitialize(initConfig) } -func initConfigAndDB() { +func initConfig() { // 1. 初始化配置文件 if err := config.InitConfig(cfgFile); err != nil { log.Fatalf("[致命错误] %v", err) @@ -53,39 +50,14 @@ func initConfigAndDB() { config.GlobalConfig.App.Debug = true } - isDebug := config.GlobalConfig.App.Debug - - // 2. 初始化数据库连接 - dbCfg := config.GlobalConfig.Database - var err error - DB, err = gorm.Open(mysql.Open(dbCfg.DSN), &gorm.Config{ - Logger: logger.Default.LogMode(logger.Warn), - }) - - if err != nil { - if isDebug { - log.Println("[Warning] 数据库连接失败,当前处于 Debug 模式,将跳过数据库操作...") - DB = nil - return - } - log.Fatalf("[致命错误] 数据库连接失败: %v", err) + // 2. 入库管线:enabled 时构建 HMAC 上报客户端,并用后台盐值初始化 hashid 编码。 + // spider 不再直连数据库,所有落库经后端 ingest 接口完成。 + ic := config.GlobalConfig.Ingest + if ic.Enabled && ic.Endpoint != "" { + ingest.InitHashID(ic.HashIDSecret) + IngestClient = ingest.NewClient(ic.Endpoint, ic.Secret) + log.Printf("[Info] 入库管线已启用 -> %s", ic.Endpoint) + } else { + IngestClient = nil } - - // 3. 设置连接池参数 - sqlDB, err := DB.DB() - if err != nil || sqlDB.Ping() != nil { - if isDebug { - log.Println("[Warning] 数据库 Ping 失败,当前处于 Debug 模式,将跳过数据库操作...") - DB = nil - return - } - log.Fatalf("[致命错误] 数据库不可用: %v", err) - } - - sqlDB.SetMaxOpenConns(dbCfg.MaxOpenConns) - sqlDB.SetMaxIdleConns(dbCfg.MaxIdleConns) - sqlDB.SetConnMaxLifetime(time.Duration(dbCfg.ConnMaxLifetime) * time.Minute) - sqlDB.SetConnMaxIdleTime(time.Duration(dbCfg.ConnMaxIdleTime) * time.Minute) - - log.Printf("[Info] MySQL 连接池初始化完成 | MaxOpen: %d | MaxIdle: %d", dbCfg.MaxOpenConns, dbCfg.MaxIdleConns) } diff --git a/cmd/theimpression.go b/cmd/theimpression.go index 8421652..e0d5e42 100644 --- a/cmd/theimpression.go +++ b/cmd/theimpression.go @@ -17,7 +17,7 @@ var theImpressionCmd = &cobra.Command{ } // 复用全局 -d/--debug 标志 isDebug := debug - s := spider.NewTheImpressionSpider(DB, theImpressionMaxCo) + s := spider.NewTheImpressionSpider(theImpressionMaxCo, IngestClient) s.Debug = isDebug s.Run() }, diff --git a/cmd/vogue.go b/cmd/vogue.go index 584f47a..e5daaf6 100644 --- a/cmd/vogue.go +++ b/cmd/vogue.go @@ -7,9 +7,9 @@ import ( ) var ( - brandID int64 - onlyPlatform bool - maxCo int + brandID int64 + maxCo int + maxShows int ) var vogueCmd = &cobra.Command{ @@ -17,7 +17,7 @@ var vogueCmd = &cobra.Command{ Short: "启动 Vogue 时装发布会采集器", Run: func(cmd *cobra.Command, args []string) { // 1. 修复参数数量 & 类型强转 uint(brandID) - vogueSpider := spider.NewVogueSpider(DB, maxCo, uint(brandID), onlyPlatform) + vogueSpider := spider.NewVogueSpider(maxCo, uint(brandID), maxShows, IngestClient) // 2. 修复方法名:将 Execute() 改为 Run() vogueSpider.Run() @@ -29,8 +29,10 @@ func init() { // 绑定命令行参数 vogueCmd.Flags().Int64VarP(&brandID, "brand", "b", 0, "指定抓取的品牌 ID") - vogueCmd.Flags().BoolVarP(&onlyPlatform, "only-platform", "p", false, "是否仅抓取平台关联品牌") - + // 修复:将 "c" 改为 "m"(或者 ""),避免与全局的 -c 参数冲突 vogueCmd.Flags().IntVarP(&maxCo, "max-co", "m", 5, "最大并发数限制") -} \ No newline at end of file + + // 小批量试爬:仅抓取单个品牌的前 N 场发布会(0=不限制)。用于控制入库量 / 先看效果。 + vogueCmd.Flags().IntVarP(&maxShows, "max-shows", "s", 0, "单个品牌最多抓取的发布会场次(0=不限制,小批量试爬用)") +} diff --git a/config.yml b/config.yml index 973c60a..4ea6d21 100644 --- a/config.yml +++ b/config.yml @@ -1,12 +1,14 @@ # 爬虫框架通用配置 app: max_co: 3 # 默认最大并发协程数 - debug: true # 默认为 true(开发调试),生产环境中设为 false + debug: false # 开发调试仅打印控制台;入库设为 false -# MySQL 数据库及连接池配置 -database: - dsn: "root:123456@tcp(127.0.0.1:3306)/database?charset=utf8mb4&parseTime=True&loc=Local" - max_open_conns: 50 # 连接池最大打开连接数 - max_idle_conns: 10 # 连接池最大空闲连接数 - conn_max_lifetime: 30 # 连接可复用的最长时间(分钟) - conn_max_idle_time: 10 # 连接空闲最大存活时间(分钟) \ No newline at end of file +# 入库管线:爬虫 → 后台 ingest 接口(HMAC 签名),替代直写 brand_runway +# enabled: true 时,vogue 走秀改为签名 POST 上送元数据,由后台 worker 异步下载入库。 +# spider 不再直连数据库:品牌任务经 GET /admin/internal/crawl/brands 拉取。 +# secret / hashid_secret 必须与后台 INGEST_SECRET / HASHID_SECRET 完全一致。 +ingest: + enabled: true + endpoint: "http://localhost:8092/admin/internal/ingest" + secret: "dev-ingest-secret-2026" + hashid_secret: "" \ No newline at end of file diff --git a/go.mod b/go.mod index ffe7a72..a64b5b4 100644 --- a/go.mod +++ b/go.mod @@ -7,19 +7,13 @@ require ( github.com/spf13/cobra v1.10.2 github.com/spf13/viper v1.21.0 github.com/tidwall/gjson v1.19.0 - gorm.io/driver/mysql v1.6.0 - gorm.io/gorm v1.31.2 ) require ( - filippo.io/edwards25519 v1.1.0 // indirect github.com/andybalholm/cascadia v1.3.2 // indirect github.com/fsnotify/fsnotify v1.9.0 // indirect - github.com/go-sql-driver/mysql v1.8.1 // indirect github.com/go-viper/mapstructure/v2 v2.4.0 // indirect github.com/inconshreveable/mousetrap v1.1.0 // indirect - github.com/jinzhu/inflection v1.0.0 // indirect - github.com/jinzhu/now v1.1.5 // indirect github.com/pelletier/go-toml/v2 v2.2.4 // indirect github.com/sagikazarmark/locafero v0.11.0 // indirect github.com/sourcegraph/conc v0.3.1-0.20240121214520-5f936abd7ae8 // indirect diff --git a/go.sum b/go.sum index d149529..36d5eec 100644 --- a/go.sum +++ b/go.sum @@ -1,5 +1,3 @@ -filippo.io/edwards25519 v1.1.0 h1:FNf4tywRC1HmFuKW5xopWpigGjJKiJSV0Cqo0cJWDaA= -filippo.io/edwards25519 v1.1.0/go.mod h1:BxyFTGdWcka3PhytdK4V28tE5sGfRvvvRV7EaN4VDT4= github.com/PuerkitoBio/goquery v1.9.2 h1:4/wZksC3KgkQw7SQgkKotmKljk0M6V8TUvA8Wb4yPeE= github.com/PuerkitoBio/goquery v1.9.2/go.mod h1:GHPCaP0ODyyxqcNoFGYlAprUFH81NuRPd0GX3Zu2Mvk= github.com/andybalholm/cascadia v1.3.2 h1:3Xi6Dw5lHF15JtdcmAHD3i1+T8plmv7BQ/nsViSLyss= @@ -11,24 +9,16 @@ github.com/frankban/quicktest v1.14.6 h1:7Xjx+VpznH+oBnejlPUj8oUpdxnVs4f8XU8WnHk github.com/frankban/quicktest v1.14.6/go.mod h1:4ptaffx2x8+WTWXmUCuVU6aPUX1/Mz7zb5vbUoiM6w0= github.com/fsnotify/fsnotify v1.9.0 h1:2Ml+OJNzbYCTzsxtv8vKSFD9PbJjmhYF14k/jKC7S9k= github.com/fsnotify/fsnotify v1.9.0/go.mod h1:8jBTzvmWwFyi3Pb8djgCCO5IBqzKJ/Jwo8TRcHyHii0= -github.com/go-sql-driver/mysql v1.8.1 h1:LedoTUt/eveggdHS9qUFC1EFSa8bU2+1pZjSRpvNJ1Y= -github.com/go-sql-driver/mysql v1.8.1/go.mod h1:wEBSXgmK//2ZFJyE+qWnIsVGmvmEKlqwuVSjsCm7DZg= github.com/go-viper/mapstructure/v2 v2.4.0 h1:EBsztssimR/CONLSZZ04E8qAkxNYq4Qp9LvH92wZUgs= github.com/go-viper/mapstructure/v2 v2.4.0/go.mod h1:oJDH3BJKyqBA2TXFhDsKDGDTlndYOZ6rGS0BRZIxGhM= github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8= github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw= -github.com/jinzhu/inflection v1.0.0 h1:K317FqzuhWc8YvSVlFMCCUb36O/S9MCKRDI7QkRKD/E= -github.com/jinzhu/inflection v1.0.0/go.mod h1:h+uFLlag+Qp1Va5pdKtLDYj+kHp5pxUVkryuEj+Srlc= -github.com/jinzhu/now v1.1.5 h1:/o9tlHleP7gOFmsnYNz3RGnqzefHA47wQpKrrdTIwXQ= -github.com/jinzhu/now v1.1.5/go.mod h1:d3SSVoowX0Lcu0IBviAWJpolVfI5UJVZZ7cO71lE/z8= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= -github.com/mattn/go-sqlite3 v1.14.22 h1:2gZY6PC6kBnID23Tichd1K+Z0oS6nE/XwU+Vz/5o4kU= -github.com/mattn/go-sqlite3 v1.14.22/go.mod h1:Uh1q+B4BYcTPb+yiD3kU8Ct7aC0hY9fxUwlHK0RXw+Y= github.com/pelletier/go-toml/v2 v2.2.4 h1:mye9XuhQ6gvn5h28+VilKrrPoQVanw5PMw/TB0t5Ec4= github.com/pelletier/go-toml/v2 v2.2.4/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8t9Stt+h56mCY= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= @@ -108,9 +98,3 @@ gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15 h1:YR8cESwS4TdDjEe65xsg0ogR gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= -gorm.io/driver/mysql v1.6.0 h1:eNbLmNTpPpTOVZi8MMxCi2aaIm0ZpInbORNXDwyLGvg= -gorm.io/driver/mysql v1.6.0/go.mod h1:D/oCC2GWK3M/dqoLxnOlaNKmXz8WNTfcS9y5ovaSqKo= -gorm.io/driver/sqlite v1.6.0 h1:WHRRrIiulaPiPFmDcod6prc4l2VGVWHz80KspNsxSfQ= -gorm.io/driver/sqlite v1.6.0/go.mod h1:AO9V1qIQddBESngQUKWL9yoH93HIeA1X6V633rBwyT8= -gorm.io/gorm v1.31.2 h1:3o8FXNo9v9S858gil+3LlZA1LkCOzgb4g5BL64FgaCo= -gorm.io/gorm v1.31.2/go.mod h1:XyQVbO2k6YkOis7C2437jSit3SsDK72s7n7rsSHd+Gs= diff --git a/internal/config/config.go b/internal/config/config.go index 459561d..edd7f7d 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -10,8 +10,18 @@ import ( var GlobalConfig Config type Config struct { - App AppConfig `mapstructure:"app"` - Database DatabaseConfig `mapstructure:"database"` + App AppConfig `mapstructure:"app"` + Ingest IngestConfig `mapstructure:"ingest"` +} + +// IngestConfig 入库管线配置:爬虫 → 后台 ingest 接口(HMAC 签名)。 +// spider 不直连数据库:品牌任务经 GET /admin/internal/crawl/brands 拉取, +// 抓取结果经 POST /admin/internal/ingest 上报,由后台 worker 异步落库。 +type IngestConfig struct { + Enabled bool `mapstructure:"enabled"` // 是否启用签名上报(替代直写) + Endpoint string `mapstructure:"endpoint"` // 后台 ingest 地址,如 http://localhost:8092/admin/internal/ingest + Secret string `mapstructure:"secret"` // 必须与后台 INGEST_SECRET 一致 + HashIDSecret string `mapstructure:"hashid_secret"` // 必须与后台 HASHID_SECRET 一致(用于编码 brand_uid) } type AppConfig struct { @@ -19,14 +29,6 @@ type AppConfig struct { Debug bool `mapstructure:"debug"` // 新增 Debug 字段 } -type DatabaseConfig struct { - DSN string `mapstructure:"dsn"` - MaxOpenConns int `mapstructure:"max_open_conns"` - MaxIdleConns int `mapstructure:"max_idle_conns"` - ConnMaxLifetime int `mapstructure:"conn_max_lifetime"` - ConnMaxIdleTime int `mapstructure:"conn_max_idle_time"` -} - // InitConfig 初始化并读取 YAML 配置文件 func InitConfig(cfgFile string) error { if cfgFile != "" { diff --git a/internal/model/models.go b/internal/model/models.go deleted file mode 100644 index 20816ee..0000000 --- a/internal/model/models.go +++ /dev/null @@ -1,79 +0,0 @@ -package model - -// AppBrand 品牌数据表模型 -type AppBrand struct { - ID int64 `gorm:"primaryKey;column:id"` - Name string `gorm:"column:name"` - SpiderOrigin string `gorm:"column:spider_origin"` -} - -func (AppBrand) TableName() string { - return "app_brands" // 根据实际表名调整 -} - -// BrandRunway 时尚发布会主表 -type BrandRunway struct { - ID uint `gorm:"primaryKey;column:id;autoIncrement" json:"id"` - Title string `gorm:"column:title" json:"title"` - Description string `gorm:"column:description" json:"description"` - CreatedAt int64 `gorm:"column:created_at" json:"created_at"` - UpdatedAt int64 `gorm:"column:updated_at" json:"updated_at"` - IsDeleted uint8 `gorm:"column:is_deleted" json:"is_deleted"` - ImageCount uint16 `gorm:"column:image_count" json:"image_count"` - BrandID uint `gorm:"column:brand_id" json:"brand_id"` - Year uint16 `gorm:"column:year" json:"year"` - Cover string `gorm:"column:cover" json:"cover"` - SourceURL string `gorm:"column:source_url" json:"source_url"` - CollectionType string `gorm:"column:collection_type" json:"collection_type"` - Season string `gorm:"column:season" json:"season"` - SeasonCode string `gorm:"column:season_code" json:"season_code"` -} - -// TableName 指定主表名 -func (BrandRunway) TableName() string { - return "brand_runway" -} - -// BrandRunwayImage 发布会图片表 -type BrandRunwayImage struct { - ID uint `gorm:"primaryKey;column:id;autoIncrement" json:"id"` - CreatedAt int64 `gorm:"column:created_at" json:"created_at"` - UpdatedAt int64 `gorm:"column:updated_at" json:"updated_at"` - IsDeleted uint8 `gorm:"column:is_deleted" json:"is_deleted"` - Image string `gorm:"column:image" json:"image"` - RunwayID uint `gorm:"column:runway_id" json:"runway_id"` - BrandID uint `gorm:"column:brand_id" json:"brand_id"` - Name string `gorm:"column:name" json:"name"` - SortOrder uint `gorm:"column:sort_order" json:"sort_order"` -} - -// TableName 指定图片明细表名 -func (BrandRunwayImage) TableName() string { - return "brand_runway_images" -} - -// Article The Impression 街拍文章表(对应 PHP 的 articles 模型) -// 字段与 PHP 的 ArticleModuleEnum::STREET 文章结构保持一致 -type Article struct { - ID uint `gorm:"primaryKey;column:id;autoIncrement" json:"id"` - Title string `gorm:"column:title" json:"title"` - Platform string `gorm:"column:platform" json:"platform"` - Module string `gorm:"column:module" json:"module"` - Year uint16 `gorm:"column:year" json:"year"` - Cover string `gorm:"column:cover" json:"cover"` - Images string `gorm:"column:images" json:"images"` - SourceURL string `gorm:"column:source_url" json:"source_url"` - Brand uint `gorm:"column:brand;default:0" json:"brand"` - CreatedAt int64 `gorm:"column:created_at" json:"created_at"` - UpdatedAt int64 `gorm:"column:updated_at" json:"updated_at"` -} - -// TableName 指定文章主表名 -func (Article) TableName() string { - return "articles" -} - -// ImageItem 图片条目,序列化进 Article.Images(JSON 数组) -type ImageItem struct { - Src string `json:"src"` -} \ No newline at end of file diff --git a/internal/spider/theimpression.go b/internal/spider/theimpression.go index c326799..e50f8e0 100644 --- a/internal/spider/theimpression.go +++ b/internal/spider/theimpression.go @@ -1,6 +1,7 @@ package spider import ( + "context" "encoding/json" "fmt" "io" @@ -13,9 +14,8 @@ import ( "time" "github.com/PuerkitoBio/goquery" - "gorm.io/gorm" - "my-spiders/internal/model" + "my-spiders/internal/ingest" ) const ( @@ -28,20 +28,19 @@ const ( // TheImpressionSpider The Impression 街拍采集器(对应 PHP 的 TheImpressionStreetCommand) type TheImpressionSpider struct { - Client *http.Client - DB *gorm.DB - MaxCo int // 最大并发数 - Debug bool // Debug 模式:仅输出不入库 - ForceUpdate bool // 强制更新已存在的记录 + Client *http.Client + Ingest *ingest.Client // 入库管线客户端(nil = 无法入库,仅 Debug 输出) + MaxCo int // 最大并发数 + Debug bool // Debug 模式:仅输出不入库 } -func NewTheImpressionSpider(db *gorm.DB, maxCo int) *TheImpressionSpider { +func NewTheImpressionSpider(maxCo int, ingestClient *ingest.Client) *TheImpressionSpider { if maxCo <= 0 { maxCo = 5 } return &TheImpressionSpider{ - DB: db, - MaxCo: maxCo, + Ingest: ingestClient, + MaxCo: maxCo, Client: &http.Client{ Timeout: 30 * time.Second, }, @@ -150,7 +149,8 @@ func (s *TheImpressionSpider) Run() { log.Println("[Info] The Impression 所有抓取任务已完成!") } -// getDetail 抓取并解析单个文章详情(对应 PHP 的 getDetail) +// getDetail 抓取并解析单个街拍详情(对应 PHP 的 getDetail),通过入库管线上报至后台审核。 +// 判重与入库(street_snap / street_snap_draft)全部交由后台 worker 处理,本爬虫不直连数据库。 func (s *TheImpressionSpider) getDetail(tk streetTask) { url := tk.URL title := tk.Title @@ -158,15 +158,6 @@ func (s *TheImpressionSpider) getDetail(tk streetTask) { return } - // 判重:按 title + platform 查询是否已存在 - var exist model.Article - if !s.Debug && s.DB != nil { - s.DB.Where("title = ? AND platform = ?", title, TheImpressionPlatform).First(&exist) - if exist.ID > 0 && !s.ForceUpdate { - return - } - } - body, httpCode := s.request(url) if httpCode != 200 || body == "" { log.Printf("[错误] %s 请求失败, code: %d", url, httpCode) @@ -179,75 +170,66 @@ func (s *TheImpressionSpider) getDetail(tk streetTask) { return } - // 采集 figure a img 的图片(去重) - imagesMap := make(map[string]model.ImageItem) - doc.Find("figure a img").Each(func(_ int, node *goquery.Selection) { + // 采集文章正文 figure 内图片(去重,按 DOM 出现顺序保序)。 + // 注意:theimpression.com 当前结构为
, + // 图片并不包裹在 内,故选择器须用 "figure img" 而非旧版的 "figure a img"; + // 同时兜底读取 data-src(部分懒加载主题把真图地址放在 data-src)。 + seen := make(map[string]struct{}) + var imageURLs []string + doc.Find("figure img").Each(func(_ int, node *goquery.Selection) { src, _ := node.Attr("src") - if src != "" { - if _, ok := imagesMap[src]; !ok { - log.Printf("采集图片: %s", src) - imagesMap[src] = model.ImageItem{Src: src} - } + if src == "" { + src, _ = node.Attr("data-src") } + if src == "" { + return + } + // 过滤掉占位/1px 透明图与空 srcset 噪点 + if strings.Contains(src, "data:image") || strings.HasSuffix(src, "1x1.gif") { + return + } + if _, ok := seen[src]; ok { + return + } + seen[src] = struct{}{} + log.Printf("采集图片: %s", src) + imageURLs = append(imageURLs, src) }) - if len(imagesMap) == 0 { + if len(imageURLs) == 0 { log.Printf("[Warning] %s 未采集到图片,跳过保存.", url) return } - cover := "" - var imageList []model.ImageItem - for _, v := range imagesMap { - if cover == "" { - cover = v.Src - } - imageList = append(imageList, v) - } + year := s.parseYear(title) + city := s.parseCity(title) - imagesJSON, _ := json.Marshal(imageList) - now := time.Now().Unix() - - article := model.Article{ - Title: title, - Platform: TheImpressionPlatform, - Module: TheImpressionModule, - Year: uint16(s.parseYear(title)), - Cover: cover, - Images: string(imagesJSON), - SourceURL: url, - Brand: 0, - CreatedAt: now, - UpdatedAt: now, - } - - // Debug 模式或无数据库连接时,仅控制台输出 - if s.Debug || s.DB == nil { - out, _ := json.MarshalIndent(article, "", " ") - log.Printf("\n================ [DEBUG OUTPUT] ================\n%s\n提取图片总数: %d 张\n================================================", - string(out), len(imageList)) + // Debug 模式:仅控制台输出,不上报 + if s.Debug { + log.Printf("[DEBUG] 街拍: %s | year=%d | city=%s | url=%s | images=%d", + title, year, city, url, len(imageURLs)) return } - // 入库(存在则更新,不存在则新增) - err = s.DB.Transaction(func(tx *gorm.DB) error { - if exist.ID > 0 { - article.ID = exist.ID - if err := tx.Save(&article).Error; err != nil { - return err - } - } else { - if err := tx.Create(&article).Error; err != nil { - return err - } - } - return nil - }) + // 入库分支:启用入库管线(s.Ingest != nil)时,把街拍元数据 HMAC 签名 POST 到后台 + // ingest,由后台 worker 异步下载图、传七牛、写 street_snap_draft 待审;通过后才晋升 + // street_snap 正式表(图片统一存七牛云)。判重全交后台 worker(按实体键 city+year 幂等)。 + if s.Ingest == nil { + log.Printf("[Warning] 入库管线未启用 (Ingest==nil),无法上报街拍: %s", title) + return + } - if err != nil { - log.Printf("[错误] 数据入库失败 [%s]: %v", title, err) + p := ingest.RunwayIngest{ + Kind: ingest.KindStreet, + TitleEn: title, + Year: uint16(year), + City: city, + Images: imageURLs, + } + if _, err := s.Ingest.Submit(context.Background(), p); err != nil { + log.Printf("[错误] 街拍入库上报失败 [%s]: %v", title, err) } else { - log.Printf("[成功] 保存文章: %s (共 %d 张图片)", title, len(imageList)) + log.Printf("[成功] 已上报街拍待入库: %s (共 %d 张图片)", title, len(imageURLs)) } } @@ -283,3 +265,21 @@ func (s *TheImpressionSpider) parseYear(title string) int { } return time.Now().Year() } + +// streetCities 常见时装周城市(用于从街拍标题中粗提 city 字段)。 +var streetCities = []string{ + "New York", "London", "Milan", "Paris", "Florence", "Berlin", + "Tokyo", "Seoul", "Shanghai", "Beijing", "Copenhagen", "Stockholm", + "Sydney", "Madrid", "Barcelona", "Amsterdam", "São Paulo", "Sao Paulo", +} + +// parseCity 从街拍标题中粗略识别城市名(命中时装周城市则返回,否则空串)。 +func (s *TheImpressionSpider) parseCity(title string) string { + t := strings.ToLower(title) + for _, city := range streetCities { + if strings.Contains(t, strings.ToLower(city)) { + return city + } + } + return "" +} diff --git a/internal/spider/vogue.go b/internal/spider/vogue.go index 76bf8c2..5376a12 100644 --- a/internal/spider/vogue.go +++ b/internal/spider/vogue.go @@ -1,7 +1,7 @@ package spider import ( - "encoding/json" + "context" "io" "log" "net/http" @@ -12,9 +12,8 @@ import ( "time" "github.com/tidwall/gjson" - "gorm.io/gorm" - "my-spiders/internal/model" + "my-spiders/internal/ingest" ) const ( @@ -23,24 +22,25 @@ const ( ) type VogueSpider struct { - Client *http.Client - DB *gorm.DB - MaxCo int // 最大并发数 - BrandID uint // 客户端指定的品牌ID - OnlyPlatform bool // 是否仅抓取该平台关联的品牌 - Debug bool // 供 s.Debug 使用 - ForceUpdate bool // 供 s.ForceUpdate 使用 + Client *http.Client + Ingest *ingest.Client // 入库管线客户端(nil = 离线/Mock 模式) + MaxCo int // 最大并发数 + BrandID uint // 客户端指定的品牌 ID(--brand,单品牌调试用) + MaxShows int // 小批量模式:单品牌最多抓取的发布会场次(0=不限制) + Debug bool // Debug 模式:仅输出不入库 } -func NewVogueSpider(db *gorm.DB, maxCo int, brandID uint, onlyPlatform bool) *VogueSpider { +// NewVogueSpider 创建 Vogue 爬虫。spider 不再直连数据库:品牌任务经 ingest 接口从后端拉取, +// 抓取结果经 ingest 上报后端,由后端异步下载图、写正式表。spider 是纯 HTTP 客户端。 +func NewVogueSpider(maxCo int, brandID uint, maxShows int, ingestClient *ingest.Client) *VogueSpider { if maxCo <= 0 { maxCo = 5 // 默认限制 5 个并发,防止打满带宽 } return &VogueSpider{ - DB: db, - MaxCo: maxCo, - BrandID: brandID, - OnlyPlatform: onlyPlatform, + Ingest: ingestClient, + MaxCo: maxCo, + BrandID: brandID, + MaxShows: maxShows, Client: &http.Client{ Timeout: 30 * time.Second, }, @@ -65,7 +65,7 @@ func (s *VogueSpider) Run() { brandSem <- struct{}{} wg.Add(1) - go func(t model.AppBrand) { + go func(t ingest.CrawlBrand) { defer func() { <-brandSem wg.Done() @@ -78,41 +78,44 @@ func (s *VogueSpider) Run() { log.Println("[Info] Vogue 所有抓取任务已完成!") } -// getTask 获取需要抓取的品牌任务 -func (s *VogueSpider) getTask() ([]model.AppBrand, error) { - // 1. 本地无数据库时的 Mock 调试场景 - if s.DB == nil { - log.Println("[Debug] 当前未连接数据库,启动 Mock 任务数据模式") - brandID := s.BrandID - if brandID == 0 { - brandID = 108 // 默认给一个 Chanel 示例 ID +// getTask 从后端 ingest 接口拉取要抓取的品牌任务(不再直连数据库读 brand 表)。 +// 无 ingest 客户端时退化为 Mock(给一个示例品牌),便于完全离线调试。 +func (s *VogueSpider) getTask() ([]ingest.CrawlBrand, error) { + if s.Ingest == nil { + log.Println("[Debug] 未配置入库客户端,启动 Mock 任务数据模式") + id := s.BrandID + if id == 0 { + id = 108 // 默认给一个 Chanel 示例 ID } - - return []model.AppBrand{ - {ID: int64(brandID), Name: "Alexander McQueen", SpiderOrigin: VoguePlatform}, + return []ingest.CrawlBrand{ + {BrandUID: ingest.EncodeBrand(uint32(id)), Name: "Alexander McQueen"}, }, nil } - // 2. 有数据库时的查询逻辑 - var brands []model.AppBrand - query := s.DB.Model(&model.AppBrand{}) - - if s.BrandID > 0 { - query = query.Where("id = ?", s.BrandID) - } else { - query = query.Where("id > ?", 1) - if s.OnlyPlatform { - query = query.Where("spider_origin = ?", VoguePlatform) - } - query = query.Order("id asc") + brands, err := s.Ingest.GetCrawlBrands(context.Background(), s.BrandID) + if err != nil { + return nil, err } + return brands, nil +} - err := query.Find(&brands).Error - return brands, err +// sourceURL 由列表项拼出图集详情页地址(供抓取详情页用;不再用于去重/上报)。 +func (s *VogueSpider) sourceURL(info gjson.Result) string { + return VogueBaseURL + info.Get("url").String() + "/slideshow/collection" +} + +// filterUncrawled 返回待抓取的图集列表(原样返回)。 +// +// 历史上这里会调用后端 ExistsSourceURLs 预检、在抓取前跳过已爬图集以省流量; +// 现 source_url 已从 ingestion 管线移除,后端按「实体键(品牌+季节码+系列)」去重, +// 预检既不可行也不必要——重复抓取由 worker 在晋升阶段幂等收敛,不会重复建行。 +// 因此这里直接返回完整列表,去重交给后端。 +func (s *VogueSpider) filterUncrawled(showsList []gjson.Result) []gjson.Result { + return showsList } // SpiderStart 针对单个品牌开启抓取 -func (s *VogueSpider) SpiderStart(task model.AppBrand) { +func (s *VogueSpider) SpiderStart(task ingest.CrawlBrand) { brandName := s.getTaskName(task.Name) url := VogueBaseURL + "/fashion-shows/designer/" + brandName log.Printf("[Command] brandName: %s; spiderUrl: %s", brandName, url) @@ -123,6 +126,20 @@ func (s *VogueSpider) SpiderStart(task model.AppBrand) { return } + // 抓取前不再做预检(source_url 预检已下线);直接抓取,去重由后端 worker 按实体键幂等收敛。 + showsList = s.filterUncrawled(showsList) + if len(showsList) == 0 { + log.Printf("[Info] 品牌 [%s] 的图集均已爬取过,全部跳过", brandName) + return + } + + // 小批量模式:仅抓取前 N 场发布会(用于试爬 / 控制入库量)。 + // 可多次递增 --max-shows 追加抓取;已晋升的秀由后端按实体键(品牌+季节码+系列)幂等收敛,不会重复建行。 + if s.MaxShows > 0 && len(showsList) > s.MaxShows { + log.Printf("[Info] 小批量模式:仅抓取前 %d / %d 场发布会", s.MaxShows, len(showsList)) + showsList = showsList[:s.MaxShows] + } + // 核心修复:限制该品牌下发布会详情页的抓取并发数 detailSem := make(chan struct{}, s.MaxCo) var wg sync.WaitGroup @@ -137,7 +154,7 @@ func (s *VogueSpider) SpiderStart(task model.AppBrand) { wg.Done() }() - s.getDetail(task.ID, info) + s.getDetail(task.BrandUID, info) // 频控:每次下载后稍作停顿,保护带宽并防封 IP time.Sleep(200 * time.Millisecond) @@ -165,27 +182,13 @@ func (s *VogueSpider) getShowsList(url string) []gjson.Result { } // getDetail 获取并解析单个发布会详情 -func (s *VogueSpider) getDetail(brandID int64, info gjson.Result) { +func (s *VogueSpider) getDetail(brandUID string, info gjson.Result) { hed := info.Get("hed").String() - var exist model.BrandRunway - - // 如果处于非 Debug 模式,进行数据库查询判重 - if !s.Debug && s.DB != nil { - s.DB.Where("brand_id = ? AND title = ?", brandID, hed).First(&exist) - if exist.ID > 0 && !s.ForceUpdate { - return - } - } - - // 如果非 forceUpdate 且记录已存在,跳过更新 - if exist.ID > 0 && !s.ForceUpdate { - return - } + // 判重交由后端 worker 处理(按实体键幂等去重),spider 不再直连 brand_runway 判重。 // 获取图片(保持你原本的 URL 拼接与请求逻辑) - pageURI := info.Get("url").String() - requestURL := VogueBaseURL + pageURI + "/slideshow/collection" + requestURL := s.sourceURL(info) log.Printf("正在匹配发布会详情 %s", requestURL) body, httpCode := s.request(requestURL) @@ -208,113 +211,65 @@ func (s *VogueSpider) getDetail(brandID int64, info gjson.Result) { return } - type ImageItem struct { - Src string `json:"src"` - } + // 按 Vogue 原生结构逐 look 提取:每个 gallery item = 1 张主图(image.sources.xxl.url) + // + 0..N 张细节图(details[].image.sources.xxl.url)。主图与细节图在此建立父子关系, + // 不再像旧逻辑那样把所有图扁平混排(那样会丢失「细节图属于哪张主图」)。 + var looks []ingest.RunwayLook + // 兼容旧 worker / 兜底:把主图铺平进 Images(无细节分组)。 + var mainImages []string - var saveURL []ImageItem - var detailURL []ImageItem - - // 保持你原有的图片与细节图提取逻辑 for _, img := range imagesResult.Array() { - xxlURL := img.Get("image.sources.xxl.url").String() - if xxlURL != "" { - saveURL = append(saveURL, ImageItem{Src: xxlURL}) - log.Printf("%s", xxlURL) + mainURL := img.Get("image.sources.xxl.url").String() + if mainURL == "" { + continue } + mainImages = append(mainImages, mainURL) - // 详情图片提取 + look := ingest.RunwayLook{Main: mainURL} for _, detail := range img.Get("details").Array() { detailXxlURL := detail.Get("image.sources.xxl.url").String() if detailXxlURL != "" { - detailURL = append(detailURL, ImageItem{Src: detailXxlURL}) + look.Details = append(look.Details, detailXxlURL) } } + log.Printf("[look] 主图 %s 附 %d 张细节图", mainURL, len(look.Details)) + looks = append(looks, look) } - // ---------------- 下面对接你的双表数据库结构 ---------------- - - // 1. 整理图片与封面 - cover := "" - if len(saveURL) > 0 { - cover = saveURL[0].Src + totalImages := len(mainImages) + for _, l := range looks { + totalImages += len(l.Details) } - // 汇总所有图片(主图 + 细节图)准备存入 brand_runway_images 表 - var allImages []string - for _, item := range saveURL { - allImages = append(allImages, item.Src) - } - for _, item := range detailURL { - allImages = append(allImages, item.Src) - } - - now := time.Now().Unix() - - // 2. 构建主表 BrandRunway 结构体 - runway := model.BrandRunway{ - Title: hed, - BrandID: uint(brandID), - Year: uint16(s.parseYear(hed)), - Cover: cover, - SourceURL: requestURL, - ImageCount: uint16(len(allImages)), - CreatedAt: now, - UpdatedAt: now, - } - - // 3. Debug 模式或无数据库连接时控制台输出 - if s.Debug || s.DB == nil { - out, _ := json.MarshalIndent(runway, "", " ") - log.Printf("\n================ [DEBUG OUTPUT] ================\n%s\n提取图片总数: %d 张\n================================================", string(out), len(allImages)) + // Debug 模式:仅控制台输出,不上报 + if s.Debug { + log.Printf("\n================ [DEBUG OUTPUT] ================\n标题: %s\n抓取 URL: %s\n主图(look)数: %d 张 / 含细节共 %d 张\n================================================", hed, requestURL, len(looks), totalImages) return } - // 4. 执行数据库事务保存(主表 brand_runway + 从表 brand_runway_images) - err := s.DB.Transaction(func(tx *gorm.DB) error { - // 判断是更新还是新增 - if exist.ID > 0 { - runway.ID = exist.ID - if err := tx.Save(&runway).Error; err != nil { - return err - } - // 删除旧关联图片 - if err := tx.Where("runway_id = ?", exist.ID).Delete(&model.BrandRunwayImage{}).Error; err != nil { - return err - } + // 入库分支:把走秀元数据 HMAC 签名 POST 到后台 ingest,由后台 worker 异步下载图、 + // 补 season_code、写正式表(替代直写 brand_runway)。brand_uid 直接来自任务接口。 + if s.Ingest != nil { + collectionType, season := ingest.ParseCollection(hed) + p := ingest.RunwayIngest{ + Kind: ingest.KindRunway, + BrandUID: brandUID, + TitleEn: hed, + Year: uint16(s.parseYear(hed)), + Season: season, + CollectionType: collectionType, + // Looks 优先:结构化承载「主图 + 细节图」分组;Images 仅铺平主图做兼容兜底。 + Looks: looks, + Images: mainImages, + } + if _, err := s.Ingest.Submit(context.Background(), p); err != nil { + log.Printf("[错误] 入库上报失败 [%s]: %v", hed, err) } else { - if err := tx.Create(&runway).Error; err != nil { - return err - } + log.Printf("[成功] 已上报待入库: %s (共 %d 个 look / %d 张图)", hed, len(looks), totalImages) } - - // 批量插入图片从表 - var runwayImages []model.BrandRunwayImage - for idx, imgURL := range allImages { - runwayImages = append(runwayImages, model.BrandRunwayImage{ - RunwayID: runway.ID, // 获取上面生成的自增 ID - BrandID: uint(brandID), - Image: imgURL, - SortOrder: uint(idx + 1), - CreatedAt: now, - UpdatedAt: now, - }) - } - - if len(runwayImages) > 0 { - if err := tx.Create(&runwayImages).Error; err != nil { - return err - } - } - - return nil - }) - - if err != nil { - log.Printf("[错误] 数据入库失败 [%s]: %v", hed, err) - } else { - log.Printf("[成功] 保存发布会: %s (共 %d 张图片)", hed, len(allImages)) + return } + } // getTaskName 品牌名称格式化(对应 getTaskName) diff --git a/spider b/spider index f8a4bfa..7f3e698 100644 Binary files a/spider and b/spider differ