Add sorting flag to import command.

This commit is contained in:
Ben Johnson 2017-06-05 20:24:28 -06:00
parent f066d1e125
commit 1c34ea7f8f
No known key found for this signature in database
GPG key ID: 81741CD251883081
8 changed files with 12 additions and 344 deletions

View file

@ -54,6 +54,7 @@ omitted. If it is present then its format should be YYYY-MM-DDTHH:MM.
flags.StringVarP(&Importer.Index, "index", "i", "", "Pilosa index to import into.")
flags.StringVarP(&Importer.Frame, "frame", "f", "", "Frame to import into.")
flags.IntVarP(&Importer.BufferSize, "buffer-size", "s", 10000000, "Number of bits to buffer/sort before importing.")
flags.BoolVarP(&Importer.Sort, "sort", "", false, "Enables sorting before import.")
return importCmd
}

View file

@ -1,64 +0,0 @@
// Copyright 2017 Pilosa Corp.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package cmd
import (
"context"
"fmt"
"io"
"os"
"github.com/spf13/cobra"
"github.com/pilosa/pilosa/ctl"
)
var Sorter *ctl.SortCommand
func NewSortCommand(stdin io.Reader, stdout, stderr io.Writer) *cobra.Command {
Sorter = ctl.NewSortCommand(os.Stdin, os.Stdout, os.Stderr)
sortCmd := &cobra.Command{
Use: "sort <path>",
Short: "Sort import data for optimal import performance.",
Long: `
Sorts the import data at PATH into the optimal sort order for importing.
The format of the CSV file is:
ROWID,COLUMNID
The file should contain no headers.
`,
RunE: func(cmd *cobra.Command, args []string) error {
if len(args) == 0 {
return fmt.Errorf("path required")
} else if len(args) > 1 {
return fmt.Errorf("only one path supported")
}
Sorter.Path = args[0]
if err := Sorter.Run(context.Background()); err != nil {
return err
}
return nil
},
}
return sortCmd
}
func init() {
subcommandFns["sort"] = NewSortCommand
}

View file

@ -1,43 +0,0 @@
// Copyright 2017 Pilosa Corp.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package cmd_test
import (
"strings"
"testing"
)
func TestSortHelp(t *testing.T) {
output, err := ExecNewRootCommand(t, "sort", "--help")
if !strings.Contains(output, "Usage:") ||
!strings.Contains(output, "Flags:") ||
!strings.Contains(output, "pilosa sort") || err != nil {
t.Fatalf("Command 'sort --help' not working, err: '%v', output: '%s'", err, output)
}
}
func TestSortNoPath(t *testing.T) {
output, err := ExecNewRootCommand(t, "sort")
if !strings.Contains(err.Error(), "path required") {
t.Fatalf("Command 'sort' without args should error but: err: '%v', output: '%v'", err, output)
}
}
func TestSortMultiPath(t *testing.T) {
output, err := ExecNewRootCommand(t, "sort", "one", "two")
if !strings.Contains(err.Error(), "only one path") {
t.Fatalf("Command 'sort' without args should error but: err: '%v', output: '%v'", err, output)
}
}

View file

