diff --git a/.circleci/config.yml b/.circleci/config.yml
index c5bf059e2..b42b2aa64 100644
--- a/.circleci/config.yml
+++ b/.circleci/config.yml
@@ -68,6 +68,7 @@ jobs:
docker:
- image: circleci/python:2.7-jessie
steps:
+ - run: '[[ -v CIRCLE_PR_NUMBER ]] && circleci step halt || true' # Skip job if this is a PR
- *fast-checkout
- run: sudo pip install awscli
- run: make prerelease-upload
diff --git a/cmd/server_test.go b/cmd/server_test.go
index e58f8af4d..84b2297d4 100644
--- a/cmd/server_test.go
+++ b/cmd/server_test.go
@@ -42,7 +42,7 @@ func TestServerConfig(t *testing.T) {
tests := []commandTest{
// TEST 0
{
- args: []string{"server", "--data-dir", actualDataDir, "--cluster.hosts", "localhost:10111,localhost:10110", "--bind", "localhost:10111"},
+ args: []string{"server", "--data-dir", actualDataDir, "--cluster.hosts", "localhost:10111,localhost:10110", "--bind", "localhost:10111", "--translation.map-size", "100000"},
env: map[string]string{"PILOSA_DATA_DIR": "/tmp/myEnvDatadir", "PILOSA_CLUSTER_LONG_QUERY_TIME": "1m30s", "PILOSA_MAX_WRITES_PER_REQUEST": "2000"},
cfgFileContent: `
data-dir = "/tmp/myFileDatadir"
@@ -65,13 +65,14 @@ func TestServerConfig(t *testing.T) {
v.Check(cmd.Server.Config.Cluster.Hosts, []string{"localhost:10111", "localhost:10110"})
v.Check(cmd.Server.Config.Cluster.LongQueryTime, toml.Duration(time.Second*90))
v.Check(cmd.Server.Config.MaxWritesPerRequest, 2000)
+ v.Check(cmd.Server.Config.Translation.MapSize, 100000)
return v.Error()
},
},
// TEST 1
{
args: []string{"server", "--anti-entropy.interval", "9m0s"},
- env: map[string]string{"PILOSA_CLUSTER_HOSTS": "localhost:1110,localhost:1111", "PILOSA_BIND": "localhost:1110"},
+ env: map[string]string{"PILOSA_CLUSTER_HOSTS": "localhost:1110,localhost:1111", "PILOSA_BIND": "localhost:1110", "PILOSA_TRANSLATION_MAP_SIZE": "100000"},
cfgFileContent: `
bind = "localhost:0"
data-dir = "` + actualDataDir + `"
@@ -85,12 +86,13 @@ func TestServerConfig(t *testing.T) {
v := validator{}
v.Check(cmd.Server.Config.Cluster.Hosts, []string{"localhost:1110", "localhost:1111"})
v.Check(cmd.Server.Config.AntiEntropy.Interval, toml.Duration(time.Minute*9))
+ v.Check(cmd.Server.Config.Translation.MapSize, 100000)
return v.Error()
},
},
// TEST 2
{
- args: []string{"server", "--log-path", logFile.Name(), "--cluster.disabled", "true"},
+ args: []string{"server", "--log-path", logFile.Name(), "--cluster.disabled", "true", "--translation.map-size", "100000"},
env: map[string]string{},
cfgFileContent: `
bind = "localhost:19444"
diff --git a/ctl/import_test.go b/ctl/import_test.go
index cd75a5d0a..2913005fb 100644
--- a/ctl/import_test.go
+++ b/ctl/import_test.go
@@ -279,3 +279,49 @@ func TestImportCommand_BugOverwriteValue(t *testing.T) {
t.Fatalf("Import Run with values doesn't work: %s", err)
}
}
+
+// Ensure that import into bool field runs.
+func TestImportCommand_RunBool(t *testing.T) {
+ buf := bytes.Buffer{}
+ stdin, stdout, stderr := GetIO(buf)
+ cm := NewImportCommand(stdin, stdout, stderr)
+ ctx := context.Background()
+
+ cmd := test.MustRunCluster(t, 1)[0]
+ cm.Host = cmd.API.Node().URI.HostPort()
+
+ http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+cm.Host+"/index/i", strings.NewReader("")))
+ http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+cm.Host+"/index/i/field/f", strings.NewReader(`{"options":{"type": "bool"}}`)))
+
+ cm.Index = "i"
+ cm.Field = "f"
+
+ t.Run("Valid", func(t *testing.T) {
+ file, err := ioutil.TempFile("", "import-bool.csv")
+ if err != nil {
+ t.Fatal(err)
+ }
+ file.Write([]byte("0,1\n1,2\n1,3"))
+
+ cm.Paths = []string{file.Name()}
+ err = cm.Run(ctx)
+ if err != nil {
+ t.Fatalf("Import Run to bool field doesn't work: %s", err)
+ }
+ })
+
+ // Ensure that invalid bool values return an error.
+ t.Run("Invalid", func(t *testing.T) {
+ file, err := ioutil.TempFile("", "import-invalid-bool.csv")
+ if err != nil {
+ t.Fatal(err)
+ }
+ file.Write([]byte("0,1\n1,2\n1,3\n2,4"))
+
+ cm.Paths = []string{file.Name()}
+ err = cm.Run(ctx)
+ if !strings.Contains(err.Error(), "bool field imports only support values 0 and 1") {
+ t.Fatalf("expect error: bool field imports only support values 0 and 1, actual: %s", err)
+ }
+ })
+}
diff --git a/ctl/server.go b/ctl/server.go
index 90fdfb46e..7f384ec7b 100644
--- a/ctl/server.go
+++ b/ctl/server.go
@@ -45,6 +45,7 @@ func BuildServerFlags(cmd *cobra.Command, srv *server.Command) {
// Translation
flags.StringVarP(&srv.Config.Translation.PrimaryURL, "translation.primary-url", "", srv.Config.Translation.PrimaryURL, "DEPRECATED: URL for primary translation node for replication.")
+ flags.IntVarP(&srv.Config.Translation.MapSize, "translation.map-size", "", srv.Config.Translation.MapSize, "Size in bytes of mmap to allocate for key translation.")
// Gossip
flags.StringVarP(&srv.Config.Gossip.Port, "gossip.port", "", srv.Config.Gossip.Port, "Port to which pilosa should bind for internal state sharing.")
diff --git a/docs/administration.md b/docs/administration.md
index cd4e30184..e626cc0e5 100644
--- a/docs/administration.md
+++ b/docs/administration.md
@@ -64,6 +64,20 @@ If you are using [integer](../data-model/#bsi-range-encoding) field values, the
pilosa import -i project -f stargazer-counts project-stargazer-counts.csv
```
+##### Importing Boolean Values
+
+If you are using a [boolean](../data-model/#boolean) field, the CSV file should be in the format `Boolean,Value`, where `Boolean` is either `0` (false) or `1` (true).
+
+For example, importing a file with the following contents will result in columns 3 and 9 being set in the `false` row, and columns 1, 2, 4, and 8 being set in the `true` row.
+```
+0,3
+0,9
+1,1
+1,2
+1,4
+1,8
+```
+
Note that you must first create a field. View Create Field for more details. The `-e` flag can create the necessary schema when using a field of type "set".
diff --git a/docs/api-reference.md b/docs/api-reference.md
index 890b37941..216a01bca 100644
--- a/docs/api-reference.md
+++ b/docs/api-reference.md
@@ -108,6 +108,8 @@ The request payload is in JSON, and may contain the `options` field. The `option
* `int`
* `min` (int): Minimum integer value allowed for the field.
* `max` (int): Maximum integer value allowed for the field.
+* `bool`
+ * (boolean fields take no arguments)
* `time`
* `timeQuantum` (string): [Time Quantum](../data-model/#time-quantum) for this field.
* `mutex`
diff --git a/docs/configuration.md b/docs/configuration.md
index 22a67ddcf..5750441ef 100644
--- a/docs/configuration.md
+++ b/docs/configuration.md
@@ -307,6 +307,18 @@ The config file is in the [toml format](https://github.com/toml-lang/toml) and h
skip-verify = true
```
+#### Translation Map Size
+
+* Description: Size in bytes of mmap to allocate for key translation
+* Flag: `translation.map-size`
+* Env: `PILOSA_TRANSLATION_MAP_SIZE`
+* Config:
+
+ ```toml
+ [translation]
+ map-size = 10737418240
+ ```
+
### Example Cluster Configuration
A three node cluster running on different hosts could be minimally configured as follows:
diff --git a/docs/data-model.md b/docs/data-model.md
index ba7c12024..bababdd34 100644
--- a/docs/data-model.md
+++ b/docs/data-model.md
@@ -117,7 +117,7 @@ Query operations run in parallel, and they are evenly distributed across a clust
### Field Type
-Upon creation, fields are configured to be of a certain type. Pilosa supports the following field types: `set`, `int`, `time`, and `mutex`.
+Upon creation, fields are configured to be of a certain type. Pilosa supports the following field types: `set`, `int`, `bool`, `time`, and `mutex`.
#### Set
@@ -192,3 +192,7 @@ Set(3, A=8, 2017-05-19T00:00)
#### Mutex
Mutex fields are similar to `set` fields, with the distinction of requiring the row value for each column to be mutually exclusive. In other words, each column can only have a single value for the field. If the field value for a column is updated on a `mutex` field, then the previous field value for that column will be cleared. This field type is like a field in an RDBMS table where every record contains a single value for a particular field.
+
+#### Boolean
+
+A boolean field is similar to a `mutex` field tracking only two values: `true` and `false`. Boolean fields do not maintain a sorted cache, nor do they support key values.
diff --git a/docs/glossary.md b/docs/glossary.md
index c39a0bf6f..fccbe0ae8 100644
--- a/docs/glossary.md
+++ b/docs/glossary.md
@@ -22,7 +22,7 @@ nav = []
Fragment: A Fragment is the intersection of a [field](#field) and a [shard](#shard) in an [index](#index).
-[Field](../data-model/#field): Fields are used to group [rows](#row) into different categories. Row IDs are namespaced by field such that the same row ID in a different field refers to a different row. For [ranked](#topn) fields, rows are kept in sorted order within the field. Fields are one of four types: set, [int](#bsi), time, and mutex. For more information, see [data model](../data-model/) and [Creating fields](../api-reference/#create-field).
+[Field](../data-model/#field): Fields are used to group [rows](#row) into different categories. Row IDs are namespaced by field such that the same row ID in a different field refers to a different row. For [ranked](#topn) fields, rows are kept in sorted order within the field. Fields are one of four types: set, [int](#bsi), bool, time, and mutex. For more information, see [data model](../data-model/) and [Creating fields](../api-reference/#create-field).
[Frame](../data-model/#field): Prior to Pilosa 1.0, fields were known as frames.
diff --git a/docs/tutorials.md b/docs/tutorials.md
index 040d85a1c..47344db4f 100644
--- a/docs/tutorials.md
+++ b/docs/tutorials.md
@@ -483,7 +483,7 @@ curl localhost:10101/index/patients/field/tcells \
{"success":true}
```
-Next, let's populate our fields with data. There are two ways to get data into fields: use the `SetFieldValue()` PQL function to set fields individually, or use the `pilosa import` command to import many values at once. First, let's set some field data using PQL.
+Next, let's populate our fields with data. There are two ways to get data into fields: use the `Set()` PQL function to set fields individually, or use the `pilosa import` command to import many values at once. First, let's set some field data using PQL.
The following queries set the age, weight, and t-cell count for the patient with ID `1` in our system:
``` request
diff --git a/executor.go b/executor.go
index b0083a514..ff3db909d 100644
--- a/executor.go
+++ b/executor.go
@@ -734,8 +734,7 @@ func (e *executor) executeBitmapShard(_ context.Context, index string, c *pql.Ca
rowID, rowOK, rowErr := c.UintArg(fieldName)
if rowErr != nil {
return nil, fmt.Errorf("Row() error with arg for row: %v", rowErr)
- }
- if !rowOK {
+ } else if !rowOK {
return nil, fmt.Errorf("Row() must specify %v", rowLabel)
}
@@ -1140,7 +1139,6 @@ func (e *executor) executeClearBitField(ctx context.Context, index string, c *pq
// executeSet executes a Set() call.
func (e *executor) executeSet(ctx context.Context, index string, c *pql.Call, opt *execOptions) (bool, error) {
-
// Read colID.
colID, ok, err := c.UintArg("_" + columnLabel)
if err != nil {
@@ -1172,8 +1170,9 @@ func (e *executor) executeSet(ctx context.Context, index string, c *pql.Call, op
}
}
+ // Int field.
if f.Type() == FieldTypeInt {
- // Read remaining fields using labels.
+ // Read row value.
rowVal, ok, err := c.IntArg(fieldName)
if err != nil {
return false, fmt.Errorf("reading Set() row: %v", err)
@@ -1182,27 +1181,27 @@ func (e *executor) executeSet(ctx context.Context, index string, c *pql.Call, op
}
return e.executeSetValueField(ctx, index, c, f, colID, rowVal, opt)
- } else {
- // Read remaining fields using labels.
- rowID, ok, err := c.UintArg(fieldName)
- if err != nil {
- return false, fmt.Errorf("reading Set() row: %v", err)
- } else if !ok {
- return false, fmt.Errorf("Set() row argument '%v' required", rowLabel)
- }
-
- var timestamp *time.Time
- sTimestamp, ok := c.Args["_timestamp"].(string)
- if ok {
- t, err := time.Parse(TimeFormat, sTimestamp)
- if err != nil {
- return false, fmt.Errorf("invalid date: %s", sTimestamp)
- }
- timestamp = &t
- }
-
- return e.executeSetBitField(ctx, index, c, f, colID, rowID, timestamp, opt)
}
+
+ // Read row ID.
+ rowID, ok, err := c.UintArg(fieldName)
+ if err != nil {
+ return false, fmt.Errorf("reading Set() row: %v", err)
+ } else if !ok {
+ return false, fmt.Errorf("Set() row argument '%v' required", rowLabel)
+ }
+
+ var timestamp *time.Time
+ sTimestamp, ok := c.Args["_timestamp"].(string)
+ if ok {
+ t, err := time.Parse(TimeFormat, sTimestamp)
+ if err != nil {
+ return false, fmt.Errorf("invalid date: %s", sTimestamp)
+ }
+ timestamp = &t
+ }
+
+ return e.executeSetBitField(ctx, index, c, f, colID, rowID, timestamp, opt)
}
// executeSetBitField executes a Set() call for a specific field.
@@ -1675,7 +1674,21 @@ func (e *executor) translateCall(index string, idx *Index, c *pql.Call) error {
// will raise an error downstream when it's used.
return nil
}
- if field.keys() {
+
+ // Bool field keys do not use the translator because there
+ // are only two possible values. Instead, they are handled
+ // directly.
+ if field.Type() == FieldTypeBool {
+ boolVal, err := callArgBool(c, rowKey)
+ if err != nil {
+ return errors.Wrap(err, "getting bool key")
+ }
+ rowID := falseRowID
+ if boolVal {
+ rowID = trueRowID
+ }
+ c.Args[rowKey] = rowID
+ } else if field.keys() {
if c.Args[rowKey] != nil && !isString(c.Args[rowKey]) {
return errors.New("row value must be a string when field 'keys' option enabled")
}
@@ -1832,6 +1845,18 @@ func (vc *ValCount) larger(other ValCount) ValCount {
}
}
+func callArgBool(call *pql.Call, key string) (bool, error) {
+ value, ok := call.Args[key]
+ if !ok {
+ return false, errors.New("missing bool argument")
+ }
+ b, ok := value.(bool)
+ if !ok {
+ return false, fmt.Errorf("invalid bool argument type: %T", value)
+ }
+ return b, nil
+}
+
func callArgString(call *pql.Call, key string) string {
value, ok := call.Args[key]
if !ok {
diff --git a/executor_test.go b/executor_test.go
index ee919f001..828a40b09 100644
--- a/executor_test.go
+++ b/executor_test.go
@@ -376,6 +376,78 @@ func TestExecutor_Execute_SetBit(t *testing.T) {
})
}
+// Ensure a set query can be executed on a bool field.
+func TestExecutor_Execute_SetBool(t *testing.T) {
+ t.Run("Basic", func(t *testing.T) {
+ c := test.MustRunCluster(t, 1)
+ defer c.Close()
+ hldr := test.Holder{Holder: c[0].Server.Holder()}
+
+ // Create fields.
+ index := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{})
+ if _, err := index.CreateFieldIfNotExists("f", pilosa.OptFieldTypeBool()); err != nil {
+ t.Fatal(err)
+ }
+
+ // Set a true bit.
+ if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(100, f=true)`}); err != nil {
+ t.Fatal(err)
+ } else if !res.Results[0].(bool) {
+ t.Fatalf("expected column changed")
+ }
+
+ // Set the same bit to true again verify nothing changed.
+ if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(100, f=true)`}); err != nil {
+ t.Fatal(err)
+ } else if res.Results[0].(bool) {
+ t.Fatalf("expected column to be unchanged")
+ }
+
+ // Set the same bit to false.
+ if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(100, f=false)`}); err != nil {
+ t.Fatal(err)
+ } else if !res.Results[0].(bool) {
+ t.Fatalf("expected column changed")
+ }
+
+ // Ensure that the false row is set.
+ if result, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Row(f=false)`}); err != nil {
+ t.Fatal(err)
+ } else if columns := result.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(columns, []uint64{100}) {
+ t.Fatalf("unexpected colums: %+v", columns)
+ }
+
+ // Ensure that the true row is empty.
+ if result, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Row(f=true)`}); err != nil {
+ t.Fatal(err)
+ } else if columns := result.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(columns, []uint64{}) {
+ t.Fatalf("unexpected colums: %+v", columns)
+ }
+ })
+ t.Run("Error", func(t *testing.T) {
+ c := test.MustRunCluster(t, 1)
+ defer c.Close()
+ hldr := test.Holder{Holder: c[0].Server.Holder()}
+
+ // Create fields.
+ index := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{})
+ if _, err := index.CreateFieldIfNotExists("f", pilosa.OptFieldTypeBool()); err != nil {
+ t.Fatal(err)
+ }
+
+ // Set bool using a string value.
+ if _, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(100, f="true")`}); err == nil {
+ t.Fatalf("expected invalid bool type error")
+ }
+
+ // Set bool using an integer.
+ if _, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(100, f=1)`}); err == nil {
+ t.Fatalf("expected invalid bool type error")
+ }
+
+ })
+}
+
// Ensure old PQL syntax doesn't break anything too badly.
func TestExecutor_Execute_OldPQL(t *testing.T) {
c := test.MustRunCluster(t, 1)
@@ -397,7 +469,7 @@ func TestExecutor_Execute_SetValue(t *testing.T) {
defer c.Close()
hldr := test.Holder{Holder: c[0].Server.Holder()}
- // Create felds.
+ // Create fields.
index := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{})
if _, err := index.CreateFieldIfNotExists("f", pilosa.OptFieldTypeInt(0, 50)); err != nil {
t.Fatal(err)
diff --git a/field.go b/field.go
index d038b606c..822522bc4 100644
--- a/field.go
+++ b/field.go
@@ -52,6 +52,7 @@ const (
FieldTypeInt = "int"
FieldTypeTime = "time"
FieldTypeMutex = "mutex"
+ FieldTypeBool = "bool"
)
// Field represents a container for views.
@@ -155,6 +156,16 @@ func OptFieldTypeMutex(cacheType string, cacheSize uint32) FieldOption {
}
}
+func OptFieldTypeBool() FieldOption {
+ return func(fo *FieldOptions) error {
+ if fo.Type != "" {
+ return errors.Errorf("field type is already set to: %s", fo.Type)
+ }
+ fo.Type = FieldTypeBool
+ return nil
+ }
+}
+
// NewField returns a new instance of field.
func NewField(path, index, name string, opts FieldOption) (*Field, error) {
err := validateName(name)
@@ -441,6 +452,14 @@ func (f *Field) applyOptions(opt FieldOptions) error {
f.options.Max = 0
f.options.TimeQuantum = ""
f.options.Keys = opt.Keys
+ case FieldTypeBool:
+ f.options.Type = FieldTypeBool
+ f.options.CacheType = CacheTypeNone
+ f.options.CacheSize = 0
+ f.options.Min = 0
+ f.options.Max = 0
+ f.options.TimeQuantum = ""
+ f.options.Keys = false
default:
return errors.New("invalid field type")
}
@@ -681,6 +700,10 @@ func (f *Field) deleteView(name string) error {
}
// Row returns a row of the standard view.
+// It seems this method is only being used by the test
+// package, and the fact that it's only allowed on
+// `set` fields is odd. This may be considered for
+// deprecation in a future version.
func (f *Field) Row(rowID uint64) (*Row, error) {
if f.Type() != FieldTypeSet {
return nil, errors.Errorf("row method unsupported for field type: %s", f.Type())
@@ -954,10 +977,18 @@ func (f *Field) Import(rowIDs, columnIDs []uint64, timestamps []*time.Time) erro
return errors.New("time quantum not set in field")
}
+ fieldType := f.Type()
+
// Split import data by fragment.
dataByFragment := make(map[importKey]importData)
for i := range rowIDs {
rowID, columnID := rowIDs[i], columnIDs[i]
+
+ // Bool-specific data validation.
+ if fieldType == FieldTypeBool && rowID > 1 {
+ return errors.New("bool field imports only support values 0 and 1")
+ }
+
var timestamp *time.Time
if len(timestamps) > i {
timestamp = timestamps[i]
@@ -1191,6 +1222,12 @@ func (o *FieldOptions) MarshalJSON() ([]byte, error) {
o.CacheSize,
o.Keys,
})
+ case FieldTypeBool:
+ return json.Marshal(struct {
+ Type string `json:"type"`
+ }{
+ o.Type,
+ })
}
return nil, errors.New("invalid field type")
}
diff --git a/fragment.go b/fragment.go
index 7f651b135..312b4af10 100644
--- a/fragment.go
+++ b/fragment.go
@@ -72,6 +72,10 @@ const (
// defaultFragmentMaxOpN is the default value for Fragment.MaxOpN.
defaultFragmentMaxOpN = 2000
+
+ // Row ids used for boolean fields.
+ falseRowID = uint64(0)
+ trueRowID = uint64(1)
)
// fragment represents the intersection of a field and shard in an index.
@@ -1330,16 +1334,29 @@ func (f *fragment) bulkImport(rowIDs, columnIDs []uint64) error {
return fmt.Errorf("mismatch of row/column len: %d != %d", len(rowIDs), len(columnIDs))
}
+ if f.mutexVector != nil {
+ return f.bulkImportMutex(rowIDs, columnIDs)
+ }
+ return f.bulkImportStandard(rowIDs, columnIDs)
+}
+
+// bulkImportStandard performs a bulk import on a standard fragment.
+func (f *fragment) bulkImportStandard(rowIDs, columnIDs []uint64) error {
// Create a temporary bitmap which will be populated by rowIDs and columnIDs
// and then merged into the existing fragment's bitmap.
localBitmap := roaring.NewBitmap()
+
// Disconnect op writer so we don't append updates.
localBitmap.OpWriter = nil
- // Process every bit.
- // If an error occurs then reopen the storage.
- lastID := uint64(0)
+ // rowSet maintains the set of rowIDs present in this import.
+ // It allows the cache to be updated once per row, instead of once
+ // per bit.
rowSet := make(map[uint64]struct{})
+ lastRowID := uint64(0)
+
+ // Process every bit by writing to a local bitmap,
+ // to be merged with fragment storage next.
for i := range rowIDs {
rowID, columnID := rowIDs[i], columnIDs[i]
@@ -1349,18 +1366,18 @@ func (f *fragment) bulkImport(rowIDs, columnIDs []uint64) error {
return err
}
- // Write to storage.
+ // Write to local storage.
_, err = localBitmap.Add(pos)
if err != nil {
return err
}
+
// Reduce the StatsD rate for high volume stats
f.stats.Count("ImportBit", 1, 0.0001)
- // import optimization to avoid linear foreach calls
- // slight risk of concurrent cache counter being off but
- // no real danger
- if i == 0 || rowID != lastID {
- lastID = rowID
+
+ // Add row to rowSet.
+ if i == 0 || rowID != lastRowID {
+ lastRowID = rowID
rowSet[rowID] = struct{}{}
}
@@ -1389,6 +1406,96 @@ func (f *fragment) bulkImport(rowIDs, columnIDs []uint64) error {
return unprotectedWriteToFragment(f, results)
}
+// bulkImportMutex performs a bulk import on a fragment while ensuring
+// mutex restrictions. Because the mutex requirements must be checked
+// against storage, this method must acquire a write lock on the fragment
+// during the entire process, and it handles every bit independently.
+func (f *fragment) bulkImportMutex(rowIDs, columnIDs []uint64) error {
+ f.mu.Lock()
+ defer f.mu.Unlock()
+
+ // Disconnect op writer so we don't append updates.
+ f.storage.OpWriter = nil
+
+ // If an error occurs then reopen the storage.
+ if err := func() error {
+ // rowSet maintains the set of rowIDs present in this import.
+ // It allows the cache to be updated once per row, instead of once
+ // per bit.
+ rowSet := make(map[uint64]struct{})
+ lastRowID := uint64(0)
+
+ // Process every bit.
+ for i := range rowIDs {
+ rowID, columnID := rowIDs[i], columnIDs[i]
+
+ // Handle mutex vector (i.e. clear an existing row).
+ if existingRowID, found := f.mutexVector.Get(columnID); found && existingRowID != rowID {
+ // Determine the position of the bit in the storage.
+ pos, err := f.pos(existingRowID, columnID)
+ if err != nil {
+ return err
+ }
+
+ // Clear storage.
+ _, err = f.storage.Remove(pos)
+ if err != nil {
+ return err
+ }
+
+ rowSet[existingRowID] = struct{}{}
+ }
+
+ // Determine the position of the bit in the storage.
+ pos, err := f.pos(rowID, columnID)
+ if err != nil {
+ return err
+ }
+
+ // Write to storage.
+ _, err = f.storage.Add(pos)
+ if err != nil {
+ return err
+ }
+
+ // Reduce the StatsD rate for high volume stats
+ f.stats.Count("ImportBit", 1, 0.0001)
+
+ // Add row to rowSet.
+ if i == 0 || rowID != lastRowID {
+ lastRowID = rowID
+ rowSet[rowID] = struct{}{}
+ }
+
+ // Invalidate block checksum.
+ delete(f.checksums, int(rowID/HashBlockSize))
+ }
+
+ // Update cache counts for all rows.
+ for rowID := range rowSet {
+ // Import should ALWAYS have row() load a new bm from fragment.storage
+ // because the row that's in rowCache hasn't been updated with
+ // this import's data.
+ f.cache.BulkAdd(rowID, f.unprotectedRow(rowID).Count())
+ }
+
+ f.cache.Invalidate()
+
+ return nil
+ }(); err != nil {
+ _ = f.closeStorage()
+ _ = f.openStorage()
+ return err
+ }
+
+ // Write the storage to disk and reload.
+ if err := f.snapshot(); err != nil {
+ return err
+ }
+
+ return nil
+}
+
// importValue bulk imports a set of range-encoded values.
func (f *fragment) importValue(columnIDs, values []uint64, bitDepth uint) error {
f.mu.Lock()
@@ -2103,3 +2210,33 @@ func (v *rowsVector) Get(colID uint64) (uint64, bool) {
// Set is not used for rowsVector.
func (v *rowsVector) Set(colID, rowID uint64) {}
+
+// boolVector implements the vector interface by looking
+// at data in rows 0 and 1.
+type boolVector struct {
+ f *fragment
+}
+
+// newBoolVector returns a boolVector for a given fragment.
+func newBoolVector(f *fragment) *boolVector {
+ return &boolVector{
+ f: f,
+ }
+}
+
+// Get returns the rowID associated to the given colID.
+// Additionally, it returns true if a value was found,
+// otherwise it returns false.
+func (v *boolVector) Get(colID uint64) (uint64, bool) {
+ rows := v.f.rowsForColumn(colID)
+ if len(rows) == 1 {
+ switch rows[0] {
+ case falseRowID, trueRowID:
+ return rows[0], true
+ }
+ }
+ return 0, false
+}
+
+// Set is not used for boolVector.
+func (v *boolVector) Set(colID, rowID uint64) {}
diff --git a/fragment_internal_test.go b/fragment_internal_test.go
index dd2d00b41..87f4e6137 100644
--- a/fragment_internal_test.go
+++ b/fragment_internal_test.go
@@ -1178,6 +1178,146 @@ func TestFragment_SetMutex(t *testing.T) {
}
}
+// Ensure a fragment can import mutually exclusive values.
+func TestFragment_ImportMutex(t *testing.T) {
+ tests := []struct {
+ rowIDs []uint64
+ colIDs []uint64
+ exp map[uint64][]uint64
+ }{
+ {
+ []uint64{1, 1, 1, 1},
+ []uint64{0, 1, 2, 3},
+ map[uint64][]uint64{
+ 1: {0, 1, 2, 3},
+ },
+ },
+ {
+ []uint64{1, 1, 1, 1, 2, 2, 2, 2},
+ []uint64{0, 1, 2, 3, 0, 1, 2, 3},
+ map[uint64][]uint64{
+ 1: {},
+ 2: {0, 1, 2, 3},
+ },
+ },
+ {
+ []uint64{1, 1, 1, 1, 2},
+ []uint64{0, 1, 2, 3, 1},
+ map[uint64][]uint64{
+ 1: {0, 2, 3},
+ 2: {1},
+ },
+ },
+ {
+ []uint64{1, 1, 1, 1, 2, 2, 1},
+ []uint64{0, 1, 2, 3, 1, 8, 1},
+ map[uint64][]uint64{
+ 1: {0, 1, 2, 3},
+ 2: {8},
+ },
+ },
+ {
+ []uint64{1, 2, 3},
+ []uint64{8, 8, 8},
+ map[uint64][]uint64{
+ 1: {},
+ 2: {},
+ 3: {8},
+ },
+ },
+ }
+
+ for i, test := range tests {
+ t.Run(fmt.Sprintf("importmutex%d", i), func(t *testing.T) {
+ f := mustOpenMutexFragment("i", "f", viewStandard, 0, "")
+ defer f.Close()
+
+ err := f.bulkImport(test.rowIDs, test.colIDs)
+ if err != nil {
+ t.Fatalf("bulk importing ids: %v", err)
+ }
+
+ // Check for expected results.
+ for k, v := range test.exp {
+ cols := f.row(k).Columns()
+ if !reflect.DeepEqual(cols, v) {
+ t.Fatalf("expected: %v, but got: %v", v, cols)
+ }
+ }
+ })
+ }
+}
+
+// Ensure a fragment can import bool values.
+func TestFragment_ImportBool(t *testing.T) {
+ tests := []struct {
+ rowIDs []uint64
+ colIDs []uint64
+ exp map[uint64][]uint64
+ }{
+ {
+ []uint64{1, 1, 1, 1},
+ []uint64{0, 1, 2, 3},
+ map[uint64][]uint64{
+ 1: {0, 1, 2, 3},
+ },
+ },
+ {
+ []uint64{0, 0, 0, 0, 1, 1, 1, 1},
+ []uint64{0, 1, 2, 3, 0, 1, 2, 3},
+ map[uint64][]uint64{
+ 0: {},
+ 1: {0, 1, 2, 3},
+ },
+ },
+ {
+ []uint64{0, 0, 0, 0, 1},
+ []uint64{0, 1, 2, 3, 1},
+ map[uint64][]uint64{
+ 0: {0, 2, 3},
+ 1: {1},
+ },
+ },
+ {
+ []uint64{1, 1, 1, 1, 0, 0, 1},
+ []uint64{0, 1, 2, 3, 1, 8, 1},
+ map[uint64][]uint64{
+ 0: {8},
+ 1: {0, 1, 2, 3},
+ },
+ },
+ {
+ []uint64{0, 1, 2},
+ []uint64{8, 8, 8},
+ map[uint64][]uint64{
+ 0: {},
+ 1: {}, // This isn't {8} because fragment doesn't validate bool values.
+ 2: {8},
+ },
+ },
+ }
+
+ for i, test := range tests {
+ t.Run(fmt.Sprintf("importmutex%d", i), func(t *testing.T) {
+ f := mustOpenBoolFragment("i", "f", viewStandard, 0, "")
+ defer f.Close()
+
+ err := f.bulkImport(test.rowIDs, test.colIDs)
+ if err != nil {
+ t.Fatalf("bulk importing ids: %v", err)
+ }
+
+ // Check for expected results.
+ for k, v := range test.exp {
+ cols := f.row(k).Columns()
+ if !reflect.DeepEqual(cols, v) {
+ t.Fatalf("expected: %v, but got: %v", v, cols)
+ }
+ }
+ })
+ }
+}
+
func BenchmarkFragment_Snapshot(b *testing.B) {
if *FragmentPath == "" {
b.Skip("no fragment specified")
@@ -1302,6 +1442,13 @@ func mustOpenMutexFragment(index, field, view string, shard uint64, cacheType st
return frag
}
+// mustOpenBoolFragment returns a new instance of Fragment for a bool field.
+func mustOpenBoolFragment(index, field, view string, shard uint64, cacheType string) *fragment {
+ frag := mustOpenFragment(index, field, view, shard, cacheType)
+ frag.mutexVector = newBoolVector(frag)
+ return frag
+}
+
// Reopen closes the fragment and reopens it as a new instance.
func (f *fragment) reopen() error {
if err := f.Close(); err != nil {
diff --git a/http/handler.go b/http/handler.go
index bc108568a..820bc6521 100644
--- a/http/handler.go
+++ b/http/handler.go
@@ -687,6 +687,8 @@ func (h *Handler) handlePostField(w http.ResponseWriter, r *http.Request) {
fos = append(fos, pilosa.OptFieldTypeTime(*req.Options.TimeQuantum))
case pilosa.FieldTypeMutex:
fos = append(fos, pilosa.OptFieldTypeMutex(*req.Options.CacheType, *req.Options.CacheSize))
+ case pilosa.FieldTypeBool:
+ fos = append(fos, pilosa.OptFieldTypeBool())
}
if req.Options.Keys != nil {
if *req.Options.Keys {
@@ -778,6 +780,20 @@ func (o *fieldOptions) validate() error {
} else if o.TimeQuantum != nil {
return pilosa.NewBadRequestError(errors.New("timeQuantum does not apply to field type mutex"))
}
+ case pilosa.FieldTypeBool:
+ if o.CacheType != nil {
+ return pilosa.NewBadRequestError(errors.New("cacheType does not apply to field type bool"))
+ } else if o.CacheSize != nil {
+ return pilosa.NewBadRequestError(errors.New("cacheSize does not apply to field type bool"))
+ } else if o.Min != nil {
+ return pilosa.NewBadRequestError(errors.New("min does not apply to field type bool"))
+ } else if o.Max != nil {
+ return pilosa.NewBadRequestError(errors.New("max does not apply to field type bool"))
+ } else if o.TimeQuantum != nil {
+ return pilosa.NewBadRequestError(errors.New("timeQuantum does not apply to field type bool"))
+ } else if o.Keys != nil {
+ return pilosa.NewBadRequestError(errors.New("keys does not apply to field type bool"))
+ }
default:
return errors.Errorf("invalid field type: %s", o.Type)
}
diff --git a/pql/ast.go b/pql/ast.go
index cd483e030..0e2ff6aef 100644
--- a/pql/ast.go
+++ b/pql/ast.go
@@ -248,6 +248,23 @@ func (c *Call) FieldArg() (string, error) {
return "", fmt.Errorf("No field argument specified")
}
+// BoolArg is for reading the value at key from call.Args as a bool. If the
+// key is not in Call.Args, the value of the returned bool will be false, and
+// the error will be nil. The value is assumed to be a bool. An error is
+// returned if the value is not a bool.
+func (c *Call) BoolArg(key string) (bool, bool, error) {
+ val, ok := c.Args[key]
+ if !ok {
+ return false, false, nil
+ }
+ switch tval := val.(type) {
+ case bool:
+ return tval, true, nil
+ default:
+ return false, true, fmt.Errorf("could not convert %v of type %T to bool in Call.BoolArg", tval, tval)
+ }
+}
+
// UintArg is for reading the value at key from call.Args as a uint64. If the
// key is not in Call.Args, the value of the returned bool will be false, and
// the error will be nil. The value is assumed to be a uint64 or an int64 and
diff --git a/server.go b/server.go
index cde1dd406..ce6319f45 100644
--- a/server.go
+++ b/server.go
@@ -234,6 +234,13 @@ func OptServerClusterHasher(h Hasher) ServerOption {
}
}
+func OptServerTranslateFileMapSize(mapSize int) ServerOption {
+ return func(s *Server) error {
+ s.holder.translateFile = NewTranslateFile(OptTranslateFileMapSize(mapSize))
+ return nil
+ }
+}
+
// NewServer returns a new instance of Server.
func NewServer(opts ...ServerOption) (*Server, error) {
s := &Server{
diff --git a/server/config.go b/server/config.go
index 8367eb835..d255e4bb6 100644
--- a/server/config.go
+++ b/server/config.go
@@ -71,8 +71,9 @@ type Config struct {
// Gossip config is based around memberlist.Config.
Gossip gossip.Config `toml:"gossip"`
- // DEPRECATED: Translation config supports translation store replication.
Translation struct {
+ MapSize int `toml:"map-size"`
+ // DEPRECATED: Translation config supports translation store replication.
PrimaryURL string `toml:"primary-url"`
} `toml:"translation"`
diff --git a/server/server.go b/server/server.go
index 4fdd711fe..b4ac595fa 100644
--- a/server/server.go
+++ b/server/server.go
@@ -283,6 +283,14 @@ func (m *Command) SetupServer() error {
coordinatorOpt,
}
+ if m.Config.Translation.MapSize > 0 {
+ serverOptions = append(
+ serverOptions,
+ pilosa.OptServerTranslateFileMapSize(
+ m.Config.Translation.MapSize,
+ ),
+ )
+ }
serverOptions = append(serverOptions, m.serverOptions...)
m.Server, err = pilosa.NewServer(serverOptions...)
diff --git a/server/server_test.go b/server/server_test.go
index 67440c21e..ed73975d7 100644
--- a/server/server_test.go
+++ b/server/server_test.go
@@ -404,6 +404,7 @@ func TestClusteringNodesReplica1(t *testing.T) {
// Create new main with the same config.
config := cluster[2].Command.Config
+ config.Translation.MapSize = 100000
// config.Bind = cluster[2].API.Node().URI.HostPort()
// this isn't necessary, but makes the test run way faster
@@ -477,6 +478,7 @@ func TestClusteringNodesReplica2(t *testing.T) {
// Create new main with the same config.
config := cluster[2].Command.Config
+ config.Translation.MapSize = 100000
// config.Bind = cluster[2].API.Node().URI.HostPort()
// this isn't necessary, but makes the test run way faster
@@ -497,6 +499,7 @@ func TestClusteringNodesReplica2(t *testing.T) {
// Create new main with the same config.
config = cluster[1].Command.Config
// config.Bind = cluster[1].API.Node().URI.HostPort()
+ config.Translation.MapSize = 100000
// this isn't necessary, but makes the test run way faster
config.Gossip.Port = strconv.Itoa(int(cluster[1].Command.GossipTransport().URI.Port))
diff --git a/test/pilosa.go b/test/pilosa.go
index 9dbd78057..f8f60eebe 100644
--- a/test/pilosa.go
+++ b/test/pilosa.go
@@ -55,11 +55,20 @@ func newCommand(opts ...server.CommandOption) *Command {
panic(err)
}
- // set aggressive close timeout by default to avoid hanging tests. This was
+ // Set aggressive close timeout by default to avoid hanging tests. This was
// a problem with PDK tests which used go-pilosa as well. We put it at the
// beginning of the option slice so that it can be overridden by user-passed
// options.
- opts = append([]server.CommandOption{server.OptCommandCloseTimeout(time.Millisecond * 2)}, opts...)
+ // Also set TranslateFile MapSize to a smaller number so memory allocation
+ // does not fail on 32-bit systems.
+ opts = append([]server.CommandOption{
+ server.OptCommandCloseTimeout(time.Millisecond * 2),
+ server.OptCommandServerOptions(
+ pilosa.OptServerTranslateFileMapSize(
+ 2 << 25,
+ ),
+ ),
+ }, opts...)
m := &Command{commandOptions: opts}
m.Command = server.NewCommand(bytes.NewReader(nil), ioutil.Discard, ioutil.Discard, opts...)
m.Config.DataDir = path
diff --git a/translate.go b/translate.go
index 41f3b4852..2ef289f45 100644
--- a/translate.go
+++ b/translate.go
@@ -5,7 +5,6 @@ import (
"bytes"
"context"
"encoding/binary"
- "errors"
"fmt"
"io"
"io/ioutil"
@@ -17,6 +16,7 @@ import (
"time"
"github.com/cespare/xxhash"
+ "github.com/pkg/errors"
)
const (
@@ -81,9 +81,29 @@ type TranslateFile struct {
replicationRetryInterval time.Duration
}
+// TranslateFileOption is a functional option type for pilosa.TranslateFile
+type TranslateFileOption func(f *TranslateFile) error
+
+func OptTranslateFileMapSize(mapSize int) TranslateFileOption {
+ return func(f *TranslateFile) error {
+ f.mapSize = mapSize
+ return nil
+ }
+}
+
// NewTranslateFile returns a new instance of TranslateFile.
-func NewTranslateFile() *TranslateFile {
- return &TranslateFile{
+func NewTranslateFile(opts ...TranslateFileOption) *TranslateFile {
+ var defaultMapSize64 int64 = 10 * (1 << 30)
+ var defaultMapSize int
+
+ if ^uint(0)>>32 > 0 {
+ // 10GB default map size
+ defaultMapSize = int(defaultMapSize64)
+ } else {
+ // Use 2GB default map size on 32-bit systems
+ defaultMapSize = (1 << 31) - 1
+ }
+ f := &TranslateFile{
writeNotify: make(chan struct{}),
closing: make(chan struct{}),
cols: make(map[string]*index),
@@ -96,6 +116,16 @@ func NewTranslateFile() *TranslateFile {
replicationRetryInterval: defaultReplicationRetryInterval,
}
+
+ for _, opt := range opts {
+ err := opt(f)
+ if err != nil {
+ // TODO (2.0): Change func signature to return error
+ panic(errors.Wrap(err, "applying option"))
+ }
+ }
+
+ return f
}
func (s *TranslateFile) Open() (err error) {
diff --git a/translate_mapsize_386.go b/translate_mapsize_386.go
deleted file mode 100644
index 259b735ac..000000000
--- a/translate_mapsize_386.go
+++ /dev/null
@@ -1,5 +0,0 @@
-package pilosa
-
-// defaultMapSize is the default size of mapped memory for the translate store.
-// It is passed as an int to syscall.Mmap and so must be < 2^31
-const defaultMapSize = (1 << 31) - 1 // 2GB
diff --git a/translate_mapsize_all64bitsystems.go b/translate_mapsize_all64bitsystems.go
deleted file mode 100644
index 605d6270a..000000000
--- a/translate_mapsize_all64bitsystems.go
+++ /dev/null
@@ -1,8 +0,0 @@
-// +build !386
-
-package pilosa
-
-// defaultMapSize is the default size of mapped memory for the translate store.
-// It is passed as an int to syscall.Mmap and so can only be larger than 2^31 on
-// 64bit systems.
-const defaultMapSize = 10 * (1 << 30) // 10GB
diff --git a/translate_test.go b/translate_test.go
index 706fe365f..1aea26cfe 100644
--- a/translate_test.go
+++ b/translate_test.go
@@ -367,7 +367,7 @@ func TestPrintTranslateFile(t *testing.T) {
}
f.Close()
- s := pilosa.NewTranslateFile()
+ s := pilosa.NewTranslateFile(pilosa.OptTranslateFileMapSize(2 << 25))
s.Path = f.Name()
err = s.Open()
if err != nil {
@@ -809,7 +809,7 @@ func NewTranslateFile() *TranslateFile {
}
f.Close()
- s := &TranslateFile{TranslateFile: pilosa.NewTranslateFile()}
+ s := &TranslateFile{TranslateFile: pilosa.NewTranslateFile(pilosa.OptTranslateFileMapSize(2 << 25))}
s.Path = f.Name()
return s
}
@@ -847,7 +847,7 @@ func (s *TranslateFile) Reopen() error {
}
s.lock.Lock()
- s.TranslateFile = pilosa.NewTranslateFile()
+ s.TranslateFile = pilosa.NewTranslateFile(pilosa.OptTranslateFileMapSize(2 << 25))
s.lock.Unlock()
s.Path = prev.Path
s.SetPrimaryStore("restored-primary", prev.PrimaryTranslateStore)
diff --git a/view.go b/view.go
index 2feac8b44..5fc0f5eb8 100644
--- a/view.go
+++ b/view.go
@@ -244,6 +244,8 @@ func (v *view) newFragment(path string, shard uint64) *fragment {
frag.stats = v.stats.WithTags(fmt.Sprintf("shard:%d", shard))
if v.fieldType == FieldTypeMutex {
frag.mutexVector = newRowsVector(frag)
+ } else if v.fieldType == FieldTypeBool {
+ frag.mutexVector = newBoolVector(frag)
}
return frag
}