Files
backend_v2/internal/repository/ingest_repository.go
toom1996 6c31b18e0f update
2026-09-17 00:38:01 +08:00

440 lines
16 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 repository
import (
"context"
"errors"
"time"
"fashionapi/internal/model"
"gorm.io/gorm"
)
// IngestRepository 爬虫入库管线专属仓储:任务队列(ingest_jobs)+ nonce 防重放
// (ingest_nces)+ 走秀正式表写入(brand_runway / brand_runway_images)。
//
// 去重由 worker 按「实体键」完成:走秀 = brand_id + season_code + collection_type,
// 街拍 = city + year(与 review_repository 的 SaveRunwayFromDraft / SaveStreetSnapFromDraft
// 晋升时的幂等键一致),因此一个秀/街拍只会在正式表留一行,多来源图片在晋升阶段聚合。
//
// 入队与领取用同一张 ingest_jobs 表,领取靠 MySQL 的 FOR UPDATE SKIP LOCKED
// 实现「多 worker 安全并发」——同一条任务只会被一个 worker 拿到,其它 worker 跳过它。
type IngestRepository interface {
// Enqueue 写入一条待处理任务(payload 为原始 JSON)。
Enqueue(ctx context.Context, job *model.IngestJob) error
// Claim 原子领取最多 limit 条 pending 任务并置为 processing,返回这些任务。
// 用 SKIP LOCKED 保证多 worker 不抢同一条。
Claim(ctx context.Context, limit int) ([]model.IngestJob, error)
// MarkDone 标记任务成功。
MarkDone(ctx context.Context, id uint32) error
// MarkFailed 标记任务失败并记录错误(attempts 自增)。
MarkFailed(ctx context.Context, id uint32, errMsg string) error
// ReserveNonce 写入一次性随机串;若已存在(重放)返回 ok=false。
ReserveNonce(ctx context.Context, nonce string) (ok bool, err error)
// RunwayIDByEntity 按实体键(brand_id + season_code + collection_type)查是否已存在走秀正式表;
// 返回 (id, found)。worker 据此跳过已晋升秀的重爬,避免重复下载。
RunwayIDByEntity(ctx context.Context, brandID uint32, seasonCode, collectionType string) (uint32, bool, error)
// CreateRunwayDraft 插入走秀草稿行(status=pending),返回自增主键。
CreateRunwayDraft(ctx context.Context, d *model.BrandRunwayDraft) (uint32, error)
// CreateRunwayDraftImages 批量插入草稿图片行。
CreateRunwayDraftImages(ctx context.Context, imgs []model.BrandRunwayDraftImage) error
// StreetSnapIDByEntity 按实体键(city + year)查是否已存在街拍正式表;返回 (id, found)。
StreetSnapIDByEntity(ctx context.Context, city string, year uint16) (uint32, bool, error)
// CreateStreetSnapDraft 插入街拍草稿行(status=pending),返回自增主键。
CreateStreetSnapDraft(ctx context.Context, d *model.StreetSnapDraft) (uint32, error)
// CreateStreetSnapDraftImages 批量插入街拍草稿图片行。
CreateStreetSnapDraftImages(ctx context.Context, imgs []model.StreetSnapDraftImage) error
// ListJobs 按 id 倒序分页列出入库任务(用于后台监控页)。offset/limit 控制分页。
ListJobs(ctx context.Context, offset, limit int) ([]model.IngestJob, error)
// CountJobs 返回 ingest_jobs 总条数(用于分页计算总页数)。
CountJobs(ctx context.Context) (int64, error)
// RetryJob 把一条 failed 任务重置回 pending,清 last_error/locked_at,等待 worker 重新处理。
RetryJob(ctx context.Context, id uint32) error
// EnqueueMediaCleanup 写入一条「清理七牛孤儿图」任务(payload 为待清理 key 的 JSON)。
// 由删除图集的服务调用,把同步的七牛删除改为异步队列,避免阻塞删除请求。
EnqueueMediaCleanup(ctx context.Context, payload string) error
// ScheduleRetry 失败时调用:attempts+1,未达上限则退避后重置 pending,达上限则置 failed。
// 用于临时失败(网络抖动 / 单图下载失败)的自动重试,区别于永久失败(payload 解析错等)直接 MarkFailed。
ScheduleRetry(ctx context.Context, id uint32, errMsg string) error
// ListImagePHashes 返回全部已晋升图片(runway + street)的 (id, phash, kind),
// 跳过 is_deleted 与 phash=0/NULL(存量未计算)。供入库时与新增图做全局汉明比对(近似去重)。
ListImagePHashes(ctx context.Context) ([]ImagePHash, error)
}
// ImagePHash 已晋升图片的感知哈希摘要,供入库时全局近似去重比对。
type ImagePHash struct {
ID uint32 // 图片数字主键
Phash uint64 // 感知哈希(SQL 已过滤 0/NULL)
Kind string // "runway" | "street"(决定 dup_of 的 hashid 类型)
}
type ingestRepository struct {
db *gorm.DB
}
// NewIngestRepository 创建入库管线仓储。
func NewIngestRepository(db *gorm.DB) IngestRepository {
return &ingestRepository{db: db}
}
func (r *ingestRepository) Enqueue(ctx context.Context, job *model.IngestJob) error {
now := uint32(time.Now().Unix())
job.CreatedAt = now
job.UpdatedAt = now
job.Status = model.IngestStatusPending
return r.db.WithContext(ctx).Create(job).Error
}
// Claim 在事务内 SELECT ... FOR UPDATE SKIP LOCKED 锁定 pending 行,
// 立即置为 processing,再返回这些行,保证领取与状态变更原子、且不被其它 worker 重复领取。
//
// 事务内先做两件事:
// 1. 回收卡死的 processing 任务——worker 崩溃/被杀会留下 processing 孤儿永久卡死,
// 锁定超时(IngestStuckTimeoutSec)后重置回 pending 并立即可领(next_attempt_at=now)。
// 2. 仅领取「已到重试时间」的 pending(next_attempt_at <= now),未到退避点的暂不领。
func (r *ingestRepository) Claim(ctx context.Context, limit int) ([]model.IngestJob, error) {
if limit < 1 {
limit = 10
}
now := uint32(time.Now().Unix())
tx := r.db.WithContext(ctx).Begin()
if tx.Error != nil {
return nil, tx.Error
}
defer func() {
if tx.Error != nil {
tx.Rollback()
}
}()
// 1) 回收卡死的 processing 任务(仅当时间足够大,避免服务器启动初期把刚领取的任务误回收)。
if now > model.IngestStuckTimeoutSec {
if err := tx.Model(&model.IngestJob{}).
Where("status = ? AND locked_at > 0 AND locked_at < ?",
model.IngestStatusProcessing, now-model.IngestStuckTimeoutSec).
Updates(map[string]any{
"status": model.IngestStatusPending,
"locked_at": 0,
"next_attempt_at": now,
"updated_at": now,
}).Error; err != nil {
return nil, err
}
}
// 2) 领取 pending 且已到重试时间的任务(含刚回收的 + 新入队的 + 退避已到期的)。
// 优先处理 media_cleanup(删除图集时异步清理七牛),避免被大批 crawl 任务排到后面、清理迟迟不触发。
var ids []uint32
if err := tx.Raw(
"SELECT id FROM ingest_jobs WHERE status = ? AND next_attempt_at <= ? "+
"ORDER BY CASE kind WHEN ? THEN 0 ELSE 1 END, id ASC LIMIT ? FOR UPDATE SKIP LOCKED",
model.IngestStatusPending, now, model.IngestKindMediaCleanup, limit,
).Scan(&ids).Error; err != nil {
return nil, err
}
if len(ids) == 0 {
tx.Commit()
return nil, nil
}
if err := tx.Model(&model.IngestJob{}).
Where("id IN ?", ids).
Updates(map[string]any{
"status": model.IngestStatusProcessing,
"locked_at": now,
"updated_at": now,
}).Error; err != nil {
return nil, err
}
if err := tx.Commit().Error; err != nil {
return nil, err
}
var jobs []model.IngestJob
if err := r.db.WithContext(ctx).Where("id IN ?", ids).Find(&jobs).Error; err != nil {
return nil, err
}
return jobs, nil
}
func (r *ingestRepository) MarkDone(ctx context.Context, id uint32) error {
return r.db.WithContext(ctx).
Model(&model.IngestJob{}).
Where("id = ?", id).
Updates(map[string]any{
"status": model.IngestStatusDone,
"last_error": "",
"updated_at": uint32(time.Now().Unix()),
}).Error
}
func (r *ingestRepository) MarkFailed(ctx context.Context, id uint32, errMsg string) error {
return r.db.WithContext(ctx).
Model(&model.IngestJob{}).
Where("id = ?", id).
Updates(map[string]any{
"status": model.IngestStatusFailed,
"attempts": gorm.Expr("attempts + 1"),
"last_error": errMsg,
"updated_at": uint32(time.Now().Unix()),
}).Error
}
// ScheduleRetry 失败时调度自动重试:attempts+1,未达上限(IngestMaxAttempts)则按指数退避
// 重置为 pending 并写入 next_attempt_at(到点才可被 Claim 领取);达上限则置 failed,需人工处理。
// 用于临时失败(网络抖动 / 单图下载失败),区别于永久失败(payload 解析错等)直接 MarkFailed。
func (r *ingestRepository) ScheduleRetry(ctx context.Context, id uint32, errMsg string) error {
now := uint32(time.Now().Unix())
var job model.IngestJob
if err := r.db.WithContext(ctx).Select("attempts").Where("id = ?", id).First(&job).Error; err != nil {
return err
}
attempts := int(job.Attempts) + 1
if attempts >= model.IngestMaxAttempts {
return r.db.WithContext(ctx).
Model(&model.IngestJob{}).
Where("id = ?", id).
Updates(map[string]any{
"status": model.IngestStatusFailed,
"attempts": attempts,
"last_error": errMsg,
"updated_at": now,
}).Error
}
delay := model.IngestRetryBackoff(attempts)
return r.db.WithContext(ctx).
Model(&model.IngestJob{}).
Where("id = ?", id).
Updates(map[string]any{
"status": model.IngestStatusPending,
"attempts": attempts,
"last_error": errMsg,
"next_attempt_at": now + uint32(delay),
"locked_at": 0,
"updated_at": now,
}).Error
}
// EnqueueMediaCleanup 写入一条「清理七牛孤儿图」任务(kind=media_cleanup),
// payload 为待清理 key 的 JSON,由 worker 的 processMediaCleanup 按引用计数判定真孤儿后删除。
func (r *ingestRepository) EnqueueMediaCleanup(ctx context.Context, payload string) error {
now := uint32(time.Now().Unix())
job := &model.IngestJob{
Kind: model.IngestKindMediaCleanup,
Payload: payload,
CreatedAt: now,
UpdatedAt: now,
Status: model.IngestStatusPending,
}
return r.db.WithContext(ctx).Create(job).Error
}
// ReserveNonce 写入一次性随机串;依赖 ingest_nces.nonce 主键唯一约束,
// 重复插入触发 DuplicateEntry → 视为重放,返回 ok=false。
func (r *ingestRepository) ReserveNonce(ctx context.Context, nonce string) (bool, error) {
err := r.db.WithContext(ctx).Create(&model.IngestNonce{
Nonce: nonce,
CreatedAt: uint32(time.Now().Unix()),
}).Error
if err != nil {
// 唯一键冲突 → 重放
if errors.Is(err, gorm.ErrDuplicatedKey) || isDuplicateKey(err) {
return false, nil
}
return false, err
}
return true, nil
}
// RunwayIDByEntity 按实体键查走秀正式表是否已存在(与 SaveRunwayFromDraft 晋升键一致)。
func (r *ingestRepository) RunwayIDByEntity(ctx context.Context, brandID uint32, seasonCode, collectionType string) (uint32, bool, error) {
var row struct {
ID uint32 `gorm:"column:id"`
}
err := r.db.WithContext(ctx).
Model(&model.BrandRunway{}).
Select("id").
Where("brand_id = ? AND season_code = ? AND collection_type = ? AND is_deleted = 0", brandID, seasonCode, collectionType).
Limit(1).
Scan(&row).Error
if err != nil {
return 0, false, err
}
if row.ID == 0 {
return 0, false, nil
}
return row.ID, true, nil
}
func (r *ingestRepository) CreateRunwayDraft(ctx context.Context, d *model.BrandRunwayDraft) (uint32, error) {
now := uint32(time.Now().Unix())
d.CreatedAt = now
d.UpdatedAt = now
if d.Status == "" {
d.Status = model.DraftStatusPending
}
if err := r.db.WithContext(ctx).Create(d).Error; err != nil {
return 0, err
}
return d.ID, nil
}
func (r *ingestRepository) CreateRunwayDraftImages(ctx context.Context, imgs []model.BrandRunwayDraftImage) error {
if len(imgs) == 0 {
return nil
}
now := uint32(time.Now().Unix())
for i := range imgs {
imgs[i].CreatedAt = now
imgs[i].UpdatedAt = now
}
return r.db.WithContext(ctx).Create(&imgs).Error
}
// StreetSnapIDByEntity 按实体键(city + year)查街拍正式表是否已存在(与 SaveStreetSnapFromDraft 晋升键一致)。
func (r *ingestRepository) StreetSnapIDByEntity(ctx context.Context, city string, year uint16) (uint32, bool, error) {
var row struct {
ID uint32 `gorm:"column:id"`
}
err := r.db.WithContext(ctx).
Model(&model.StreetSnap{}).
Select("id").
Where("city = ? AND year = ? AND is_deleted = 0", city, year).
Limit(1).
Scan(&row).Error
if err != nil {
return 0, false, err
}
if row.ID == 0 {
return 0, false, nil
}
return row.ID, true, nil
}
func (r *ingestRepository) CreateStreetSnapDraft(ctx context.Context, d *model.StreetSnapDraft) (uint32, error) {
now := uint32(time.Now().Unix())
d.CreatedAt = now
d.UpdatedAt = now
if d.Status == "" {
d.Status = model.DraftStatusPending
}
if err := r.db.WithContext(ctx).Create(d).Error; err != nil {
return 0, err
}
return d.ID, nil
}
func (r *ingestRepository) CreateStreetSnapDraftImages(ctx context.Context, imgs []model.StreetSnapDraftImage) error {
if len(imgs) == 0 {
return nil
}
now := uint32(time.Now().Unix())
for i := range imgs {
imgs[i].CreatedAt = now
imgs[i].UpdatedAt = now
}
return r.db.WithContext(ctx).Create(&imgs).Error
}
// ListJobs 按 id 倒序分页列出入库任务(后台监控页用)。offset/limit 控制分页区间。
func (r *ingestRepository) ListJobs(ctx context.Context, offset, limit int) ([]model.IngestJob, error) {
if limit < 1 {
limit = 50
}
if offset < 0 {
offset = 0
}
var jobs []model.IngestJob
if err := r.db.WithContext(ctx).Order("id DESC").Offset(offset).Limit(limit).Find(&jobs).Error; err != nil {
return nil, err
}
return jobs, nil
}
// CountJobs 返回 ingest_jobs 总条数(后台监控页分页用)。
func (r *ingestRepository) CountJobs(ctx context.Context) (int64, error) {
var n int64
if err := r.db.WithContext(ctx).Model(&model.IngestJob{}).Count(&n).Error; err != nil {
return 0, err
}
return n, nil
}
// RetryJob 把一条 failed 任务重置回 pending,清掉 last_error / locked_at / next_attempt_at,
// 让 worker(每 3s 扫一次 pending)立即重新拉起处理。仅对 failed 生效,其它状态原样不动。
func (r *ingestRepository) RetryJob(ctx context.Context, id uint32) error {
now := uint32(time.Now().Unix())
return r.db.WithContext(ctx).
Model(&model.IngestJob{}).
Where("id = ? AND status = ?", id, model.IngestStatusFailed).
Updates(map[string]any{
"status": model.IngestStatusPending,
"locked_at": 0,
"last_error": "",
"next_attempt_at": 0,
"updated_at": now,
}).Error
}
// ListImagePHashes 返回全部已晋升图片(runway + street)的感知哈希摘要,供入库时全局近似去重。
// 跳过 is_deleted 与 phash=0/NULL(存量未计算)。结果合并 runway + street 两类,
// 用 Kind 标注类型,调用方据此把 dup_of 编码成对应 hashid 类型。
//
// 注:每次入库任务都会全量拉取一次(图片量当前为千级,可接受);若后续图片量到十万级,
// 可改为按 runway_id/snap_id 分批或加内存缓存 + 定时刷新,避免每 job 一次全表扫描。
func (r *ingestRepository) ListImagePHashes(ctx context.Context) ([]ImagePHash, error) {
type phRow struct {
ID uint32 `gorm:"column:id"`
Phash uint64 `gorm:"column:phash"`
}
out := make([]ImagePHash, 0, 64)
const whereActive = "is_deleted = 0 AND phash IS NOT NULL AND phash <> 0"
var rw []phRow
if err := r.db.WithContext(ctx).Model(&model.BrandRunwayImage{}).
Select("id, phash").Where(whereActive).Scan(&rw).Error; err != nil {
return nil, err
}
for _, x := range rw {
out = append(out, ImagePHash{ID: x.ID, Phash: x.Phash, Kind: "runway"})
}
var sn []phRow
if err := r.db.WithContext(ctx).Model(&model.StreetSnapImage{}).
Select("id, phash").Where(whereActive).Scan(&sn).Error; err != nil {
return nil, err
}
for _, x := range sn {
out = append(out, ImagePHash{ID: x.ID, Phash: x.Phash, Kind: "street"})
}
return out, nil
}
// isDuplicateKey 兜底:gorm 的 ErrDuplicatedKey 在不同驱动下的封装不一定一致,
// 直接命中 MySQL 1062 错误号更稳。
func isDuplicateKey(err error) bool {
if err == nil {
return false
}
msg := err.Error()
return containsAny(msg, "Duplicate entry", "1062", "UNIQUE constraint failed")
}
func containsAny(s string, subs ...string) bool {
for _, sub := range subs {
if len(sub) > 0 && indexOf(s, sub) >= 0 {
return true
}
}
return false
}
func indexOf(s, sub string) int {
for i := 0; i+len(sub) <= len(s); i++ {
if s[i:i+len(sub)] == sub {
return i
}
}
return -1
}