mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-10 12:57:54 +00:00
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.
This commit is contained in:
parent
3793719485
commit
fd96d3a02e
4 changed files with 86 additions and 122 deletions
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
107
ctl/import.go
107
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)
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue