package esaudiobook import ( "fmt" "time" "91porn-server/common/elastic" "91porn-server/common/log" "91porn-server/common/timeutil/timerange" "91porn-server/models/s/statrecordmod" "91porn-server/models/v/audiobookmod" "91porn-server/skd/skdg" ) var es *elastic.Client func RunSync(defaultDate time.Time) { es = skdg.VideoES log.Info("Elastic Search audiobook sync start...") defer log.Info("Elastic Search audiobook sync end.") // 获取上次数据同步的记录时间 lastSyncAt, err := statrecordmod.LastRecordTime(statrecordmod.ElasticSyncAudioBookJob, defaultDate) if err != nil { log.Error("Elastic sync audiobook data RunSync get last record time failed", log.E(err)) return } var now = time.Now().Add(-time.Minute * 20) var recentMinute = timerange.RecentMinute(now, statrecordmod.FiveMinuteScale) for recentMinute.After(lastSyncAt) { var end = lastSyncAt.Add(statrecordmod.FiveMinuteScale * time.Minute) if end.After(recentMinute) { end = recentMinute } var subTimeRange = timerange.TimeRange{ Head: lastSyncAt, Tail: end, } data, err := audiobookmod.GetListByUpdateTimeRange(subTimeRange.Head, subTimeRange.Tail) if err != nil { return } if err = commit(data, subTimeRange.Tail); err != nil { log.Error("Elastic sync audiobook data RunSync commit err! ", log.E(err)) return } lastSyncAt = end } } func commit(data []audiobookmod.AudioBookBase, syncTime time.Time) error { if err := SyncDataToES(data); err != nil { return fmt.Errorf("Elastic sync commit audiobook data SyncDataToES failed!, err:%+v\n", err) } if err := statrecordmod.UpsertOneTrans(nil, statrecordmod.ElasticSyncAudioBookJob, syncTime); err != nil { return fmt.Errorf("Elastic sync commit audiobook data UpsertOneTrans failed!, err:%+v\n", err) } return nil }