diff --git a/internal/model/runway_draft.go b/internal/model/runway_draft.go index aa5ef17..f956899 100644 --- a/internal/model/runway_draft.go +++ b/internal/model/runway_draft.go @@ -32,16 +32,19 @@ func (BrandRunwayDraft) TableName() string { return "brand_runway_draft" } // BrandRunwayDraftImage 草稿图片(对应 brand_runway_draft_images)。 type BrandRunwayDraftImage struct { - ID uint32 `gorm:"primaryKey;column:id" json:"id"` - DraftID uint32 `gorm:"column:draft_id" json:"draft_id"` - Image string `gorm:"column:image" json:"image"` - Name string `gorm:"column:name" json:"name"` - SortOrder uint32 `gorm:"column:sort_order" json:"sort_order"` - LookIndex uint32 `gorm:"column:look_index" json:"look_index"` // 细节图归属的主图序号 - IsDetail uint8 `gorm:"column:is_detail" json:"is_detail"` // 0=主图 1=细节图 - IsDeleted uint8 `gorm:"column:is_deleted" json:"is_deleted"` - CreatedAt uint32 `gorm:"column:created_at" json:"created_at"` - UpdatedAt uint32 `gorm:"column:updated_at" json:"updated_at"` + ID uint32 `gorm:"primaryKey;column:id" json:"id"` + DraftID uint32 `gorm:"column:draft_id" json:"draft_id"` + Image string `gorm:"column:image" json:"image"` + Name string `gorm:"column:name" json:"name"` + SortOrder uint32 `gorm:"column:sort_order" json:"sort_order"` + LookIndex uint32 `gorm:"column:look_index" json:"look_index"` // 细节图归属的主图序号 + IsDetail uint8 `gorm:"column:is_detail" json:"is_detail"` // 0=主图 1=细节图 + IsDeleted uint8 `gorm:"column:is_deleted" json:"is_deleted"` + CreatedAt uint32 `gorm:"column:created_at" json:"created_at"` + UpdatedAt uint32 `gorm:"column:updated_at" json:"updated_at"` + Phash uint64 `gorm:"column:phash" json:"phash"` // 感知哈希 64-bit;0=未计算 + IsDuplicate uint8 `gorm:"column:is_duplicate" json:"is_duplicate"` // 近似重复标记:0=否 1=是(命中留痕不删) + DupOf string `gorm:"column:dup_of" json:"dup_of"` // 近似重复指向的图 uid(hashid),''=非重复 } // TableName 指定草稿图片表名。 diff --git a/internal/model/runway_image.go b/internal/model/runway_image.go index 430d4f1..5a09e64 100644 --- a/internal/model/runway_image.go +++ b/internal/model/runway_image.go @@ -2,17 +2,20 @@ package model // BrandRunwayImage 走秀图片。 type BrandRunwayImage struct { - ID uint32 `gorm:"primaryKey;column:id" json:"id"` - CreatedAt uint32 `gorm:"column:created_at" json:"created_at"` - UpdatedAt uint32 `gorm:"column:updated_at" json:"updated_at"` - IsDeleted uint8 `gorm:"column:is_deleted" json:"is_deleted"` - Image string `gorm:"column:image" json:"image"` - RunwayID uint32 `gorm:"column:runway_id" json:"runway_id"` - BrandID uint32 `gorm:"column:brand_id" json:"brand_id"` - Name string `gorm:"column:name" json:"name"` - SortOrder uint32 `gorm:"column:sort_order" json:"sort_order"` // 拖拽排序用,由迁移脚本新增 - LookIndex uint32 `gorm:"column:look_index" json:"look_index"` // 细节图归属的主图序号(主图=该 look;细节图=所属主图的 look) - IsDetail uint8 `gorm:"column:is_detail" json:"is_detail"` // 0=主图/look 图(默认展示) 1=细节图 + ID uint32 `gorm:"primaryKey;column:id" json:"id"` + CreatedAt uint32 `gorm:"column:created_at" json:"created_at"` + UpdatedAt uint32 `gorm:"column:updated_at" json:"updated_at"` + IsDeleted uint8 `gorm:"column:is_deleted" json:"is_deleted"` + Image string `gorm:"column:image" json:"image"` + RunwayID uint32 `gorm:"column:runway_id" json:"runway_id"` + BrandID uint32 `gorm:"column:brand_id" json:"brand_id"` + Name string `gorm:"column:name" json:"name"` + SortOrder uint32 `gorm:"column:sort_order" json:"sort_order"` // 拖拽排序用,由迁移脚本新增 + LookIndex uint32 `gorm:"column:look_index" json:"look_index"` // 细节图归属的主图序号(主图=该 look;细节图=所属主图的 look) + IsDetail uint8 `gorm:"column:is_detail" json:"is_detail"` // 0=主图/look 图(默认展示) 1=细节图 + Phash uint64 `gorm:"column:phash" json:"phash"` // 感知哈希 64-bit;0=未计算(存量) + IsDuplicate uint8 `gorm:"column:is_duplicate" json:"is_duplicate"` // 近似重复标记:0=否 1=是(命中留痕不删) + DupOf string `gorm:"column:dup_of" json:"dup_of"` // 近似重复指向的图 uid(hashid),''=非重复 } // TableName 指定表名。 diff --git a/internal/model/street_snap.go b/internal/model/street_snap.go index 35e7217..4210b5b 100644 --- a/internal/model/street_snap.go +++ b/internal/model/street_snap.go @@ -23,14 +23,17 @@ func (StreetSnap) TableName() string { return "street_snap" } // // 仿 brand_runway_images,但把 runway_id 改名为 snap_id,去掉 brand_id(街拍无品牌关联)。 type StreetSnapImage struct { - ID uint32 `gorm:"primaryKey;column:id" json:"id"` - SnapID uint32 `gorm:"column:snap_id" json:"snap_id"` - Image string `gorm:"column:image" json:"image"` - Name string `gorm:"column:name" json:"name"` - SortOrder uint32 `gorm:"column:sort_order" json:"sort_order"` - IsDeleted uint8 `gorm:"column:is_deleted" json:"is_deleted"` - CreatedAt uint32 `gorm:"column:created_at" json:"created_at"` - UpdatedAt uint32 `gorm:"column:updated_at" json:"updated_at"` + ID uint32 `gorm:"primaryKey;column:id" json:"id"` + SnapID uint32 `gorm:"column:snap_id" json:"snap_id"` + Image string `gorm:"column:image" json:"image"` + Name string `gorm:"column:name" json:"name"` + SortOrder uint32 `gorm:"column:sort_order" json:"sort_order"` + IsDeleted uint8 `gorm:"column:is_deleted" json:"is_deleted"` + CreatedAt uint32 `gorm:"column:created_at" json:"created_at"` + UpdatedAt uint32 `gorm:"column:updated_at" json:"updated_at"` + Phash uint64 `gorm:"column:phash" json:"phash"` // 感知哈希 64-bit;0=未计算(存量) + IsDuplicate uint8 `gorm:"column:is_duplicate" json:"is_duplicate"` // 近似重复标记:0=否 1=是(命中留痕不删) + DupOf string `gorm:"column:dup_of" json:"dup_of"` // 近似重复指向的图 uid(hashid),''=非重复 } // TableName 指定图片明细表名。 diff --git a/internal/model/street_snap_draft.go b/internal/model/street_snap_draft.go index 442040e..30f0155 100644 --- a/internal/model/street_snap_draft.go +++ b/internal/model/street_snap_draft.go @@ -28,14 +28,17 @@ func (StreetSnapDraft) TableName() string { return "street_snap_draft" } // StreetSnapDraftImage 街拍草稿图片(对应 street_snap_draft_images)。 type StreetSnapDraftImage struct { - ID uint32 `gorm:"primaryKey;column:id" json:"id"` - DraftID uint32 `gorm:"column:draft_id" json:"draft_id"` - Image string `gorm:"column:image" json:"image"` - Name string `gorm:"column:name" json:"name"` - SortOrder uint32 `gorm:"column:sort_order" json:"sort_order"` - IsDeleted uint8 `gorm:"column:is_deleted" json:"is_deleted"` - CreatedAt uint32 `gorm:"column:created_at" json:"created_at"` - UpdatedAt uint32 `gorm:"column:updated_at" json:"updated_at"` + ID uint32 `gorm:"primaryKey;column:id" json:"id"` + DraftID uint32 `gorm:"column:draft_id" json:"draft_id"` + Image string `gorm:"column:image" json:"image"` + Name string `gorm:"column:name" json:"name"` + SortOrder uint32 `gorm:"column:sort_order" json:"sort_order"` + IsDeleted uint8 `gorm:"column:is_deleted" json:"is_deleted"` + CreatedAt uint32 `gorm:"column:created_at" json:"created_at"` + UpdatedAt uint32 `gorm:"column:updated_at" json:"updated_at"` + Phash uint64 `gorm:"column:phash" json:"phash"` // 感知哈希 64-bit;0=未计算 + IsDuplicate uint8 `gorm:"column:is_duplicate" json:"is_duplicate"` // 近似重复标记:0=否 1=是(命中留痕不删) + DupOf string `gorm:"column:dup_of" json:"dup_of"` // 近似重复指向的图 uid(hashid),''=非重复 } // TableName 指定草稿图片表名。 diff --git a/internal/pkg/imgurl/vip_check_test.go b/internal/pkg/imgurl/vip_check_test.go index 4d29360..5b35a07 100644 --- a/internal/pkg/imgurl/vip_check_test.go +++ b/internal/pkg/imgurl/vip_check_test.go @@ -10,7 +10,7 @@ import ( // - 查询串模式(style 含 =):base_url/key?w=..&q=.. // - /uploads/ 与 http(s) 原样返回 func TestDisplayStyle(t *testing.T) { - c := New("https://mybucket.s3.bitiful.net", "high", "thumb", "") + c := New("https://mybucket.s3.bitiful.net", "high", "thumb") got := c.Compose("runway/abc123.jpg") want := "https://mybucket.s3.bitiful.net/runway/abc123.jpg!style:high" if got != want { @@ -18,13 +18,13 @@ func TestDisplayStyle(t *testing.T) { } // 缺省回落 high - c0 := New("https://mybucket.s3.bitiful.net", "", "", "") + c0 := New("https://mybucket.s3.bitiful.net", "", "") if c0.Compose("runway/abc123.jpg") != "https://mybucket.s3.bitiful.net/runway/abc123.jpg!style:high" { t.Fatal("style 缺省应回落 high") } // 查询串模式 - cq := New("https://mybucket.s3.bitiful.net", "w=1080&q=80&fmt=webp", "", "") + cq := New("https://mybucket.s3.bitiful.net", "w=1080&q=80&fmt=webp", "") gq := cq.Compose("runway/abc123.jpg") wq := "https://mybucket.s3.bitiful.net/runway/abc123.jpg?w=1080&q=80&fmt=webp" if gq != wq { @@ -45,7 +45,7 @@ func TestDisplayStyle(t *testing.T) { // TestComposeThumb 验证列表缩略图拼装:走 StyleThumb,兜底逻辑与 Compose 一致。 func TestComposeThumb(t *testing.T) { - c := New("https://mybucket.s3.bitiful.net", "high", "thumb", "") + c := New("https://mybucket.s3.bitiful.net", "high", "thumb") got := c.ComposeThumb("runway/abc123.jpg") want := "https://mybucket.s3.bitiful.net/runway/abc123.jpg!style:thumb" if got != want { @@ -53,7 +53,7 @@ func TestComposeThumb(t *testing.T) { } // 未配置 styleThumb 时应回落展示样式,不得拼出空样式(key!style:) - c0 := New("https://mybucket.s3.bitiful.net", "high", "", "") + c0 := New("https://mybucket.s3.bitiful.net", "high", "") if got := c0.ComposeThumb("runway/abc123.jpg"); got != "https://mybucket.s3.bitiful.net/runway/abc123.jpg!style:high" { t.Fatalf("styleThumb 缺省应回落 styleDisplay,实际 %s", got) } @@ -72,7 +72,7 @@ func TestComposeThumb(t *testing.T) { // TestDisplayStyleURLSafe 确认展示 URL 不含任何签名/过期参数(匿名可访问)。 func TestDisplayStyleURLSafe(t *testing.T) { - c := New("https://mybucket.s3.bitiful.net", "high", "thumb", "") + c := New("https://mybucket.s3.bitiful.net", "high", "thumb") u := c.Compose("runway/abc123.jpg") for _, bad := range []string{"X-Amz", "sign", "e="} { if strings.Contains(u, bad) { diff --git a/internal/repository/ingest_repository.go b/internal/repository/ingest_repository.go index 9bb2245..29a6f9e 100644 --- a/internal/repository/ingest_repository.go +++ b/internal/repository/ingest_repository.go @@ -56,6 +56,16 @@ type IngestRepository interface { // 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 { @@ -366,6 +376,40 @@ func (r *ingestRepository) RetryJob(ctx context.Context, id uint32) error { }).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 { diff --git a/internal/repository/review_repository.go b/internal/repository/review_repository.go index 2f6ea10..9835419 100644 --- a/internal/repository/review_repository.go +++ b/internal/repository/review_repository.go @@ -245,15 +245,18 @@ func (r *reviewRepository) SaveRunwayFromDraft(ctx context.Context, draftID uint rows := make([]model.BrandRunwayImage, 0, len(imgs)) for i, im := range imgs { rows = append(rows, model.BrandRunwayImage{ - RunwayID: runwayID, - BrandID: draft.BrandID, - Image: im.Image, - Name: im.Name, - SortOrder: uint32(i + 1), - LookIndex: im.LookIndex, - IsDetail: im.IsDetail, - CreatedAt: now, - UpdatedAt: now, + RunwayID: runwayID, + BrandID: draft.BrandID, + Image: im.Image, + Name: im.Name, + SortOrder: uint32(i + 1), + LookIndex: im.LookIndex, + IsDetail: im.IsDetail, + Phash: im.Phash, + IsDuplicate: im.IsDuplicate, + DupOf: im.DupOf, + CreatedAt: now, + UpdatedAt: now, }) } if cErr := tx.Create(&rows).Error; cErr != nil { @@ -437,12 +440,15 @@ func (r *reviewRepository) SaveStreetSnapFromDraft(ctx context.Context, draftID rows := make([]model.StreetSnapImage, 0, len(imgs)) for i, im := range imgs { rows = append(rows, model.StreetSnapImage{ - SnapID: snapID, - Image: im.Image, - Name: im.Name, - SortOrder: uint32(i + 1), - CreatedAt: now, - UpdatedAt: now, + SnapID: snapID, + Image: im.Image, + Name: im.Name, + SortOrder: uint32(i + 1), + Phash: im.Phash, + IsDuplicate: im.IsDuplicate, + DupOf: im.DupOf, + CreatedAt: now, + UpdatedAt: now, }) } if cErr := tx.Create(&rows).Error; cErr != nil { diff --git a/internal/service/ingest_cleanup_test.go b/internal/service/ingest_cleanup_test.go index 5f48fd2..bae1e24 100644 --- a/internal/service/ingest_cleanup_test.go +++ b/internal/service/ingest_cleanup_test.go @@ -30,7 +30,7 @@ func TestFetchImagesCleansUpOnFailure(t *testing.T) { defer srv.Close() urls := []string{srv.URL + "/ok1.jpg", srv.URL + "/ok2.jpg", srv.URL + "/bad.jpg"} - _, out, keys, failed := s.fetchImages(context.Background(), urls, "runway") + _, out, _, keys, failed := s.fetchImages(context.Background(), urls, "runway") if !failed { t.Fatalf("expected failed=true when one image errors") } diff --git a/internal/service/ingest_service.go b/internal/service/ingest_service.go index a3231e3..025f1e9 100644 --- a/internal/service/ingest_service.go +++ b/internal/service/ingest_service.go @@ -15,6 +15,7 @@ import ( "fashionapi/internal/dto" "fashionapi/internal/model" "fashionapi/internal/pkg/hashid" + "fashionapi/internal/pkg/phash" "fashionapi/internal/pkg/season" "fashionapi/internal/pkg/storage" "fashionapi/internal/repository" @@ -234,17 +235,24 @@ func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob, // 4) 下载图片并上传到存储(结构化 Looks 优先:主图+细节图分组;否则回退 Images 全部视为主图)。 // 内容哈希(sha1)key 保证重爬不产生孤儿文件:失败回滚删本批 key 即可。 + // 同时算出每张图的 pHash(downloadOne 内基于原始字节),供后续全局近似去重标记。 var cover string var draftImages []model.BrandRunwayDraftImage var keys []string var imgFailed bool var imageCount uint16 + var phashList []uint64 if len(p.Looks) > 0 { cover, draftImages, keys, imgFailed = s.fetchLookImages(ctx, p.Looks, "runway") imageCount = uint16(len(p.Looks)) + phashList = make([]uint64, len(draftImages)) + for i, d := range draftImages { + phashList[i] = d.Phash + } } else { var imgs []string - cover, imgs, keys, imgFailed = s.fetchImages(ctx, p.Images, "runway") + var phs []uint64 + cover, imgs, phs, keys, imgFailed = s.fetchImages(ctx, p.Images, "runway") imageCount = uint16(len(imgs)) for i, img := range imgs { draftImages = append(draftImages, model.BrandRunwayDraftImage{ @@ -253,8 +261,10 @@ func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob, SortOrder: uint32(i + 1), LookIndex: uint32(i + 1), IsDetail: 0, + Phash: phs[i], }) } + phashList = phs } if imgFailed { // 单图失败=整任务失败:先回滚本批已上传的图(S4 + 本地兜底),避免孤儿文件永远堆在存储里, @@ -264,6 +274,13 @@ func (s *IngestService) processRunway(ctx context.Context, job model.IngestJob, return } + // 4.5) 全局近似去重标记:与「已晋升图片」比对汉明距离,命中则留痕(is_duplicate=1 + dup_of), + // 不丢弃、交后台人工裁决(只拦新增,不碰存量)。比对失败仅告警,不阻断入库。 + s.tagDuplicates(ctx, phashList, func(i int, dupOf string) { + draftImages[i].IsDuplicate = 1 + draftImages[i].DupOf = dupOf + }) + // 5) 写草稿表(status=pending),等待后台审核通过后再晋升正式表 draft := &model.BrandRunwayDraft{ JobID: job.ID, @@ -309,7 +326,7 @@ func (s *IngestService) fetchLookImages(ctx context.Context, looks []dto.RunwayL lookIdx := li + 1 if look.Main != "" { order++ - url, key, err := s.downloadOne(ctx, look.Main, order, prefix) + url, key, ph, err := s.downloadOne(ctx, look.Main, order, prefix) if err != nil { failed = true } else { @@ -323,6 +340,7 @@ func (s *IngestService) fetchLookImages(ctx context.Context, looks []dto.RunwayL SortOrder: uint32(order), LookIndex: uint32(lookIdx), IsDetail: 0, + Phash: ph, }) } } @@ -331,7 +349,7 @@ func (s *IngestService) fetchLookImages(ctx context.Context, looks []dto.RunwayL continue } order++ - url, key, err := s.downloadOne(ctx, d, order, prefix) + url, key, ph, err := s.downloadOne(ctx, d, order, prefix) if err != nil { failed = true continue @@ -343,6 +361,7 @@ func (s *IngestService) fetchLookImages(ctx context.Context, looks []dto.RunwayL SortOrder: uint32(order), LookIndex: uint32(lookIdx), IsDetail: 1, + Phash: ph, }) } } @@ -358,8 +377,8 @@ func (s *IngestService) processStreet(ctx context.Context, job model.IngestJob, return } - // 2) 下载图片并上传到存储(七牛优先,失败兜底本地) - cover, imgs, keys, imgFailed := s.fetchImages(ctx, p.Images, "street") + // 2) 下载图片并上传到存储(七牛优先,失败兜底本地);同时算 pHash 供近似去重。 + cover, imgs, phs, keys, imgFailed := s.fetchImages(ctx, p.Images, "street") if len(p.Images) > 0 && imgFailed { // 单图失败=整任务失败:先回滚本批已上传的图(七牛 + 本地兜底),避免孤儿文件永远堆在存储里, // 然后按指数退避自动重试,达上限才置 failed 等后台手动重试。 @@ -390,8 +409,16 @@ func (s *IngestService) processStreet(ctx context.Context, job model.IngestJob, Image: img, Name: fmt.Sprintf("Look %d", i+1), SortOrder: uint32(i + 1), + Phash: phs[i], }) } + + // 3.5) 全局近似去重标记(与走秀同逻辑,只拦新增、留痕不删)。 + s.tagDuplicates(ctx, phs, func(i int, dupOf string) { + rows[i].IsDuplicate = 1 + rows[i].DupOf = dupOf + }) + if err := s.repo.CreateStreetSnapDraftImages(ctx, rows); err != nil { s.failOrRetry(ctx, job.ID, "create street draft images: "+err.Error()) return @@ -399,16 +426,18 @@ func (s *IngestService) processStreet(ctx context.Context, job model.IngestJob, _ = s.repo.MarkDone(ctx, job.ID) } -// fetchImages 下载图片并上传到存储,返回 (cover 地址, 全部图片地址, 已成功上传对象的 key 列表, 是否有任意一张失败)。 +// fetchImages 下载图片并上传到存储,返回 (cover 地址, 全部图片地址, 各图 pHash, 已成功上传对象的 key 列表, 是否有任意一张失败)。 // prefix 为七牛 key 前缀(runway/ 或 street/)。只要任意一张下载/上传失败,failed 即置 true, // 调用方据此把整条任务判为失败(不再写草稿),并拿 keys 回滚本批已上传的对象,符合「单图失败=整任务失败」策略。 -func (s *IngestService) fetchImages(ctx context.Context, urls []string, prefix string) (string, []string, []string, bool) { +// 返回的 pHash 与图片地址按索引对齐(phash=0 表示解码失败未计算,比对时跳过)。 +func (s *IngestService) fetchImages(ctx context.Context, urls []string, prefix string) (string, []string, []uint64, []string, bool) { cover := "" out := make([]string, 0, len(urls)) + phashes := make([]uint64, 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) + url, key, ph, err := s.downloadOne(ctx, u, i, prefix) if err != nil { failed = true continue @@ -417,9 +446,10 @@ func (s *IngestService) fetchImages(ctx context.Context, urls []string, prefix s cover = url } out = append(out, url) + phashes = append(phashes, ph) keys = append(keys, key) } - return cover, out, keys, failed + return cover, out, phashes, keys, failed } // cleanupUploads 删除一批本批次成功上传的对象(七牛 + 本地兜底),用于任务失败回滚: @@ -442,19 +472,20 @@ func (s *IngestService) cleanupUploads(ctx context.Context, keys []string) { } // downloadOne 把单张远程图下载后上传到存储(七牛优先,失败兜底本地), -// 返回可直接写入数据库的访问地址(七牛为完整 https URL,本地为相对 /uploads 路径)。 -func (s *IngestService) downloadOne(ctx context.Context, u string, idx int, prefix string) (string, string, error) { +// 返回 (访问地址, 对象 key, 感知哈希, error)。访问地址可直接写入数据库 +// (七牛为完整 https URL,本地为相对 /uploads 路径);感知哈希基于原始字节算一次,供入库时近似去重。 +func (s *IngestService) downloadOne(ctx context.Context, u string, idx int, prefix string) (string, string, uint64, error) { resp, err := s.httpClient.Get(u) if err != nil { - return "", "", err + return "", "", 0, err } defer resp.Body.Close() if resp.StatusCode != http.StatusOK { - return "", "", fmt.Errorf("status %d", resp.StatusCode) + return "", "", 0, fmt.Errorf("status %d", resp.StatusCode) } data, err := io.ReadAll(resp.Body) if err != nil { - return "", "", err + return "", "", 0, err } ext := path.Ext(u) if ext == "" || len(ext) > 5 { @@ -467,16 +498,74 @@ func (s *IngestService) downloadOne(ctx context.Context, u string, idx int, pref key := fmt.Sprintf("%s/%s%s", prefix, hash, ext) ct := resp.Header.Get("Content-Type") + // 感知哈希:同一份字节在此算一次(decode 失败时 ph=0,比对时跳过,不阻断入库)。 + ph := phash.Of(data) + if ph == 0 { + log.Printf("[warn] ingest pHash 未计算(解码失败或非 jpeg/png/gif 格式) url=%s", u) + } + // 主上传器(七牛) if url, err := s.uploader.Upload(ctx, key, data, ct); err == nil { - return url, key, nil + return url, key, ph, nil } else if s.local != nil { // 兜底本地,避免图片完全丢失 if lurl, lerr := s.local.Upload(ctx, key, data, ct); lerr == nil { - return lurl, key, nil + return lurl, key, ph, nil } else { - return "", "", err + return "", "", 0, err } } - return "", "", err + return "", "", 0, err +} + +// tagDuplicates 对一批图片(phashList 与 apply 下标对齐)做全局近似去重标记: +// 与「已晋升图片」的 phash 库比汉明距离(≤ phash.DefaultThreshold 即近似重复), +// 命中则调用 apply(i, dupOf) 由调用方把第 i 张图标 IsDuplicate=1、DupOf=命中图 uid。 +// 不丢弃任何图、仅留痕(交后台人工裁决),契合「只拦新增、全局跨所有图、标记不删」。 +// 比对库拉取失败仅告警并跳过(不阻断入库);phash=0 的图不比对、不误杀。 +func (s *IngestService) tagDuplicates(ctx context.Context, phashList []uint64, apply func(i int, dupOf string)) { + uids, err := s.matchDuplicates(ctx, phashList) + if err != nil { + log.Printf("[warn] ingest 近似去重比对失败 err=%v(跳过标记,不阻断入库)", err) + return + } + for i, u := range uids { + if u != "" { + apply(i, u) + } + } +} + +// matchDuplicates 返回与 phashList 等长的 dup_of uid 切片(""=未命中近似重复)。 +// 对每张新图,在已晋升图片库里找汉明距离 ≤ 阈值的命中,取距离最近者, +// 按其 kind(runway/street)编码成对应 hashid 类型作为 dup_of。 +func (s *IngestService) matchDuplicates(ctx context.Context, phashList []uint64) ([]string, error) { + refs, err := s.repo.ListImagePHashes(ctx) + if err != nil { + return nil, err + } + out := make([]string, len(phashList)) + for i, ph := range phashList { + if ph == 0 { + continue // 未计算 phash 的图不比对(避免误杀) + } + bestUID := "" + bestDist := phash.DefaultThreshold + 1 + for _, ref := range refs { + if ref.Phash == 0 { + continue + } + d := phash.Hamming(ph, ref.Phash) + if d <= phash.DefaultThreshold && d < bestDist { + var typ byte = hashid.TypeRunwayImage + if ref.Kind == "street" { + typ = hashid.TypeSnapImage + } + bestUID = hashid.EncodeWithType(ref.ID, typ) + bestDist = d + } + } + out[i] = bestUID + } + return out, nil } diff --git a/scripts/pgvector/setup_pg_docker.sh b/scripts/pgvector/setup_pg_docker.sh new file mode 100644 index 0000000..2104a46 --- /dev/null +++ b/scripts/pgvector/setup_pg_docker.sh @@ -0,0 +1,161 @@ +#!/usr/bin/env bash +# +# setup_pg_docker.sh — 用 Docker 起 PostgreSQL 16 + pgvector 本地开发环境 +# +# 用法(在 WSL 终端里,不要用 Windows 的 CMD/PowerShell): +# sudo bash setup_pg_docker.sh +# +# 行为: +# 1) 先卸载上次用 apt 装失败的 PostgreSQL 残留(清理 pgdg 源 / 坏 key / 包),幂等。 +# 2) 确保 docker 可用(缺失则尝试 apt 装 docker.io 并起守护进程;仍不行则给出 Docker Desktop 指引)。 +# 3) docker compose up -d 起 pgvector 容器(幂等:已存在则不变)。 +# 4) 等就绪后验证 version() 与 vector 扩展版本。 +# +# 仅用于本地开发 / 迁移验证。生产请改强密码并单独评审配置。 +# +set -euo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +COMPOSE_FILE="$SCRIPT_DIR/docker-compose.yml" +DB_USER=fashion +DB_PASS="fashion_dev_2026" +DB_NAME=fashion + +log() { echo -e "\033[32m==>\033[0m $*"; } +warn() { echo -e "\033[33m⚠️ \033[0m $*"; } + +if [ "$(id -u)" -ne 0 ]; then + echo "请以 root 运行: sudo bash $0"; exit 1 +fi + +# ---------- 1) 卸载上次 apt 版 PostgreSQL 的残留 ---------- +if command -v psql >/dev/null 2>&1 || [ -f /etc/apt/sources.list.d/pgdg.list ]; then + log "卸载上次 apt 版 PostgreSQL 残留(清理失败残留,幂等)" + pg_ctlcluster "$(ls /etc/postgresql 2>/dev/null | head -1)" main stop 2>/dev/null \ + || service postgresql stop 2>/dev/null \ + || true + apt-get remove --purge -y 'postgresql-*' 2>/dev/null || true + rm -f /etc/apt/sources.list.d/pgdg.list + rm -f /usr/share/postgresql-common/pgdg/apt.postgresql.org.asc \ + /usr/share/postgresql-common/pgdg/apt.postgresql.org.gpg + apt-get autoremove -y 2>/dev/null || true + log "apt 版 PostgreSQL 残留已清理" +else + log "未发现 apt 版 PostgreSQL,跳过卸载" +fi + +# ---------- 2) 确保 docker 可用 ---------- +if command -v docker >/dev/null 2>&1 && docker info >/dev/null 2>&1; then + log "docker 可用" +else + if command -v docker >/dev/null 2>&1; then + warn "docker 已安装但守护进程未起,尝试启动" + service docker start 2>/dev/null || (dockerd >/var/log/docker.log 2>&1 &) || true + sleep 3 + fi + if ! command -v docker >/dev/null 2>&1 || ! docker info >/dev/null 2>&1; then + warn "docker 不可用,尝试 apt 安装 docker.io" + apt-get update -y + if apt-get install -y docker.io; then + service docker start 2>/dev/null || (dockerd >/var/log/docker.log 2>&1 &) || true + sleep 3 + fi + fi + if ! command -v docker >/dev/null 2>&1 || ! docker info >/dev/null 2>&1; then + echo "❌ 仍无法使用 docker。请二选一:" + echo " A) 装 Docker Desktop (Windows),安装时勾选 'Use WSL 2'," + echo " 装好后在其 Settings > Resources > WSL Integration 里启用你的 WSL 发行版,重启 WSL 后重试;" + echo " B) 或在本 WSL 内: sudo apt-get install -y docker.io && sudo service docker start" + exit 1 + fi +fi + +# 选 docker compose 命令(v2 插件优先,回退 v1) +if docker compose version >/dev/null 2>&1; then + DC=(docker compose) +elif command -v docker-compose >/dev/null 2>&1; then + DC=(docker-compose) +else + echo "❌ 未找到 docker compose / docker-compose,请安装 Docker Compose 后重试"; exit 1 +fi + +# ---------- 2.5) 配置国内镜像源加速(避免直连 Docker Hub 慢) ---------- +MIRRORS='["https://docker.m.daocloud.io","https://hub-mirror.c.163.com","https://mirror.baidubce.com"]' +if docker context ls 2>/dev/null | grep -qw "desktop-linux" || [[ "${DOCKER_HOST:-}" == *"npipe"* ]]; then + # Docker Desktop(守护进程在 Windows 侧):WSL 内的 daemon.json 不生效,必须在 GUI 配 + warn "检测到 Docker Desktop:WSL 内的 daemon.json 不影响 Windows 侧守护进程。" + echo " 请在 Windows 的 Docker Desktop > Settings > Docker Engine 的 JSON 中加入:" + echo " \"registry-mirrors\": $MIRRORS" + echo " 然后点 Apply & Restart;重启后再直接跑: docker compose -f $COMPOSE_FILE up -d" + echo " (本脚本不再自动拉起,避免与 Docker Desktop 守护进程冲突)" + exit 1 +elif ! docker info 2>/dev/null | grep -q "docker.m.daocloud.io"; then + # docker.io 跑在 WSL 内:直接写 daemon.json 并重启守护进程 + if [ ! -f /etc/docker/daemon.json ] || ! grep -q registry-mirrors /etc/docker/daemon.json; then + log "配置国内镜像源加速(写入 /etc/docker/daemon.json 并重启 docker 守护进程)" + cat > /etc/docker/daemon.json </dev/null || systemctl restart docker 2>/dev/null \ + || (dockerd --registry-mirror=https://docker.m.daocloud.io >/var/log/docker.log 2>&1 &) || true + sleep 3 + else + log "daemon.json 已含 registry-mirrors,跳过" + fi + # 等守护进程重新就绪 + for i in $(seq 1 15); do + docker info >/dev/null 2>&1 && break + sleep 1 + done +fi + +# ---------- 3) 起容器(幂等) ---------- +log "启动 PostgreSQL + pgvector 容器(docker compose up -d)" +"${DC[@]}" -f "$COMPOSE_FILE" up -d + +# ---------- 4) 等就绪 ---------- +log "等待数据库就绪(最多 30s)" +ready=0 +for i in $(seq 1 30); do + if docker exec pgvector pg_isready -U "$DB_USER" -d "$DB_NAME" >/dev/null 2>&1; then + ready=1; break + fi + sleep 1 +done +if [ "$ready" -ne 1 ]; then + echo "❌ 数据库未就绪,查看日志: docker logs pgvector"; exit 1 +fi + +# ---------- 5) 验证 ---------- +log "验证安装" +docker exec -i pgvector psql -U "$DB_USER" -d "$DB_NAME" \ + -c "SELECT version();" \ + -c "SELECT extname, extversion FROM pg_extension WHERE extname='vector';" + +cat < '[...]'::vector) AS sim + FROM image_embedding + WHERE kind = 0 + ORDER BY embedding <=> '[...]'::vector + LIMIT 20; +EOF