This commit is contained in:
toom1996
2026-09-13 00:52:35 +08:00
parent 10d8a96e8c
commit dbee704c95
44 changed files with 755 additions and 1148 deletions

View File

@ -62,32 +62,28 @@ func (s *IngestService) Submit(ctx context.Context, payload dto.RunwayIngest) (u
if err != nil {
return 0, NewError(http.StatusBadRequest, "payload invalid: "+err.Error())
}
if payload.SourceURL == "" {
return 0, NewError(http.StatusBadRequest, "source_url required")
}
// runway 必须带品牌;street 不关联品牌,brand_uid 可选
if payload.Kind == dto.IngestKindRunway && payload.BrandUID == "" {
return 0, NewError(http.StatusBadRequest, "brand_uid required")
}
job := &model.IngestJob{SourceURL: payload.SourceURL, Payload: string(raw), Kind: model.IngestKindCrawl}
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
}
// ExistsSourceURLs 批量判断 source_url 是否已爬取过(爬虫抓取前的预检)。
//
// 返回命中集合:key 为已存在的 source_url,value 恒为 true。爬虫拿到结果后应跳过这些图集,
// 不必再抓详情页、提取图片 URL 并上报——这一整套动作在判重命中时本来也会被 worker 丢弃。
// 判定条件与 worker 判重完全一致(正式表未删除 / 草稿表 pending),因此不会漏抓也不会重复抓。
func (s *IngestService) ExistsSourceURLs(ctx context.Context, urls []string) (map[string]bool, error) {
return s.repo.ExistsSourceURLs(ctx, urls)
}
// ListJobs 列出最近的入库任务(后台监控页用)。
func (s *IngestService) ListJobs(ctx context.Context, limit int) ([]model.IngestJob, error) {
return s.repo.ListJobs(ctx, limit)
// 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 重新处理。
@ -225,24 +221,43 @@ func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob,
return
}
// 2) 按 source_url 去重(幂等):每来源一条草稿。多来源爬同品牌同季时各建一条草稿,
// 晋升阶段(SaveRunwayFromDraft)按实体键(品牌+季节码+系列)合并为正式表一条并聚合图片。
if _, found, err := s.repo.RunwayIDBySourceURL(ctx, p.SourceURL); err == nil && found {
_ = s.repo.MarkDone(ctx, job.ID)
return
}
if _, found, err := s.repo.DraftIDBySourceURL(ctx, p.SourceURL); err == nil && found {
_ = s.repo.MarkDone(ctx, job.ID)
return
}
// 3) 补 season_code(vogue.go 历史漏填的 bug,统一在此兜底)
// 2) 按实体键(brand_id + season_code + collection_type)去重:正式表已存在该秀则跳过,
// 避免重爬重复下载。注意:只判正式表、不判 pending 草稿——否则多来源(Vogue + theImpression)
// 爬同一场秀时第二个来源会被误判重复而丢弃,破坏晋升阶段的图片聚合。
seasonCode := season.Derive(p.Year, p.CollectionType, p.Season)
if _, found, err := s.repo.RunwayIDByEntity(ctx, brandID, seasonCode, p.CollectionType); err == nil && found {
_ = s.repo.MarkDone(ctx, job.ID)
return
}
// 4) 下载图片并上传到存储(七牛优先,失败兜底本地)
cover, imgs, keys, imgFailed := s.fetchImages(ctx, p.Images, "runway")
if len(p.Images) > 0 && imgFailed {
// 单图失败=整任务失败:先回滚本批已上传的图(七牛 + 本地兜底),避免孤儿文件永远堆在存储里,
// 3) season_code 已在去重前补齐(vogue.go 历史漏填的 bug,统一在此兜底)
// 4) 下载图片并上传到存储(结构化 Looks 优先:主图+细节图分组;否则回退 Images 全部视为主图)。
// 内容哈希(sha1)key 保证重爬不产生孤儿文件:失败回滚删本批 key 即可。
var cover string
var draftImages []model.BrandRunwayDraftImage
var keys []string
var imgFailed bool
var imageCount uint16
if len(p.Looks) > 0 {
cover, draftImages, keys, imgFailed = s.fetchLookImages(ctx, p.Looks, "runway")
imageCount = uint16(len(p.Looks))
} else {
var imgs []string
cover, imgs, keys, imgFailed = s.fetchImages(ctx, p.Images, "runway")
imageCount = uint16(len(imgs))
for i, img := range imgs {
draftImages = append(draftImages, model.BrandRunwayDraftImage{
Image: img,
Name: fmt.Sprintf("Look %d", i+1),
SortOrder: uint32(i + 1),
LookIndex: uint32(i + 1),
IsDetail: 0,
})
}
}
if imgFailed {
// 单图失败=整任务失败:先回滚本批已上传的图(S4 + 本地兜底),避免孤儿文件永远堆在存储里,
// 然后按指数退避自动重试,达上限才置 failed 等后台手动重试。
s.cleanupUploads(ctx, keys)
s.failOrRetry(ctx, job.ID, "image download failed")
@ -262,8 +277,7 @@ func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob,
CollectionType: p.CollectionType,
SeasonCode: seasonCode,
Cover: cover,
SourceURL: p.SourceURL,
ImageCount: uint16(len(imgs)),
ImageCount: imageCount,
Status: model.DraftStatusPending,
}
id, err := s.repo.CreateRunwayDraft(ctx, draft)
@ -271,31 +285,75 @@ func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob,
s.failOrRetry(ctx, job.ID, "create draft: "+err.Error())
return
}
rows := make([]model.BrandRunwayDraftImage, 0, len(imgs))
for i, img := range imgs {
rows = append(rows, model.BrandRunwayDraftImage{
DraftID: id,
Image: img,
Name: fmt.Sprintf("Look %d", i+1),
SortOrder: uint32(i + 1),
})
for i := range draftImages {
draftImages[i].DraftID = id
}
if err := s.repo.CreateRunwayDraftImages(ctx, rows); err != nil {
if err := s.repo.CreateRunwayDraftImages(ctx, draftImages); err != nil {
s.failOrRetry(ctx, job.ID, "create draft images: "+err.Error())
return
}
_ = s.repo.MarkDone(ctx, job.ID)
}
// fetchLookImages 按 Looks 结构下载主图+细节图并上传到存储,返回可直接落草稿的
// BrandRunwayDraftImage 行(带 look_index / is_detail 分组)。任意一张下载/上传失败即把
// failed 置 true,调用方据此把整条任务判失败并回滚本批已上传的 key,符合「单图失败=整任务失败」策略。
// cover 取首个成功下载的主图;image_count(主图数)由调用方按 len(Looks) 计,不在此返回。
func (s *IngestService) fetchLookImages(ctx context.Context, looks []dto.RunwayLook, prefix string) (string, []model.BrandRunwayDraftImage, []string, bool) {
cover := ""
rows := make([]model.BrandRunwayDraftImage, 0)
keys := make([]string, 0)
failed := false
order := 0
for li, look := range looks {
lookIdx := li + 1
if look.Main != "" {
order++
url, key, err := s.downloadOne(ctx, look.Main, order, prefix)
if err != nil {
failed = true
} else {
if cover == "" {
cover = url
}
keys = append(keys, key)
rows = append(rows, model.BrandRunwayDraftImage{
Image: url,
Name: fmt.Sprintf("Look %d", lookIdx),
SortOrder: uint32(order),
LookIndex: uint32(lookIdx),
IsDetail: 0,
})
}
}
for di, d := range look.Details {
if d == "" {
continue
}
order++
url, key, err := s.downloadOne(ctx, d, order, prefix)
if err != nil {
failed = true
continue
}
keys = append(keys, key)
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,
})
}
}
return cover, rows, keys, failed
}
// processStreet 街拍入库:去重 → 下载图 → 写 street_snap_draft(无品牌)。
func (s *IngestService) processStreet(ctx context.Context, job model.IngestJob, p dto.RunwayIngest) {
// 1) 按 source_url 去重(幂等):每来源一条草稿。多来源爬同城市同年时各建一条草稿,
// 晋升阶段(SaveStreetSnapFromDraft)按实体键(城市+年份)合并为正式表一条并聚合图片。
if _, found, err := s.repo.StreetSnapIDBySourceURL(ctx, p.SourceURL); err == nil && found {
_ = s.repo.MarkDone(ctx, job.ID)
return
}
if _, found, err := s.repo.DraftStreetIDBySourceURL(ctx, p.SourceURL); err == nil && found {
// 1) 按实体键(city + year)去重:正式表已存在该街拍则跳过,避免重爬重复下载。
// 只判正式表、不判 pending 草稿,保留多来源街拍图片在晋升阶段聚合。
if _, found, err := s.repo.StreetSnapIDByEntity(ctx, p.City, p.Year); err == nil && found {
_ = s.repo.MarkDone(ctx, job.ID)
return
}
@ -317,7 +375,6 @@ func (s *IngestService) processStreet(ctx context.Context, job model.IngestJob,
Year: p.Year,
City: p.City,
Cover: cover,
SourceURL: p.SourceURL,
ImageCount: uint16(len(imgs)),
Status: model.DraftStatusPending,
}