81 lines
2.3 KiB
Go
81 lines
2.3 KiB
Go
package esmedia
|
|
|
|
import (
|
|
"fmt"
|
|
"time"
|
|
|
|
"91porn-server/common/elastic"
|
|
"91porn-server/common/log"
|
|
"91porn-server/models/s/statrecordmod"
|
|
"91porn-server/models/v/mediamod"
|
|
"91porn-server/skd/skdg"
|
|
|
|
"go.mongodb.org/mongo-driver/bson/primitive"
|
|
)
|
|
|
|
var es *elastic.Client
|
|
|
|
func RunSync(defaultDate time.Time) {
|
|
es = skdg.VideoES
|
|
log.Info("Elastic Search ACG media sync start...")
|
|
defer log.Info("Elastic Search ACG media sync end.")
|
|
// MongoDB 日期精度为毫秒;本轮开始后的修改必须留给下一轮再次同步。
|
|
startedAt := time.Now().Truncate(time.Millisecond)
|
|
lastSyncAt, err := statrecordmod.LastRecordTime(statrecordmod.ElasticSyncMediaJob, defaultDate)
|
|
if err != nil {
|
|
log.Error("Elastic sync ACG media get last record time failed", log.E(err))
|
|
return
|
|
}
|
|
if !startedAt.After(lastSyncAt) {
|
|
return
|
|
}
|
|
// 固定 ID 上界,避免持续新增数据使本轮无法结束。
|
|
maxID, err := mediamod.GetLatestIDForESSync()
|
|
if err != nil {
|
|
log.Error("Elastic sync ACG media get latest ID failed", log.E(err))
|
|
return
|
|
}
|
|
sync := mediaSync{
|
|
read: mediamod.GetListForESSync,
|
|
write: SyncDataToES,
|
|
checkpoint: func(at time.Time) error {
|
|
return statrecordmod.UpsertOneTrans(nil, statrecordmod.ElasticSyncMediaJob, at)
|
|
},
|
|
}
|
|
if err = sync.run(lastSyncAt, startedAt, maxID); err != nil {
|
|
log.Error("Elastic sync ACG media failed; checkpoint not confirmed", log.E(err))
|
|
}
|
|
}
|
|
|
|
type mediaSync struct {
|
|
read func(time.Time, primitive.ObjectID, primitive.ObjectID, int) ([]*mediamod.Media, error)
|
|
write func([]*mediamod.Media) error
|
|
checkpoint func(time.Time) error
|
|
}
|
|
|
|
func (s mediaSync) run(since, startedAt time.Time, maxID primitive.ObjectID) error {
|
|
const size = 1000
|
|
var after primitive.ObjectID
|
|
for !maxID.IsZero() {
|
|
data, err := s.read(since, after, maxID, size)
|
|
if err != nil {
|
|
return fmt.Errorf("read media batch: %w", err)
|
|
}
|
|
if len(data) == 0 {
|
|
break
|
|
}
|
|
if err = s.write(data); err != nil {
|
|
return fmt.Errorf("write media batch: %w", err)
|
|
}
|
|
after = data[len(data)-1].ID
|
|
if len(data) < size || after == maxID {
|
|
break
|
|
}
|
|
}
|
|
// 所有批次及每条 ES 写入均成功后才推进;失败时下一轮从原进度幂等重试。
|
|
if err := s.checkpoint(startedAt); err != nil {
|
|
return fmt.Errorf("save media checkpoint: %w", err)
|
|
}
|
|
return nil
|
|
}
|