From 923fd648792486415c5c789fd3729df2e1b59e5e Mon Sep 17 00:00:00 2001 From: Rogee Date: Wed, 30 Sep 2026 12:50:10 +0800 Subject: [PATCH] =?UTF-8?q?feat(creator):=20=E8=B4=A6=E5=8F=B7=E7=94=BB?= =?UTF-8?q?=E5=83=8F=E8=A1=A5=E4=BA=92=E5=85=B3=E6=8C=87=E6=A0=87=E5=B9=B6?= =?UTF-8?q?=E8=AE=A9=E7=9B=91=E6=8E=A7=E5=88=97=E8=A1=A8=E6=90=BA=E5=B8=A6?= =?UTF-8?q?=E6=9C=80=E6=96=B0=E5=BF=AB=E7=85=A7=E4=B8=8E=E8=AF=84=E8=AE=BA?= =?UTF-8?q?=E8=81=9A=E5=90=88?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - migration 1044:creator_account_metric 加 friend_count(self profile 接口 user.friend_count,即互关数) - douyin connector 解析 friend_count(负值校验同其他计数),快照落库透传 - monitor-views 视图新增:每账号最新非空画像计数(粉丝/关注/获赞/作品/互关,按采集时间回填缺失维度) - 评论口径:自有作品最近指标 comments_count 求和(comment_total),无账号级评论采集源 - 详情页指标时序 AccountMetricPoint 带回 friend_count --- internal/controlplane/api/creator.go | 2 +- internal/creator/accounts.go | 79 +++++++++++++++++-- internal/creator/content.go | 19 +++-- internal/creator/integration_test.go | 44 +++++++++++ internal/creator/models.go | 2 + .../environment/migration043_probe_test.go | 4 +- internal/environment/migration_test.go | 4 +- ...44_creator_account_metric_friend_count.sql | 2 + internal/environment/store.go | 5 +- internal/platform/douyin/connector.go | 5 +- 10 files changed, 146 insertions(+), 20 deletions(-) create mode 100644 internal/environment/migrations/1044_creator_account_metric_friend_count.sql diff --git a/internal/controlplane/api/creator.go b/internal/controlplane/api/creator.go index d53431f..6a52203 100644 --- a/internal/controlplane/api/creator.go +++ b/internal/controlplane/api/creator.go @@ -1663,7 +1663,7 @@ func recordCreatorAccountMetricSnapshot(ctx context.Context, store *creator.Stor return store.RecordAccountMetric(ctx, creator.AccountMetricInput{ AccountID: accountID, CollectedAt: time.Now().UTC(), FollowerCount: profile.FollowerCount, FollowingCount: profile.FollowingCount, - TotalFavorited: profile.TotalFavorited, AwemeCount: profile.AwemeCount, + TotalFavorited: profile.TotalFavorited, AwemeCount: profile.AwemeCount, FriendCount: profile.FriendCount, }) } diff --git a/internal/creator/accounts.go b/internal/creator/accounts.go index d691c30..360d365 100644 --- a/internal/creator/accounts.go +++ b/internal/creator/accounts.go @@ -105,11 +105,17 @@ func (s *Store) ListAccountProfiles(ctx context.Context) ([]AccountProfile, erro return profiles, rows.Err() } -// AccountMonitorView 自有账号监控列表视图:账号画像 + 作品聚合统计(作品数/最近发布)。 +// AccountMonitorView 自有账号监控列表视图:账号画像 + 作品聚合统计(作品数/最近发布)+ 最新画像指标快照 + 自有作品评论累计。 type AccountMonitorView struct { AccountProfile WorkCount int64 `json:"work_count"` LatestPublishedAt *time.Time `json:"latest_published_at,omitempty"` + FollowerCount *int64 `json:"follower_count,omitempty"` + FollowingCount *int64 `json:"following_count,omitempty"` + TotalFavorited *int64 `json:"total_favorited,omitempty"` + AwemeCount *int64 `json:"aweme_count,omitempty"` + FriendCount *int64 `json:"friend_count,omitempty"` + CommentTotal int64 `json:"comment_total"` } // ListAccountMonitorViews 自有账号监控列表:画像 + owned 作品聚合,形态对齐竞品的 ListCompetitorsWithProfile。 @@ -122,11 +128,19 @@ func (s *Store) ListAccountMonitorViews(ctx context.Context) ([]AccountMonitorVi if err != nil { return nil, err } + latest, err := s.latestAccountMetrics(ctx) + if err != nil { + return nil, err + } views := make([]AccountMonitorView, 0, len(profiles)) for _, profile := range profiles { view := AccountMonitorView{AccountProfile: profile} if stat, exists := stats[profile.ID]; exists { - view.WorkCount, view.LatestPublishedAt = stat.WorkCount, stat.LatestPublishedAt + view.WorkCount, view.LatestPublishedAt, view.CommentTotal = stat.WorkCount, stat.LatestPublishedAt, stat.CommentTotal + } + if metric, exists := latest[profile.ID]; exists { + view.FollowerCount, view.FollowingCount, view.TotalFavorited, view.AwemeCount, view.FriendCount = + metric.FollowerCount, metric.FollowingCount, metric.TotalFavorited, metric.AwemeCount, metric.FriendCount } views = append(views, view) } @@ -136,12 +150,67 @@ func (s *Store) ListAccountMonitorViews(ctx context.Context) ([]AccountMonitorVi type ownedWorkStat struct { WorkCount int64 LatestPublishedAt *time.Time + CommentTotal int64 } -// ownedWorkStats 按 source_id 聚合 owned 作品数与最近发布时间。 +// latestAccountMetricRow 最新快照内一行非空计数(无对应快照时不返回行)。 +type latestAccountMetricRow struct { + FollowerCount *int64 + FollowingCount *int64 + TotalFavorited *int64 + AwemeCount *int64 + FriendCount *int64 +} + +// latestAccountMetrics 每账号取最新非空画像计数:按采集时间倒序回填,缺失维度保留上一轮的值(详情页同口径)。 +func (s *Store) latestAccountMetrics(ctx context.Context) (map[string]latestAccountMetricRow, error) { + rows, err := s.db.QueryContext(ctx, ` + SELECT account.account_id, metric.follower_count, metric.following_count, metric.total_favorited, metric.aweme_count, metric.friend_count + FROM creator_account_metric metric + JOIN social_account account ON account.id = metric.account_id + ORDER BY account.account_id, metric.collected_at`) + if err != nil { + return nil, databaseError(err) + } + defer rows.Close() + latest := make(map[string]latestAccountMetricRow) + for rows.Next() { + var accountID string + var follower, following, favorited, aweme, friend sql.NullInt64 + if err := rows.Scan(&accountID, &follower, &following, &favorited, &aweme, &friend); err != nil { + return nil, err + } + current := latest[accountID] + if follower.Valid { + value := follower.Int64 + current.FollowerCount = &value + } + if following.Valid { + value := following.Int64 + current.FollowingCount = &value + } + if favorited.Valid { + value := favorited.Int64 + current.TotalFavorited = &value + } + if aweme.Valid { + value := aweme.Int64 + current.AwemeCount = &value + } + if friend.Valid { + value := friend.Int64 + current.FriendCount = &value + } + latest[accountID] = current + } + return latest, rows.Err() +} + +// AccountCollectionStatus 自有账号采集状态(checkpoint 形态,对齐竞品的 sync_status 展示语义)。 func (s *Store) ownedWorkStats(ctx context.Context) (map[string]ownedWorkStat, error) { rows, err := s.db.QueryContext(ctx, ` - SELECT w.source_id, COUNT(*) AS work_count, MAX(w.published_at) AS latest_published_at + SELECT w.source_id, COUNT(*) AS work_count, MAX(w.published_at) AS latest_published_at, + COALESCE(SUM(w.comments_count), 0) AS comment_total FROM creator_work w WHERE w.source_type = 'owned' GROUP BY w.source_id`) @@ -154,7 +223,7 @@ func (s *Store) ownedWorkStats(ctx context.Context) (map[string]ownedWorkStat, e var stat ownedWorkStat var sourceID string var latest sql.NullTime - if err := rows.Scan(&sourceID, &stat.WorkCount, &latest); err != nil { + if err := rows.Scan(&sourceID, &stat.WorkCount, &latest, &stat.CommentTotal); err != nil { return nil, err } stat.LatestPublishedAt = nullableTime(latest) diff --git a/internal/creator/content.go b/internal/creator/content.go index 972aaf5..e7995a3 100644 --- a/internal/creator/content.go +++ b/internal/creator/content.go @@ -623,16 +623,18 @@ func nullableArg(value time.Time) any { // RecordAccountMetric 落库账号画像快照。同一 (account_id, collected_at) 重复采集覆盖更新。 func (s *Store) RecordAccountMetric(ctx context.Context, input AccountMetricInput) error { if input.AccountID == "" || input.CollectedAt.IsZero() || input.FollowerCount != nil && *input.FollowerCount < 0 || - input.FollowingCount != nil && *input.FollowingCount < 0 || input.TotalFavorited != nil && *input.TotalFavorited < 0 || input.AwemeCount != nil && *input.AwemeCount < 0 { + input.FollowingCount != nil && *input.FollowingCount < 0 || input.TotalFavorited != nil && *input.TotalFavorited < 0 || + input.AwemeCount != nil && *input.AwemeCount < 0 || input.FriendCount != nil && *input.FriendCount < 0 { return ErrInvalid } result, err := s.db.ExecContext(ctx, ` - INSERT INTO creator_account_metric (account_id, collected_at, follower_count, following_count, total_favorited, aweme_count) - SELECT account.id, $2, $3, $4, $5, $6 FROM social_account account WHERE account.account_id = $1 + INSERT INTO creator_account_metric (account_id, collected_at, follower_count, following_count, total_favorited, aweme_count, friend_count) + SELECT account.id, $2, $3, $4, $5, $6, $7 FROM social_account account WHERE account.account_id = $1 ON CONFLICT (account_id, collected_at) DO UPDATE SET follower_count = EXCLUDED.follower_count, following_count = EXCLUDED.following_count, - total_favorited = EXCLUDED.total_favorited, aweme_count = EXCLUDED.aweme_count`, - input.AccountID, input.CollectedAt.UTC(), input.FollowerCount, input.FollowingCount, input.TotalFavorited, input.AwemeCount) + total_favorited = EXCLUDED.total_favorited, aweme_count = EXCLUDED.aweme_count, + friend_count = EXCLUDED.friend_count`, + input.AccountID, input.CollectedAt.UTC(), input.FollowerCount, input.FollowingCount, input.TotalFavorited, input.AwemeCount, input.FriendCount) if err != nil { return databaseError(err) } @@ -651,7 +653,7 @@ func (s *Store) ListAccountMetrics(ctx context.Context, accountID string) ([]Acc return nil, ErrInvalid } rows, err := s.db.QueryContext(ctx, ` - SELECT collected_at, follower_count, following_count, total_favorited, aweme_count + SELECT collected_at, follower_count, following_count, total_favorited, aweme_count, friend_count FROM creator_account_metric metric JOIN social_account account ON account.id = metric.account_id WHERE account.account_id = $1 ORDER BY metric.collected_at`, accountID) if err != nil { return nil, databaseError(err) @@ -660,13 +662,14 @@ func (s *Store) ListAccountMetrics(ctx context.Context, accountID string) ([]Acc result := make([]AccountMetricPoint, 0) for rows.Next() { var point AccountMetricPoint - var follower, following, favorited, aweme sql.NullInt64 - if err := rows.Scan(&point.CollectedAt, &follower, &following, &favorited, &aweme); err != nil { + var follower, following, favorited, aweme, friend sql.NullInt64 + if err := rows.Scan(&point.CollectedAt, &follower, &following, &favorited, &aweme, &friend); err != nil { return nil, err } point.CollectedAt = point.CollectedAt.UTC() point.FollowerCount, point.FollowingCount = nullableInt64(follower), nullableInt64(following) point.TotalFavorited, point.AwemeCount = nullableInt64(favorited), nullableInt64(aweme) + point.FriendCount = nullableInt64(friend) result = append(result, point) } return result, rows.Err() diff --git a/internal/creator/integration_test.go b/internal/creator/integration_test.go index 3f03524..f37a847 100644 --- a/internal/creator/integration_test.go +++ b/internal/creator/integration_test.go @@ -552,6 +552,13 @@ func TestCreatorPostgresAccountMonitorViewAndCollectionStatus(t *testing.T) { if view == nil || view.WorkCount != 0 || view.LatestPublishedAt != nil { t.Fatalf("monitor view without works: %+v", view) } + // 无指标快照:画像指标为空(不伪造成 0)。 + if view.FollowerCount != nil || view.FollowingCount != nil || view.TotalFavorited != nil || view.AwemeCount != nil { + t.Fatalf("monitor view without metrics must leave profile counters nil: %+v", view) + } + if view.CommentTotal != 0 { + t.Fatalf("monitor view without works must have zero comments: %+v", view) + } // 写入 checkpoint(succeeded + 窗口)与 owned 作品后:状态带下次窗口,视图聚合作品数。 lease, err := store.beginCheckpoint(ctx, SourceOwned, accountID, "works", now.Add(-24*time.Hour), now) @@ -585,6 +592,43 @@ func TestCreatorPostgresAccountMonitorViewAndCollectionStatus(t *testing.T) { t.Fatalf("monitor view aggregates: %+v", view) } } + + // 画像指标快照:写入两轮后视图取最新一轮;评论聚合来自 owned 作品最近指标之和。 + first, second := int64(100), int64(200) + following := int64(30) + favorited := int64(5000) + aweme := int64(7) + mutual := int64(9) + if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now.Add(-2 * time.Hour), FollowerCount: &first, FollowingCount: &following, TotalFavorited: &favorited, AwemeCount: &aweme, FriendCount: &mutual}); err != nil { + t.Fatal(err) + } + if err := store.RecordAccountMetric(ctx, AccountMetricInput{AccountID: accountID, CollectedAt: now.Add(-time.Hour), FollowerCount: &second}); err != nil { + t.Fatal(err) + } + monitorWork, _, err := store.UpsertWork(ctx, WorkInput{Platform: PlatformDouyin, WorkKey: "monitor-view-" + stamp, SourceType: SourceOwned, SourceID: accountID, Title: "MonitorView", PublishedAt: &published, PublishedAtStatus: "verified"}, now) + if err != nil { + t.Fatal(err) + } + if _, err := store.RecordMetric(ctx, MetricInput{WorkID: monitorWork.ID, CollectedAt: now, CommentsCount: &favorited}, Settings{LookbackDays: 30, NewWorkIntervalSeconds: 1800, MetricInitialIntervalSeconds: 3600, MetricMultiplier: 2, MetricMaxIntervalSeconds: 86400, MetricAgeSeconds: 30 * 24 * 60 * 60}, now); err != nil { + t.Fatal(err) + } + views, err = store.ListAccountMonitorViews(ctx) + if err != nil { + t.Fatal(err) + } + view = viewFor(accountID) + if view == nil || view.FollowerCount == nil || *view.FollowerCount != second { + t.Fatalf("monitor view must carry the latest metric snapshot: %+v", view) + } + if view.FollowingCount == nil || *view.FollowingCount != following || view.TotalFavorited == nil || *view.TotalFavorited != favorited || view.AwemeCount == nil || *view.AwemeCount != aweme { + t.Fatalf("monitor view must fall back to the latest non-null counters: %+v", view) + } + if view.FriendCount == nil || *view.FriendCount != mutual { + t.Fatalf("monitor view must carry the latest mutual friend count: %+v", view) + } + if view.CommentTotal != favorited { + t.Fatalf("monitor view must aggregate owned work comments: got %d want %d", view.CommentTotal, favorited) + } // failed checkpoint 带错误信息(普通错误落在 failed,不是 blocked)。 lease, err = store.beginCheckpoint(ctx, SourceOwned, accountID, "comments", now.Add(-24*time.Hour), now) if err != nil { diff --git a/internal/creator/models.go b/internal/creator/models.go index dc19528..3828d2b 100644 --- a/internal/creator/models.go +++ b/internal/creator/models.go @@ -259,6 +259,7 @@ type AccountMetricInput struct { FollowingCount *int64 TotalFavorited *int64 AwemeCount *int64 + FriendCount *int64 } type AccountMetricPoint struct { @@ -267,6 +268,7 @@ type AccountMetricPoint struct { FollowingCount *int64 `json:"following_count"` TotalFavorited *int64 `json:"total_favorited"` AwemeCount *int64 `json:"aweme_count"` + FriendCount *int64 `json:"friend_count"` } type MaterialJob struct { diff --git a/internal/environment/migration043_probe_test.go b/internal/environment/migration043_probe_test.go index 2dff345..c2e09fd 100644 --- a/internal/environment/migration043_probe_test.go +++ b/internal/environment/migration043_probe_test.go @@ -70,6 +70,6 @@ func TestMigration043ConsolidatedSchemaShape(t *testing.T) { // audit_event 任务列已删 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-1043、43 全部登记 - assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 50) + // 统一登记表:1-38(除 36)、1017-1044、43 全部登记 + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 51) } diff --git a/internal/environment/migration_test.go b/internal/environment/migration_test.go index f8cf13e..b6e8f9d 100644 --- a/internal/environment/migration_test.go +++ b/internal/environment/migration_test.go @@ -30,7 +30,7 @@ func TestUnifiedAccountMigration(t *testing.T) { t.Fatal(err) } defer db.Close() - assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 50) + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 51) 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`, 50) + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 51) }) t.Run("legacy migration 013 without account secrets is repaired forward", func(t *testing.T) { diff --git a/internal/environment/migrations/1044_creator_account_metric_friend_count.sql b/internal/environment/migrations/1044_creator_account_metric_friend_count.sql new file mode 100644 index 0000000..53de441 --- /dev/null +++ b/internal/environment/migrations/1044_creator_account_metric_friend_count.sql @@ -0,0 +1,2 @@ +-- 自有账号画像快照补互关数:self profile 接口 user.friend_count(实测字段存在),与粉丝/关注同节奏落库。 +ALTER TABLE creator_account_metric ADD COLUMN IF NOT EXISTS friend_count bigint; diff --git a/internal/environment/store.go b/internal/environment/store.go index 9ac92c5..bcfe6d4 100644 --- a/internal/environment/store.go +++ b/internal/environment/store.go @@ -170,6 +170,9 @@ var migration1042 string //go:embed migrations/1043_gateway_health.sql var migration1043 string +//go:embed migrations/1044_creator_account_metric_friend_count.sql +var migration1044 string + var ( ErrConflict = errors.New("resource conflicts with existing state") ErrInvalid = errors.New("invalid environment input") @@ -309,7 +312,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}, + {1041, migration1041}, {1042, migration1042}, {1043, migration1043}, {1044, migration1044}, {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 { diff --git a/internal/platform/douyin/connector.go b/internal/platform/douyin/connector.go index cbb7198..22cedf6 100644 --- a/internal/platform/douyin/connector.go +++ b/internal/platform/douyin/connector.go @@ -319,6 +319,7 @@ type identityEnvelope struct { FollowerCount *int64 `json:"follower_count"` FollowingCount *int64 `json:"following_count"` AwemeCount *int64 `json:"aweme_count"` + FriendCount *int64 `json:"friend_count"` TotalFavorited *flexibleInt `json:"total_favorited"` } `json:"user"` } @@ -356,7 +357,7 @@ func parseIdentity(body []byte) (identityEnvelope, bool) { (identity.User.UniqueID != "" && !keyPattern.MatchString(identity.User.UniqueID)) { return identityEnvelope{}, false } - for _, value := range []*int64{identity.User.FollowerCount, identity.User.FollowingCount, identity.User.AwemeCount} { + for _, value := range []*int64{identity.User.FollowerCount, identity.User.FollowingCount, identity.User.AwemeCount, identity.User.FriendCount} { if value != nil && *value < 0 { return identityEnvelope{}, false } @@ -375,6 +376,7 @@ type SelfProfile struct { FollowingCount *int64 TotalFavorited *int64 AwemeCount *int64 + FriendCount *int64 } // FetchSelfProfile 经网关受限 fetch 拉取登录账号自己的画像(复用 self profile 身份接口,返回体自带粉丝/关注/获赞/作品数)。 @@ -413,6 +415,7 @@ func parseSelfProfile(body []byte) (SelfProfile, bool) { return SelfProfile{ UID: identity.User.UID, SecUID: identity.User.SecUID, UniqueID: identity.User.UniqueID, Nickname: identity.User.Nickname, AvatarURL: avatarURL, FollowerCount: identity.User.FollowerCount, FollowingCount: identity.User.FollowingCount, TotalFavorited: totalFavorited, AwemeCount: identity.User.AwemeCount, + FriendCount: identity.User.FriendCount, }, true }