Files
backend_v2/internal/service/ingest_service.go
toom1996 10d8a96e8c update
2026-09-07 00:04:01 +08:00

426 lines
16 KiB
Go
Raw Permalink 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"
"encoding/hex"
"encoding/json"
"fmt"
"io"
"log"
"net/http"
"path"
"time"
"fashionapi/internal/dto"
"fashionapi/internal/model"
"fashionapi/internal/pkg/hashid"
"fashionapi/internal/pkg/season"
"fashionapi/internal/pkg/storage"
"fashionapi/internal/repository"
)
// IngestService 爬虫入库管线的业务逻辑层:Submit 入队,worker 异步处理。
//
// 队列复用一张 ingest_jobs 表,靠 kind 列分流两类任务:
// - crawl:爬虫上报的走秀/街拍入库(下载图 → 写草稿)
// - media_cleanup:删除图集时异步清理七牛孤儿图(引用计数归零才真删)
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 // 七牛失败时的兜底落地
httpClient *http.Client
}
// NewIngestService 创建入库服务。
//
// uploader: 主上传器(七牛或本地)
// del: 删除器(七牛或本地),用于 media_cleanup 真删七牛/本地孤儿文件
// local: 本地兜底上传器(七牛上传失败时回退,避免图片完全丢失)
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},
}
}
// 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())
}
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}
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)
}
// 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) {
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
}
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)
}
}
// MediaCleanupPayload 是「清理七牛孤儿图」任务的 payload:待清理的图片 key 列表。
type MediaCleanupPayload struct {
Keys []string `json:"keys"`
}
// EnqueueMediaCleanup 把一批待清理的七牛 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 清理七牛孤儿图任务:解析 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_runway_draft。
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) 按 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,统一在此兜底)
seasonCode := season.Derive(p.Year, p.CollectionType, p.Season)
// 4) 下载图片并上传到存储(七牛优先,失败兜底本地)
cover, imgs, keys, imgFailed := s.fetchImages(ctx, p.Images, "runway")
if len(p.Images) > 0 && imgFailed {
// 单图失败=整任务失败:先回滚本批已上传的图(七牛 + 本地兜底),避免孤儿文件永远堆在存储里,
// 然后按指数退避自动重试,达上限才置 failed 等后台手动重试。
s.cleanupUploads(ctx, keys)
s.failOrRetry(ctx, job.ID, "image download failed")
return
}
// 5) 写草稿表(status=pending),等待后台审核通过后再晋升正式表
draft := &model.BrandRunwayDraft{
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,
SourceURL: p.SourceURL,
ImageCount: uint16(len(imgs)),
Status: model.DraftStatusPending,
}
id, err := s.repo.CreateRunwayDraft(ctx, draft)
if err != nil {
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),
})
}
if err := s.repo.CreateRunwayDraftImages(ctx, rows); err != nil {
s.failOrRetry(ctx, job.ID, "create draft images: "+err.Error())
return
}
_ = s.repo.MarkDone(ctx, job.ID)
}
// 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 {
_ = s.repo.MarkDone(ctx, job.ID)
return
}
// 2) 下载图片并上传到存储(七牛优先,失败兜底本地)
cover, imgs, keys, imgFailed := s.fetchImages(ctx, p.Images, "street")
if len(p.Images) > 0 && imgFailed {
// 单图失败=整任务失败:先回滚本批已上传的图(七牛 + 本地兜底),避免孤儿文件永远堆在存储里,
// 然后按指数退避自动重试,达上限才置 failed 等后台手动重试。
s.cleanupUploads(ctx, keys)
s.failOrRetry(ctx, job.ID, "image download failed")
return
}
// 3) 写草稿表(status=pending),等待后台审核通过后再晋升 street_snap 正式表
draft := &model.StreetSnapDraft{
JobID: job.ID,
Title: p.TitleEn, // 街拍单标题,爬虫优先填 title_en
Year: p.Year,
City: p.City,
Cover: cover,
SourceURL: p.SourceURL,
ImageCount: uint16(len(imgs)),
Status: model.DraftStatusPending,
}
id, err := s.repo.CreateStreetSnapDraft(ctx, draft)
if err != nil {
s.failOrRetry(ctx, job.ID, "create street draft: "+err.Error())
return
}
rows := make([]model.StreetSnapDraftImage, 0, len(imgs))
for i, img := range imgs {
rows = append(rows, model.StreetSnapDraftImage{
DraftID: id,
Image: img,
Name: fmt.Sprintf("Look %d", i+1),
SortOrder: uint32(i + 1),
})
}
if err := s.repo.CreateStreetSnapDraftImages(ctx, rows); err != nil {
s.failOrRetry(ctx, job.ID, "create street draft images: "+err.Error())
return
}
_ = s.repo.MarkDone(ctx, job.ID)
}
// fetchImages 下载图片并上传到存储,返回 (cover 地址, 全部图片地址, 已成功上传对象的 key 列表, 是否有任意一张失败)。
// prefix 为七牛 key 前缀(runway/ 或 street/)。只要任意一张下载/上传失败,failed 即置 true,
// 调用方据此把整条任务判为失败(不再写草稿),并拿 keys 回滚本批已上传的对象,符合「单图失败=整任务失败」策略。
func (s *IngestService) fetchImages(ctx context.Context, urls []string, prefix string) (string, []string, []string, bool) {
cover := ""
out := make([]string, 0, len(urls))
keys := make([]string, 0, len(urls))
failed := false
for i, u := range urls {
url, key, err := s.downloadOne(ctx, u, i, prefix)
if err != nil {
failed = true
continue
}
if i == 0 {
cover = url
}
out = append(out, url)
keys = append(keys, key)
}
return cover, out, keys, failed
}
// cleanupUploads 删除一批本批次成功上传的对象(七牛 + 本地兜底),用于任务失败回滚:
// 「单图失败=整任务失败」时,前面已成功上传的图若放任不管就成了孤儿,永远堆在存储里。
// 对两种存储都尝试删除(任一不存在即按幂等成功处理);删除失败仅告警,不阻断任务置失败。
// 注意: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: 七牛删除失败 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)
}
}
}
}
// downloadOne 把单张远程图下载后上传到存储(七牛优先,失败兜底本地),
// 返回可直接写入数据库的访问地址(七牛为完整 https URL,本地为相对 /uploads 路径)。
func (s *IngestService) downloadOne(ctx context.Context, u string, idx int, prefix string) (string, string, error) {
resp, err := s.httpClient.Get(u)
if err != nil {
return "", "", err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return "", "", fmt.Errorf("status %d", resp.StatusCode)
}
data, err := io.ReadAll(resp.Body)
if err != nil {
return "", "", err
}
ext := path.Ext(u)
if ext == "" || len(ext) > 5 {
ext = ".jpg"
}
// 内容哈希作 key:同一张图(无论来自哪篇文章/job)永远得到相同 key,
// 七牛覆盖写即天然幂等,不会因重复采集、崩溃重试、reject 重爬而累积孤儿文件。
h := sha1.Sum(data)
hash := hex.EncodeToString(h[:])
key := fmt.Sprintf("%s/%s%s", prefix, hash, ext)
ct := resp.Header.Get("Content-Type")
// 主上传器(七牛)
if url, err := s.uploader.Upload(ctx, key, data, ct); err == nil {
return url, key, nil
} else if s.local != nil {
// 兜底本地,避免图片完全丢失
if lurl, lerr := s.local.Upload(ctx, key, data, ct); lerr == nil {
return lurl, key, nil
} else {
return "", "", err
}
}
return "", "", err
}