es.go 3.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129
  1. package dao
  2. import (
  3. "context"
  4. "encoding/json"
  5. "fmt"
  6. "strings"
  7. "go-common/app/service/main/search/model"
  8. "go-common/library/ecode"
  9. "go-common/library/log"
  10. elastic "gopkg.in/olivere/elastic.v5"
  11. )
  12. // UpdateMapBulk .
  13. func (d *Dao) UpdateMapBulk(c context.Context, esName string, bulkData []BulkMapItem) (err error) {
  14. if _, ok := d.esPool[esName]; !ok {
  15. PromError(fmt.Sprintf("es:集群不存在%s", esName), "s.dao.searchResult indexName:%s", esName)
  16. err = ecode.SearchUpdateIndexFailed
  17. return
  18. }
  19. bulkRequest := d.esPool[esName].Bulk()
  20. for _, b := range bulkData {
  21. request := elastic.NewBulkUpdateRequest().Index(b.IndexName()).Type(b.IndexType()).Id(b.IndexID()).Doc(b.PField()).DocAsUpsert(true)
  22. bulkRequest.Add(request)
  23. }
  24. if _, err = bulkRequest.Do(context.TODO()); err != nil {
  25. log.Error("esName(%s) bulk error(%v)", esName, err)
  26. }
  27. return
  28. }
  29. func (d *Dao) UpdateBulk(c context.Context, esName string, bulkData []BulkItem) (err error) {
  30. if _, ok := d.esPool[esName]; !ok {
  31. PromError(fmt.Sprintf("es:集群不存在%s", esName), "s.dao.searchResult indexName:%s", esName)
  32. err = ecode.SearchUpdateIndexFailed
  33. return
  34. }
  35. bulkRequest := d.esPool[esName].Bulk()
  36. for _, b := range bulkData {
  37. request := elastic.NewBulkUpdateRequest().Index(b.IndexName()).Type(b.IndexType()).Id(b.IndexID()).Doc(b).DocAsUpsert(true)
  38. //fmt.Println(request)
  39. bulkRequest.Add(request)
  40. }
  41. if bulkRequest.NumberOfActions() == 0 {
  42. return
  43. }
  44. if _, err = bulkRequest.Do(context.TODO()); err != nil {
  45. log.Error("esName(%s) bulk error(%v)", esName, err)
  46. }
  47. return
  48. }
  49. // searchResult get result from ES.
  50. func (d *Dao) searchResult(c context.Context, esClusterName, indexName string, query elastic.Query, bsp *model.BasicSearchParams) (res *model.SearchResult, err error) {
  51. res = &model.SearchResult{Debug: ""}
  52. if bsp.Debug {
  53. var src interface{}
  54. if src, err = query.Source(); err == nil {
  55. var data []byte
  56. if data, err = json.Marshal(src); err == nil {
  57. res = &model.SearchResult{Debug: string(data)}
  58. }
  59. }
  60. }
  61. if _, ok := d.esPool[esClusterName]; !ok {
  62. PromError(fmt.Sprintf("es:集群不存在%s", esClusterName), "s.dao.searchResult indexName:%s", indexName)
  63. res = &model.SearchResult{Debug: fmt.Sprintf("es:集群不存在%s, %s", esClusterName, res.Debug)}
  64. return
  65. }
  66. // multi sort
  67. sorterSlice := []elastic.Sorter{}
  68. if bsp.KW != "" {
  69. sorterSlice = append(sorterSlice, elastic.NewScoreSort().Desc())
  70. }
  71. for i, d := range bsp.Order {
  72. if len(bsp.Sort) < i+1 {
  73. if bsp.Sort[0] == "desc" {
  74. sorterSlice = append(sorterSlice, elastic.NewFieldSort(d).Desc())
  75. } else {
  76. sorterSlice = append(sorterSlice, elastic.NewFieldSort(d).Asc())
  77. }
  78. } else {
  79. if bsp.Sort[i] == "desc" {
  80. sorterSlice = append(sorterSlice, elastic.NewFieldSort(d).Desc())
  81. } else {
  82. sorterSlice = append(sorterSlice, elastic.NewFieldSort(d).Asc())
  83. }
  84. }
  85. }
  86. fsc := elastic.NewFetchSourceContext(true).Include(bsp.Source...)
  87. searchResult, err := d.esPool[esClusterName].
  88. Search().Index(indexName).
  89. Query(query).
  90. SortBy(sorterSlice...).
  91. From((bsp.Pn - 1) * bsp.Ps).
  92. Size(bsp.Ps).
  93. Pretty(true).
  94. FetchSourceContext(fsc).
  95. Do(c)
  96. if err != nil {
  97. PromError(fmt.Sprintf("es:执行查询失败%s ", esClusterName), "%v", err)
  98. res = &model.SearchResult{Debug: res.Debug + "es:执行查询失败"}
  99. return
  100. }
  101. var data []json.RawMessage
  102. for _, hit := range searchResult.Hits.Hits {
  103. var t json.RawMessage
  104. e := json.Unmarshal(*hit.Source, &t)
  105. if e != nil {
  106. PromError(fmt.Sprintf("es:%s 索引有脏数据", esClusterName), "s.dao.SearchArchiveCheck(%d,%d) error(%v) ", bsp.Pn*bsp.Ps, bsp.Ps, e)
  107. continue
  108. }
  109. data = append(data, t)
  110. }
  111. res = &model.SearchResult{
  112. Order: strings.Join(bsp.Order, ","),
  113. Sort: strings.Join(bsp.Sort, ","),
  114. Result: data,
  115. Debug: res.Debug,
  116. Page: &model.Page{
  117. Pn: bsp.Pn,
  118. Ps: bsp.Ps,
  119. Total: searchResult.Hits.TotalHits,
  120. },
  121. }
  122. return
  123. }