feat(creator): 监控账号画像采集入库——粉丝/关注/获赞/作品数与作者介绍

抖音 profile/other 响应本就携带画像数据但被丢弃,同步链路也无回填。
- migration 1046: creator_competitor 增加 signature、total_favorited 列及非负约束
- douyin 解析层: TargetProfile 增加签名与计数,total_favorited 兼容字符串大数
- store: UpdateCompetitorProfile 改为画像回填(nil 计数保留原值)
- 同步链路: 作品采集成功后拉取 profile 回填,失败仅记日志不阻断
This commit is contained in:
2026-09-30 16:37:56 +08:00
parent a9358f5c96
commit 3e809c4fc7
11 changed files with 176 additions and 25 deletions
+23
View File
@@ -1285,6 +1285,24 @@ func previewDouyinCompetitor(ctx context.Context, store *creator.Store, phaseASt
}, nil
}
// refreshCompetitorProfile 拉取抖音公开主页画像并回填监控账号画像字段。
// profile/other 响应自带粉丝/关注/获赞/作品数与签名,此前未回填导致详情页画像恒为空。
func refreshCompetitorProfile(ctx context.Context, store *creator.Store, browser creatorGatewayBrowser, competitor creator.Competitor) error {
profile, err := (douyin.CreatorCollector{Browser: browser}).ResolveTarget(ctx, competitor.PlatformAccountKey)
if err != nil {
return err
}
_, err = store.UpdateCompetitorProfile(ctx, competitor.ID, creator.CompetitorProfileUpdate{
AvatarURL: profile.AvatarURL,
Signature: profile.Signature,
FollowerCount: profile.FollowerCount,
FollowingCount: profile.FollowingCount,
TotalFavorited: profile.TotalFavorited,
AwemeCount: profile.AwemeCount,
})
return err
}
func syncCreatorCompetitor(ctx context.Context, store *creator.Store, hubStore *hub.Store, competitorID string) (creator.CollectionReport, error) {
return syncCreatorCompetitorWithClaim(ctx, store, hubStore, competitorID, true)
}
@@ -1349,6 +1367,11 @@ func syncCreatorCompetitorWithClaim(ctx context.Context, store *creator.Store, h
}
report, collectErr := store.CollectSource(ctx, competitor.Platform, creator.SourceCompetitor, competitor.ID, collector, collectionNow)
if collectErr == nil {
// 采集成功后回填监控账号画像(粉丝/关注/获赞/作品数/作者介绍);
// profile 拉取失败只记日志,不影响本次作品同步结果。
if profileErr := refreshCompetitorProfile(ctx, store, creatorGatewayBrowser{gateway: lease.gateway, environment: lease.environment}, competitor); profileErr != nil {
logrus.WithError(profileErr).WithField("competitor_id", competitorID).Warn("creator competitor profile refresh failed")
}
// 采集成功后把缺封面的作品封面缓存到本地(失败只记日志,不影响本次同步结果)。
if coverErr := cacheCompetitorWorkCovers(ctx, store, creatorGatewayBrowser{gateway: lease.gateway, environment: lease.environment}, competitor.ID); coverErr != nil {
logrus.WithError(coverErr).WithField("competitor_id", competitorID).Warn("creator competitor cover cache failed")
+35 -14
View File
@@ -72,20 +72,21 @@ func (s *Store) CreateCompetitor(ctx context.Context, input CompetitorInput) (Co
// competitorColumns 是 creator_competitor 的完整列清单(含画像属性),各查询共用。
const competitorColumns = `competitor_id, platform, platform_account_key, unique_id, nickname, avatar_url, homepage_url, tags,
enabled, follower_count, following_count, aweme_count, sync_status, sync_cursor, sync_error, sync_lease_until,
signature, enabled, follower_count, following_count, total_favorited, aweme_count, sync_status, sync_cursor, sync_error, sync_lease_until,
last_sync_at, next_sync_at, created_at, updated_at`
type competitorScan struct {
competitor Competitor
tags pgtype.FlatArray[string]
leaseUntil, lastSync, nextSync sql.NullTime
follower, following, aweme sql.NullInt64
competitor Competitor
tags pgtype.FlatArray[string]
leaseUntil, lastSync, nextSync sql.NullTime
follower, following, favorited sql.NullInt64
aweme sql.NullInt64
}
func competitorScanDestinations(scan *competitorScan) []any {
return []any{&scan.competitor.ID, &scan.competitor.Platform, &scan.competitor.PlatformAccountKey, &scan.competitor.UniqueID, &scan.competitor.Nickname,
&scan.competitor.AvatarURL, &scan.competitor.HomepageURL, pgtype.NewMap().SQLScanner(&scan.tags), &scan.competitor.Enabled,
&scan.follower, &scan.following, &scan.aweme, &scan.competitor.SyncStatus, &scan.competitor.SyncCursor,
&scan.competitor.AvatarURL, &scan.competitor.HomepageURL, pgtype.NewMap().SQLScanner(&scan.tags), &scan.competitor.Signature, &scan.competitor.Enabled,
&scan.follower, &scan.following, &scan.favorited, &scan.aweme, &scan.competitor.SyncStatus, &scan.competitor.SyncCursor,
&scan.competitor.SyncError, &scan.leaseUntil, &scan.lastSync, &scan.nextSync, &scan.competitor.CreatedAt, &scan.competitor.UpdatedAt}
}
@@ -93,6 +94,7 @@ func (scan *competitorScan) materialize() Competitor {
scan.competitor.Tags = []string(scan.tags)
scan.competitor.FollowerCount = nullableInt64(scan.follower)
scan.competitor.FollowingCount = nullableInt64(scan.following)
scan.competitor.TotalFavorited = nullableInt64(scan.favorited)
scan.competitor.AwemeCount = nullableInt64(scan.aweme)
scan.competitor.SyncLeaseUntil = nullableTime(scan.leaseUntil)
scan.competitor.LastSyncAt = nullableTime(scan.lastSync)
@@ -166,18 +168,37 @@ func scanCompetitorView(scanner interface{ Scan(...any) error }) (CompetitorView
return result, nil
}
// UpdateCompetitorProfile 更新监控账号画像属性(头像与粉丝/关注/作品计数)。
// 当前仅提供数据通道(字段定义与接口返回),采集侧回填由后续接入。
func (s *Store) UpdateCompetitorProfile(ctx context.Context, id, avatarURL string, followerCount, followingCount, awemeCount int64) (Competitor, error) {
// CompetitorProfileUpdate 监控账号画像回填值;nil 计数表示本次未采集到,保留原值。
type CompetitorProfileUpdate struct {
AvatarURL string
Signature string
FollowerCount *int64
FollowingCount *int64
TotalFavorited *int64
AwemeCount *int64
}
// UpdateCompetitorProfile 回填监控账号画像(头像/签名/粉丝/关注/获赞/作品计数)。
func (s *Store) UpdateCompetitorProfile(ctx context.Context, id string, update CompetitorProfileUpdate) (Competitor, error) {
id = strings.TrimSpace(id)
avatarURL = strings.TrimSpace(avatarURL)
if id == "" || utf8.RuneCountInString(avatarURL) > 2000 || followerCount < 0 || followingCount < 0 || awemeCount < 0 {
avatarURL := strings.TrimSpace(update.AvatarURL)
signature := strings.TrimSpace(update.Signature)
if id == "" || utf8.RuneCountInString(avatarURL) > 2000 || utf8.RuneCountInString(signature) > 2000 {
return Competitor{}, ErrInvalid
}
for _, value := range []*int64{update.FollowerCount, update.FollowingCount, update.TotalFavorited, update.AwemeCount} {
if value != nil && *value < 0 {
return Competitor{}, ErrInvalid
}
}
result, err := s.db.ExecContext(ctx, `
UPDATE creator_competitor
SET avatar_url = $2, follower_count = $3, following_count = $4, aweme_count = $5, updated_at = now()
WHERE competitor_id = $1`, id, avatarURL, followerCount, followingCount, awemeCount)
SET avatar_url = CASE WHEN $2 = '' THEN avatar_url ELSE $2 END,
signature = CASE WHEN $3 = '' THEN signature ELSE $3 END,
follower_count = COALESCE($4, follower_count), following_count = COALESCE($5, following_count),
total_favorited = COALESCE($6, total_favorited), aweme_count = COALESCE($7, aweme_count),
updated_at = now()
WHERE competitor_id = $1`, id, avatarURL, signature, update.FollowerCount, update.FollowingCount, update.TotalFavorited, update.AwemeCount)
if err != nil {
return Competitor{}, databaseError(err)
}
+45 -1
View File
@@ -171,6 +171,49 @@ func TestCompetitorShareJobLifecycle(t *testing.T) {
}
// 强制同步抢锁失败的分类:租约被占 → ErrSyncInProgress;竞品停用 → ErrSyncDisabled。
func TestCompetitorProfileUpdateRoundTrip(t *testing.T) {
store, _, ctx := openCreatorIntegrationStore(t)
stamp := fmt.Sprintf("%d", time.Now().UnixNano())
competitor, err := store.UpsertCompetitor(ctx, CompetitorInput{
Platform: PlatformDouyin,
PlatformAccountKey: "sec_uid_profile_" + stamp,
Nickname: "Profile Target",
HomepageURL: "https://www.douyin.com/user/sec_uid_profile_" + stamp,
})
if err != nil {
t.Fatal(err)
}
follower, following, favorited, aweme := int64(123), int64(45), int64(98765), int64(7)
updated, err := store.UpdateCompetitorProfile(ctx, competitor.ID, CompetitorProfileUpdate{
AvatarURL: "https://example.invalid/avatar",
Signature: "第一行介绍\n第二行介绍",
FollowerCount: &follower,
FollowingCount: &following,
TotalFavorited: &favorited,
AwemeCount: &aweme,
})
if err != nil {
t.Fatal(err)
}
if updated.Signature != "第一行介绍\n第二行介绍" || updated.FollowerCount == nil || *updated.FollowerCount != 123 ||
updated.FollowingCount == nil || *updated.FollowingCount != 45 || updated.TotalFavorited == nil || *updated.TotalFavorited != 98765 ||
updated.AwemeCount == nil || *updated.AwemeCount != 7 {
t.Fatalf("profile update: %+v", updated)
}
// nil 计数保留原值,不置 NULL。
reread, err := store.UpdateCompetitorProfile(ctx, competitor.ID, CompetitorProfileUpdate{Signature: "新签名"})
if err != nil || reread.Signature != "新签名" || reread.FollowerCount == nil || *reread.FollowerCount != 123 || reread.TotalFavorited == nil || *reread.TotalFavorited != 98765 {
t.Fatalf("profile partial update: %+v err=%v", reread, err)
}
negative := int64(-1)
if _, err := store.UpdateCompetitorProfile(ctx, competitor.ID, CompetitorProfileUpdate{FollowerCount: &negative}); err == nil {
t.Fatal("negative follower count accepted")
}
if _, err := store.UpdateCompetitorProfile(ctx, "missing_"+stamp, CompetitorProfileUpdate{}); err == nil {
t.Fatal("missing competitor accepted")
}
}
func TestCompetitorSyncClaimFailureClassification(t *testing.T) {
store, _, ctx := openCreatorIntegrationStore(t)
stamp := fmt.Sprintf("%d", time.Now().UnixNano())
@@ -378,7 +421,8 @@ func TestCreatorPostgresContentAndWorkflow(t *testing.T) {
if views[0].WorkCount != 1 || views[0].LatestPublishedAt == nil {
t.Fatalf("competitor view aggregates: work_count=%d latest=%v", views[0].WorkCount, views[0].LatestPublishedAt)
}
if _, err := store.UpdateCompetitorProfile(ctx, competitor.ID, "https://example.invalid/avatar.jpg", 1200, 300, 88); err != nil {
follower, following, aweme := int64(1200), int64(300), int64(88)
if _, err := store.UpdateCompetitorProfile(ctx, competitor.ID, CompetitorProfileUpdate{AvatarURL: "https://example.invalid/avatar.jpg", FollowerCount: &follower, FollowingCount: &following, AwemeCount: &aweme}); err != nil {
t.Fatal(err)
}
views, err = store.ListCompetitorsWithProfile(ctx, "")
+2
View File
@@ -103,9 +103,11 @@ type Competitor struct {
AvatarURL string `json:"avatar_url"`
HomepageURL string `json:"homepage_url"`
Tags []string `json:"tags"`
Signature string `json:"signature"`
Enabled bool `json:"enabled"`
FollowerCount *int64 `json:"follower_count,omitempty"`
FollowingCount *int64 `json:"following_count,omitempty"`
TotalFavorited *int64 `json:"total_favorited,omitempty"`
AwemeCount *int64 `json:"aweme_count,omitempty"`
SyncStatus string `json:"sync_status"`
SyncCursor string `json:"sync_cursor,omitempty"`
@@ -71,5 +71,5 @@ func TestMigration043ConsolidatedSchemaShape(t *testing.T) {
assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema()
AND table_name = 'audit_event' AND column_name IN ('confirmation_id','confirmation_version','attempt_id','task_id','runtime_instance_id')`, 0)
// 统一登记表:1-38(除 36)、1017-1045、43 全部登记
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 52)
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 53)
}
+2 -2
View File
@@ -30,7 +30,7 @@ func TestUnifiedAccountMigration(t *testing.T) {
t.Fatal(err)
}
defer db.Close()
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 52)
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 53)
assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.tables WHERE table_schema = current_schema()
AND table_name IN ('social_account', 'browser_env', 'network_exit', 'environment_binding')`, 3)
assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.tables WHERE table_schema = current_schema() AND table_name = 'browser_image'`, 0)
@@ -42,7 +42,7 @@ func TestUnifiedAccountMigration(t *testing.T) {
store = openFullyMigratedHub(t, ctx, testURL)
store.Close()
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 52)
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 53)
})
t.Run("legacy migration 013 without account secrets is repaired forward", func(t *testing.T) {
@@ -0,0 +1,13 @@
-- 监控账号画像补全:作者介绍与获赞总数(同步时经抖音 profile/other 接口回填)。
ALTER TABLE creator_competitor
ADD COLUMN IF NOT EXISTS signature text NOT NULL DEFAULT '',
ADD COLUMN IF NOT EXISTS total_favorited bigint;
ALTER TABLE creator_competitor
ADD CONSTRAINT creator_competitor_counts_non_negative CHECK (
(follower_count IS NULL OR follower_count >= 0)
AND (following_count IS NULL OR following_count >= 0)
AND (total_favorited IS NULL OR total_favorited >= 0)
AND (aweme_count IS NULL OR aweme_count >= 0)
);
+4 -1
View File
@@ -176,6 +176,9 @@ var migration1044 string
//go:embed migrations/1045_drop_authorization_status.sql
var migration1045 string
//go:embed migrations/1046_competitor_profile_signature.sql
var migration1046 string
var (
ErrConflict = errors.New("resource conflicts with existing state")
ErrInvalid = errors.New("invalid environment input")
@@ -315,7 +318,7 @@ func (s *Store) migrate(ctx context.Context) error {
{1023, migration1023}, {1024, migration1024}, {1025, migration1025}, {1026, migration1026}, {1027, migration1027}, {1028, migration1028},
{1029, migration1029}, {1030, migration1030}, {1031, migration1031}, {1032, migration1032}, {1033, migration1033}, {1034, migration1034},
{1035, migration1035}, {1036, migration1036}, {1037, migration1037}, {1038, migration1038}, {1039, migration1039}, {1040, migration1040},
{1041, migration1041}, {1042, migration1042}, {1043, migration1043}, {1044, migration1044}, {1045, migration1045},
{1041, migration1041}, {1042, migration1042}, {1043, migration1043}, {1044, migration1044}, {1045, migration1045}, {1046, migration1046},
{43, migration043}, {44, migration044}} {
var applied bool
if err := tx.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM schema_migration WHERE version = $1)`, migration.version).Scan(&applied); err != nil {
+9
View File
@@ -316,6 +316,7 @@ type identityEnvelope struct {
AvatarThumb *struct {
URLList []string `json:"url_list"`
} `json:"avatar_thumb"`
Signature string `json:"signature"`
FollowerCount *int64 `json:"follower_count"`
FollowingCount *int64 `json:"following_count"`
AwemeCount *int64 `json:"aweme_count"`
@@ -327,6 +328,14 @@ type identityEnvelope struct {
// flexibleInt 兼容抖音返回的数字与字符串两种计数形态(total_favorited 实测为大数字符串)。
type flexibleInt int64
func flexibleIntPtr(value *flexibleInt) *int64 {
if value == nil {
return nil
}
converted := int64(*value)
return &converted
}
func (value *flexibleInt) UnmarshalJSON(data []byte) error {
text := strings.TrimSpace(string(data))
if text == "null" {
+15 -6
View File
@@ -27,11 +27,16 @@ type CreatorCollector struct {
}
type TargetProfile struct {
UID string
SecUID string
UniqueID string
Nickname string
AvatarURL string
UID string
SecUID string
UniqueID string
Nickname string
AvatarURL string
Signature string
FollowerCount *int64
FollowingCount *int64
TotalFavorited *int64
AwemeCount *int64
}
// WorkStatistics 是作品详情接口返回的指标快照,用于增长曲线时序采集。
@@ -174,7 +179,11 @@ func (c CreatorCollector) ResolveTarget(ctx context.Context, expectedKey string)
if identity.User.AvatarThumb != nil && len(identity.User.AvatarThumb.URLList) > 0 {
avatarURL = identity.User.AvatarThumb.URLList[0]
}
return TargetProfile{UID: identity.User.UID, SecUID: identity.User.SecUID, UniqueID: identity.User.UniqueID, Nickname: identity.User.Nickname, AvatarURL: avatarURL}, nil
return TargetProfile{
UID: identity.User.UID, SecUID: identity.User.SecUID, UniqueID: identity.User.UniqueID, Nickname: identity.User.Nickname, AvatarURL: avatarURL,
Signature: identity.User.Signature, FollowerCount: identity.User.FollowerCount, FollowingCount: identity.User.FollowingCount,
TotalFavorited: flexibleIntPtr(identity.User.TotalFavorited), AwemeCount: identity.User.AwemeCount,
}, nil
}
func (c CreatorCollector) ResolveWork(ctx context.Context, workKey string) (TargetProfile, error) {
@@ -37,6 +37,33 @@ func TestResolveTargetReturnsVerifiedProfileFields(t *testing.T) {
}
}
func TestResolveTargetExtractsPublicProfileMetrics(t *testing.T) {
browser := &collectorBrowser{response: Response{Status: 200, Body: []byte(`{"status_code":0,"user":{"uid":"2328120603967913","sec_uid":"MS4wLjABAAAA9f_a7k0bzVizLYXlpC7R61EIaqJ8Ordug7yp7AB8fGKuuF8Fzqk5_DM-eutXnPIK","unique_id":"96332518739","nickname":"目标","signature":"第一行介绍\n第二行介绍","follower_count":123,"following_count":45,"aweme_count":7,"total_favorited":"98765"}}`)}}
profile, err := (CreatorCollector{Browser: browser}).ResolveTarget(context.Background(), "MS4wLjABAAAA9f_a7k0bzVizLYXlpC7R61EIaqJ8Ordug7yp7AB8fGKuuF8Fzqk5_DM-eutXnPIK")
if err != nil {
t.Fatal(err)
}
if profile.Signature != "第一行介绍\n第二行介绍" {
t.Fatalf("signature: %q", profile.Signature)
}
if profile.FollowerCount == nil || *profile.FollowerCount != 123 || profile.FollowingCount == nil || *profile.FollowingCount != 45 ||
profile.AwemeCount == nil || *profile.AwemeCount != 7 {
t.Fatalf("counts: %+v", profile)
}
// total_favorited 实测为大数字符串,必须经 flexibleInt 兼容解析。
if profile.TotalFavorited == nil || *profile.TotalFavorited != 98765 {
t.Fatalf("total_favorited: %+v", profile.TotalFavorited)
}
}
func TestResolveTargetToleratesMissingPublicMetrics(t *testing.T) {
browser := &collectorBrowser{response: Response{Status: 200, Body: []byte(`{"status_code":0,"user":{"uid":"2328120603967913","sec_uid":"MS4wLjABAAAA9f_a7k0bzVizLYXlpC7R61EIaqJ8Ordug7yp7AB8fGKuuF8Fzqk5_DM-eutXnPIK","nickname":"目标"}}`)}}
profile, err := (CreatorCollector{Browser: browser}).ResolveTarget(context.Background(), "MS4wLjABAAAA9f_a7k0bzVizLYXlpC7R61EIaqJ8Ordug7yp7AB8fGKuuF8Fzqk5_DM-eutXnPIK")
if err != nil || profile.Signature != "" || profile.FollowerCount != nil || profile.TotalFavorited != nil {
t.Fatalf("target profile without metrics: %+v err=%v", profile, err)
}
}
func TestCanonicalTargetSecUIDRejectsUniqueIDLookup(t *testing.T) {
browser := &collectorBrowser{}
_, err := (CreatorCollector{Browser: browser}).CanonicalTargetSecUID(context.Background(), "creator_handle")