Files
backend_v2/internal/service/ingest_service.go
toom1996 d15d2a4701 update
2026-09-28 10:52:50 +08:00

862 lines
34 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package service
import (
"context"
"crypto/sha1"
"database/sql"
"encoding/hex"
"encoding/json"
"fmt"
"io"
"log"
"net/http"
"net/url"
"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, downloadProxy string) *IngestService {
// 下载代理:download_proxy(或 INGEST_DOWNLOAD_PROXY)非空时,下载客户端经该代理访问源站;
// 仅作用于图片下载,S4 上传走独立 s3 客户端、pgx 数据库连接不走 http,均不受影响。
// 留空回退 http.ProxyFromEnvironment(默认直连,仍可被系统 HTTP(S)_PROXY 影响)。
dlTransport := &http.Transport{
MaxIdleConns: imageFetchConcurrency * 4,
MaxIdleConnsPerHost: imageFetchConcurrency * 2,
IdleConnTimeout: 90 * time.Second,
TLSHandshakeTimeout: 10 * time.Second,
}
if downloadProxy != "" {
if pu, perr := url.Parse(downloadProxy); perr == nil {
dlTransport.Proxy = http.ProxyURL(pu)
log.Printf("[ingest] 图片下载走代理: %s", downloadProxy)
} else {
log.Printf("[ingest] download_proxy 非法,已忽略并回退 ProxyFromEnvironment: %q", downloadProxy)
dlTransport.Proxy = http.ProxyFromEnvironment
}
} else {
dlTransport.Proxy = http.ProxyFromEnvironment
}
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: dlTransport,
},
}
}
// 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) 实体键查重并决定分支(单表模型)
// 标题缺年份(p.Year=0)时,用发布日期的年份补全年份,使 year / season_code 不再为空,
// 避免这些秀在「按年份/季节筛选」与热门排序(按 year + season_code)中丢失。
effYear := p.Year
if effYear == 0 && p.PublishedAt != 0 {
effYear = uint16(time.Unix(int64(p.PublishedAt), 0).UTC().Year())
}
seasonCode := season.Derive(effYear, p.CollectionType, p.Season)
reuseID, existStatus, exist, isDeleted, 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
}
// 删除即永久:已软删的实体键不再因重爬复活(原设计靠部分唯一索引释放键、允许重爬,已改为禁止)。
if exist && isDeleted {
log.Printf("[ingest] job=%d runway 实体已删除(is_deleted=1),跳过重爬", job.ID)
_ = s.repo.MarkDone(ctx, job.ID)
return
}
if !exist && p.Year == 0 {
// 标题无年份的历史行(year=0、season_code 为空)待回填:按旧键再查一次,
// 命中后走跳过分支按主键 id 回填(回填后season_code变为effYear码,后续重跑即可直接命中,幂等)。
var rid uint32
rid, existStatus, exist, isDeleted, err = s.repo.RunwayEntityState(ctx, brandID, "", p.CollectionType)
if err != nil {
s.failOrRetry(ctx, job.ID, "entity lookup: "+err.Error())
return
}
reuseID = rid
if exist && isDeleted {
log.Printf("[ingest] job=%d runway 实体已删除(is_deleted=1),跳过重爬", job.ID)
_ = s.repo.MarkDone(ctx, job.ID)
return
}
}
action := decideIngestAction(exist, existStatus)
if action == ingestSkip {
log.Printf("[ingest] job=%d runway 实体已存在(status=%s),跳过(回填 published_at)", job.ID, existStatus)
// 临时回填:仅更新发布时间与缺失年份,不重下图片(图片已在库)。
if p.PublishedAt != 0 {
if err := s.repo.BackfillRunwayPublishedAt(ctx, reuseID, p.PublishedAt, effYear, seasonCode); err != nil {
s.failOrRetry(ctx, job.ID, "backfill published_at: "+err.Error())
return
}
}
_ = 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: effYear,
Season: p.Season,
CollectionType: p.CollectionType,
SeasonCode: seasonCode,
PublishedAt: p.PublishedAt,
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, isDeleted, 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
}
// 删除即永久:已软删的实体键不再因重爬复活(与 runway 同源)。
if exist && isDeleted {
log.Printf("[ingest] job=%d street 实体已删除(is_deleted=1),跳过重爬", job.ID)
_ = s.repo.MarkDone(ctx, job.ID)
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
}
}