forked from olivere/elastic
-
Notifications
You must be signed in to change notification settings - Fork 0
/
msearch.go
132 lines (115 loc) · 3.38 KB
/
msearch.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
// Copyright 2012-present Oliver Eilhard. All rights reserved.
// Use of this source code is governed by a MIT-license.
// See http://olivere.mit-license.org/license.txt for details.
package elastic
import (
"context"
"encoding/json"
"fmt"
"net/url"
"strings"
)
// MultiSearch executes one or more searches in one roundtrip.
type MultiSearchService struct {
client *Client
requests []*SearchRequest
indices []string
pretty bool
maxConcurrentRequests *int
preFilterShardSize *int
restTotalHitsAsInt *bool
}
func NewMultiSearchService(client *Client) *MultiSearchService {
builder := &MultiSearchService{
client: client,
}
return builder
}
func (s *MultiSearchService) Add(requests ...*SearchRequest) *MultiSearchService {
s.requests = append(s.requests, requests...)
return s
}
func (s *MultiSearchService) Index(indices ...string) *MultiSearchService {
s.indices = append(s.indices, indices...)
return s
}
func (s *MultiSearchService) Pretty(pretty bool) *MultiSearchService {
s.pretty = pretty
return s
}
func (s *MultiSearchService) MaxConcurrentSearches(max int) *MultiSearchService {
s.maxConcurrentRequests = &max
return s
}
func (s *MultiSearchService) PreFilterShardSize(size int) *MultiSearchService {
s.preFilterShardSize = &size
return s
}
// RestTotalHitsAsInt is a flag that is temporarily available for ES 7.x
// servers to return total hits as an int64 instead of a response structure.
//
// Warning: Using it indicates that you are using elastic.v6 with ES 7.x,
// which is an unsupported scenario. Use at your own risk.
// This option will also be removed with ES 8.x.
// See https://www.elastic.co/guide/en/elasticsearch/reference/7.0/breaking-changes-7.0.html#hits-total-now-object-search-response.
func (s *MultiSearchService) RestTotalHitsAsInt(v bool) *MultiSearchService {
s.restTotalHitsAsInt = &v
return s
}
func (s *MultiSearchService) Do(ctx context.Context) (*MultiSearchResult, error) {
// Build url
path := "/_msearch"
// Parameters
params := make(url.Values)
if s.pretty {
params.Set("pretty", fmt.Sprint(s.pretty))
}
if v := s.maxConcurrentRequests; v != nil {
params.Set("max_concurrent_searches", fmt.Sprint(*v))
}
if v := s.preFilterShardSize; v != nil {
params.Set("pre_filter_shard_size", fmt.Sprint(*v))
}
if v := s.restTotalHitsAsInt; v != nil {
params.Set("rest_total_hits_as_int", fmt.Sprint(*v))
}
// Set body
var lines []string
for _, sr := range s.requests {
// Set default indices if not specified in the request
if !sr.HasIndices() && len(s.indices) > 0 {
sr = sr.Index(s.indices...)
}
header, err := json.Marshal(sr.header())
if err != nil {
return nil, err
}
body, err := sr.Body()
if err != nil {
return nil, err
}
lines = append(lines, string(header))
lines = append(lines, body)
}
body := strings.Join(lines, "\n") + "\n" // add trailing \n
// Get response
res, err := s.client.PerformRequest(ctx, PerformRequestOptions{
Method: "GET",
Path: path,
Params: params,
Body: body,
})
if err != nil {
return nil, err
}
// Return result
ret := new(MultiSearchResult)
if err := s.client.decoder.Decode(res.Body, ret); err != nil {
return nil, err
}
return ret, nil
}
// MultiSearchResult is the outcome of running a multi-search operation.
type MultiSearchResult struct {
Responses []*SearchResult `json:"responses,omitempty"`
}