lottery.go 7.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325
  1. package service
  2. import (
  3. "bytes"
  4. "context"
  5. "sort"
  6. "strconv"
  7. "sync"
  8. "time"
  9. "go-common/app/job/main/growup/model"
  10. "go-common/library/ecode"
  11. "go-common/library/log"
  12. "golang.org/x/sync/errgroup"
  13. )
  14. type bubbleType int
  15. const (
  16. _bubbleLottery bubbleType = 1
  17. _bubbleVote bubbleType = 2
  18. )
  19. type bubbleMeta struct {
  20. avid int64
  21. aType bubbleType
  22. }
  23. // use var map[] + register map ???
  24. func (s *Service) syncBubbleMetaFuncM() map[bubbleType]func(c context.Context, start, end time.Time) ([]*bubbleMeta, error) {
  25. return map[bubbleType]func(c context.Context, start, end time.Time) ([]*bubbleMeta, error){
  26. _bubbleLottery: s.syncLotteryAvs,
  27. _bubbleVote: s.syncVoteBIZAvs,
  28. }
  29. }
  30. func (s *Service) syncLotteryAvs(c context.Context, start, end time.Time) (data []*bubbleMeta, err error) {
  31. var (
  32. offset int64
  33. hasMore = 1
  34. avs []int64
  35. )
  36. for hasMore == 1 {
  37. var info *model.LotteryRes
  38. info, err = s.dao.GetLotteryRIDs(c, start.Unix(), end.Unix()-1, offset)
  39. if err != nil {
  40. log.Error("s.dao.GetLotteryRIDs error(%v)", err)
  41. return
  42. }
  43. if len(info.RIDs) > 0 {
  44. avs = append(avs, info.RIDs...)
  45. }
  46. hasMore = info.HasMore
  47. offset = info.Offset
  48. }
  49. data = make([]*bubbleMeta, 0, len(avs))
  50. for _, av := range avs {
  51. data = append(data, &bubbleMeta{
  52. avid: av,
  53. aType: _bubbleLottery,
  54. })
  55. }
  56. return
  57. }
  58. func (s *Service) syncVoteBIZAvs(c context.Context, start, end time.Time) (data []*bubbleMeta, err error) {
  59. if end.After(start.AddDate(0, 6, 0)) {
  60. err = ecode.Error(ecode.RequestErr, "同步投票视频时间区间不能超过半年")
  61. return
  62. }
  63. data = make([]*bubbleMeta, 0)
  64. var voteAVs []*model.VoteBIZArchive
  65. voteAVs, err = s.dao.VoteBIZArchive(c, start.Unix(), end.Unix())
  66. if err != nil {
  67. return
  68. }
  69. if len(voteAVs) == 0 {
  70. return
  71. }
  72. for _, av := range voteAVs {
  73. data = append(data, &bubbleMeta{
  74. avid: av.Aid,
  75. aType: _bubbleVote,
  76. })
  77. }
  78. return
  79. }
  80. func (s *Service) saveBubbleMeta(c context.Context, metas []*bubbleMeta, date time.Time) (rows int64, err error) {
  81. return s.dao.InsertBubbleMeta(c, assembleBubbleMeta(metas, date))
  82. }
  83. func assembleBubbleMeta(data []*bubbleMeta, syncDate time.Time) (values string) {
  84. var buf bytes.Buffer
  85. for _, v := range data {
  86. buf.WriteString("(")
  87. buf.WriteString(strconv.FormatInt(v.avid, 10))
  88. buf.WriteByte(',')
  89. buf.WriteString("'" + syncDate.Format(_layout) + "'")
  90. buf.WriteByte(',')
  91. buf.WriteString(strconv.Itoa(int(v.aType)))
  92. buf.WriteString(")")
  93. buf.WriteByte(',')
  94. }
  95. if buf.Len() > 0 {
  96. buf.Truncate(buf.Len() - 1)
  97. }
  98. values = buf.String()
  99. buf.Reset()
  100. return
  101. }
  102. // SyncIncomeBubbleMeta .
  103. func (s *Service) SyncIncomeBubbleMeta(c context.Context, start, end time.Time, tp int) (rows int64, err error) {
  104. syncFuncM := s.syncBubbleMetaFuncM()
  105. syncFunc, ok := syncFuncM[bubbleType(tp)]
  106. if !ok {
  107. err = ecode.Errorf(ecode.ReqParamErr, "illegal bubbleType(%d)", tp)
  108. return
  109. }
  110. metas, err := syncFunc(c, start, end)
  111. if err != nil {
  112. log.Error("s.syncVoteBIZAvs err(%v)", err)
  113. return
  114. }
  115. rows, err = s.saveBubbleMeta(c, metas, time.Now())
  116. if err != nil {
  117. log.Error("s.saveBubbleMeta err(%v)", err)
  118. }
  119. return
  120. }
  121. // SyncIncomeBubbleMetaTask sync meta data for income bubble to growup
  122. func (s *Service) SyncIncomeBubbleMetaTask(c context.Context, date time.Time) (err error) {
  123. defer func() {
  124. GetTaskService().SetTaskStatus(c, TaskBubbleMeta, date.Format(_layout), err)
  125. }()
  126. var (
  127. lock sync.Mutex
  128. group errgroup.Group
  129. syncMetaFuncM = s.syncBubbleMetaFuncM()
  130. start, end = date, date.AddDate(0, 0, 1)
  131. metas []*bubbleMeta
  132. )
  133. for _, syncMetaFunc := range syncMetaFuncM {
  134. group.Go(func() error {
  135. var data []*bubbleMeta
  136. data, err = syncMetaFunc(c, start, end)
  137. if err != nil {
  138. return err
  139. }
  140. lock.Lock()
  141. metas = append(metas, data...)
  142. lock.Unlock()
  143. return nil
  144. })
  145. }
  146. err = group.Wait()
  147. if err != nil {
  148. return
  149. }
  150. _, err = s.saveBubbleMeta(c, metas, date)
  151. if err != nil {
  152. log.Error("s.saveBubbleMeta err(%v)", err)
  153. }
  154. return
  155. }
  156. func (s *Service) avToBType(bubbleMeta map[int64][]int) map[int64]int {
  157. var (
  158. res = make(map[int64]int)
  159. typeToRatio = make(map[int]float64)
  160. chooseBType = func(bTypes []int) (bType int) {
  161. if len(bTypes) == 1 {
  162. return bTypes[0]
  163. }
  164. sort.Slice(bTypes, func(i, j int) bool {
  165. bti, btj := bTypes[i], bTypes[j]
  166. if typeToRatio[bti] == typeToRatio[btj] {
  167. return bti < btj
  168. }
  169. return typeToRatio[bti] < typeToRatio[btj]
  170. })
  171. return bTypes[0]
  172. }
  173. )
  174. for _, v := range s.conf.Bubble.BRatio {
  175. typeToRatio[v.BType] = v.Ratio
  176. }
  177. for avID, bTypes := range bubbleMeta {
  178. res[avID] = chooseBType(bTypes)
  179. }
  180. return res
  181. }
  182. // SnapshotBubbleIncomeTask get income lottery
  183. func (s *Service) SnapshotBubbleIncomeTask(c context.Context, date time.Time) (err error) {
  184. defer func() {
  185. GetTaskService().SetTaskStatus(c, TaskSnapshotBubbleIncome, date.Format("2006-01-02"), err)
  186. }()
  187. err = GetTaskService().TaskReady(c, date.Format("2006-01-02"), TaskCreativeIncome)
  188. if err != nil {
  189. return
  190. }
  191. bubbleMeta, err := s.getBubbleMeta(c)
  192. if err != nil {
  193. log.Error("s.getBubbleAVs error(%v)", err)
  194. return
  195. }
  196. avToBType := s.avToBType(bubbleMeta)
  197. var (
  198. eg errgroup.Group
  199. incomeCh = make(chan []*model.IncomeInfo, 2000)
  200. snapshotCh = make(chan []*model.IncomeInfo, 2000)
  201. )
  202. // get av income by date
  203. eg.Go(func() (err error) {
  204. defer close(incomeCh)
  205. var from int64
  206. for {
  207. var infos []*model.IncomeInfo
  208. infos, err = s.dao.GetAvIncome(c, date, from, int64(_dbLimit))
  209. if err != nil {
  210. log.Error("dao.GetAvIncome error(%v)", err)
  211. return
  212. }
  213. if len(infos) == 0 {
  214. break
  215. }
  216. incomeCh <- infos
  217. from = infos[len(infos)-1].ID
  218. }
  219. return
  220. })
  221. eg.Go(func() (err error) {
  222. defer close(snapshotCh)
  223. for income := range incomeCh {
  224. lotteryIncome := make([]*model.IncomeInfo, 0)
  225. for _, av := range income {
  226. if bType, ok := avToBType[av.AVID]; ok {
  227. av.BType = bType
  228. lotteryIncome = append(lotteryIncome, av)
  229. }
  230. }
  231. if len(lotteryIncome) > 0 {
  232. snapshotCh <- lotteryIncome
  233. }
  234. }
  235. return
  236. })
  237. eg.Go(func() (err error) {
  238. for income := range snapshotCh {
  239. _, err = s.dao.InsertBubbleIncome(c, assembleLotteryIncome(income))
  240. if err != nil {
  241. return
  242. }
  243. }
  244. return
  245. })
  246. if err = eg.Wait(); err != nil {
  247. log.Error("eg.Wait error(%v)", err)
  248. }
  249. return
  250. }
  251. func (s *Service) getBubbleMeta(c context.Context) (data map[int64][]int, err error) {
  252. var id int64
  253. data = make(map[int64][]int)
  254. for {
  255. var meta map[int64][]int
  256. meta, id, err = s.income.GetBubbleMeta(c, id, int64(_dbLimit))
  257. if err != nil {
  258. return
  259. }
  260. if len(meta) == 0 {
  261. break
  262. }
  263. for avID, bTypes := range meta {
  264. data[avID] = append(data[avID], bTypes...)
  265. }
  266. }
  267. return
  268. }
  269. func assembleLotteryIncome(as []*model.IncomeInfo) (values string) {
  270. var buf bytes.Buffer
  271. for _, a := range as {
  272. buf.WriteString("(")
  273. buf.WriteString(strconv.FormatInt(a.AVID, 10))
  274. buf.WriteByte(',')
  275. buf.WriteString(strconv.FormatInt(a.MID, 10))
  276. buf.WriteByte(',')
  277. buf.WriteString(strconv.FormatInt(a.TagID, 10))
  278. buf.WriteByte(',')
  279. buf.WriteString("'" + a.UploadTime.Format("2006-01-02 15:04:05") + "'")
  280. buf.WriteByte(',')
  281. buf.WriteString(strconv.FormatInt(a.TotalIncome, 10))
  282. buf.WriteByte(',')
  283. buf.WriteString(strconv.FormatInt(a.Income, 10))
  284. buf.WriteByte(',')
  285. buf.WriteString(strconv.FormatInt(a.TaxMoney, 10))
  286. buf.WriteByte(',')
  287. buf.WriteString("'" + a.Date.Format(_layout) + "'")
  288. buf.WriteByte(',')
  289. buf.WriteString(strconv.FormatInt(a.BaseIncome, 10))
  290. buf.WriteByte(',')
  291. buf.WriteString(strconv.Itoa(a.BType))
  292. buf.WriteString(")")
  293. buf.WriteByte(',')
  294. }
  295. if buf.Len() > 0 {
  296. buf.Truncate(buf.Len() - 1)
  297. }
  298. values = buf.String()
  299. buf.Reset()
  300. return
  301. }