@ -12,6 +12,7 @@ import (
"net/http"
"path"
"strconv"
"sync"
"time"
"fashionapi/internal/dto"
@ -27,31 +28,44 @@ import (
//
// 队列复用一张 ingest_jobs 表,靠 kind 列分流两类任务:
// - crawl: 爬虫上报的走秀/街拍入库(下载图 → 写草稿)
// - media_cleanup: 删除图集时异步清理七牛 孤儿图(引用计数归零才真删)
// - media_cleanup: 删除图集时异步清理S4 孤儿图(引用计数归零才真删)
type IngestService struct {
repo repository . IngestRepository
brandRepo repository . BrandRepository
media repository . MediaRepository // 跨表图片引用计数( media_cleanup 删孤儿用)
uploader storage . Uploader // 主上传器(七牛 启用时为七牛 ,否则本地)
del storage . Deleter // 删除器(七牛 或本地, media_cleanup 真删用)
local * storage . LocalUploader // 七牛 失败时的兜底落地
uploader storage . Uploader // 主上传器(S4 启用时为S4 ,否则本地)
del storage . Deleter // 删除器(S4 或本地, media_cleanup 真删用)
local * storage . LocalUploader // S4 失败时的兜底落地
httpClient * http . Client
}
// NewIngestService 创建入库服务。
//
// uploader: 主上传器(七牛 或本地)
// del: 删除器(七牛 或本地),用于 media_cleanup 真删七牛 /本地孤儿文件
// local: 本地兜底上传器(七牛 上传失败时回退,避免图片完全丢失)
// 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 ,
httpClient : & http . Client { Timeout : 30 * time . Second } ,
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 ,
} ,
} ,
}
}
@ -129,6 +143,7 @@ func (s *IngestService) drain(ctx context.Context, batch int) {
// 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 )
@ -149,6 +164,7 @@ func (s *IngestService) process(ctx context.Context, job model.IngestJob) {
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 )
@ -158,15 +174,17 @@ func (s *IngestService) process(ctx context.Context, job model.IngestJob) {
// 未知 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 是「清理七牛 孤儿图」任务的 payload: 待清理的图片 key 列表。
// MediaCleanupPayload 是「清理S4 孤儿图」任务的 payload: 待清理的图片 key 列表。
type MediaCleanupPayload struct {
Keys [ ] string ` json:"keys" `
}
// EnqueueMediaCleanup 把一批待清理的七牛 key 异步入队;真正删除由 worker 的
// EnqueueMediaCleanup 把一批待清理的S4 key 异步入队;真正删除由 worker 的
// processMediaCleanup 按引用计数判定,仅当 key 在所有图集/草稿表中引用归零才真删
// (内容寻址共享 key 不会被误删)。
func ( s * IngestService ) EnqueueMediaCleanup ( ctx context . Context , keys [ ] string ) error {
@ -180,7 +198,7 @@ func (s *IngestService) EnqueueMediaCleanup(ctx context.Context, keys []string)
return s . repo . EnqueueMediaCleanup ( ctx , string ( raw ) )
}
// processMediaCleanup 清理七牛 孤儿图任务:解析 key 列表,按跨表引用计数删除真孤儿。
// processMediaCleanup 清理S4 孤儿图任务:解析 key 列表,按跨表引用计数删除真孤儿。
// 删除/计数失败一律保守跳过(宁可留文件),因此本任务几乎总是成功置 done。
func ( s * IngestService ) processMediaCleanup ( ctx context . Context , job model . IngestJob ) {
if s . media == nil || s . del == nil {
@ -211,7 +229,7 @@ func (s *IngestService) failOrRetry(ctx context.Context, id uint32, errMsg strin
}
}
// processRunway 走秀入库:品牌校验 → 去重 → 补季节码 → 下载图 → 写 brand_runway_draft。
// processRunway 走秀入库:品牌校验 → 去重 → 补季节码 → 下载图 → 写 brand_runway_drafts 。
func ( s * IngestService ) processRunway ( ctx context . Context , job model . IngestJob , p dto . RunwayIngest ) {
// 1) 品牌必须存在(爬虫负责先建/复用品牌)
brandID , err := hashid . Decode ( p . BrandUID )
@ -228,26 +246,31 @@ func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob,
// 避免重爬重复下载。注意:只判正式表、不判 pending 草稿——否则多来源( Vogue + theImpression)
// 爬同一场秀时第二个来源会被误判重复而丢弃,破坏晋升阶段的图片聚合。
seasonCode := season . Derive ( p . Year , p . CollectionType , p . Season )
dedupStart := time . Now ( )
if _ , found , err := s . repo . RunwayIDByEntity ( ctx , brandID , seasonCode , p . CollectionType ) ; err == nil && found {
log . Printf ( "[ingest] job=%d runway 实体去重命中(season=%s type=%s),跳过 耗时=%v" ,
job . ID , seasonCode , p . CollectionType , time . Since ( dedupStart ) . Round ( time . Millisecond ) )
_ = s . repo . MarkDone ( ctx , job . ID )
return
}
log . Printf ( "[ingest] job=%d runway 实体去重 耗时=%v( 未命中, 继续下载) " , job . ID , time . Since ( dedupStart ) . Round ( time . Millisecond ) )
// 3) season_code 已在去重前补齐( vogue.go 历史漏填的 bug, 统一在此兜底)
// 4) 下载图片并上传到存储(结构化 Looks 优先:主图+细节图分组;否则回退 Images 全部视为主图)。
// 内容哈希( sha1) key 保证重爬不产生孤儿文件:失败回滚删本批 key 即可。
var cover string
var cover string
var draftImages [ ] model . BrandRunwayDraftImage
var keys [ ] string
var imgFailed bool
var imageCount uint16
var timing * fetchTiming
if len ( p . Looks ) > 0 {
cover , draftImages , keys , imgFailed = s . fetchLookImages ( ctx , p . Looks , "runway" )
cover , draftImages , keys , imgFailed , timing = s . fetchLookImages ( ctx , p . Looks , "runway" )
imageCount = uint16 ( len ( p . Looks ) )
} else {
var fimgs [ ] fetchedImage
cover , fimgs , keys , imgFailed = s . fetchImages ( ctx , p . Images , "runway" , dto . IngestKindRunway )
cover , fimgs , keys , imgFailed , timing = s . fetchImages ( ctx , p . Images , "runway" , dto . IngestKindRunway )
imageCount = uint16 ( len ( fimgs ) )
for i , fi := range fimgs {
draftImages = append ( draftImages , model . BrandRunwayDraftImage {
@ -256,13 +279,13 @@ func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob,
SortOrder : uint32 ( i + 1 ) ,
LookIndex : uint32 ( i + 1 ) ,
IsDetail : 0 ,
ContentSha1 : fi . sha1 ,
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 等后台手动重试。
@ -271,8 +294,8 @@ func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob,
return
}
// 5) 写草稿表( status=pending) , 等待后台审核通过后再晋升正式表
writeStart := time . Now ( )
draft := & model . BrandRunwayDraft {
JobID : job . ID ,
BrandID : brandID ,
@ -300,6 +323,7 @@ func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob,
s . failOrRetry ( ctx , job . ID , "create draft images: " + err . Error ( ) )
return
}
log . Printf ( "[ingest] job=%d runway 写草稿 耗时=%v 图片行=%d" , job . ID , time . Since ( writeStart ) . Round ( time . Millisecond ) , len ( draftImages ) )
_ = s . repo . MarkDone ( ctx , job . ID )
}
@ -307,102 +331,109 @@ func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob,
// BrandRunwayDraftImage 行(带 look_index / is_detail 分组)。任意一张下载/上传失败即把
// failed 置 true, 调用方据此把整条任务判失败并回滚本批已上传的 key, 符合「单图失败=整任务失败」策略。
// cover 取首个成功下载的主图; image_count( 主图数) 由调用方按 len(Looks) 计,不在此返回。
// 每张图入库前做去重:content_sha1 已存在则整行跳过(精确重复); 命中 dHash 近重复则仍入库但标记留痕。
func ( s * IngestService ) fetchLookImages ( ctx context . Context , looks [ ] dto . RunwayLook , prefix string ) ( string , [ ] model . BrandRunwayDraftImage , [ ] string , bool ) {
// 每张图入库前做去重:命中 dHash 近重复则仍入库但标记留痕。
//
// 实现:先把全部「主图 + 细节图」按原始顺序摊平成任务列表,用 concurrentFetch 并发完成
// 「下载 → phash → 上传」这段 IO 密集操作;随后再按下标顺序串行去重与组装。
// 这样既拿到并发收益,又保证 cover / sort_order / 批次内去重与原串行实现一致。
func ( s * IngestService ) fetchLookImages ( ctx context . Context , looks [ ] dto . RunwayLook , prefix string ) ( string , [ ] model . BrandRunwayDraftImage , [ ] string , bool , * fetchTiming ) {
start := time . Now ( )
timing := & fetchTiming { }
cover := ""
rows := make ( [ ] model . BrandRunwayDraftImage , 0 )
keys := make ( [ ] string , 0 )
seen := make ( map [ string ] bool )
failed := false
order := 0
// 摊平:顺序 = 逐个 look 先主图、再其细节图,与原串行遍历顺序完全一致。
tasks := make ( [ ] lookImageTask , 0 , len ( looks ) )
for li , look := range looks {
lookIdx := li + 1
if look . Main != "" {
url , key , sha1h , ph , err := s . downloadOne ( ctx , look . Main , order , prefix )
if err != nil {
failed = true
} else {
keys = append ( keys , key )
if seen [ sha1h ] {
continue
}
seen [ sha1h ] = true
skip , dupID , isDup := s . dedupImage ( ctx , dto . IngestKindRunway , sha1h , ph )
if skip {
continue
}
order ++
if cover == "" {
cover = url
}
rows = append ( rows , model . BrandRunwayDraftImage {
Image : url ,
Name : fmt . Sprintf ( "Look %d" , lookIdx ) ,
SortOrder : uint32 ( order ) ,
LookIndex : uint32 ( lookIdx ) ,
IsDetail : 0 ,
ContentSha1 : sha1h ,
Phash : sqlNull ( ph ) ,
IsDuplicate : isDup ,
DupOf : strconv . FormatUint ( uint64 ( dupID ) , 10 ) ,
} )
}
tasks = append ( tasks , lookImageTask { url : look . Main , lookIdx : lookIdx } )
}
for di , d := range look . Details {
if d == "" {
continue
}
url , key , sha1h , ph , err := s . downloadOne ( ctx , d , order , prefix )
if err != nil {
failed = true
continue
}
keys = append ( keys , key )
if seen [ sha1h ] {
continue
}
seen [ sha1h ] = true
skip , dupID , isDup := s . dedupImage ( ctx , dto . IngestKindRunway , sha1h , ph )
if skip {
continue
}
order ++
rows = append ( rows , model . BrandRunwayDraftImage {
Image : url ,
Name : fmt . Sprintf ( "Look %d — Detail %d" , lookIdx , di + 1 ) ,
SortOrder : uint32 ( order ) ,
LookIndex : uint32 ( lookIdx ) ,
IsDetail : 1 ,
ContentSha1 : sha1h ,
Phash : sqlNull ( ph ) ,
IsDuplicate : isDup ,
DupOf : strconv . FormatUint ( uint64 ( dupID ) , 10 ) ,
} )
tasks = append ( tasks , lookImageTask { url : d , lookIdx : lookIdx , isDetail : true , detailNo : di + 1 } )
}
}
return cover , rows , keys , failed
// 并发段:下载 + 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 . BrandRunwayDraftImage {
Image : r . url ,
Name : name ,
SortOrder : uint32 ( order ) ,
LookIndex : uint32 ( task . lookIdx ) ,
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
}
// processStreet 街拍入库:去重 → 下载图 → 写 street_snap_draft( 无品牌) 。
func ( s * IngestService ) processStreet ( ctx context . Context , job model . IngestJob , p dto . RunwayIngest ) {
// 1) 按实体键( city + year) 去重: 正式表已存在该街拍则跳过, 避免重爬重复下载。
// 只判正式表、不判 pending 草稿,保留多来源街拍图片在晋升阶段聚合。
dedupStart := time . Now ( )
if _ , found , err := s . repo . StreetSnapIDByEntity ( ctx , p . City , p . Year ) ; err == nil && found {
log . Printf ( "[ingest] job=%d street 实体去重命中(city=%s year=%d),跳过 耗时=%v" ,
job . ID , p . City , p . Year , time . Since ( dedupStart ) . Round ( time . Millisecond ) )
_ = s . repo . MarkDone ( ctx , job . ID )
return
}
log . Printf ( "[ingest] job=%d street 实体去重 耗时=%v( 未命中, 继续下载) " , job . ID , time . Since ( dedupStart ) . Round ( time . Millisecond ) )
// 2) 下载图片并上传到存储(七牛 优先,失败兜底本地)。
cover , fimgs , keys , imgFailed := s . fetchImages ( ctx , p . Images , "street" , dto . IngestKindStreet )
// 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) , 等待后台审核通过后再晋升 street_snap 正式表
// 3) 写草稿表( status=pending) , 等待后台审核通过后再晋升 street_snaps 正式表
writeStart := time . Now ( )
draft := & model . StreetSnapDraft {
JobID : job . ID ,
Title : p . TitleEn , // 街拍单标题,爬虫优先填 title_en
@ -424,66 +455,73 @@ func (s *IngestService) processStreet(ctx context.Context, job model.IngestJob,
Image : fi . url ,
Name : fmt . Sprintf ( "Look %d" , i + 1 ) ,
SortOrder : uint32 ( i + 1 ) ,
ContentSha1 : fi . sha1 ,
Phash : sqlNull ( fi . phash ) ,
IsDuplicate : fi . isDup ,
DupOf : strconv . FormatUint ( uint64 ( fi . dupID ) , 10 ) ,
} )
}
if err := s . repo . CreateStreetSnapDraftImages ( ctx , rows ) ; err != nil {
s . failOrRetry ( ctx , job . ID , "create street draft images: " + err . Error ( ) )
return
}
log . Printf ( "[ingest] job=%d street 写草稿 耗时=%v 图片行=%d" , job . ID , time . Since ( writeStart ) . Round ( time . Millisecond ) , len ( rows ) )
_ = s . repo . MarkDone ( ctx , job . ID )
}
// fetchImages 下载图片并上传到存储,返回 (cover 地址, 全部图片地址, 已成功上传对象的 key 列表, 是否有任意一张失败)。
// prefix 为七牛 key 前缀( runway/ 或 street/)。只要任意一张下载/上传失败, failed 即置 true,
// 调用方据此把整条任务判为失败(不再写草稿),并拿 keys 回滚本批已上传的对象,符合「单图失败=整任务失败」策略。
// fetchImages 下载图片并上传到存储,返回 (cover 地址, 已下载图结构, 已成功上传对象的 key 列表, 是否有任意一张失败)。
// prefix 为七牛 key 前缀( runway/ 或 street/) 。kind 用于选择去重比对表。
// prefix 为S4 key 前缀( runway/ 或 street/) 。kind 用于选择去重比对表。
// 只要任意一张下载/上传失败, failed 即置 true, 调用方据此把整条任务判为失败并回滚本批已上传的对象,
// 符合「单图失败=整任务失败」策略。每张图入库前做去重(精确跳过 + 近重复标记)。
func ( s * IngestService ) fetchImages ( ctx context . Context , urls [ ] string , prefix , kind string ) ( string , [ ] fetchedImage , [ ] string , bool ) {
// 符合「单图失败=整任务失败」策略。每张图入库前做去重(近重复标记)。
//
// 实现:先用 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
for _ , u := range urls {
url , key , sha1h , ph , err := s . downloadOne ( ctx , u , 0 , prefix )
if err != nil {
// 并发段:下载 + 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 , key )
if seen [ sha1h ] {
continue
}
seen [ sha1h ] = true
skip , dupID , isDup := s . dedupImage ( ctx , kind , sha1h , ph )
if skip {
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 = url
cover = r . url
}
out = append ( out , fetchedImage {
url : url ,
key : key ,
sha1 : sha1h ,
phash : ph ,
url : r . url,
key : r . key,
phash : r . phashBits ,
dupID : dupID ,
isDup : isDup ,
} )
}
return cover , out , keys , failed
timing . images = len ( out )
timing . total = time . Since ( start )
return cover , out , keys , failed , timing
}
// cleanupUploads 删除一批本批次成功上传的对象(七牛 + 本地兜底),用于任务失败回滚:
// cleanupUploads 删除一批本批次成功上传的对象(S4 + 本地兜底),用于任务失败回滚:
// 「单图失败=整任务失败」时,前面已成功上传的图若放任不管就成了孤儿,永远堆在存储里。
// 对两种存储都尝试删除(任一不存在即按幂等成功处理);删除失败仅告警,不阻断任务置失败。
// 注意: key 为 sha1 内容寻址,若恰与其他已晋升图集共享同一内容哈希会被一并移除,重试会重新上传补齐。
@ -491,7 +529,7 @@ 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: 七牛 删除失败 key=%s err=%v" , k , err )
log . Printf ( "[warn] ingest cleanup: S4 删除失败 key=%s err=%v" , k , err )
}
}
if s . local != nil {
@ -506,25 +544,64 @@ func (s *IngestService) cleanupUploads(ctx context.Context, keys []string) {
type fetchedImage struct {
url string
key string
sha1 string // 内容 sha1( 与存储 key 同源)
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 判断单张图是否重复:
// - 精确重复( content_sha1 已在库)→ 返回 skip=true( 不入库该行) ;
// - 近重复( dHash 汉明距离 ≤ 阈值)→ 仍入库,但标记 is_duplicate + dup_of;
// - 否则正常入库。
func ( s * IngestService ) dedupImage ( ctx context . Context , kind , sha1hash , phashBits string ) ( skip bool , dupID uint32 , isDup uint8 ) {
// 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 false , 0 , 0
return 0 , 0
}
var tables [ ] string
switch kind {
@ -533,61 +610,135 @@ func (s *IngestService) dedupImage(ctx context.Context, kind, sha1hash, phashBit
case dto . IngestKindStreet :
tables = [ ] string { "street_snap_draft_images" , "street_snap_images" }
default :
return false , 0 , 0
}
if exists , err := s . repo . ImageExistsBySha1 ( ctx , tables , sha1hash ) ; err == nil && exists {
return true , 0 , 0
return 0 , 0
}
if phashBits != "" {
if id , found , err := s . repo . FindNearDuplicateImage ( ctx , tables , phashBits , phash . DefaultThreshold ) ; err == nil && found {
return false , id , 1
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 false , 0 , 0
return 0 , 0
}
// downloadOne 把单张远程图下载后上传到存储(七牛优先,失败兜底本地),
// 返回 (访问地址, 对象 key, 内容 sha1, dHash 向量串, error)。访问地址可直接写入数据库
// (七牛为完整 https URL, 本地为相对 /uploads 路径) 。
func ( s * IngestService ) downloadOne ( ctx context . Context , u string , idx int , prefix string ) ( string , string , string , string , error ) {
// 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 / LookIndex / 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 {
return "" , "" , "" , "" , err
res . err = err
return res
}
defer resp . Body . Close ( )
if resp . StatusCode != http . StatusOK {
return "" , "" , "" , "" , fmt . Errorf ( "status %d" , resp . StatusCode )
res . err = fmt . Errorf ( "status %d" , resp . StatusCode )
return res
}
data , err := io . ReadAll ( resp . Body )
if err != nil {
return "" , "" , "" , "" , err
res . err = err
return res
}
res . download = time . Since ( dlStart )
ext := path . Ext ( u )
if ext == "" || len ( ext ) > 5 {
ext = ".jpg"
}
// 内容哈希作 key: 同一张图( 无论来自哪篇文章/job) 永远得到相同 key,
// 七牛 覆盖写即天然幂等, 不会因重复采集、崩溃重试、reject 重爬而累积孤儿文件。
// S4 覆盖写即天然幂等, 不会因重复采集、崩溃重试、reject 重爬而累积孤儿文件。
h := sha1 . Sum ( data )
hash : = hex . EncodeToString ( h [ : ] )
key : = fmt . Sprintf ( "%s/%s%s" , prefix , hash , ext )
res . sha1 = hex . EncodeToString ( h [ : ] )
res . key = fmt . Sprintf ( "%s/%s%s" , prefix , res . sha1 , ext )
ct := resp . Header . Get ( "Content-Type" )
// 计算 dHash 指纹(仅用于近重复检索;解码失败返回空串 → NULL, 不参与检索) 。
ph := phash . Of ( data )
phashBits : = phash . ToVectorBits ( ph )
// ② dHash 指纹(仅用于近重复检索;解码失败返回空串 → NULL, 不参与检索) 。
phStart := time . Now ( )
res . phashBits = phash . ToVectorBits ( phash . Of ( data ) )
res . phash = time . Since ( phStart )
// 主上传器(七牛 )
if url , err := s . uploader . Upload ( ctx , key , data , ct ) ; err == nil {
return url , key , hash , phashBits , nil
// ③ 上传(主存储 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 , key , data , ct ) ; lerr == nil {
return lurl , key , hash , phashBits , ni l
if lurl , lerr := s . local . Upload ( ctx , res . key, data , ct ) ; lerr == nil {
res . url = lur l
} else {
return "" , "" , "" , "" , err
res . err = u err
}
} else {
res . err = uerr
}
return "" , "" , "" , "" , err
}
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
}