Compare commits
2 Commits
main
...
a939848e84
| Author | SHA1 | Date | |
|---|---|---|---|
| a939848e84 | |||
| d8f667e391 |
58
cmd/root.go
58
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)
|
||||
}
|
||||
|
||||
@ -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()
|
||||
},
|
||||
|
||||
16
cmd/vogue.go
16
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, "最大并发数限制")
|
||||
}
|
||||
|
||||
// 小批量试爬:仅抓取单个品牌的前 N 场发布会(0=不限制)。用于控制入库量 / 先看效果。
|
||||
vogueCmd.Flags().IntVarP(&maxShows, "max-shows", "s", 0, "单个品牌最多抓取的发布会场次(0=不限制,小批量试爬用)")
|
||||
}
|
||||
|
||||
18
config.yml
18
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 # 连接空闲最大存活时间(分钟)
|
||||
# 入库管线:爬虫 → 后台 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: ""
|
||||
6
go.mod
6
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
|
||||
|
||||
16
go.sum
16
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=
|
||||
|
||||
@ -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 != "" {
|
||||
|
||||
171
internal/ingest/client.go
Normal file
171
internal/ingest/client.go
Normal file
@ -0,0 +1,171 @@
|
||||
// Package ingest 爬虫 → 后台 ingest 管线的客户端。
|
||||
//
|
||||
// 设计(与 backend_v2/internal/middleware/ingest_auth.go + handler/ingest_handler.go 对齐):
|
||||
// - spider 不再直连数据库:品牌任务经 GET /admin/internal/crawl/brands 拉取,
|
||||
// 抓取结果经 POST /admin/internal/ingest 上报,由后台 worker 异步下载图、补 season_code、写正式表。
|
||||
// - 两个接口都走 HMAC-SHA256 验签(X-Signature / X-Timestamp / X-Nonce),
|
||||
// 签名串 = HMAC_SHA256(secret, timestamp + "." + nonce + "." + bodyRaw)。
|
||||
// - 时间戳容忍窗口默认 ±300s;nonce 由后台一次性去重防重放。
|
||||
package ingest
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/hmac"
|
||||
"crypto/rand"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Kind 取值(与 backend_v2/internal/dto/ingest.go 一致)。
|
||||
const (
|
||||
KindRunway = "runway" // 走秀(默认,需 brand_uid)
|
||||
KindStreet = "street" // 街拍(无品牌,需 city/title)
|
||||
)
|
||||
|
||||
// RunwayIngest 上送后端的单场载荷(JSON 字段与 dto.RunwayIngest 完全对齐)。
|
||||
type RunwayIngest struct {
|
||||
Kind string `json:"kind"` // runway | street,缺省 runway
|
||||
BrandUID string `json:"brand_uid"` // 品牌编码 id(hashid),runway 必填
|
||||
TitleEn string `json:"title_en"` // 英文标题(runway 优先;street 作为单标题)
|
||||
TitleCn string `json:"title_cn"` // 中文标题(可空)
|
||||
DescriptionEn string `json:"description_en"` // 英文描述(可空)
|
||||
DescriptionCn string `json:"description_cn"` // 中文描述(可空)
|
||||
Year uint16 `json:"year"` // 年份,如 2026
|
||||
Season string `json:"season"` // spring / fall
|
||||
CollectionType string `json:"collection_type"` // rtw / menswear / couture / resort / pre_fall
|
||||
City string `json:"city"` // 地区/城市(street 用)
|
||||
Images []string `json:"images"` // 铺平的主图 URL 列表(回退用)
|
||||
Looks []RunwayLook `json:"looks,omitempty"` // 结构化「主图 + 细节图」分组(优先)
|
||||
}
|
||||
|
||||
// RunwayLook 一场秀中的一个 look:1 张主图 + 0..N 张细节图。
|
||||
type RunwayLook struct {
|
||||
Main string `json:"main"` // 主图(look)原始 URL
|
||||
Details []string `json:"details"` // 细节图原始 URL 列表
|
||||
}
|
||||
|
||||
// CrawlBrand 爬虫取任务接口返回的单个品牌。
|
||||
type CrawlBrand struct {
|
||||
BrandUID string `json:"brand_uid"`
|
||||
Name string `json:"name"`
|
||||
}
|
||||
|
||||
// Client 入库管线客户端。
|
||||
type Client struct {
|
||||
endpoint string
|
||||
secret string
|
||||
http *http.Client
|
||||
}
|
||||
|
||||
// NewClient 创建入库管线客户端。endpoint 形如 http://localhost:8092/admin/internal/ingest。
|
||||
func NewClient(endpoint, secret string) *Client {
|
||||
return &Client{
|
||||
endpoint: endpoint,
|
||||
secret: secret,
|
||||
http: &http.Client{Timeout: 30 * time.Second},
|
||||
}
|
||||
}
|
||||
|
||||
// Submit 把一场走秀/街拍元数据 HMAC 签名后 POST 到后台 ingest,返回 job_id(202 Accepted)。
|
||||
func (c *Client) Submit(ctx context.Context, p RunwayIngest) (string, error) {
|
||||
body, err := json.Marshal(p)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("marshal payload: %w", err)
|
||||
}
|
||||
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.endpoint, bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("new request: %w", err)
|
||||
}
|
||||
c.signRequest(req, body)
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
|
||||
resp, err := c.http.Do(req)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("do request: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
respBody, _ := io.ReadAll(resp.Body)
|
||||
|
||||
if resp.StatusCode != http.StatusAccepted {
|
||||
return "", fmt.Errorf("ingest submit unexpected status %d: %s", resp.StatusCode, string(respBody))
|
||||
}
|
||||
|
||||
var out struct {
|
||||
JobID string `json:"job_id"`
|
||||
Status string `json:"status"`
|
||||
}
|
||||
_ = json.Unmarshal(respBody, &out)
|
||||
return out.JobID, nil
|
||||
}
|
||||
|
||||
// GetCrawlBrands 从后台拉取可抓取品牌任务(GET /admin/internal/crawl/brands)。
|
||||
// brandID>0 时仅返回该品牌(单品牌调试)。
|
||||
func (c *Client) GetCrawlBrands(ctx context.Context, brandID uint) ([]CrawlBrand, error) {
|
||||
url := strings.Replace(c.endpoint, "/ingest", "/crawl/brands", 1)
|
||||
if brandID > 0 {
|
||||
url += "?brand=" + strconv.FormatUint(uint64(brandID), 10)
|
||||
}
|
||||
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("new request: %w", err)
|
||||
}
|
||||
c.signRequest(req, nil)
|
||||
|
||||
resp, err := c.http.Do(req)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("do request: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
respBody, _ := io.ReadAll(resp.Body)
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return nil, fmt.Errorf("crawl brands unexpected status %d: %s", resp.StatusCode, string(respBody))
|
||||
}
|
||||
|
||||
var out struct {
|
||||
Brands []CrawlBrand `json:"brands"`
|
||||
}
|
||||
if err := json.Unmarshal(respBody, &out); err != nil {
|
||||
return nil, fmt.Errorf("decode crawl brands: %w", err)
|
||||
}
|
||||
return out.Brands, nil
|
||||
}
|
||||
|
||||
// signRequest 对请求体做 HMAC-SHA256 签名并写入验签头。body 为 nil 时按空串处理(GET)。
|
||||
func (c *Client) signRequest(req *http.Request, body []byte) {
|
||||
ts := strconv.FormatInt(time.Now().Unix(), 10)
|
||||
nonce := newNonce()
|
||||
sig := sign(c.secret, ts, nonce, string(body))
|
||||
|
||||
req.Header.Set("X-Signature", sig)
|
||||
req.Header.Set("X-Timestamp", ts)
|
||||
req.Header.Set("X-Nonce", nonce)
|
||||
}
|
||||
|
||||
// sign 复刻 backend_v2/internal/pkg/hmac.Sign:HMAC_SHA256(secret, ts + "." + nonce + "." + body)。
|
||||
func sign(secret, ts, nonce, body string) string {
|
||||
mac := hmac.New(sha256.New, []byte(secret))
|
||||
mac.Write([]byte(ts))
|
||||
mac.Write([]byte("."))
|
||||
mac.Write([]byte(nonce))
|
||||
mac.Write([]byte("."))
|
||||
mac.Write([]byte(body))
|
||||
return hex.EncodeToString(mac.Sum(nil))
|
||||
}
|
||||
|
||||
// newNonce 生成 16 字节随机十六进制串作为一次性 nonce。
|
||||
func newNonce() string {
|
||||
b := make([]byte, 16)
|
||||
_, _ = rand.Read(b)
|
||||
return hex.EncodeToString(b)
|
||||
}
|
||||
87
internal/ingest/hashid.go
Normal file
87
internal/ingest/hashid.go
Normal file
@ -0,0 +1,87 @@
|
||||
package ingest
|
||||
|
||||
import (
|
||||
"crypto/sha256"
|
||||
"encoding/binary"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// 本文件是 backend_v2/internal/pkg/hashid 中「无类型编码(品牌等)」的忠实移植,
|
||||
// 仅用于 spider 离线 Mock 模式下生成 brand_uid(真实上报路径下 brand_uid 直接来自后台接口,无需自编码)。
|
||||
// 算法:SHA256(盐) 派生 4 个子密钥 → 32-bit 平衡 Feistel 网络 → base62 定长(minLen=8) 序列化。
|
||||
// 默认盐与后台一致:hashid_secret 为空时后台走 "fashion-archive-default-salt-change-me"。
|
||||
|
||||
const (
|
||||
hidAlphabet = "0123456789abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ"
|
||||
hidBase = 62
|
||||
hidMinLen = 8
|
||||
)
|
||||
|
||||
var hidKeys [4]uint32
|
||||
|
||||
// InitHashID 用部署级盐值初始化混淆密钥。盐为空时使用与后台相同的内置默认盐。
|
||||
func InitHashID(secret string) {
|
||||
if secret == "" {
|
||||
secret = "fashion-archive-default-salt-change-me"
|
||||
}
|
||||
h := sha256.Sum256([]byte(secret))
|
||||
for i := 0; i < 4; i++ {
|
||||
hidKeys[i] = binary.BigEndian.Uint32(h[i*4 : i*4+4])
|
||||
}
|
||||
}
|
||||
|
||||
// feistel 32-bit 平衡 Feistel 网络(全局密钥版,品牌等无类型实体用)。
|
||||
func feistel(v uint32, encrypt bool) uint32 {
|
||||
return feistelK(v, encrypt, hidKeys)
|
||||
}
|
||||
|
||||
func feistelK(v uint32, encrypt bool, k [4]uint32) uint32 {
|
||||
const rounds = 8
|
||||
l, r := uint16(v>>16), uint16(v&0xffff)
|
||||
for i := 0; i < rounds; i++ {
|
||||
idx := i
|
||||
if !encrypt {
|
||||
idx = rounds - 1 - i
|
||||
}
|
||||
round := func(h uint16) uint16 {
|
||||
f := uint32(h)*0x9E3779B1 + k[idx%4]
|
||||
return uint16((f ^ (f >> 16)) & 0xffff)
|
||||
}
|
||||
if encrypt {
|
||||
nl := r
|
||||
nr := l ^ round(r)
|
||||
l, r = nl, nr
|
||||
} else {
|
||||
nl := r ^ round(l)
|
||||
nr := l
|
||||
l, r = nl, nr
|
||||
}
|
||||
}
|
||||
return uint32(l)<<16 | uint32(r)
|
||||
}
|
||||
|
||||
// numToBase62 把 32-bit 值序列化为定长(minLen)base62 串,高位在前。
|
||||
func numToBase62(x uint32) string {
|
||||
var sb strings.Builder
|
||||
for x > 0 {
|
||||
sb.WriteByte(hidAlphabet[x%hidBase])
|
||||
x /= hidBase
|
||||
}
|
||||
if sb.Len() == 0 {
|
||||
sb.WriteByte(hidAlphabet[0])
|
||||
}
|
||||
runes := []rune(sb.String())
|
||||
for i, j := 0, len(runes)-1; i < j; i, j = i+1, j-1 {
|
||||
runes[i], runes[j] = runes[j], runes[i]
|
||||
}
|
||||
out := string(runes)
|
||||
if len(out) < hidMinLen {
|
||||
out = strings.Repeat("0", hidMinLen-len(out)) + out
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// EncodeBrand 把数字品牌主键编码为对外 hashid 串(与后台 hashid.Encode 同算法、同盐)。
|
||||
func EncodeBrand(id uint32) string {
|
||||
return numToBase62(feistel(id, true))
|
||||
}
|
||||
45
internal/ingest/parse.go
Normal file
45
internal/ingest/parse.go
Normal file
@ -0,0 +1,45 @@
|
||||
package ingest
|
||||
|
||||
import "strings"
|
||||
|
||||
// ParseCollection 从走秀标题解析出 collection_type 与 season。
|
||||
//
|
||||
// 仅输出后端白名单枚举(与 backend_v2/internal/pkg/season.Derive 一致):
|
||||
// collection_type ∈ {rtw, menswear, couture, resort, pre_fall}
|
||||
// season ∈ {spring, fall}(resort/pre_fall 无 spring/fall 词,season 留空,由 Derive 按 collection 前缀补码)
|
||||
//
|
||||
// 常见 Vogue 标题形态:
|
||||
// "Fall 2024 Ready-to-Wear" → (rtw, fall)
|
||||
// "Spring 2024 Menswear" → (menswear, spring)
|
||||
// "Fall 2024 Couture" → (couture, fall)
|
||||
// "Resort 2024" → (resort, "")
|
||||
// "Pre-Fall 2024" → (pre_fall, "")
|
||||
// "Cruise" 视作 resort 的同义表述(Vogue 两词混用);未识别到系列词时兜底为 rtw(成衣)。
|
||||
func ParseCollection(title string) (collectionType, season string) {
|
||||
t := strings.ToLower(title)
|
||||
|
||||
// season:仅 spring / fall 两种
|
||||
switch {
|
||||
case strings.Contains(t, "spring"):
|
||||
season = "spring"
|
||||
case strings.Contains(t, "fall"):
|
||||
season = "fall"
|
||||
}
|
||||
|
||||
// collection_type:按关键字优先级匹配,最后兜底 rtw
|
||||
switch {
|
||||
case strings.Contains(t, "menswear"), strings.Contains(t, "men's"):
|
||||
collectionType = "menswear"
|
||||
case strings.Contains(t, "couture"):
|
||||
collectionType = "couture"
|
||||
case strings.Contains(t, "resort"), strings.Contains(t, "cruise"):
|
||||
collectionType = "resort"
|
||||
case strings.Contains(t, "pre-fall"), strings.Contains(t, "pre fall"), strings.Contains(t, "prefall"):
|
||||
collectionType = "pre_fall"
|
||||
case strings.Contains(t, "ready-to-wear"), strings.Contains(t, "ready to wear"), strings.Contains(t, "rtw"):
|
||||
collectionType = "rtw"
|
||||
default:
|
||||
collectionType = "rtw"
|
||||
}
|
||||
return collectionType, season
|
||||
}
|
||||
@ -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"`
|
||||
}
|
||||
@ -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 class="wp-block-image"><img>,
|
||||
// 图片并不包裹在 <a> 内,故选择器须用 "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 ""
|
||||
}
|
||||
|
||||
@ -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)
|
||||
|
||||
Reference in New Issue
Block a user