From fd96d3a02ef2285aedad638bdfac50470a27a64c Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Tue, 14 Aug 2018 12:34:17 -0500 Subject: [PATCH] This commit ensures that keyed imports are sent to the coordinator node (as opposed to sending to shard0, which may or may not be the coordinator). It adds a `Nodes()` method to the `InternalClient` which is used by the importer to determine which node is the coordinator. --- client.go | 4 ++ ctl/import.go | 107 ++++++------------------------------------------ http/client.go | 80 ++++++++++++++++++++++++------------ http/handler.go | 17 ++++++++ 4 files changed, 86 insertions(+), 122 deletions(-) diff --git a/client.go b/client.go index ccdec31c1..0f64ba8fc 100644 --- a/client.go +++ b/client.go @@ -34,6 +34,7 @@ type InternalClient interface { Schema(ctx context.Context) ([]*IndexInfo, error) CreateIndex(ctx context.Context, index string, opt IndexOptions) error FragmentNodes(ctx context.Context, index string, shard uint64) ([]*Node, error) + Nodes(ctx context.Context) ([]*Node, error) Query(ctx context.Context, index string, queryRequest *QueryRequest) (*QueryResponse, error) QueryNode(ctx context.Context, uri *URI, index string, queryRequest *QueryRequest) (*QueryResponse, error) Import(ctx context.Context, index, field string, shard uint64, bits []Bit) error @@ -89,6 +90,9 @@ func (n nopInternalClient) CreateIndex(ctx context.Context, index string, opt In func (n nopInternalClient) FragmentNodes(ctx context.Context, index string, shard uint64) ([]*Node, error) { return nil, nil } +func (n nopInternalClient) Nodes(ctx context.Context) ([]*Node, error) { + return nil, nil +} func (n nopInternalClient) Query(ctx context.Context, index string, queryRequest *QueryRequest) (*QueryResponse, error) { return nil, nil } diff --git a/ctl/import.go b/ctl/import.go index 95f7b6223..4a744fbf1 100644 --- a/ctl/import.go +++ b/ctl/import.go @@ -224,7 +224,7 @@ func (cmd *ImportCommand) bufferBits(ctx context.Context, useColumnKeys, useRowK // If we've reached the buffer size then import bits. if len(a) == cmd.BufferSize { - if err := cmd.importBits(ctx, a); err != nil { + if err := cmd.importBits(ctx, useColumnKeys, useRowKeys, a); err != nil { return err } a = a[:0] @@ -232,13 +232,22 @@ func (cmd *ImportCommand) bufferBits(ctx context.Context, useColumnKeys, useRowK } // If there are still bits in the buffer then flush them. - return cmd.importBits(ctx, a) + return cmd.importBits(ctx, useColumnKeys, useRowKeys, a) } // importBits sends batches of bits to the server. -func (cmd *ImportCommand) importBits(ctx context.Context, bits []pilosa.Bit) error { +func (cmd *ImportCommand) importBits(ctx context.Context, useColumnKeys, useRowKeys bool, bits []pilosa.Bit) error { logger := log.New(cmd.Stderr, "", log.LstdFlags) + // If keys are used, all bits are sent to the primary translate store (i.e. coordinator). + if useColumnKeys || useRowKeys { + logger.Printf("importing keys: n=%d", len(bits)) + if err := cmd.client.ImportK(ctx, cmd.Index, cmd.Field, bits); err != nil { + return errors.Wrap(err, "importing keys") + } + return nil + } + // Group bits by shard. logger.Printf("grouping %d bits", len(bits)) bitsByShard := http.Bits(bits).GroupByShard() @@ -258,98 +267,6 @@ func (cmd *ImportCommand) importBits(ctx context.Context, bits []pilosa.Bit) err return nil } -// bufferBitsK buffers slices of keys to be imported as a batch. -func (cmd *ImportCommand) bufferBitsK(ctx context.Context, path string) error { - a := make([]pilosa.Bit, 0, cmd.BufferSize) - - var r *csv.Reader - - if path != "-" { - // Open file for reading. - f, err := os.Open(path) - if err != nil { - return errors.Wrap(err, "opening file") - } - defer f.Close() - - // Read rows as bits. - r = csv.NewReader(f) - } else { - r = csv.NewReader(cmd.Stdin) - } - - r.FieldsPerRecord = -1 - rnum := 0 - for { - rnum++ - - // Read CSV row. - record, err := r.Read() - if err == io.EOF { - break - } else if err != nil { - return errors.Wrap(err, "reading") - } - - // Ignore blank rows. - if record[0] == "" { - continue - } else if len(record) < 2 { - return fmt.Errorf("bad column count on row %d: col=%d", rnum, len(record)) - } - - var bit pilosa.Bit - - // Parse row key. - if record[0] == "" { - return fmt.Errorf("invalid row key on row %d: %q", rnum, record[0]) - } - bit.RowKey = record[0] - - // Parse column key. - if record[1] == "" { - return fmt.Errorf("invalid column id on row %d: %q", rnum, record[1]) - } - bit.ColumnKey = record[1] - - // Parse time, if exists. - if len(record) > 2 && record[2] != "" { - t, err := time.Parse(pilosa.TimeFormat, record[2]) - if err != nil { - return fmt.Errorf("invalid timestamp on row %d: %q", rnum, record[2]) - } - bit.Timestamp = t.UnixNano() - } - - a = append(a, bit) - - // If we've reached the buffer size then import bits. - if len(a) == cmd.BufferSize { - if err := cmd.importBitsK(ctx, a); err != nil { - return err - } - a = a[:0] - } - } - - // If there are still bitKs in the buffer then flush them. - return cmd.importBitsK(ctx, a) -} - -// importBitsK sends batches of bitKs to the server. -func (cmd *ImportCommand) importBitsK(ctx context.Context, bits []pilosa.Bit) error { - logger := log.New(cmd.Stderr, "", log.LstdFlags) - - // TODO: does it help to sort the rowKeys? - - logger.Printf("importing keys: n=%d", len(bits)) - if err := cmd.client.ImportK(ctx, cmd.Index, cmd.Field, bits); err != nil { - return errors.Wrap(err, "importing keys") - } - - return nil -} - // bufferValues buffers slices of FieldValues to be imported as a batch. func (cmd *ImportCommand) bufferValues(ctx context.Context, useColumnKeys bool, path string) error { a := make([]pilosa.FieldValue, 0, cmd.BufferSize) diff --git a/http/client.go b/http/client.go index 4ce77041a..5889f0589 100644 --- a/http/client.go +++ b/http/client.go @@ -207,6 +207,37 @@ func (c *InternalClient) FragmentNodes(ctx context.Context, index string, shard return a, nil } +// Nodes returns a list of all nodes. +func (c *InternalClient) Nodes(ctx context.Context) ([]*pilosa.Node, error) { + // Execute request against the host. + u := uriPathToURL(c.defaultURI, "/internal/nodes") + + // Build request. + req, err := http.NewRequest("GET", u.String(), nil) + if err != nil { + return nil, errors.Wrap(err, "creating request") + } + + req.Header.Set("User-Agent", "pilosa/"+pilosa.Version) + req.Header.Set("Accept", "application/json") + + // Execute request. + resp, err := c.httpClient.Do(req.WithContext(ctx)) + if err != nil { + return nil, errors.Wrap(err, "executing request") + } + defer resp.Body.Close() + + var a []*pilosa.Node + if resp.StatusCode != http.StatusOK { + return nil, fmt.Errorf("http: status=%d", resp.StatusCode) + } else if err := json.NewDecoder(resp.Body).Decode(&a); err != nil { + return nil, fmt.Errorf("json decode: %s", err) + } + + return a, nil +} + // Query executes query against the index. func (c *InternalClient) Query(ctx context.Context, index string, queryRequest *pilosa.QueryRequest) (*pilosa.QueryResponse, error) { return c.QueryNode(ctx, c.defaultURI, index, queryRequest) @@ -291,26 +322,42 @@ func (c *InternalClient) Import(ctx context.Context, index, field string, shard return nil } +func getCoordinatorNode(nodes []*pilosa.Node) *pilosa.Node { + for _, node := range nodes { + if node.IsCoordinator { + return node + } + } + return nil +} + // ImportK bulk imports bits specified by string keys to a host. -func (c *InternalClient) ImportK(ctx context.Context, index, field string, columns []pilosa.Bit) error { +func (c *InternalClient) ImportK(ctx context.Context, index, field string, bits []pilosa.Bit) error { if index == "" { return pilosa.ErrIndexRequired } else if field == "" { return pilosa.ErrFieldRequired } - buf, err := c.marshalImportPayloadK(index, field, columns) + buf, err := c.marshalImportPayload(index, field, 0, bits) if err != nil { return fmt.Errorf("Error Creating Payload: %s", err) } - node := &pilosa.Node{ - URI: *c.defaultURI, + // Get the coordinator node; all bits are sent to the + // primary translate store (i.e. coordinator). + nodes, err := c.Nodes(ctx) + if err != nil { + return fmt.Errorf("getting nodes: %s", err) + } + coord := getCoordinatorNode(nodes) + if coord == nil { + return fmt.Errorf("could not find coordinator node") } // Import to node. - if err := c.importNode(ctx, node, index, field, buf); err != nil { - return fmt.Errorf("import node: host=%s, err=%s", node.URI, err) + if err := c.importNode(ctx, coord, index, field, buf); err != nil { + return fmt.Errorf("import node: host=%s, err=%s", coord.URI, err) } return nil @@ -358,27 +405,6 @@ func (c *InternalClient) marshalImportPayload(index, field string, shard uint64, return buf, nil } -// marshalImportPayloadK marshalls the import parameters into a protobuf byte slice. -func (c *InternalClient) marshalImportPayloadK(index, field string, bits []pilosa.Bit) ([]byte, error) { - // Separate row and column IDs to reduce allocations. - rowKeys := Bits(bits).RowKeys() - columnKeys := Bits(bits).ColumnKeys() - timestamps := Bits(bits).Timestamps() - - // Marshal data to protobuf. - buf, err := c.serializer.Marshal(&pilosa.ImportRequest{ - Index: index, - Field: field, - RowKeys: rowKeys, - ColumnKeys: columnKeys, - Timestamps: timestamps, - }) - if err != nil { - return nil, fmt.Errorf("marshal import request: %s", err) - } - return buf, nil -} - // importNode sends a pre-marshaled import request to a node. func (c *InternalClient) importNode(ctx context.Context, node *pilosa.Node, index, field string, buf []byte) error { // Create URL & HTTP request. diff --git a/http/handler.go b/http/handler.go index 80d1563a0..6462f793a 100644 --- a/http/handler.go +++ b/http/handler.go @@ -232,6 +232,7 @@ func newRouter(handler *Handler) *mux.Router { router.HandleFunc("/internal/fragment/nodes", handler.handleGetFragmentNodes).Methods("GET").Name("GetFragmentNodes") router.HandleFunc("/internal/index/{index}/attr/diff", handler.handlePostIndexAttrDiff).Methods("POST") router.HandleFunc("/internal/index/{index}/field/{field}/attr/diff", handler.handlePostFieldAttrDiff).Methods("POST") + router.HandleFunc("/internal/nodes", handler.handleGetNodes).Methods("GET").Name("GetNodes") router.HandleFunc("/internal/shards/max", handler.handleGetShardsMax).Methods("GET") // TODO: deprecate, but it's being used by the client router.HandleFunc("/internal/translate/data", handler.handleGetTranslateData).Methods("GET") @@ -1062,6 +1063,22 @@ func (h *Handler) handleGetFragmentNodes(w http.ResponseWriter, r *http.Request) } } +// handleGetNodes handles /internal/nodes requests. +func (h *Handler) handleGetNodes(w http.ResponseWriter, r *http.Request) { + if !validHeaderAcceptJSON(r.Header) { + http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable) + return + } + + // Retrieve all nodes. + nodes := h.api.Hosts(r.Context()) + + // Write to response. + if err := json.NewEncoder(w).Encode(nodes); err != nil { + h.logger.Printf("json write error: %s", err) + } +} + // handleGetFragmentBlockData handles GET /internal/fragment/block/data requests. func (h *Handler) handleGetFragmentBlockData(w http.ResponseWriter, r *http.Request) { buf, err := h.api.FragmentBlockData(r.Context(), r.Body)