@ -22,6 +22,7 @@ import (
"io"
"log"
"os"
"sort"
"strconv"
"time"
@ -43,6 +44,9 @@ type ImportCommand struct {
// Size of buffer used to chunk import.
BufferSize int `json:"bufferSize"`
// Enables sorting of data file before import.
Sort bool `json:"sort"`
// Reusable client.
Client *pilosa.Client `json:"-"`
@ -185,6 +189,10 @@ func (cmd *ImportCommand) importBits(ctx context.Context, bits []pilosa.Bit) err
// Parse path into bits.
for slice, bits := range bitsBySlice {
if cmd.Sort {
sort.Sort(pilosa.BitsByPos(bits))
}
logger.Printf("importing slice: %d, n=%d", slice, len(bits))
if err := cmd.Client.Import(ctx, cmd.Index, cmd.Frame, slice, bits); err != nil {
return err

View file

@ -1,148 +0,0 @@
// Copyright 2017 Pilosa Corp.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package ctl
import (
"bufio"
"context"
"encoding/csv"
"errors"
"fmt"
"io"
"os"
"sort"
"strconv"
"time"
"github.com/pilosa/pilosa"
)
// SortCommand represents a command for sorting import data.
type SortCommand struct {
// Filename to sort
Path string
// Standard input/output
*pilosa.CmdIO
}
// NewSortCommand returns a new instance of SortCommand.
func NewSortCommand(stdin io.Reader, stdout, stderr io.Writer) *SortCommand {
return &SortCommand{
CmdIO: pilosa.NewCmdIO(stdin, stdout, stderr),
}
}
// Run executes the sort command.
func (cmd *SortCommand) Run(ctx context.Context) error {
// Open file for reading.
f, err := os.Open(cmd.Path)
if err != nil {
return err
}
defer f.Close()
// Read rows as bits.
r := csv.NewReader(f)
r.FieldsPerRecord = -1
a := make([]pilosa.Bit, 0, 1000000)
for {
rowID, columnID, timestamp, err := readCSVRow(r)
if err == io.EOF {
break
} else if err == errBlank {
continue
} else if err != nil {
return err
}
a = append(a, pilosa.Bit{RowID: rowID, ColumnID: columnID, Timestamp: timestamp})
}
// Sort bits by position.
sort.Sort(pilosa.BitsByPos(a))
// Rewrite to STDOUT.
w := bufio.NewWriter(cmd.Stdout)
buf := make([]byte, 0, 1024)
for _, bit := range a {
// Write CSV to buffer.
buf = buf[:0]
buf = strconv.AppendUint(buf, bit.RowID, 10)
buf = append(buf, ',')
buf = strconv.AppendUint(buf, bit.ColumnID, 10)
if bit.Timestamp != 0 {
buf = append(buf, ',')
buf = append(buf, time.Unix(0, bit.Timestamp).UTC().Format(pilosa.TimeFormat)...)
}
buf = append(buf, '\n')
// Write to output.
if _, err := w.Write(buf); err != nil {
return err
}
}
// Ensure buffer is flushed before exiting.
if err := w.Flush(); err != nil {
return err
}
return nil
}
// readCSVRow reads a row/column pair from a CSV row.
func readCSVRow(r *csv.Reader) (rowID, columnID uint64, timestamp int64, err error) {
// Read CSV row.
record, err := r.Read()
if err != nil {
return 0, 0, 0, err
}
// Ignore blank rows.
if record[0] == "" {
return 0, 0, 0, errBlank
} else if len(record) < 2 {
return 0, 0, 0, fmt.Errorf("bad column count: %d", len(record))
}
// Parse row id.
rowID, err = strconv.ParseUint(record[0], 10, 64)
if err != nil {
return 0, 0, 0, fmt.Errorf("invalid row id: %q", record[0])
}
// Parse column id.
columnID, err = strconv.ParseUint(record[1], 10, 64)
if err != nil {
return 0, 0, 0, fmt.Errorf("invalid column id: %q", record[1])
}
// Parse timestamp, if available.
if len(record) > 2 && record[2] != "" {
t, err := time.Parse(pilosa.TimeFormat, record[2])
if err != nil {
return 0, 0, 0, fmt.Errorf("invalid timestamp: %q", record[2])
}
timestamp = t.UnixNano()
}
return rowID, columnID, timestamp, nil
}
// errBlank indicates a blank row in a CSV file.
var errBlank = errors.New("blank row")

View file

@ -1,83 +0,0 @@
// Copyright 2017 Pilosa Corp.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package ctl
import (
"bytes"
"golang.org/x/net/context"
"io"
"io/ioutil"
"os"
"strings"
"testing"
)
func TestSortCommand_Run(t *testing.T) {
file, _ := ioutil.TempFile("", "file.csv")
content := "3,3\n1,2\n2,4"
file.Write([]byte(content))
file.Close()
rder := []byte{}
stdin := bytes.NewReader(rder)
r, w, _ := os.Pipe()
cm := NewSortCommand(stdin, w, w)
cm.Path = file.Name()
err := cm.Run(context.Background())
w.Close()
var buf bytes.Buffer
io.Copy(&buf, r)
if err != nil {
t.Fatal(err)
} else if !strings.Contains(buf.String(), "1,2\n2,4\n3,3") {
t.Fatalf("File is not sorted, actual result: %s", buf.String())
}
}
func TestSortCommand_InvalidFile(t *testing.T) {
buf := bytes.Buffer{}
stdin, stdout, stderr := GetIO(buf)
file, _ := ioutil.TempFile("", "file.csv")
file.Write([]byte("3,3\na,8\n2,4"))
file.Close()
cm := NewSortCommand(stdin, stdout, stderr)
cm.Path = file.Name()
err := cm.Run(context.Background())
if !strings.Contains(err.Error(), "invalid row id") {
t.Fatalf("expect err: invalid row id, actual: %s", err)
}
file, _ = ioutil.TempFile("", "file.csv")
file.Write([]byte("3,3\n1,a\n2,4"))
file.Close()
cm.Path = file.Name()
err = cm.Run(context.Background())
if !strings.Contains(err.Error(), "invalid column id") {
t.Fatalf("expect err: invalid column id, actual: %s", err)
}
file, _ = ioutil.TempFile("", "file.csv")
file.Write([]byte("3,3,1234\n1,2,34345\n2,4"))
file.Close()
cm.Path = file.Name()
err = cm.Run(context.Background())
if !strings.Contains(err.Error(), "invalid timestamp") {
t.Fatalf("expect err: invalid timestamp, actual: %s", err)
}
}

View file

@ -36,9 +36,10 @@ While Pilosa does have some high system requirements it is not a best practice t
The import API expects a csv of RowID,ColumnID's.
When importing large datasets remember it is much faster to pre sort the data by RowID and then by ColumnID in ascending order. You can use `pilosa sort CSV_FILE` to do that. Also, avoid querying Pilosa until the import is complete, otherwise you will experience inconsistent results.
When importing large datasets remember it is much faster to pre sort the data by RowID and then by ColumnID in ascending order. You can use the `--sort` flag to do that. Also, avoid querying Pilosa until the import is complete, otherwise you will experience inconsistent results.
```
pilosa import -d project -f stargazer project-stargazer.csv
pilosa import --sort -d project -f stargazer project-stargazer.csv
```
#### Exporting

View file

@ -61,7 +61,6 @@ There are three ways to install Pilosa on MacOS: download the binary (recommende
inspect Get stats on a pilosa data file.
restore Restore data to pilosa from a backup file.
server Run Pilosa.
sort Sort import data for optimal import performance.
Flags:
-c, --config string Configuration file to read from.
@ -121,7 +120,6 @@ There are three ways to install Pilosa on MacOS: download the binary (recommende
inspect Get stats on a pilosa data file.
restore Restore data to pilosa from a backup file.
server Run Pilosa.
sort Sort import data for optimal import performance.
Flags:
-c, --config string Configuration file to read from.
@ -211,7 +209,6 @@ There are three ways to install Pilosa on Linux: download the binary (recommende
inspect Get stats on a pilosa data file.
restore Restore data to pilosa from a backup file.
server Run Pilosa.
sort Sort import data for optimal import performance.
Flags:
-c, --config string Configuration file to read from.
@ -271,7 +268,6 @@ There are three ways to install Pilosa on Linux: download the binary (recommende
inspect Get stats on a pilosa data file.
restore Restore data to pilosa from a backup file.
server Run Pilosa.
sort Sort import data for optimal import performance.
Flags:
-c, --config string Configuration file to read from.