package service import ( "context" "crypto/sha1" "database/sql" "encoding/hex" "encoding/json" "fmt" "io" "log" "net/http" "path" "strconv" "sync" "time" "fashionapi/internal/dto" "fashionapi/internal/model" "fashionapi/internal/pkg/hashid" "fashionapi/internal/pkg/phash" "fashionapi/internal/pkg/season" "fashionapi/internal/pkg/storage" "fashionapi/internal/repository" ) // IngestService 爬虫入库管线的业务逻辑层:Submit 入队,worker 异步处理。 // // 队列复用一张 ingest_jobs 表,靠 kind 列分流两类任务: // - crawl:爬虫上报的走秀/街拍入库(下载图 → 写草稿) // - media_cleanup:删除图集时异步清理S4孤儿图(引用计数归零才真删) type IngestService struct { repo repository.IngestRepository brandRepo repository.BrandRepository media repository.MediaRepository // 跨表图片引用计数(media_cleanup 删孤儿用) uploader storage.Uploader // 主上传器(S4启用时为S4,否则本地) del storage.Deleter // 删除器(S4或本地,media_cleanup 真删用) local *storage.LocalUploader // S4失败时的兜底落地 httpClient *http.Client } // NewIngestService 创建入库服务。 // // uploader: 主上传器(S4或本地) // del: 删除器(S4或本地),用于 media_cleanup 真删S4/本地孤儿文件 // local: 本地兜底上传器(S4上传失败时回退,避免图片完全丢失) func NewIngestService(repo repository.IngestRepository, brandRepo repository.BrandRepository, media repository.MediaRepository, uploader storage.Uploader, del storage.Deleter, local *storage.LocalUploader) *IngestService { return &IngestService{ repo: repo, brandRepo: brandRepo, media: media, uploader: uploader, del: del, local: local, // 图片下载是并发执行的(见 imageFetchConcurrency),连接池须匹配并发度: // Go 默认 MaxIdleConnsPerHost=2,多余连接会在每轮下载时反复重建、白付 TLS 握手开销。 httpClient: &http.Client{ // 超时须显著大于单张图的正常下载耗时(实测 4~22s,且随源站波动剧烈)。 // 原 30s 余量太薄:一次抖动就会让「单图失败=整任务失败」把整批打回重试。 Timeout: 60 * time.Second, Transport: &http.Transport{ Proxy: http.ProxyFromEnvironment, MaxIdleConns: imageFetchConcurrency * 4, MaxIdleConnsPerHost: imageFetchConcurrency * 2, IdleConnTimeout: 90 * time.Second, TLSHandshakeTimeout: 10 * time.Second, }, }, } } // Submit 把一条上报写入队列(payload 原样存 JSON),立即返回任务 id。 // 真正处理由 worker 异步完成,因此本方法本身很快、不阻塞爬虫。 func (s *IngestService) Submit(ctx context.Context, payload dto.RunwayIngest) (uint32, error) { if payload.Kind == "" { payload.Kind = dto.IngestKindRunway } raw, err := json.Marshal(payload) if err != nil { return 0, NewError(http.StatusBadRequest, "payload invalid: "+err.Error()) } // runway 必须带品牌;street 不关联品牌,brand_uid 可选 if payload.Kind == dto.IngestKindRunway && payload.BrandUID == "" { return 0, NewError(http.StatusBadRequest, "brand_uid required") } job := &model.IngestJob{Payload: string(raw), Kind: model.IngestKindCrawl} if err := s.repo.Enqueue(ctx, job); err != nil { return 0, internalErr(err.Error()) } return job.ID, nil } // ListJobs 分页列出入库任务(后台监控页用)。返回本页任务、总条数与错误。 func (s *IngestService) ListJobs(ctx context.Context, offset, limit int) ([]model.IngestJob, int64, error) { jobs, err := s.repo.ListJobs(ctx, offset, limit) if err != nil { return nil, 0, err } total, err := s.repo.CountJobs(ctx) if err != nil { return nil, 0, err } return jobs, total, nil } // RetryJob 把一条 failed 任务重新入队,等待 worker 重新处理。 func (s *IngestService) RetryJob(ctx context.Context, id uint32) error { return s.repo.RetryJob(ctx, id) } // Run 启动 worker 循环:定时领取 pending 任务并处理,直到 ctx 取消。 // 多实例安全:Claim 用 SKIP LOCKED,互不抢同一条。 func (s *IngestService) Run(ctx context.Context, batch int, interval time.Duration) { if batch < 1 { batch = 2 } if interval <= 0 { interval = 3 * time.Second } ticker := time.NewTicker(interval) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: s.drain(ctx, batch) } } } // drain 一次性领取并处理一批任务。 func (s *IngestService) drain(ctx context.Context, batch int) { jobs, err := s.repo.Claim(ctx, batch) if err != nil { return } for i := range jobs { s.process(ctx, jobs[i]) } } // process 处理单条任务:先按 job.Kind(DB 列)分流到 crawl / media_cleanup 两条管线。 // 旧爬虫任务(迁移前)kind 为空,回落 crawl 走秀/街拍入库管线。 func (s *IngestService) process(ctx context.Context, job model.IngestJob) { start := time.Now() switch job.Kind { case model.IngestKindMediaCleanup: s.processMediaCleanup(ctx, job) return case "", model.IngestKindCrawl: // 走秀/街拍入库管线,按 payload 内 Kind 再分流 runway / street。 default: // 未知 job.Kind 绝不静默处理,标记失败避免脏数据。 _ = s.repo.MarkFailed(ctx, job.ID, "unknown job kind: "+job.Kind) return } var p dto.RunwayIngest if err := json.Unmarshal([]byte(job.Payload), &p); err != nil { _ = s.repo.MarkFailed(ctx, job.ID, "payload parse: "+err.Error()) return } if p.Kind == "" { p.Kind = dto.IngestKindRunway } log.Printf("[ingest] job=%d kind=%s 开始处理(图片=%d 张)", job.ID, p.Kind, len(p.Images)) switch p.Kind { case dto.IngestKindRunway: s.processRunway(ctx, job, p) case dto.IngestKindStreet: s.processStreet(ctx, job, p) default: // 未知 kind 绝不静默当成走秀处理;所有爬虫数据都必须落到已注册的审核模块, // 否则标记任务失败,避免出现「未审核就入库」的脏数据。 _ = s.repo.MarkFailed(ctx, job.ID, "unknown ingest kind: "+p.Kind) return } log.Printf("[ingest] job=%d kind=%s 处理结束 总耗时=%v", job.ID, p.Kind, time.Since(start).Round(time.Millisecond)) } // MediaCleanupPayload 是「清理S4孤儿图」任务的 payload:待清理的图片 key 列表。 type MediaCleanupPayload struct { Keys []string `json:"keys"` } // EnqueueMediaCleanup 把一批待清理的S4 key 异步入队;真正删除由 worker 的 // processMediaCleanup 按引用计数判定,仅当 key 在所有图集/草稿表中引用归零才真删 // (内容寻址共享 key 不会被误删)。 func (s *IngestService) EnqueueMediaCleanup(ctx context.Context, keys []string) error { if len(keys) == 0 { return nil } raw, err := json.Marshal(MediaCleanupPayload{Keys: keys}) if err != nil { return err } return s.repo.EnqueueMediaCleanup(ctx, string(raw)) } // processMediaCleanup 清理S4孤儿图任务:解析 key 列表,按跨表引用计数删除真孤儿。 // 删除/计数失败一律保守跳过(宁可留文件),因此本任务几乎总是成功置 done。 func (s *IngestService) processMediaCleanup(ctx context.Context, job model.IngestJob) { if s.media == nil || s.del == nil { _ = s.repo.MarkFailed(ctx, job.ID, "media cleanup not configured") return } var p MediaCleanupPayload if err := json.Unmarshal([]byte(job.Payload), &p); err != nil { _ = s.repo.MarkFailed(ctx, job.ID, "media cleanup payload parse: "+err.Error()) return } if len(p.Keys) == 0 { _ = s.repo.MarkDone(ctx, job.ID) return } purgeOrphanImages(ctx, s.media, s.del, p.Keys) _ = s.repo.MarkDone(ctx, job.ID) } // failOrRetry 把临时失败(如网络抖动 / 单图下载失败)按指数退避自动重试: // attempts 未达上限则重置为 pending 并推迟 next_attempt_at,worker 到点再领; // 达上限才置 failed,等待后台手动「重试」。永久失败(payload 解析错 / 未知 kind)直接 // MarkFailed,不进重试。 func (s *IngestService) failOrRetry(ctx context.Context, id uint32, errMsg string) { if err := s.repo.ScheduleRetry(ctx, id, errMsg); err != nil { // 调度失败兜底硬失败,避免任务卡在 processing 无人收。 _ = s.repo.MarkFailed(ctx, id, errMsg) } } // processRunway 走秀入库:品牌校验 → 实体键分三支(新建 / 放弃 / 复用驳回)→ 下载图 → 直写 brand_runways。 func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob, p dto.RunwayIngest) { // 1) 品牌必须存在(爬虫负责先建/复用品牌) brandID, err := hashid.Decode(p.BrandUID) if err != nil { _ = s.repo.MarkFailed(ctx, job.ID, "unknown brand_uid") return } if _, err := s.brandRepo.FindByID(ctx, brandID); err != nil { _ = s.repo.MarkFailed(ctx, job.ID, "brand not found") return } // 2) 实体键查重并决定分支(单表模型) seasonCode := season.Derive(p.Year, p.CollectionType, p.Season) reuseID, existStatus, exist, err := s.repo.RunwayEntityState(ctx, brandID, seasonCode, p.CollectionType) if err != nil { // 读失败绝不能当成「未命中」:误判为新建会走一次注定冲突的写入——实体键唯一索引 //(2026-09-22-03)会把它拒成一次可重试的写库错误,重试时实体键已能查到该行(pending) // 而直接跳过。但唯一索引是最后一道防线,不该靠它来替代这里的前置判断, // 故仍 fail-closed + 退避重试,不做无谓的写尝试。 s.failOrRetry(ctx, job.ID, "entity lookup: "+err.Error()) return } action := decideIngestAction(exist, existStatus) if action == ingestSkip { log.Printf("[ingest] job=%d runway 实体已存在(status=%s),跳过", job.ID, existStatus) _ = s.repo.MarkDone(ctx, job.ID) return } // 3) season_code 已在去重前补齐(vogue.go 历史漏填的 bug,统一在此兜底) // 4) 下载图片并上传到存储(结构化 Looks 优先:主图+细节图分组;否则回退 Images 全部视为主图)。 // 内容哈希(sha1)key 保证重爬不产生孤儿文件:失败回滚删本批 key 即可。 var cover string var images []model.BrandRunwayImage var keys []string var imgFailed bool var imageCount uint16 var timing *fetchTiming if len(p.Looks) > 0 { cover, images, keys, imgFailed, timing = s.fetchLookImages(ctx, p.Looks, "runway") // image_count 必须等于**实际落库的主图行数**,不能是 len(p.Looks): // fetchLookImages 会跳过 look.Main == "" 的 look 与 sha1 重复的行,实际主图行可能更少。 // 旧流程靠晋升时按实际行重算自愈,该路径已随单表发布模型删除 —— 一旦用 len(Looks), // 漂移会永久留在公开卡片的「N 张」角标与 FeaturedIDs 的 SUM(image_count) 热度上。 imageCount = countMainImages(images) } else { var fimgs []fetchedImage cover, fimgs, keys, imgFailed, timing = s.fetchImages(ctx, p.Images, "runway", dto.IngestKindRunway) imageCount = uint16(len(fimgs)) for i, fi := range fimgs { images = append(images, model.BrandRunwayImage{ Image: fi.url, Name: fmt.Sprintf("Look %d", i+1), SortOrder: uint32(i + 1), IsDetail: 0, Phash: sqlNull(fi.phash), IsDuplicate: fi.isDup, DupOf: strconv.FormatUint(uint64(fi.dupID), 10), }) } } log.Printf("[ingest] job=%d runway 拉图 总耗时=%v %s", job.ID, timing.total.Round(time.Millisecond), timing) if imgFailed { // 单图失败=整任务失败:先回滚本批已上传的图(S4 + 本地兜底),避免孤儿文件永远堆在存储里, // 然后按指数退避自动重试,达上限才置 failed 等后台手动重试。 s.cleanupUploads(ctx, keys) s.failOrRetry(ctx, job.ID, "image download failed") return } // 5) 直写正式表(status=pending),等待后台审核通过(审核只改状态,不重建图片) writeStart := time.Now() rw := &model.BrandRunway{ JobID: job.ID, BrandID: brandID, TitleEn: p.TitleEn, TitleCn: p.TitleCn, DescriptionEn: p.DescriptionEn, DescriptionCn: p.DescriptionCn, Year: p.Year, Season: p.Season, CollectionType: p.CollectionType, SeasonCode: seasonCode, Cover: cover, ImageCount: imageCount, Status: model.StatusPending, } for i := range images { images[i].BrandID = brandID } var saveErr error if action == ingestReuseRejected { saveErr = s.repo.ReuseRejectedRunway(ctx, reuseID, rw, images) } else { _, saveErr = s.repo.CreateRunwayWithImages(ctx, rw, images) } if saveErr != nil { s.failOrRetry(ctx, job.ID, "save runway: "+saveErr.Error()) return } // 入库后按 sort_order 回填细节图的 parent_image_id,使新走秀自动成组(与街拍一致)。 runwayID := rw.ID if runwayID == 0 { runwayID = reuseID } if rerr := s.repo.RepairRunwayDetailParents(ctx, runwayID); rerr != nil { log.Printf("[ingest] job=%d runway 回填 parent 失败: %v", job.ID, rerr) } log.Printf("[ingest] job=%d runway 写正式表 耗时=%v 图片行=%d 复用驳回=%v", job.ID, time.Since(writeStart).Round(time.Millisecond), len(images), action == ingestReuseRejected) _ = s.repo.MarkDone(ctx, job.ID) } // fetchLookImages 按 Looks 结构下载主图+细节图并上传到存储,返回可直接落正式图片表的 // BrandRunwayImage 行(带 parent_image_id / is_detail 分组)。任意一张下载/上传失败即把 // failed 置 true,调用方据此把整条任务判失败并回滚本批已上传的 key,符合「单图失败=整任务失败」策略。 // cover 取首个成功下载的主图;image_count 由调用方按返回行中 is_detail == 0 的数量计 //(countMainImages),而不是 len(Looks):跳过空主图 / sha1 重复后实际行数可能更少。 // 每张图入库前做去重:命中 dHash 近重复则仍入库但标记留痕。 // // 实现:先把全部「主图 + 细节图」按原始顺序摊平成任务列表,用 concurrentFetch 并发完成 // 「下载 → phash → 上传」这段 IO 密集操作;随后再按下标顺序串行去重与组装。 // 这样既拿到并发收益,又保证 cover / sort_order / 批次内去重与原串行实现一致。 func (s *IngestService) fetchLookImages(ctx context.Context, looks []dto.RunwayLook, prefix string) (string, []model.BrandRunwayImage, []string, bool, *fetchTiming) { start := time.Now() timing := &fetchTiming{} cover := "" rows := make([]model.BrandRunwayImage, 0) keys := make([]string, 0) seen := make(map[string]bool) failed := false // 摊平:顺序 = 逐个 look 先主图、再其细节图,与原串行遍历顺序完全一致。 tasks := make([]lookImageTask, 0, len(looks)) for li, look := range looks { lookIdx := li + 1 if look.Main != "" { tasks = append(tasks, lookImageTask{url: look.Main, lookIdx: lookIdx}) } for di, d := range look.Details { if d == "" { continue } tasks = append(tasks, lookImageTask{url: d, lookIdx: lookIdx, isDetail: true, detailNo: di + 1}) } } // 并发段:下载 + phash + 上传(耗时大头)。 results := concurrentFetch(ctx, len(tasks), func(c context.Context, i int) downloadResult { return s.downloadAndUpload(c, tasks[i].url, prefix) }) // 串行段:按原始顺序累计耗时、做批次内去重、组装行。 order := 0 for i, r := range results { task := tasks[i] timing.download += r.download timing.phash += r.phash timing.upload += r.upload if r.err != nil { failed = true continue } keys = append(keys, r.key) if seen[r.sha1] { continue } seen[r.sha1] = true dupID, isDup := s.dedupImage(ctx, dto.IngestKindRunway, r.phashBits, timing) order++ if cover == "" && !task.isDetail { cover = r.url } name := fmt.Sprintf("Look %d", task.lookIdx) isDetail := uint8(0) if task.isDetail { name = fmt.Sprintf("Look %d — Detail %d", task.lookIdx, task.detailNo) isDetail = 1 } rows = append(rows, model.BrandRunwayImage{ Image: r.url, Name: name, SortOrder: uint32(order), IsDetail: isDetail, Phash: sqlNull(r.phashBits), IsDuplicate: isDup, DupOf: strconv.FormatUint(uint64(dupID), 10), }) } timing.images = len(rows) timing.total = time.Since(start) return cover, rows, keys, failed, timing } // countMainImages 数一批待落库图片行中的主图(is_detail = 0)数量,即 image_count 的口径。 // 与 SoftDeleteRunwayImage 审核期重算的口径一致(只计主图),保证入库值与后续重算值同源。 func countMainImages(images []model.BrandRunwayImage) uint16 { var n uint16 for i := range images { if images[i].IsDetail == 0 { n++ } } return n } // processStreet 街拍入库:实体键分三支(新建 / 放弃 / 复用驳回)→ 下载图 → 直写 street_snaps(无品牌)。 func (s *IngestService) processStreet(ctx context.Context, job model.IngestJob, p dto.RunwayIngest) { // 1) 实体键查重并决定分支(单表模型,同 runway)。 // 实体键含 title:与下方 snap.Title 用的是同一个归一化后的 title,保证查重口径与落库内容一致, // 同城同年的不同专题(Day 2 / Day 3)因此不会被误判为重复。 // 标题先经 NormalizeStreetTitle 归一化(空白折叠),消除「纯空白微调」造成的重复实体(设计规格 §12.4)。 title := repository.NormalizeStreetTitle(p.TitleEn) reuseID, existStatus, exist, err := s.repo.StreetSnapEntityState(ctx, p.City, p.Year, title) if err != nil { // 同 runway:读失败不得当「未命中」,否则会给同一实体留两行。fail-closed + 退避重试。 s.failOrRetry(ctx, job.ID, "entity lookup: "+err.Error()) return } action := decideIngestAction(exist, existStatus) if action == ingestSkip { log.Printf("[ingest] job=%d street 实体已存在(status=%s),跳过", job.ID, existStatus) _ = s.repo.MarkDone(ctx, job.ID) return } // 2) 下载图片并上传到存储(S4优先,失败兜底本地)。 cover, fimgs, keys, imgFailed, timing := s.fetchImages(ctx, p.Images, "street", dto.IngestKindStreet) log.Printf("[ingest] job=%d street 拉图 总耗时=%v %s", job.ID, timing.total.Round(time.Millisecond), timing) if len(p.Images) > 0 && imgFailed { // 单图失败=整任务失败:先回滚本批已上传的图(S4 + 本地兜底),避免孤儿文件永远堆在存储里, // 然后按指数退避自动重试,达上限才置 failed 等后台手动重试。 s.cleanupUploads(ctx, keys) s.failOrRetry(ctx, job.ID, "image download failed") return } // 3) 直写正式表(status=pending) writeStart := time.Now() snap := &model.StreetSnap{ JobID: job.ID, Title: title, // 街拍单标题,爬虫优先填 title_en(已归一化,与查重口径一致) Year: p.Year, City: p.City, Cover: cover, ImageCount: uint16(len(fimgs)), Status: model.StatusPending, } rows := make([]model.StreetSnapImage, 0, len(fimgs)) for i, fi := range fimgs { rows = append(rows, model.StreetSnapImage{ Image: fi.url, Name: fmt.Sprintf("Look %d", i+1), SortOrder: uint32(i + 1), Phash: sqlNull(fi.phash), IsDuplicate: fi.isDup, DupOf: strconv.FormatUint(uint64(fi.dupID), 10), }) } var saveErr error if action == ingestReuseRejected { saveErr = s.repo.ReuseRejectedStreetSnap(ctx, reuseID, snap, rows) } else { _, saveErr = s.repo.CreateStreetSnapWithImages(ctx, snap, rows) } if saveErr != nil { s.failOrRetry(ctx, job.ID, "save street snap: "+saveErr.Error()) return } log.Printf("[ingest] job=%d street 写正式表 耗时=%v 图片行=%d 复用驳回=%v", job.ID, time.Since(writeStart).Round(time.Millisecond), len(rows), action == ingestReuseRejected) _ = s.repo.MarkDone(ctx, job.ID) } // fetchImages 下载图片并上传到存储,返回 (cover 地址, 已下载图结构, 已成功上传对象的 key 列表, 是否有任意一张失败)。 // prefix 为S4 key 前缀(runway/ 或 street/)。kind 用于选择去重比对表。 // 只要任意一张下载/上传失败,failed 即置 true,调用方据此把整条任务判为失败并回滚本批已上传的对象, // 符合「单图失败=整任务失败」策略。每张图入库前做去重(近重复标记)。 // // 实现:先用 concurrentFetch 并发完成全部图片的「下载 → phash → 上传」,再按下标顺序串行 // 做去重与组装,保证 cover / 顺序 / 批次内去重与原串行实现一致。 func (s *IngestService) fetchImages(ctx context.Context, urls []string, prefix, kind string) (string, []fetchedImage, []string, bool, *fetchTiming) { start := time.Now() timing := &fetchTiming{} cover := "" out := make([]fetchedImage, 0, len(urls)) keys := make([]string, 0, len(urls)) seen := make(map[string]bool) failed := false // 并发段:下载 + phash + 上传(耗时大头)。 results := concurrentFetch(ctx, len(urls), func(c context.Context, i int) downloadResult { return s.downloadAndUpload(c, urls[i], prefix) }) // 串行段:按原始顺序累计耗时、做批次内去重、组装结果。 for _, r := range results { timing.download += r.download timing.phash += r.phash timing.upload += r.upload if r.err != nil { failed = true continue } keys = append(keys, r.key) if seen[r.sha1] { continue } seen[r.sha1] = true dupID, isDup := s.dedupImage(ctx, kind, r.phashBits, timing) if cover == "" { cover = r.url } out = append(out, fetchedImage{ url: r.url, key: r.key, phash: r.phashBits, dupID: dupID, isDup: isDup, }) } timing.images = len(out) timing.total = time.Since(start) return cover, out, keys, failed, timing } // cleanupUploads 删除一批本批次成功上传的对象(S4 + 本地兜底),用于任务失败回滚: // 「单图失败=整任务失败」时,前面已成功上传的图若放任不管就成了孤儿,永远堆在存储里。 // 对两种存储都尝试删除(任一不存在即按幂等成功处理);删除失败仅告警,不阻断任务置失败。 // 注意:key 为 sha1 内容寻址,若恰与其他已晋升图集共享同一内容哈希会被一并移除,重试会重新上传补齐。 func (s *IngestService) cleanupUploads(ctx context.Context, keys []string) { for _, k := range keys { if d, ok := s.uploader.(storage.Deleter); ok { if err := d.Delete(ctx, k); err != nil { log.Printf("[warn] ingest cleanup: S4删除失败 key=%s err=%v", k, err) } } if s.local != nil { if err := s.local.Delete(ctx, k); err != nil { log.Printf("[warn] ingest cleanup: 本地删除失败 key=%s err=%v", k, err) } } } } // fetchedImage 是单张已下载图片的暂存结构,携带去重所需的指纹信息。 type fetchedImage struct { url string key string phash string // dHash 的 pgvector 二进制向量串,空串表示无法解码 dupID uint32 // 命中近重复时的参考图 id isDup uint8 // 是否标记为近重复(供审核留痕) } // slowImageThreshold 单张图各阶段合计超过该阈值时单独打一条告警,便于从大量图里定位异常慢图。 const slowImageThreshold = 2 * time.Second // imageFetchConcurrency 单条入库任务内「下载 + 上传」的并发度。 // // 实测结论(theimpression.com,59 张图): // - 并发 1:单张下载 ~3.7s,整批 59 张约 3m54s,任务成功。 // - 并发 5:单张下载涨到 ~23s(约 6 倍),5m24s 仅完成 12/59 张便因 30s 超时失败。 // // 即源站对同 IP 的并发连接有强限流:并发不仅不提速,反而把吞吐从 ~0.27 张/秒 // 打到 ~0.04 张/秒并直接拖垮任务。因此默认取 1(等效串行)。 // 换到不限流的源站时,可把这个值调大(机制已就绪),但务必先做小批量实测。 const imageFetchConcurrency = 1 // fetchTiming 汇总一次图集拉取各阶段的累计耗时,用于定位「慢在哪一步」。 // 下载 / 上传是网络 IO,phash 是 CPU(解码 + 哈希),近重复是 DB(HNSW 检索)—— // 三类瓶颈成因不同,分开计时才能对症下药。 type fetchTiming struct { images int // 成功入行的图片数 total time.Duration // 整个拉取过程的墙钟耗时 download time.Duration // 各张 HTTP 下载耗时之和 phash time.Duration // 各张 dHash 指纹计算耗时之和 upload time.Duration // 各张上传存储耗时之和 dedup time.Duration // 近重复查询累计 } // String 输出各阶段耗时(毫秒精度),一条日志看清瓶颈。 // 注意:下载 / phash / 上传是「各张图耗时之和」,并发执行下其总和会大于墙钟总耗时, // 因此三者与「总耗时」的比值不再等于时间占比——它们用于横向对比哪个阶段最重。 func (t *fetchTiming) String() string { if t == nil { return "" } return fmt.Sprintf("图片=%d 并发=%d 下载=%v phash=%v 上传=%v 近重复=%v(下载/phash/上传为各张累计值)", t.images, imageFetchConcurrency, t.download.Round(time.Millisecond), t.phash.Round(time.Millisecond), t.upload.Round(time.Millisecond), t.dedup.Round(time.Millisecond), ) } // sqlNull 把 phash 字符串转成可空向量字段:空串 → NULL(不参与近邻检索)。 func sqlNull(ph string) sql.NullString { return sql.NullString{String: ph, Valid: ph != ""} } // dedupImage 判断单张图是否近重复(dHash 汉明距离 ≤ 阈值):命中则仍入库,但标记 is_duplicate + dup_of。 func (s *IngestService) dedupImage(ctx context.Context, kind, phashBits string, timing *fetchTiming) (dupID uint32, isDup uint8) { // repo 未注入(如离线单测)时跳过去重,不阻断主流程。生产环境 repo 必不为空。 if s.repo == nil { return 0, 0 } var tables []string switch kind { case dto.IngestKindRunway: tables = []string{"brand_runway_images"} case dto.IngestKindStreet: tables = []string{"street_snap_images"} default: return 0, 0 } if phashBits != "" { start := time.Now() id, found, err := s.repo.FindNearDuplicateImage(ctx, tables, phashBits, phash.DefaultThreshold) if timing != nil { timing.dedup += time.Since(start) } if err == nil && found { return id, 1 } } return 0, 0 } // downloadResult 是单张图「下载 → phash → 上传」的结果,供并发执行后按下标串行组装。 // 各阶段耗时随结果一并带回、由调用方在串行段汇总(而非直接累加到共享 timing), // 这样并发路径上无需加锁即可保证计时准确。 type downloadResult struct { url string // 存储访问地址(失败为空) key string // 对象 key(内容寻址) sha1 string // 内容 sha1,用于批次内去重 phashBits string // dHash 向量串,空串表示无法解码 download time.Duration // 本张下载耗时 phash time.Duration // 本张指纹计算耗时 upload time.Duration // 本张上传耗时 err error // 下载或上传失败原因(非空即该张失败) } // concurrentFetch 以 imageFetchConcurrency 路并发执行 fn(ctx, i)(i ∈ [0,n)),返回长度 n、 // 下标与入参一一对齐的结果切片。 // // 下标对齐是刻意的:调用方随后按原始顺序串行组装,使 cover / sort_order / 批次内去重结果 // 与串行实现完全一致,只把「下载+上传」这段 IO 并行化。 // 每个下标只由一个 goroutine 写入,故无需加锁。 func concurrentFetch(ctx context.Context, n int, fn func(context.Context, int) downloadResult) []downloadResult { results := make([]downloadResult, n) if n == 0 { return results } sem := make(chan struct{}, imageFetchConcurrency) var wg sync.WaitGroup for i := 0; i < n; i++ { wg.Add(1) go func(i int) { defer wg.Done() sem <- struct{}{} // 占坑:超过并发度即在此排队 defer func() { <-sem }() // 释放坑位 results[i] = fn(ctx, i) }(i) } wg.Wait() return results } // lookImageTask 是 fetchLookImages 摊平后的一张待下载任务:记录它属于哪个 look、 // 是主图还是第几张细节图,供并发下载完成后重建 Name / 主副图归属 / IsDetail。 type lookImageTask struct { url string lookIdx int // 1-based look 序号 isDetail bool // true=细节图 detailNo int // 细节图序号(1-based),主图为 0 } // downloadAndUpload 把单张远程图下载后上传到存储(S4优先,失败兜底本地)。 // 内容 sha1 用于生成内容寻址的存储 key(重爬幂等)与批次内去重;访问地址可直接写入数据库 // (S4为完整 https URL,本地为相对 /uploads 路径)。 // // 本函数是并发调用点:每次调用各自返回一份独立的 downloadResult,不共享可变状态; // 耗时汇总与去重留痕由调用方在串行段统一完成。 func (s *IngestService) downloadAndUpload(ctx context.Context, u, prefix string) downloadResult { res := downloadResult{} // ① 下载 dlStart := time.Now() resp, err := s.httpClient.Get(u) if err != nil { res.err = err return res } defer resp.Body.Close() if resp.StatusCode != http.StatusOK { res.err = fmt.Errorf("status %d", resp.StatusCode) return res } data, err := io.ReadAll(resp.Body) if err != nil { res.err = err return res } res.download = time.Since(dlStart) ext := path.Ext(u) if ext == "" || len(ext) > 5 { ext = ".jpg" } // 内容哈希作 key:同一张图(无论来自哪篇文章/job)永远得到相同 key, // S4覆盖写即天然幂等,不会因重复采集、崩溃重试、reject 重爬而累积孤儿文件。 h := sha1.Sum(data) res.sha1 = hex.EncodeToString(h[:]) res.key = fmt.Sprintf("%s/%s%s", prefix, res.sha1, ext) ct := resp.Header.Get("Content-Type") // ② dHash 指纹(仅用于近重复检索;解码失败返回空串 → NULL,不参与检索)。 phStart := time.Now() res.phashBits = phash.ToVectorBits(phash.Of(data)) res.phash = time.Since(phStart) // ③ 上传(主存储 S4,失败兜底本地) upStart := time.Now() if url, uerr := s.uploader.Upload(ctx, res.key, data, ct); uerr == nil { res.url = url } else if s.local != nil { // 兜底本地,避免图片完全丢失 if lurl, lerr := s.local.Upload(ctx, res.key, data, ct); lerr == nil { res.url = lurl } else { res.err = uerr } } else { res.err = uerr } res.upload = time.Since(upStart) // 单图异常慢(合计超阈值)单独告警:便于从几十上百张图里一眼锁定是哪张、卡在哪一段。 if total := res.download + res.phash + res.upload; total > slowImageThreshold { log.Printf("[ingest] 慢图 url=%s 合计=%v 下载=%v phash=%v 上传=%v", u, total.Round(time.Millisecond), res.download.Round(time.Millisecond), res.phash.Round(time.Millisecond), res.upload.Round(time.Millisecond)) } return res } // ingestActionKind 入库分支决策结果(单表发布模型)。 type ingestActionKind int const ( ingestCreate ingestActionKind = iota // 实体键未命中 → 新建记录 ingestSkip // 命中 pending/published → 放弃本次入库 ingestReuseRejected // 命中 rejected → 复用该行重审 ) // decideIngestAction 按实体键命中情况与既有行状态决定入库分支。 // 抽成纯函数,是为了让「驳回不是永久黑名单」这条规则有独立的、不需要数据库的测试。 func decideIngestAction(found bool, status string) ingestActionKind { switch { case !found: return ingestCreate case status == model.StatusRejected: return ingestReuseRejected default: return ingestSkip } }