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

623 lines
22 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 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
}