perf: 重构 FTS 索引为 rowid 寻址+触发器维护,根治重启 CPU 占满与扫描卡死

- FTS5 普通列(含 UNINDEXED)不支持索引查找,旧版按 media_id 做
  NOT EXISTS/DELETE 全是整表扫描:启动回填 O(N^2) 烧 CPU 数小时,
  扫描时每个 upsert 一次全表扫描,均隔着全局写锁拖死登录(125s 超时)
- v2 布局:FTS rowid 与 media.rowid 对齐,由 INSERT/UPDATE/DELETE
  触发器实时维护;回填只剩 rowid 点查兜底,且去掉 ORDER BY
- 顺带修复刮削直写 Updates() 后新标题搜不到的问题(触发器覆盖)
- 写门改为 context 感知:长写语句不再让登录请求无限期排队
- warmMediaSearchIndex 延迟 30s 错峰启动
- library_scan 周期任务跳过云盘库(由 cloud_sync 夜间窗口负责),
  首轮延迟从 15s 改为等满一个周期,重启不再立即全量扫描风暴
- prune 改为只取 id/path 并按 500 条批量删除,缩短写锁占用
This commit is contained in:
ShukeBta
2026-06-12 06:41:14 +00:00
parent 8d1d1f18a8
commit 1568ae127a
6 changed files with 187 additions and 77 deletions
+95 -24
View File
@@ -3,9 +3,9 @@
package database
import (
"context"
"fmt"
"path/filepath"
"sync"
"github.com/glebarez/sqlite"
"go.uber.org/zap"
@@ -57,12 +57,23 @@ func installSQLiteWriteGate(db *gorm.DB) {
if db == nil {
return
}
gate := &sqliteWriteGate{}
const lockedKey = "mediastation:sqlite_write_locked"
gate := newSQLiteWriteGate()
lock := func(tx *gorm.DB) {
gate.Lock()
ctx := context.Background()
if tx.Statement != nil && tx.Statement.Context != nil {
ctx = tx.Statement.Context
}
if err := gate.Lock(ctx); err != nil {
_ = tx.AddError(err)
return
}
tx.InstanceSet(lockedKey, struct{}{})
}
unlock := func(tx *gorm.DB) {
gate.Unlock()
if _, ok := tx.InstanceGet(lockedKey); ok {
gate.Unlock()
}
}
_ = db.Callback().Create().Before("gorm:create").Register("mediastation:sqlite_write_lock", lock)
_ = db.Callback().Create().After("gorm:create").Register("mediastation:sqlite_write_unlock", unlock)
@@ -74,16 +85,41 @@ func installSQLiteWriteGate(db *gorm.DB) {
_ = db.Callback().Raw().After("gorm:raw").Register("mediastation:sqlite_write_unlock", unlock)
}
// sqliteWriteGate 串行化进程内的 SQLite 写操作,避免多连接写竞争触发
// SQLITE_BUSY。Lock 尊重语句自身的 context:此前用 sync.Mutex 时,一条
// 长写语句(如 FTS 回填批次)会让登录等关键写操作无限期排队——客户端
// 早已超时断开,goroutine 还挂在互斥锁上。现在等待方可随 context 取消
// 及时失败,不再把整个进程的写路径拖死。
type sqliteWriteGate struct {
mu sync.Mutex
ch chan struct{}
}
func (g *sqliteWriteGate) Lock() {
g.mu.Lock()
func newSQLiteWriteGate() *sqliteWriteGate {
return &sqliteWriteGate{ch: make(chan struct{}, 1)}
}
func (g *sqliteWriteGate) Lock(ctx context.Context) error {
select {
case g.ch <- struct{}{}:
return nil
default:
}
if ctx == nil {
ctx = context.Background()
}
select {
case g.ch <- struct{}{}:
return nil
case <-ctx.Done():
return ctx.Err()
}
}
func (g *sqliteWriteGate) Unlock() {
g.mu.Unlock()
select {
case <-g.ch:
default:
}
}
func buildDSN(cfg *config.Config) string {
@@ -139,9 +175,31 @@ func ensurePerformanceIndexes(db *gorm.DB) error {
return nil
}
// mediaSearchIndexSchemaVersion 标识 FTS 索引的物理布局版本。
// v2:FTS 行的 rowid 与 media.rowid 对齐,并由触发器实时维护。
const mediaSearchIndexSchemaVersion = 2
func ensureMediaSearchIndex(db *gorm.DB) error {
if mediaSearchIndexNeedsRebuild(db) {
_ = db.Exec(`DROP TABLE IF EXISTS media_search_fts`).Error
if err := db.Exec(`CREATE TABLE IF NOT EXISTS media_search_meta (id INTEGER PRIMARY KEY CHECK (id = 1), version INTEGER NOT NULL)`).Error; err != nil {
return nil
}
var version int
_ = db.Raw(`SELECT version FROM media_search_meta WHERE id = 1`).Scan(&version).Error
if version != mediaSearchIndexSchemaVersion {
// 旧版(v1)FTS 表按 UNINDEXED 的 media_id 寻址。FTS5 的普通列
// 不支持索引查找,按 media_id 的 DELETE / NOT EXISTS 都是整表
// 扫描:十几万行的库每次启动回填要做上百亿次行访问,纯 Go
// sqlite 直接把 CPU 钉满数小时,并隔着全局写锁拖死登录。
// v2 起 FTS 行的 rowid 与 media.rowid 对齐,所有寻址走 rowid
// 点查,索引一致性交给下方触发器维护。
for _, stmt := range []string{
`DROP TRIGGER IF EXISTS media_search_fts_ai`,
`DROP TRIGGER IF EXISTS media_search_fts_au`,
`DROP TRIGGER IF EXISTS media_search_fts_ad`,
`DROP TABLE IF EXISTS media_search_fts`,
} {
_ = db.Exec(stmt).Error
}
}
if err := db.Exec(`CREATE VIRTUAL TABLE IF NOT EXISTS media_search_fts USING fts5(media_id UNINDEXED, title, original_name, path, genres, tokenize='trigram')`).Error; err != nil {
if fallbackErr := db.Exec(`CREATE VIRTUAL TABLE IF NOT EXISTS media_search_fts USING fts5(media_id UNINDEXED, title, original_name, path, genres, tokenize='unicode61')`).Error; fallbackErr != nil {
@@ -151,22 +209,35 @@ func ensureMediaSearchIndex(db *gorm.DB) error {
return nil
}
}
return nil
}
func mediaSearchIndexNeedsRebuild(db *gorm.DB) bool {
var cols []struct {
Name string
}
if err := db.Raw(`PRAGMA table_info(media_search_fts)`).Scan(&cols).Error; err != nil || len(cols) == 0 {
return false
}
for _, col := range cols {
if col.Name == "genres" {
return false
// 触发器让 FTS 与 media 行保持同步(新增/标题刮削改写/软删/恢复/
// 硬删全覆盖),应用层不再需要按 media_id 手工刷新索引——也顺带
// 修复了刮削直写 Updates() 后新标题搜不到的问题。
for _, stmt := range []string{
`CREATE TRIGGER IF NOT EXISTS media_search_fts_ai AFTER INSERT ON media WHEN new.deleted_at IS NULL BEGIN
DELETE FROM media_search_fts WHERE rowid = new.rowid;
INSERT INTO media_search_fts(rowid, media_id, title, original_name, path, genres)
VALUES (new.rowid, new.id, COALESCE(new.title, ''), COALESCE(new.original_name, ''), COALESCE(new.path, ''), COALESCE(new.genres, ''));
END`,
`CREATE TRIGGER IF NOT EXISTS media_search_fts_au AFTER UPDATE OF title, original_name, path, genres, deleted_at ON media BEGIN
DELETE FROM media_search_fts WHERE rowid = old.rowid;
INSERT INTO media_search_fts(rowid, media_id, title, original_name, path, genres)
SELECT new.rowid, new.id, COALESCE(new.title, ''), COALESCE(new.original_name, ''), COALESCE(new.path, ''), COALESCE(new.genres, '')
WHERE new.deleted_at IS NULL;
END`,
`CREATE TRIGGER IF NOT EXISTS media_search_fts_ad AFTER DELETE ON media BEGIN
DELETE FROM media_search_fts WHERE rowid = old.rowid;
END`,
} {
if err := db.Exec(stmt).Error; err != nil {
return err
}
}
return true
if version != mediaSearchIndexSchemaVersion {
if err := db.Exec(`INSERT INTO media_search_meta(id, version) VALUES (1, ?) ON CONFLICT(id) DO UPDATE SET version = excluded.version`, mediaSearchIndexSchemaVersion).Error; err != nil {
return err
}
}
return nil
}
func enforceTelegramBindingOneToOne(db *gorm.DB) error {
+10 -30
View File
@@ -323,7 +323,6 @@ func (r *MediaRepository) upsert(ctx context.Context, m *model.Media) error {
m.ScrapeStatus = "pending"
}
if createErr := r.db.WithContext(ctx).Create(m).Error; createErr == nil {
_ = r.refreshSearchIndex(ctx, m.ID)
return nil
} else if retryErr := r.db.WithContext(ctx).Unscoped().Where("path = ?", m.Path).First(&existing).Error; retryErr != nil {
return createErr
@@ -426,7 +425,6 @@ func (r *MediaRepository) upsert(ctx context.Context, m *model.Media) error {
Where("id = ?", existing.ID).Updates(updates).Error; err != nil {
return err
}
_ = r.refreshSearchIndex(ctx, existing.ID)
// 回写 ID / 不可变字段,让 caller 拿到完整的现有行。
*m = existing
return nil
@@ -518,7 +516,7 @@ func (r *MediaRepository) searchFilteredFTS(ctx context.Context, query string, o
var items []model.Media
q := r.db.WithContext(ctx).
Table("media").
Joins("JOIN media_search_fts ON media_search_fts.media_id = media.id").
Joins("JOIN media_search_fts ON media_search_fts.rowid = media.rowid").
Where("media.deleted_at IS NULL").
Where("media_search_fts MATCH ?", ftsQuery)
q = applyQualifiedMediaQueryFilter(q, filter)
@@ -625,23 +623,6 @@ func escapeLike(value string) string {
return value
}
func (r *MediaRepository) refreshSearchIndex(ctx context.Context, mediaID string) error {
if strings.TrimSpace(mediaID) == "" {
return nil
}
if !r.searchIndexEnabled(ctx) {
return nil
}
tx := r.db.WithContext(ctx)
_ = tx.Exec(`DELETE FROM media_search_fts WHERE media_id = ?`, mediaID).Error
return tx.Exec(`
INSERT INTO media_search_fts(media_id, title, original_name, path, genres)
SELECT id, COALESCE(title, ''), COALESCE(original_name, ''), COALESCE(path, ''), COALESCE(genres, '')
FROM media
WHERE id = ? AND deleted_at IS NULL
`, mediaID).Error
}
func (r *MediaRepository) BackfillSearchIndex(ctx context.Context, batchLimit int) (int64, error) {
if batchLimit <= 0 {
batchLimit = 1000
@@ -649,15 +630,19 @@ func (r *MediaRepository) BackfillSearchIndex(ctx context.Context, batchLimit in
if !r.searchIndexEnabled(ctx) {
return 0, nil
}
// 关键性能点:FTS5 普通列(含 UNINDEXED)不支持索引查找,按
// media_id 做 NOT EXISTS 是对 FTS 表的整表扫描,再叠加 ORDER BY
// 后每个批次都要对全部 media 行探测一遍——大库一次启动回填等于
// 上百亿次行访问,曾把 CPU 钉满数小时。v2 布局下 FTS 行 rowid 与
// media.rowid 对齐,NOT EXISTS 走 rowid 点查,且无需排序。
res := r.db.WithContext(ctx).Exec(`
INSERT INTO media_search_fts(media_id, title, original_name, path, genres)
SELECT m.id, COALESCE(m.title, ''), COALESCE(m.original_name, ''), COALESCE(m.path, ''), COALESCE(m.genres, '')
INSERT INTO media_search_fts(rowid, media_id, title, original_name, path, genres)
SELECT m.rowid, m.id, COALESCE(m.title, ''), COALESCE(m.original_name, ''), COALESCE(m.path, ''), COALESCE(m.genres, '')
FROM media AS m
WHERE m.deleted_at IS NULL
AND NOT EXISTS (
SELECT 1 FROM media_search_fts AS f WHERE f.media_id = m.id
SELECT 1 FROM media_search_fts AS f WHERE f.rowid = m.rowid
)
ORDER BY m.created_at DESC
LIMIT ?
`, batchLimit)
return res.RowsAffected, res.Error
@@ -679,18 +664,13 @@ func (r *MediaRepository) searchIndexEnabled(ctx context.Context) bool {
// DeleteByLibrary purges all media tied to a library.
func (r *MediaRepository) DeleteByLibrary(ctx context.Context, libraryID string) error {
if r.searchIndexEnabled(ctx) {
_ = r.db.WithContext(ctx).Exec(`DELETE FROM media_search_fts WHERE media_id IN (SELECT id FROM media WHERE library_id = ?)`, libraryID).Error
}
// FTS 行由 media 表上的触发器同步清理(软删/硬删都覆盖)。
return r.db.WithContext(ctx).Where("library_id = ?", libraryID).Delete(&model.Media{}).Error
}
// PurgeByLibrary permanently removes media tied to a library. Used for virtual
// cloud mounts where "remove mount" must not populate the recycle bin.
func (r *MediaRepository) PurgeByLibrary(ctx context.Context, libraryID string) error {
if r.searchIndexEnabled(ctx) {
_ = r.db.WithContext(ctx).Exec(`DELETE FROM media_search_fts WHERE media_id IN (SELECT id FROM media WHERE library_id = ?)`, libraryID).Error
}
return r.db.WithContext(ctx).Unscoped().Where("library_id = ?", libraryID).Delete(&model.Media{}).Error
}
+22 -6
View File
@@ -153,12 +153,17 @@ func TestMediaSearchIndexBackfillRunsInBatches(t *testing.T) {
}).Error; err != nil {
t.Fatal(err)
}
var before int64
if err := repos.DB.Raw(`SELECT COUNT(*) FROM media_search_fts`).Scan(&before).Error; err != nil {
// 插入触发器应当同步维护 FTS 行。
var indexed int64
if err := repos.DB.Raw(`SELECT COUNT(*) FROM media_search_fts`).Scan(&indexed).Error; err != nil {
t.Fatal(err)
}
if before != 0 {
t.Fatalf("startup migrate should not synchronously backfill FTS, got %d rows", before)
if indexed != 1 {
t.Fatalf("insert trigger should index new media, got %d rows", indexed)
}
// 清空 FTS 模拟旧库升级后索引缺失,回填应按批补齐且 rowid 对齐。
if err := repos.DB.Exec(`DELETE FROM media_search_fts`).Error; err != nil {
t.Fatal(err)
}
n, err := repos.Media.BackfillSearchIndex(t.Context(), 1)
if err != nil {
@@ -167,11 +172,22 @@ func TestMediaSearchIndexBackfillRunsInBatches(t *testing.T) {
if n != 1 {
t.Fatalf("backfilled rows = %d, want 1", n)
}
var aligned int64
if err := repos.DB.Raw(`SELECT COUNT(*) FROM media_search_fts f JOIN media m ON f.rowid = m.rowid AND f.media_id = m.id`).Scan(&aligned).Error; err != nil {
t.Fatal(err)
}
if aligned != 1 {
t.Fatalf("fts rows aligned with media rowid = %d, want 1", aligned)
}
// 软删除后触发器应清理对应 FTS 行,避免搜索命中已删媒体。
if err := repos.DB.Delete(&model.Media{}, "id = ?", "m-backfill").Error; err != nil {
t.Fatal(err)
}
var after int64
if err := repos.DB.Raw(`SELECT COUNT(*) FROM media_search_fts`).Scan(&after).Error; err != nil {
t.Fatal(err)
}
if after != 1 {
t.Fatalf("fts rows = %d, want 1", after)
if after != 0 {
t.Fatalf("soft delete should drop fts row, got %d", after)
}
}
+38 -16
View File
@@ -1491,13 +1491,20 @@ func (s *ScannerService) mediaPathExists(ctx context.Context, path string) bool
}
func (s *ScannerService) pruneMissingMedia(ctx context.Context, libraryID string, seen map[string]struct{}) (int64, error) {
var rows []model.Media
// 只取 id/path,并把删除按批提交:此前整表载入完整 Media 结构体、
// 每行一条 DELETE,大库 prune 既费内存又长期占用写锁。
var rows []struct {
ID string
Path string
}
if err := s.repo.DB.WithContext(ctx).
Model(&model.Media{}).
Select("id, path").
Where("library_id = ?", libraryID).
Find(&rows).Error; err != nil {
return 0, err
}
var removed int64
stale := make([]string, 0)
for _, row := range rows {
if row.Path == "" {
continue
@@ -1510,9 +1517,26 @@ func (s *ScannerService) pruneMissingMedia(ctx context.Context, libraryID string
} else if !os.IsNotExist(err) {
continue
}
res := s.repo.DB.WithContext(ctx).
Where("id = ?", row.ID).
Delete(&model.Media{})
stale = append(stale, row.ID)
}
return s.deleteMediaByIDs(ctx, stale, false)
}
// deleteMediaByIDs removes media rows in fixed-size batches so each write
// transaction stays short and the global write gate is released frequently.
func (s *ScannerService) deleteMediaByIDs(ctx context.Context, ids []string, hard bool) (int64, error) {
const batch = 500
var removed int64
for i := 0; i < len(ids); i += batch {
end := i + batch
if end > len(ids) {
end = len(ids)
}
q := s.repo.DB.WithContext(ctx)
if hard {
q = q.Unscoped()
}
res := q.Where("id IN ?", ids[i:end]).Delete(&model.Media{})
if res.Error != nil {
return removed, res.Error
}
@@ -1522,27 +1546,25 @@ func (s *ScannerService) pruneMissingMedia(ctx context.Context, libraryID string
}
func (s *ScannerService) pruneMissingCloudMedia(ctx context.Context, libraryID string, seen map[string]struct{}) (int64, error) {
var rows []model.Media
var rows []struct {
ID string
Path string
}
if err := s.repo.DB.WithContext(ctx).
Model(&model.Media{}).
Select("id, path").
Where("library_id = ? AND path LIKE ?", libraryID, "cloud://%").
Find(&rows).Error; err != nil {
return 0, err
}
var removed int64
stale := make([]string, 0)
for _, row := range rows {
if _, ok := seen[row.Path]; ok {
continue
}
res := s.repo.DB.WithContext(ctx).
Unscoped().
Where("id = ?", row.ID).
Delete(&model.Media{})
if res.Error != nil {
return removed, res.Error
}
removed += res.RowsAffected
stale = append(stale, row.ID)
}
return removed, nil
return s.deleteMediaByIDs(ctx, stale, true)
}
func parseCloudLibraryPath(raw string) (typ, dirID string, ok bool) {
+15 -1
View File
@@ -134,7 +134,14 @@ func (s *SchedulerService) Start(ctx context.Context) {
},
}
for _, j := range s.jobs {
go s.loop(ctx, j)
initialDelay := 15 * time.Second
if j.name == "library_scan" {
// 重启后不立即整库重扫:更新/重启窗口恰是登录高峰,启动
// 15 秒即全量扫描曾把 CPU/磁盘打满导致无法登录。首轮等满
// 一个完整周期再跑,平时的每小时节奏不变。
initialDelay = j.interval
}
go s.loopWithInitialDelay(ctx, j, initialDelay)
}
}
@@ -259,6 +266,13 @@ func (s *SchedulerService) jobScanLibraries(ctx context.Context) error {
if !l.Enabled {
continue
}
if _, ok := ParseCloudLibraryMount(l.Path); ok {
// 云盘库由 cloud_sync 任务在夜间窗口低频同步;周期性整库
// 重扫只面向本地磁盘库。否则十几个云盘库每小时全量遍历
// 会把 CPU/网络长期吃满,还会占住唯一的云扫描槽位,让
// 手动扫描看起来一直"卡死"在排队。
continue
}
if _, err := s.scanner.ScanLibrary(ctx, l.ID); err != nil {
s.log.Warn("scheduled scan failed",
zap.String("library", l.ID), zap.Error(err))
+7
View File
@@ -268,6 +268,13 @@ func (c *Container) warmMediaSearchIndex(ctx context.Context) {
if c == nil || c.Repo == nil || c.Repo.Media == nil {
return
}
// 错峰:FTS 正常由 media 表触发器实时维护,回填只是升级或异常后的
// 兜底。先让登录、首页等关键路径跑起来,再开始后台补索引。
select {
case <-ctx.Done():
return
case <-time.After(30 * time.Second):
}
const batchSize = 1000
const pause = 100 * time.Millisecond
total := int64(0)