This commit is contained in:
toom1996
2026-09-07 00:04:01 +08:00
parent 6a5a4378ab
commit 10d8a96e8c
29 changed files with 5349 additions and 0 deletions

View File

@ -0,0 +1,622 @@
package repository
import (
"context"
"errors"
"time"
"fashionapi/internal/model"
"gorm.io/gorm"
)
// IngestRepository 爬虫入库管线专属仓储:任务队列(ingest_jobs)+ nonce 防重放
// (ingest_nces)+ 走秀正式表写入(brand_runway / brand_runway_images,按 source_url 去重)。
//
// 入队与领取用同一张 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)
// RunwayIDBySourceURL 按来源链接查是否已存在走秀;返回 (id, found)。
RunwayIDBySourceURL(ctx context.Context, sourceURL string) (uint32, bool, error)
// CreateRunway 插入走秀正式行,返回自增主键。
CreateRunway(ctx context.Context, r *model.BrandRunway) (uint32, error)
// CreateRunwayImages 批量插入走秀图片行。
CreateRunwayImages(ctx context.Context, imgs []model.BrandRunwayImage) error
// CreateRunwayDraft 插入走秀草稿行(status=pending),返回自增主键。
CreateRunwayDraft(ctx context.Context, d *model.BrandRunwayDraft) (uint32, error)
// CreateRunwayDraftImages 批量插入草稿图片行。
CreateRunwayDraftImages(ctx context.Context, imgs []model.BrandRunwayDraftImage) error
// DraftIDBySourceURL 按来源链接查是否已有 pending 草稿;返回 (id, found)。
DraftIDBySourceURL(ctx context.Context, sourceURL string) (uint32, bool, error)
// StreetSnapIDBySourceURL 按来源链接查是否已存在街拍正式表;返回 (id, found)。
StreetSnapIDBySourceURL(ctx context.Context, sourceURL string) (uint32, bool, error)
// DraftStreetIDBySourceURL 按来源链接查是否已有 pending 街拍草稿;返回 (id, found)。
DraftStreetIDBySourceURL(ctx context.Context, sourceURL string) (uint32, bool, error)
// RunwayIDByUnique 按实体键(品牌+季节码+系列)查是否已存在正式走秀;返回 (id, found)。
// 多来源爬同品牌同季时应合并为一条,故去重键是实体而非 source_url。
RunwayIDByUnique(ctx context.Context, brandID uint32, seasonCode, collectionType string) (uint32, bool, error)
// DraftRunwayIDByUnique 按实体键查是否已有 pending 走秀草稿;返回 (id, found)。
DraftRunwayIDByUnique(ctx context.Context, brandID uint32, seasonCode, collectionType string) (uint32, bool, error)
// StreetSnapIDByUnique 按实体键(城市+年份)查是否已存在街拍正式表;返回 (id, found)。
StreetSnapIDByUnique(ctx context.Context, city string, year uint16) (uint32, bool, error)
// DraftStreetSnapIDByUnique 按实体键查是否已有 pending 街拍草稿;返回 (id, found)。
DraftStreetSnapIDByUnique(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 倒序列出最近的入库任务(用于后台监控页)。
ListJobs(ctx context.Context, limit int) ([]model.IngestJob, 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
// ExistsSourceURLs 批量判断一批 source_url 是否已爬取过,返回命中集合(true = 已存在,无需再抓)。
//
// 语义必须与 worker 判重(processRunway / processStreet)严格一致,否则会出现
// 「预检说没有、worker 又判重命中」的重复抓取,或「预检说有、实际已被拒/已删」的漏抓:
// - 正式表(brand_runway / street_snap):source_url 命中且 is_deleted = 0
// - 草稿表(brand_runway_draft / street_snap_draft):status = pending 且 is_deleted = 0
// —— 被审核拒绝(rejected)的草稿不算已存在,允许重新抓取。
ExistsSourceURLs(ctx context.Context, urls []string) (map[string]bool, error)
}
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
}
func (r *ingestRepository) RunwayIDBySourceURL(ctx context.Context, sourceURL string) (uint32, bool, error) {
if sourceURL == "" {
return 0, false, nil
}
var row struct {
ID uint32 `gorm:"column:id"`
}
err := r.db.WithContext(ctx).
Model(&model.BrandRunway{}).
Select("id").
Where("source_url = ? AND is_deleted = 0", sourceURL).
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) CreateRunway(ctx context.Context, rw *model.BrandRunway) (uint32, error) {
now := uint32(time.Now().Unix())
rw.CreatedAt = now
rw.UpdatedAt = now
if rw.IsDeleted == 0 {
rw.IsDeleted = 0
}
if err := r.db.WithContext(ctx).Create(rw).Error; err != nil {
return 0, err
}
return rw.ID, nil
}
func (r *ingestRepository) CreateRunwayImages(ctx context.Context, imgs []model.BrandRunwayImage) error {
if len(imgs) == 0 {
return nil
}
return r.db.WithContext(ctx).Create(&imgs).Error
}
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
}
func (r *ingestRepository) DraftIDBySourceURL(ctx context.Context, sourceURL string) (uint32, bool, error) {
if sourceURL == "" {
return 0, false, nil
}
var row struct {
ID uint32 `gorm:"column:id"`
}
err := r.db.WithContext(ctx).
Model(&model.BrandRunwayDraft{}).
Select("id").
Where("source_url = ? AND status = ? AND is_deleted = 0", sourceURL, model.DraftStatusPending).
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) StreetSnapIDBySourceURL(ctx context.Context, sourceURL string) (uint32, bool, error) {
if sourceURL == "" {
return 0, false, nil
}
var row struct {
ID uint32 `gorm:"column:id"`
}
err := r.db.WithContext(ctx).
Model(&model.StreetSnap{}).
Select("id").
Where("source_url = ? AND is_deleted = 0", sourceURL).
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) DraftStreetIDBySourceURL(ctx context.Context, sourceURL string) (uint32, bool, error) {
if sourceURL == "" {
return 0, false, nil
}
var row struct {
ID uint32 `gorm:"column:id"`
}
err := r.db.WithContext(ctx).
Model(&model.StreetSnapDraft{}).
Select("id").
Where("source_url = ? AND status = ? AND is_deleted = 0", sourceURL, model.DraftStatusPending).
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
}
// RunwayIDByUnique 按实体键(品牌+季节码+系列)查正式走秀,多来源同实体合并为一条。
func (r *ingestRepository) RunwayIDByUnique(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
}
// DraftRunwayIDByUnique 按实体键查 pending 走秀草稿。
func (r *ingestRepository) DraftRunwayIDByUnique(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.BrandRunwayDraft{}).
Select("id").
Where("brand_id = ? AND season_code = ? AND collection_type = ? AND status = ? AND is_deleted = 0", brandID, seasonCode, collectionType, model.DraftStatusPending).
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
}
// StreetSnapIDByUnique 按实体键(城市+年份)查街拍正式表;城市为空则不去重(避免空城市互并)。
func (r *ingestRepository) StreetSnapIDByUnique(ctx context.Context, city string, year uint16) (uint32, bool, error) {
if city == "" {
return 0, false, nil
}
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
}
// DraftStreetSnapIDByUnique 按实体键查 pending 街拍草稿;城市为空则不去重。
func (r *ingestRepository) DraftStreetSnapIDByUnique(ctx context.Context, city string, year uint16) (uint32, bool, error) {
if city == "" {
return 0, false, nil
}
var row struct {
ID uint32 `gorm:"column:id"`
}
err := r.db.WithContext(ctx).
Model(&model.StreetSnapDraft{}).
Select("id").
Where("city = ? AND year = ? AND status = ? AND is_deleted = 0", city, year, model.DraftStatusPending).
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 倒序列出最近的入库任务(后台监控页用)。
func (r *ingestRepository) ListJobs(ctx context.Context, limit int) ([]model.IngestJob, error) {
if limit < 1 {
limit = 200
}
var jobs []model.IngestJob
if err := r.db.WithContext(ctx).Order("id DESC").Limit(limit).Find(&jobs).Error; err != nil {
return nil, err
}
return jobs, 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
}
// 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
}
// ExistsSourceURLs 批量判断 source_url 是否已爬取过(预检接口用)。
//
// 走秀与街拍的正式表、草稿表各查一次 IN,命中即标记。与 worker 判重共用同一套条件,
// 保证「爬虫预检跳过」与「worker 判重跳过」判定结果一致。
func (r *ingestRepository) ExistsSourceURLs(ctx context.Context, urls []string) (map[string]bool, error) {
out := make(map[string]bool)
if len(urls) == 0 {
return out, nil
}
// 去重 + 去空:同一批 URL 可能重复出现,避免无谓的返回行与 SQL 长度。
uniq := make([]string, 0, len(urls))
seen := make(map[string]struct{}, len(urls))
for _, u := range urls {
if u == "" {
continue
}
if _, ok := seen[u]; ok {
continue
}
seen[u] = struct{}{}
uniq = append(uniq, u)
}
if len(uniq) == 0 {
return out, nil
}
// 单次请求上限,避免超长 IN 拖慢数据库(超出部分按「未抓过」处理,最多多抓几个,不会漏判已存在的)。
const maxBatch = 500
if len(uniq) > maxBatch {
uniq = uniq[:maxBatch]
}
// 四张表:正式表只看未删除,草稿表只认 pending(被拒草稿允许重抓)。
queries := []struct {
dest any
where string
args []any
}{
{dest: &model.BrandRunway{}, where: "source_url IN ? AND is_deleted = 0"},
{dest: &model.StreetSnap{}, where: "source_url IN ? AND is_deleted = 0"},
{dest: &model.BrandRunwayDraft{}, where: "source_url IN ? AND status = ? AND is_deleted = 0", args: []any{model.DraftStatusPending}},
{dest: &model.StreetSnapDraft{}, where: "source_url IN ? AND status = ? AND is_deleted = 0", args: []any{model.DraftStatusPending}},
}
type hit struct {
SourceURL string `gorm:"column:source_url"`
}
for i := range queries {
var rows []hit
q := r.db.WithContext(ctx).
Model(queries[i].dest).
Select("source_url").
Where(queries[i].where, append([]any{uniq}, queries[i].args...)...)
if err := q.Scan(&rows).Error; err != nil {
return nil, err
}
for _, row := range rows {
out[row.SourceURL] = true
}
// 全部命中就无需再查后面的表。
if len(out) == len(uniq) {
break
}
}
return out, nil
}