Merge pull request #113 from benbjohnson/sort-import

Add import sorting command
This commit is contained in:
tgruben 2016-09-21 16:30:05 -05:00 committed by GitHub
commit 286762e804
4 changed files with 166 additions and 6 deletions

View file

@ -682,3 +682,12 @@ func (a Bits) GroupBySlice() map[uint64][]Bit {
return m
}
// BitsByPos represents a slice of bits sorted by internal position.
type BitsByPos []Bit
func (p BitsByPos) Swap(i, j int) { p[i], p[j] = p[j], p[i] }
func (p BitsByPos) Len() int { return len(p) }
func (p BitsByPos) Less(i, j int) bool {
return Pos(p[i].BitmapID, p[i].ProfileID) < Pos(p[j].BitmapID, p[j].ProfileID)
}

View file

@ -1,6 +1,7 @@
package main
import (
"bufio"
"encoding/csv"
"errors"
"flag"
@ -10,6 +11,7 @@ import (
"log"
"math/rand"
"os"
"sort"
"strconv"
"strings"
"syscall"
@ -81,6 +83,7 @@ The commands are:
config prints the default configuration
import imports data from a CSV file
export exports data to a CSV file
sort sorts a data file for optimal import speed
backup backs up a frame to an archive file
restore restores a frame from an archive file
inspect inspects fragment data files
@ -112,6 +115,8 @@ func (m *Main) ParseFlags(args []string) error {
m.Cmd = NewImportCommand(m.Stdin, m.Stdout, m.Stderr)
case "export":
m.Cmd = NewExportCommand(m.Stdin, m.Stdout, m.Stderr)
case "sort":
m.Cmd = NewSortCommand(m.Stdin, m.Stdout, m.Stderr)
case "backup":
m.Cmd = NewBackupCommand(m.Stdin, m.Stdout, m.Stderr)
case "restore":
@ -239,6 +244,7 @@ func (cmd *ImportCommand) ParseFlags(args []string) error {
fs.StringVar(&cmd.Host, "host", "localhost:15000", "host:port")
fs.StringVar(&cmd.Database, "d", "", "database")
fs.StringVar(&cmd.Frame, "f", "", "frame")
fs.IntVar(&cmd.BufferSize, "buffer-size", cmd.BufferSize, "buffer size")
if err := fs.Parse(args); err != nil {
return err
}
@ -492,6 +498,112 @@ func (cmd *ExportCommand) Run() error {
return nil
}
// SortCommand represents a command for sorting import data.
type SortCommand struct {
// Filename to sort
Path string
// Standard input/output
Stdin io.Reader
Stdout io.Writer
Stderr io.Writer
}
// NewSortCommand returns a new instance of SortCommand.
func NewSortCommand(stdin io.Reader, stdout, stderr io.Writer) *SortCommand {
return &SortCommand{
Stdin: stdin,
Stdout: stdout,
Stderr: stderr,
}
}
// ParseFlags parses command line flags from args.
func (cmd *SortCommand) ParseFlags(args []string) error {
fs := flag.NewFlagSet("pilosactl", flag.ContinueOnError)
fs.SetOutput(ioutil.Discard)
if err := fs.Parse(args); err != nil {
return err
}
// Extract the data path.
if fs.NArg() == 0 {
return errors.New("path required")
} else if fs.NArg() > 1 {
return errors.New("only one path allowed")
}
cmd.Path = fs.Arg(0)
return nil
}
// Usage returns the usage message to be printed.
func (cmd *SortCommand) Usage() string {
return strings.TrimSpace(`
usage: pilosactl sort PATH
Sorts the import data at PATH into the optimal sort order for importing.
The format of the CSV file is:
BITMAPID,PROFILEID
The file should contain no headers.
`)
}
// Run executes the main program execution.
func (cmd *SortCommand) Run() 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)
a := make([]pilosa.Bit, 0, 1000000)
for {
bitmapID, profileID, err := readCSVRow(r)
if err == io.EOF {
break
} else if err == errBlank {
continue
} else if err != nil {
return err
}
a = append(a, pilosa.Bit{BitmapID: bitmapID, ProfileID: profileID})
}
// 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.BitmapID, 10)
buf = append(buf, ',')
buf = strconv.AppendUint(buf, bit.ProfileID, 10)
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
}
// BackupCommand represents a command for backing up a frame.
type BackupCommand struct {
// Destination host and port.
@ -900,3 +1012,36 @@ func (cmd *BenchCommand) runSetBit(client *pilosa.Client) error {
return nil
}
// readCSVRow reads a bitmap/profile pair from a CSV row.
func readCSVRow(r *csv.Reader) (bitmapID, profileID uint64, err error) {
// Read CSV row.
record, err := r.Read()
if err != nil {
return 0, 0, err
}
// Ignore blank rows.
if record[0] == "" {
return 0, 0, errBlank
} else if len(record) < 2 {
return 0, 0, fmt.Errorf("bad column count: %d", len(record))
}
// Parse bitmap id.
bitmapID, err = strconv.ParseUint(record[0], 10, 64)
if err != nil {
return 0, 0, fmt.Errorf("invalid bitmap id: %q", record[0])
}
// Parse bitmap id.
profileID, err = strconv.ParseUint(record[1], 10, 64)
if err != nil {
return 0, 0, fmt.Errorf("invalid profile id: %q", record[1])
}
return bitmapID, profileID, nil
}
// errBlank indicates a blank row in a CSV file.
var errBlank = errors.New("blank row")

View file

@ -427,8 +427,7 @@ func (f *Fragment) pos(bitmapID, profileID uint64) (uint64, error) {
if profileID < minProfileID || profileID >= minProfileID+SliceWidth {
return 0, errors.New("profile out of bounds")
}
return (bitmapID * SliceWidth) + (profileID % SliceWidth), nil
return Pos(bitmapID, profileID), nil
}
// ForEachBit executes fn for every bit set in the fragment.
@ -1422,3 +1421,8 @@ func byteSlicesEqual(a [][]byte) bool {
}
return true
}
// Pos returns the bitmap position of a bitmap/profile pair.
func Pos(bitmapID, profileID uint64) uint64 {
return (bitmapID * SliceWidth) + (profileID % SliceWidth)
}

View file

@ -247,20 +247,22 @@ func (b *Bitmap) OffsetRange(offset, start, end uint64) *Bitmap {
off := highbits(offset)
hi0, hi1 := highbits(start), highbits(end)
// Find starting container.
n := len(b.containers)
i := sort.Search(n, func(i int) bool { return b.keys[i] >= hi0 })
var other Bitmap
for i, c := range b.containers {
for ; i < n; i++ {
key := b.keys[i]
// If we've exceeded the upper bound then exit.
if key >= hi1 {
break
} else if key < hi0 {
continue
}
// Otherwise append container with offset key.
other.keys = append(other.keys, off+(key-hi0))
other.containers = append(other.containers, c)
other.containers = append(other.containers, b.containers[i])
}
return &other
}