update traces, es transport with batches and fasthttp
This commit is contained in:
@@ -4,15 +4,19 @@ import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"bytes"
|
||||
"net/url"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/elastic/go-elasticsearch/v7"
|
||||
"github.com/elastic/go-elasticsearch/v7/esapi"
|
||||
"github.com/elastic/go-elasticsearch/v7/esutil"
|
||||
)
|
||||
|
||||
const (
|
||||
ElasticSearchDefaultIndex = "lotus-pubsub"
|
||||
flushInterval = 10 * time.Second
|
||||
flushBytes = 1024 * 1024 // MB
|
||||
esWorkers = 2 // TODO: hardcoded
|
||||
)
|
||||
|
||||
func NewElasticSearchTransport(connectionString string, elasticsearchIndex string) (TracerTransport, error) {
|
||||
@@ -30,10 +34,10 @@ func NewElasticSearchTransport(connectionString string, elasticsearchIndex strin
|
||||
},
|
||||
Username: username,
|
||||
Password: password,
|
||||
Transport: &FastHttpTransport{},
|
||||
}
|
||||
|
||||
es, err := elasticsearch.NewClient(cfg)
|
||||
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -45,14 +49,31 @@ func NewElasticSearchTransport(connectionString string, elasticsearchIndex strin
|
||||
esIndex = ElasticSearchDefaultIndex
|
||||
}
|
||||
|
||||
// Create the BulkIndexer to batch ES trace submission
|
||||
bi, err := esutil.NewBulkIndexer(esutil.BulkIndexerConfig{
|
||||
Index: esIndex,
|
||||
Client: es,
|
||||
NumWorkers: esWorkers,
|
||||
FlushBytes: int(flushBytes),
|
||||
FlushInterval: flushInterval,
|
||||
OnError: func(ctx context.Context, err error) {
|
||||
log.Errorf("Error persisting queries %s", err.Error())
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &elasticSearchTransport{
|
||||
cl: es,
|
||||
cl: es,
|
||||
bi: bi,
|
||||
esIndex: esIndex,
|
||||
}, nil
|
||||
}
|
||||
|
||||
type elasticSearchTransport struct {
|
||||
cl *elasticsearch.Client
|
||||
bi esutil.BulkIndexer
|
||||
esIndex string
|
||||
}
|
||||
|
||||
@@ -72,26 +93,18 @@ func (est *elasticSearchTransport) Transport(evt TracerTransportEvent) error {
|
||||
return fmt.Errorf("error while marshaling event: %s", err)
|
||||
}
|
||||
|
||||
req := esapi.IndexRequest{
|
||||
Index: est.esIndex,
|
||||
Body: strings.NewReader(string(jsonEvt)),
|
||||
Refresh: "true",
|
||||
}
|
||||
|
||||
// Perform the request with the client.
|
||||
res, err := req.Do(context.Background(), est.cl)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
err = res.Body.Close()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if res.IsError() {
|
||||
return fmt.Errorf("[%s] Error indexing document ID=%s", res.Status(), req.DocumentID)
|
||||
}
|
||||
|
||||
return nil
|
||||
return est.bi.Add(
|
||||
context.Background(),
|
||||
esutil.BulkIndexerItem{
|
||||
Action: "index",
|
||||
Body: bytes.NewReader(jsonEvt),
|
||||
OnFailure: func(ctx context.Context, item esutil.BulkIndexerItem, res esutil.BulkIndexerResponseItem, err error) {
|
||||
if err != nil {
|
||||
log.Errorf("unable to submit trace - %s", err)
|
||||
} else {
|
||||
log.Errorf("unable to submit trace %s: %s", res.Error.Type, res.Error.Reason)
|
||||
}
|
||||
},
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,74 @@
|
||||
package tracer
|
||||
|
||||
import (
|
||||
"io/ioutil"
|
||||
"net/http"
|
||||
"strings"
|
||||
|
||||
"github.com/valyala/fasthttp"
|
||||
)
|
||||
|
||||
// Transport implements the elastictransport interface with
|
||||
// the github.com/valyala/fasthttp HTTP client.
|
||||
type FastHttpTransport struct{}
|
||||
|
||||
// RoundTrip performs the request and returns a response or error
|
||||
func (t *FastHttpTransport) RoundTrip(req *http.Request) (*http.Response, error) {
|
||||
freq := fasthttp.AcquireRequest()
|
||||
defer fasthttp.ReleaseRequest(freq)
|
||||
|
||||
fres := fasthttp.AcquireResponse()
|
||||
defer fasthttp.ReleaseResponse(fres)
|
||||
|
||||
t.copyRequest(freq, req)
|
||||
|
||||
err := fasthttp.Do(freq, fres)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
res := &http.Response{Header: make(http.Header)}
|
||||
t.copyResponse(res, fres)
|
||||
|
||||
return res, nil
|
||||
}
|
||||
|
||||
// copyRequest converts a http.Request to fasthttp.Request
|
||||
func (t *FastHttpTransport) copyRequest(dst *fasthttp.Request, src *http.Request) *fasthttp.Request {
|
||||
if src.Method == "GET" && src.Body != nil {
|
||||
src.Method = "POST"
|
||||
}
|
||||
|
||||
dst.SetHost(src.Host)
|
||||
dst.SetRequestURI(src.URL.String())
|
||||
|
||||
dst.Header.SetRequestURI(src.URL.String())
|
||||
dst.Header.SetMethod(src.Method)
|
||||
|
||||
for k, vv := range src.Header {
|
||||
for _, v := range vv {
|
||||
dst.Header.Set(k, v)
|
||||
}
|
||||
}
|
||||
|
||||
if src.Body != nil {
|
||||
dst.SetBodyStream(src.Body, -1)
|
||||
}
|
||||
|
||||
return dst
|
||||
}
|
||||
|
||||
// copyResponse converts a http.Response to fasthttp.Response
|
||||
func (t *FastHttpTransport) copyResponse(dst *http.Response, src *fasthttp.Response) *http.Response {
|
||||
dst.StatusCode = src.StatusCode()
|
||||
|
||||
src.Header.VisitAll(func(k, v []byte) {
|
||||
dst.Header.Set(string(k), string(v))
|
||||
})
|
||||
|
||||
// Cast to a string to make a copy seeing as src.Body() won't
|
||||
// be valid after the response is released back to the pool (fasthttp.ReleaseResponse).
|
||||
dst.Body = ioutil.NopCloser(strings.NewReader(string(src.Body())))
|
||||
|
||||
return dst
|
||||
}
|
||||
Reference in New Issue
Block a user