From 7b27a48d7c970366cb0f9c0ef4fedbd06125dc9c Mon Sep 17 00:00:00 2001 From: Seebs Date: Tue, 27 Jul 2021 14:35:21 -0500 Subject: [PATCH 1/9] allow Molecula copyrights --- Makefile | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Makefile b/Makefile index 0ffeb96bd..2c2ac0456 100644 --- a/Makefile +++ b/Makefile @@ -19,7 +19,7 @@ DOCKER_BUILD= # set to 1 to use `docker-build` instead of `build` when creating BUILD_TAGS += shardwidth$(SHARD_WIDTH) TEST_TAGS = roaringparanoia define LICENSE_HASH_CODE - head -13 $1 | sed -e 's/Copyright 20[0-9][0-9]/Copyright 20XX/g' | shasum | cut -f 1 -d " " + head -13 $1 | sed -e 's/Copyright 20[0-9][0-9]/Copyright 20XX/' -e 's/Pilosa Corp\./Molecula Corp./' | shasum | cut -f 1 -d " " endef LICENSE_HASH=$(shell $(call LICENSE_HASH_CODE, pilosa.go)) UNAME := $(shell uname -s) From dab8dff9c29850257ae66f41d6896e7ff44e62d0 Mon Sep 17 00:00:00 2001 From: Seebs Date: Mon, 19 Jul 2021 14:58:07 -0500 Subject: [PATCH 2/9] make "subdivide list by shardwidth" available for reuse We keep wanting this, and it's shardwidth-dependent code, and we keep rewriting it. It should be in the shardwidth package. --- shardwidth/helper.go | 70 ++++++++++++++++++++++++++ shardwidth/helper_test.go | 103 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 173 insertions(+) create mode 100644 shardwidth/helper.go create mode 100644 shardwidth/helper_test.go diff --git a/shardwidth/helper.go b/shardwidth/helper.go new file mode 100644 index 000000000..c04e2e72b --- /dev/null +++ b/shardwidth/helper.go @@ -0,0 +1,70 @@ +// Copyright 2021 Molecula 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 shardwidth + +import ( + "math/bits" +) + +// FindNextShard returns the index of the first item which is not in the +// same shard as i. The index it returns may be equal to the length of the +// haystack, indicatincg that the rest of the list is in the same shard. +func FindNextShard(i int, haystack []uint64) int { + // compute the last thing that's in the same shard as haystack[i]. + if i >= len(haystack) { + return i + } + // current shard: + shard := (haystack[i] >> Exponent) + // last value in shard: + shardEnd := ((shard + 1) << Exponent) - 1 + j := i + // We want to do a binary search of the haystack. For any length of + // haystack, its topmost bit gives us a reasonable halfway point; it may + // not actually be halfway, but the number of steps it'll take to search + // it will be the same as if it were. sort.Search has interface overhead + // and makes us sad. + for incr := 1 << (bits.Len64(uint64(len(haystack) - i))); incr > 0; incr >>= 1 { + if j+incr < len(haystack) { + if haystack[j+incr] <= shardEnd { + j += incr + } + } + } + // we've found the last item that is in the same shard as i, so... + return j + 1 +} + +// FindShards finds the shards in a given haystack +func FindShards(haystack []uint64) (shards []uint64, endIndexes []int) { + if len(haystack) == 0 { + return nil, nil + } + index := 0 + // the steady state of this loop is that shards contains the current + // shard, but not its ending index; each time we find a new ending + // index, we record that index as the end for the current shard, and + // the new shard, until we reach the end and append len(haystack) + // as the last index. + shards = []uint64{haystack[index] >> Exponent} + index = FindNextShard(index, haystack) + for index < len(haystack) { + shards = append(shards, haystack[index]>>Exponent) + endIndexes = append(endIndexes, index) + index = FindNextShard(index, haystack) + } + endIndexes = append(endIndexes, index) + return shards, endIndexes +} diff --git a/shardwidth/helper_test.go b/shardwidth/helper_test.go new file mode 100644 index 000000000..cd1806cd8 --- /dev/null +++ b/shardwidth/helper_test.go @@ -0,0 +1,103 @@ +// Copyright 2021 Molecula 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 shardwidth_test + +import ( + "math/rand" + "testing" + + "github.com/molecula/featurebase/v2/shardwidth" +) + +type nextShardTestCase struct { + name string + haystack [][2]uint64 // stored as shard, offset pairs + shardIndexes []int +} + +var nextShardTestCases = []nextShardTestCase{ + { + name: "all-in-one", + haystack: [][2]uint64{ + {0, 0}, + {0, 1}, + }, + shardIndexes: []int{2}, + }, + { + name: "split", + haystack: [][2]uint64{ + {0, 0}, + {1, 1}, + }, + shardIndexes: []int{1, 2}, + }, + { + name: "two-and-one", + haystack: [][2]uint64{ + {0, 0}, + {0, 1}, + {1, 1}, + }, + shardIndexes: []int{2, 3}, + }, +} + +func TestFindShards(t *testing.T) { + for _, c := range nextShardTestCases { + haystack := make([]uint64, len(c.haystack)) + for i, h := range c.haystack { + haystack[i] = (h[0] << shardwidth.Exponent) + h[1] + } + _, indexes := shardwidth.FindShards(haystack) + if len(indexes) != len(c.shardIndexes) { + t.Fatalf("%s: expected %d, got %d", c.name, c.shardIndexes, indexes) + } + for i, expected := range c.shardIndexes { + if indexes[i] != expected { + t.Fatalf("%s: expected index %d to be %d, got %d", c.name, i, expected, indexes[i]) + } + } + } + // fake up some more test cases + for i := 0; i < 100; i++ { + haystack := make([]uint64, 100) + shard := uint64(0) + bit := uint64(0) + shardIndexes := []int{} + for j := range haystack { + if rand.Intn(30) == 0 { + if j > 0 { + shardIndexes = append(shardIndexes, j) + } + shard++ + bit = 0 + } else { + bit += uint64(rand.Intn(30)) + } + haystack[j] = (shard << shardwidth.Exponent) + bit + } + shardIndexes = append(shardIndexes, len(haystack)) + _, indexes := shardwidth.FindShards(haystack) + if len(indexes) != len(shardIndexes) { + t.Fatalf("trial %d: expected %d, got %d", i, shardIndexes, indexes) + } + for idx, expected := range shardIndexes { + if indexes[idx] != expected { + t.Fatalf("trial %d: expected index %d to be %d, got %d", i, idx, expected, indexes[idx]) + } + } + } +} From 35faa39b20fd2fa47921f6a77d7f37ad8bb9523f Mon Sep 17 00:00:00 2001 From: Seebs Date: Mon, 26 Jul 2021 16:20:47 -0500 Subject: [PATCH 3/9] don't use nil qcx A nil Qcx is a crime against existence and makes baby pandas cry. Having taken out the hack that tried to accommodate this when tests did it, we now have to fix the tests. Oh no. --- api.go | 37 ----------------------- executor_test.go | 27 +++++++++-------- test/cluster.go | 77 +++++++++++++++++++++++++++++------------------- tx_test.go | 24 +++++++++++---- 4 files changed, 79 insertions(+), 86 deletions(-) diff --git a/api.go b/api.go index c9c998848..7f22a52ed 100644 --- a/api.go +++ b/api.go @@ -1430,16 +1430,6 @@ var ErrAborted = fmt.Errorf("error: update was aborted") func (api *API) ImportAtomicRecord(ctx context.Context, qcx *Qcx, req *AtomicRecord, opts ...ImportOption) error { - // this is because some of the tests pass nil qcx for convenience. - isLocalQcx := false - if qcx == nil { - isLocalQcx = true - qcx = api.Txf().NewQcx() - defer func() { - qcx.Abort() - }() - } - simPowerLoss := false lossAfter := -1 var opt ImportOptions @@ -1491,11 +1481,6 @@ func (api *API) ImportAtomicRecord(ctx context.Context, qcx *Qcx, req *AtomicRec return errors.Wrap(err, "ImportAtomicRecord ImportWithTx") } } - - // got to the end succesfully, so commit if we made the qcx - if isLocalQcx { - return qcx.Finish() - } return nil } @@ -1521,21 +1506,10 @@ func (api *API) Import(ctx context.Context, qcx *Qcx, req *ImportRequest, opts . if req.Clear { opts = addClearToImportOptions(opts) } - isLocalQcx := false - if qcx == nil { - isLocalQcx = true - qcx = api.Txf().NewQcx() - defer func() { - qcx.Abort() - }() - } err = api.ImportWithTx(ctx, qcx, req, opts...) if err != nil { return err } - if isLocalQcx { - return qcx.Finish() - } return nil } @@ -1752,14 +1726,6 @@ func (api *API) ImportValueWithTx(ctx context.Context, qcx *Qcx, req *ImportValu // don't keep that list around since we don't need it anymore req.scratch = nil } - isLocalQcx := false - if qcx == nil { - isLocalQcx = true - qcx = api.Txf().NewQcx() - defer func() { - qcx.Abort() - }() - } // if we're importing into a specific shard if req.Shard != math.MaxUint64 { @@ -1847,9 +1813,6 @@ func (api *API) ImportValueWithTx(ctx context.Context, qcx *Qcx, req *ImportValu if err != nil { return err } - if isLocalQcx { - return qcx.Finish() - } return nil } diff --git a/executor_test.go b/executor_test.go index 834906864..46b3e00ef 100644 --- a/executor_test.go +++ b/executor_test.go @@ -492,7 +492,7 @@ func TestExecutor(t *testing.T) { Set(6, f=1, 2001-01-01T00:00) Set(7, f=1, 2002-01-01T02:00) Set(8, f=1, %s) - + Set(2, f=1, 1999-12-30T00:00) Set(2, f=1, 2002-02-01T00:00) Set(2, f=10, 2001-01-01T00:00)`, nextDayExclusive.Format("2006-01-02T15:04")) @@ -539,7 +539,7 @@ func TestExecutor(t *testing.T) { Set("five", f=1, 2000-02-01T00:00) Set("six", f=1, 2001-01-01T00:00) Set("seven", f=1, 2002-01-01T02:00) - + Set("two", f=1, 1999-12-30T00:00) Set("two", f=1, 2002-02-01T00:00) Set("two", f=10, 2001-01-01T00:00)` @@ -573,7 +573,7 @@ func TestExecutor(t *testing.T) { Set(5, f="foo", 2000-02-01T00:00) Set(6, f="foo", 2001-01-01T00:00) Set(7, f="foo", 2002-01-01T02:00) - + Set(2, f="foo", 1999-12-30T00:00) Set(2, f="foo", 2002-02-01T00:00) Set(2, f="bar", 2001-01-01T00:00)` @@ -608,7 +608,7 @@ func TestExecutor(t *testing.T) { Set("five", f="foo", 2000-02-01T00:00) Set("six", f="foo", 2001-01-01T00:00) Set("seven", f="foo", 2002-01-01T02:00) - + Set("two", f="foo", 1999-12-30T00:00) Set("two", f="foo", 2002-02-01T00:00) Set("two", f="bar", 2001-01-01T00:00)` @@ -643,7 +643,7 @@ func TestExecutor(t *testing.T) { Set(5, f=1, 2000-02-01T00:00) Set(6, f=1, 2001-01-01T00:00) Set(7, f=1, 2002-01-01T02:00) - + Set(2, f=1, 1999-12-30T00:00) Set(2, f=1, 2002-02-01T00:00) Set(2, f=10, 2001-01-01T00:00)` @@ -678,7 +678,7 @@ func TestExecutor(t *testing.T) { Set(5, f=1, 2000-02-01T00:00) Set(6, f=1, 2001-01-01T00:00) Set(7, f=1, 2002-01-01T02:00) - + Set(2, f=1, 1999-12-30T00:00) Set(2, f=1, 2002-02-01T00:00) Set(2, f=10, 2001-01-01T00:00)` @@ -724,7 +724,7 @@ func TestExecutor(t *testing.T) { Set("five", f=1, 2000-02-01T00:00) Set("six", f=1, 2001-01-01T00:00) Set("seven", f=1, 2002-01-01T02:00) - + Set("two", f=1, 1999-12-30T00:00) Set("two", f=1, 2002-02-01T00:00) Set("two", f=10, 2001-01-01T00:00)` @@ -758,7 +758,7 @@ func TestExecutor(t *testing.T) { Set(5, f="foo", 2000-02-01T00:00) Set(6, f="foo", 2001-01-01T00:00) Set(7, f="foo", 2002-01-01T02:00) - + Set(2, f="foo", 1999-12-30T00:00) Set(2, f="foo", 2002-02-01T00:00) Set(2, f="bar", 2001-01-01T00:00)` @@ -793,7 +793,7 @@ func TestExecutor(t *testing.T) { Set("five", f="foo", 2000-02-01T00:00) Set("six", f="foo", 2001-01-01T00:00) Set("seven", f="foo", 2002-01-01T02:00) - + Set("two", f="foo", 1999-12-30T00:00) Set("two", f="foo", 2002-02-01T00:00) Set("two", f="bar", 2001-01-01T00:00)` @@ -5154,10 +5154,13 @@ func TestExecutor_GroupByStrings(t *testing.T) { }); err != nil { t.Fatalf("importing: %v", err) } + m0 := c.GetNode(0) + qcx := m0.API.Txf().NewQcx() + defer qcx.Abort() var v1, v2, v3, v4, v5, v6, v7, v8, v9, v10 int64 = 1, 2, 3, 4, 5, 6, 7, 8, 9, 10 var nv1, nv2, nv3, nv4 int64 = -1, -2, -3, -4 - if err := c.GetNode(0).API.ImportValue(context.Background(), nil, &pilosa.ImportValueRequest{ + if err := m0.API.ImportValue(context.Background(), qcx, &pilosa.ImportValueRequest{ Index: "istring", Field: "v", Shard: 0, @@ -5167,7 +5170,7 @@ func TestExecutor_GroupByStrings(t *testing.T) { t.Fatalf("importing: %v", err) } - if err := c.GetNode(0).API.ImportValue(context.Background(), nil, &pilosa.ImportValueRequest{ + if err := m0.API.ImportValue(context.Background(), qcx, &pilosa.ImportValueRequest{ Index: "istring", Field: "vv", Shard: 0, @@ -5177,7 +5180,7 @@ func TestExecutor_GroupByStrings(t *testing.T) { t.Fatalf("importing: %v", err) } - if err := c.GetNode(0).API.ImportValue(context.Background(), nil, &pilosa.ImportValueRequest{ + if err := m0.API.ImportValue(context.Background(), qcx, &pilosa.ImportValueRequest{ Index: "istring", Field: "nv", Shard: 0, diff --git a/test/cluster.go b/test/cluster.go index 7ff157094..126245a46 100644 --- a/test/cluster.go +++ b/test/cluster.go @@ -23,7 +23,7 @@ import ( "testing" "time" - "github.com/molecula/featurebase/v2" + pilosa "github.com/molecula/featurebase/v2" "github.com/molecula/featurebase/v2/api/client" "github.com/molecula/featurebase/v2/disco" "github.com/molecula/featurebase/v2/logger" @@ -213,32 +213,37 @@ func (c *Cluster) ImportBitsWithTimestamp(t testing.TB, index, field string, row if com.API.Node().ID != node.ID { continue } - if len(timestamps) == 0 { - err := com.API.Import(context.Background(), nil, &pilosa.ImportRequest{ - Index: index, - Field: field, - Shard: shard, - RowIDs: rowIDs, - ColumnIDs: colIDs, - }) - if err != nil { - t.Fatalf("importing data: %v", err) - } - } else { - ts := byShardTs[shard] - err := com.API.Import(context.Background(), nil, &pilosa.ImportRequest{ - Index: index, - Field: field, - Shard: shard, - RowIDs: rowIDs, - ColumnIDs: colIDs, - Timestamps: ts, - }) - if err != nil { - t.Fatalf("importing data: %v", err) - } + func() { + qcx := com.API.Txf().NewQcx() + defer qcx.Abort() + if len(timestamps) == 0 { + err := com.API.Import(context.Background(), qcx, &pilosa.ImportRequest{ + Index: index, + Field: field, + Shard: shard, + RowIDs: rowIDs, + ColumnIDs: colIDs, + }) + if err != nil { + t.Fatalf("importing data: %v", err) + } + } else { + ts := byShardTs[shard] + err := com.API.Import(context.Background(), qcx, &pilosa.ImportRequest{ + Index: index, + Field: field, + Shard: shard, + RowIDs: rowIDs, + ColumnIDs: colIDs, + Timestamps: ts, + }) + if err != nil { + t.Fatalf("importing data: %v", err) + } + + } + }() - } } } } @@ -262,7 +267,9 @@ func (c *Cluster) ImportKeyKey(t testing.TB, index, field string, valAndRecKeys importRequest.RowKeys[i] = vk[0] importRequest.ColumnKeys[i] = vk[1] } - err := c.GetPrimary().API.Import(context.Background(), nil, importRequest) + qcx := c.GetPrimary().API.Txf().NewQcx() + defer qcx.Abort() + err := c.GetPrimary().API.Import(context.Background(), qcx, importRequest) if err != nil { t.Fatalf("importing keykey data: %v", err) } @@ -292,7 +299,9 @@ func (c *Cluster) ImportTimeQuantumKey(t testing.TB, index, field string, entrie importRequest.Timestamps[i] = entry.Ts } - err := c.GetPrimary().API.Import(context.Background(), nil, importRequest) + qcx := c.GetPrimary().API.Txf().NewQcx() + defer qcx.Abort() + err := c.GetPrimary().API.Import(context.Background(), qcx, importRequest) if err != nil { t.Fatalf("importing keykey data: %v", err) } @@ -318,7 +327,9 @@ func (c *Cluster) ImportIntKey(t testing.TB, index, field string, pairs []IntKey importRequest.Values[i] = pair.Val importRequest.ColumnKeys[i] = pair.Key } - if err := c.GetPrimary().API.ImportValue(context.Background(), nil, importRequest); err != nil { + qcx := c.GetPrimary().API.Txf().NewQcx() + defer qcx.Abort() + if err := c.GetPrimary().API.ImportValue(context.Background(), qcx, importRequest); err != nil { t.Fatalf("importing IntKey data: %v", err) } } @@ -342,7 +353,9 @@ func (c *Cluster) ImportIntID(t testing.TB, index, field string, pairs []IntID) importRequest.Values[i] = pair.Val importRequest.ColumnIDs[i] = pair.ID } - if err := c.GetPrimary().API.ImportValue(context.Background(), nil, importRequest); err != nil { + qcx := c.GetPrimary().API.Txf().NewQcx() + defer qcx.Abort() + if err := c.GetPrimary().API.ImportValue(context.Background(), qcx, importRequest); err != nil { t.Fatalf("importing IntID data: %v", err) } } @@ -367,7 +380,9 @@ func (c *Cluster) ImportIDKey(t testing.TB, index, field string, pairs []KeyID) importRequest.RowIDs[i] = pair.ID importRequest.ColumnKeys[i] = pair.Key } - err := c.GetPrimary().API.Import(context.Background(), nil, importRequest) + qcx := c.GetPrimary().API.Txf().NewQcx() + defer qcx.Abort() + err := c.GetPrimary().API.Import(context.Background(), qcx, importRequest) if err != nil { t.Fatalf("importing IDKey data: %v", err) } diff --git a/tx_test.go b/tx_test.go index e53a91dbe..756fae6e5 100644 --- a/tx_test.go +++ b/tx_test.go @@ -20,7 +20,7 @@ import ( "strings" "testing" - "github.com/molecula/featurebase/v2" + pilosa "github.com/molecula/featurebase/v2" "github.com/molecula/featurebase/v2/http" "github.com/molecula/featurebase/v2/server" "github.com/molecula/featurebase/v2/storage" @@ -164,7 +164,12 @@ func TestAPI_ImportAtomicRecord(t *testing.T) { //vv("BEFORE the first ImportAtomicRecord!") - if err := m0api.ImportAtomicRecord(ctx, nil, air); err != nil { + qcx := m0api.Txf().NewQcx() + if err := m0api.ImportAtomicRecord(ctx, qcx, air); err != nil { + qcx.Abort() + t.Fatal(err) + } + if err := qcx.Finish(); err != nil { t.Fatal(err) } @@ -196,7 +201,7 @@ func TestAPI_ImportAtomicRecord(t *testing.T) { air = createAIRUpdate(expectedBalEndingAcct0, expectedBalEndingAcct1) - qcx := m0api.Txf().NewQcx() + qcx = m0api.Txf().NewQcx() //vv("just before the SECOND ImportAtomicRecord, qcx is %p, should NOT BE NIL", qcx) err = m0api.ImportAtomicRecord(ctx, qcx, air, opt) //err = m0api.ImportAtomicRecord(ctx, nil, air, opt) @@ -223,9 +228,12 @@ func TestAPI_ImportAtomicRecord(t *testing.T) { // happy path with no power failure half-way through. - err = m0api.ImportAtomicRecord(ctx, nil, air) + qcx = m0api.Txf().NewQcx() + err = m0api.ImportAtomicRecord(ctx, qcx, air) PanicOn(err) - + if err := qcx.Finish(); err != nil { + t.Fatal(err) + } eb0, eb1 := queryBalances(m0api, acctOwnerID, fieldAcct0, fieldAcct1, index) // should have been applied this time. @@ -240,7 +248,11 @@ func TestAPI_ImportAtomicRecord(t *testing.T) { air.Ivr[1].Clear = true air.Ir[0].Clear = true - err = m0api.ImportAtomicRecord(ctx, nil, air) + qcx = m0api.Txf().NewQcx() + err = m0api.ImportAtomicRecord(ctx, qcx, air) + if err := qcx.Finish(); err != nil { + t.Fatal(err) + } PanicOn(err) eb0, eb1 = queryBalances(m0api, acctOwnerID, fieldAcct0, fieldAcct1, index) From e768fc89ea7ab84a3906631630994a0fad1f5b2a Mon Sep 17 00:00:00 2001 From: Seebs Date: Fri, 23 Jul 2021 12:26:45 -0500 Subject: [PATCH 4/9] stop using pointers to time.Time We're reading timestamps as []int64, instead of allocating a time.Time for each timestamp, just use the same logic to determine whether to use the int64 timestamp that we would have used to decide whether to allocate it. We still have to check the whole run, though, because we're providing a large list of 0s instead of "no timestamps", for Reasons. --- api.go | 12 +++++------- field.go | 15 +++++++-------- index.go | 11 ----------- 3 files changed, 12 insertions(+), 26 deletions(-) diff --git a/api.go b/api.go index 7f22a52ed..0353b0f42 100644 --- a/api.go +++ b/api.go @@ -1607,14 +1607,12 @@ func (api *API) ImportWithTx(ctx context.Context, qcx *Qcx, req *ImportRequest, return errors.Wrap(err, "validating shard ownership") } - // Convert timestamps to time.Time. - timestamps := make([]*time.Time, len(req.Timestamps)) - for i, ts := range req.Timestamps { - if ts == 0 { - continue + var timestamps []int64 + for _, v := range req.Timestamps { + if v != 0 { + timestamps = req.Timestamps + break } - t := time.Unix(0, ts).UTC() - timestamps[i] = &t } // Import columnIDs into existence field. diff --git a/field.go b/field.go index 0bb88d6ce..839999aaf 100644 --- a/field.go +++ b/field.go @@ -1437,7 +1437,7 @@ func (f *Field) Range(qcx *Qcx, name string, op pql.Token, predicate int64) (*Ro } // Import bulk imports data. -func (f *Field) Import(qcx *Qcx, rowIDs, columnIDs []uint64, timestamps []*time.Time, opts ...ImportOption) (err0 error) { +func (f *Field) Import(qcx *Qcx, rowIDs, columnIDs []uint64, timestamps []int64, opts ...ImportOption) (err0 error) { // Set up import options. options := &ImportOptions{} @@ -1450,7 +1450,7 @@ func (f *Field) Import(qcx *Qcx, rowIDs, columnIDs []uint64, timestamps []*time. // Determine quantum if timestamps are set. q := f.TimeQuantum() - if hasTime(timestamps) { + if len(timestamps) > 0 { if q == "" { return errors.New("time quantum not set in field") } else if options.Clear { @@ -1470,16 +1470,15 @@ func (f *Field) Import(qcx *Qcx, rowIDs, columnIDs []uint64, timestamps []*time. return errors.New("bool field imports only support values 0 and 1") } - var timestamp *time.Time - if len(timestamps) > i { - timestamp = timestamps[i] - } + hasTime := len(timestamps) > i && timestamps[i] != 0 var standard []string - if timestamp == nil { + if !hasTime { standard = []string{viewStandard} } else { - standard = viewsByTime(viewStandard, *timestamp, q) + // Yes, we mean `0, ts`; ts is an int64 in UnixNano units, + // time.Unix takes seconds-and-nanoseconds. + standard = viewsByTime(viewStandard, time.Unix(0, timestamps[i]).UTC(), q) if !f.options.NoStandardView { // In order to match the logic of `SetBit()`, we want bits // with timestamps to write to both time and standard views. diff --git a/index.go b/index.go index b7204007f..370a328fb 100644 --- a/index.go +++ b/index.go @@ -22,7 +22,6 @@ import ( "sort" "strconv" "sync" - "time" "github.com/molecula/featurebase/v2/disco" "github.com/molecula/featurebase/v2/roaring" @@ -841,16 +840,6 @@ type IndexOptions struct { TrackExistence bool `json:"trackExistence"` } -// hasTime returns true if a contains a non-nil time. -func hasTime(a []*time.Time) bool { - for _, t := range a { - if t != nil { - return true - } - } - return false -} - type importKey struct { View string Shard uint64 From 1745a93aee5f1342cd2d8ba42e8fbf9186d7f57d Mon Sep 17 00:00:00 2001 From: Seebs Date: Thu, 29 Jul 2021 13:08:15 -0500 Subject: [PATCH 5/9] allow "us" for microseconds in timestamp units The convention of using a "u" for "micro" is pretty well-established and some people will have trouble typing the Greek letter, accept that as a synonym. --- field.go | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/field.go b/field.go index 839999aaf..bfffc3509 100644 --- a/field.go +++ b/field.go @@ -2091,13 +2091,14 @@ const ( TimeUnitSeconds = "s" TimeUnitMilliseconds = "ms" TimeUnitMicroseconds = "µs" + TimeUnitUSeconds = "us" TimeUnitNanoseconds = "ns" ) // IsValidTimeUnit returns true if unit is valid. func IsValidTimeUnit(unit string) bool { switch unit { - case TimeUnitSeconds, TimeUnitMilliseconds, TimeUnitMicroseconds, TimeUnitNanoseconds: + case TimeUnitSeconds, TimeUnitMilliseconds, TimeUnitMicroseconds, TimeUnitUSeconds, TimeUnitNanoseconds: return true default: return false @@ -2111,7 +2112,7 @@ func TimeUnitNanos(unit string) int64 { return int64(time.Second) case TimeUnitMilliseconds: return int64(time.Millisecond) - case TimeUnitMicroseconds: + case TimeUnitMicroseconds, TimeUnitUSeconds: return int64(time.Microsecond) default: return int64(time.Nanosecond) From ab9b70d7e84bf6b7f3bb5d23733a89242fdcf195 Mon Sep 17 00:00:00 2001 From: Seebs Date: Wed, 11 Aug 2021 15:22:54 -0500 Subject: [PATCH 6/9] mapperLocal: actually leave loop on read from done channel staticcheck points out that the break is otherwise an ineffective break because it just ends the current case clause of the switch it's in, which is true. --- executor.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/executor.go b/executor.go index ec071f34f..8cd24f33e 100644 --- a/executor.go +++ b/executor.go @@ -5939,6 +5939,7 @@ func (e *executor) mapperLocal(ctx context.Context, shards []uint64, mapFn mapFu ch := make(chan mapResponse, len(shards)) expected := 0 + shardLoop: for _, shard := range shards { j := job{ shard: shard, @@ -5949,7 +5950,7 @@ func (e *executor) mapperLocal(ctx context.Context, shards []uint64, mapFn mapFu } select { case <-done: - break + break shardLoop case e.work <- j: expected++ } From 6af987a5badc10500feef0c42abdf6e93807a605 Mon Sep 17 00:00:00 2001 From: Seebs Date: Wed, 11 Aug 2021 15:29:03 -0500 Subject: [PATCH 7/9] fix error formatting/spelling Go convention is that error messages don't end with periods and don't start with capital letters. --- bolt.go | 2 +- rbf.go | 2 +- time.go | 10 +++++----- txfactory.go | 10 +++++----- 4 files changed, 12 insertions(+), 12 deletions(-) diff --git a/bolt.go b/bolt.go index 366b77391..e679b1400 100644 --- a/bolt.go +++ b/bolt.go @@ -281,7 +281,7 @@ func (w *BoltWrapper) DeleteIndex(indexName string) error { // index name in the key prefix, so we cannot allow indexNames // themselves to contain apostrophies. if strings.Contains(indexName, "/") { - return fmt.Errorf("error: bad indexName `%v` in BoltWrapper.DeleteIndex() call: indexName cannot contain '/'.", indexName) + return fmt.Errorf("error: bad indexName `%v` in BoltWrapper.DeleteIndex() call: indexName cannot contain '/'", indexName) } prefix := txkey.IndexOnlyPrefix(indexName) return w.DeletePrefix(prefix) diff --git a/rbf.go b/rbf.go index e75e7768e..64a268281 100644 --- a/rbf.go +++ b/rbf.go @@ -542,7 +542,7 @@ func (w *RbfDBWrapper) DeleteField(index, field, fieldPath string) error { func (w *RbfDBWrapper) DeleteIndex(indexName string) error { if strings.Contains(indexName, "'") { - return fmt.Errorf("error: bad indexName `%v` in RbfDBWrapper.DeleteIndex() call: indexName cannot contain apostrophes/single quotes.", indexName) + return fmt.Errorf("error: bad indexName `%v` in RbfDBWrapper.DeleteIndex() call: indexName cannot contain apostrophes/single quotes", indexName) } prefix := txkey.IndexOnlyPrefix(indexName) diff --git a/time.go b/time.go index b87774f1a..a35ba7dd4 100644 --- a/time.go +++ b/time.go @@ -257,16 +257,16 @@ func parsePartialTime(t string) (time.Time, error) { // has minutes minute, err = strconv.Atoi(subStrings[1]) if err != nil { - return -1, -1, errors.New("Invalid Time") + return -1, -1, errors.New("invalid time") } fallthrough case 1: hour, err = strconv.Atoi(subStrings[0]) if err != nil { - return -1, -1, errors.New("Invalid Time") + return -1, -1, errors.New("invalid time") } default: - return -1, -1, errors.New("Invalid Time") + return -1, -1, errors.New("invalid time") } return @@ -282,7 +282,7 @@ func parsePartialTime(t string) (time.Time, error) { return true } if len(subMatches) <= 1 { - return nil, errors.New("Invalid time") + return nil, errors.New("invalid time") } // ignore full match which is at index 0 subMatches = subMatches[1:] // ignore full match which is at index 0 @@ -295,7 +295,7 @@ func parsePartialTime(t string) (time.Time, error) { } else { // rest must be empty for date-time to be valid if !restAreEmpty(subMatches[i:]) { - return nil, errors.New("Invalid date-time") + return nil, errors.New("invalid date-time") } break } diff --git a/txfactory.go b/txfactory.go index 84742cd2e..93e9449df 100644 --- a/txfactory.go +++ b/txfactory.go @@ -1383,7 +1383,7 @@ func (f *TxFactory) greenHasData() (hasData bool, err error) { case 2: return f.dbPerShard.HasData(1) } - err = fmt.Errorf("unsupported len(f.types): %v. Must be 1 or 2.", n) + err = fmt.Errorf("unsupported len(f.types): %v; must be 1 or 2", n) PanicOn(err) return } @@ -1410,7 +1410,7 @@ func (f *TxFactory) green2blue(holder *Holder) (err0 error) { greenSrc := f.types[1] if blueDest == roaringTxn { - return fmt.Errorf("error: cannot migrate to 'roaring': not implemented.") + return fmt.Errorf("error: cannot migrate to 'roaring': not implemented") } idxs := holder.Indexes() @@ -1427,13 +1427,13 @@ func (f *TxFactory) green2blue(holder *Holder) (err0 error) { return errors.Wrap(err, "TxFactory.green2blue f.greenHasData()") } if !blueHasData && !greenHasData { - holder.Logger.Infof("no data in blue or green. No migration or verification to do.") + holder.Logger.Infof("no data in blue or green. No migration or verification to do") return nil } // INVAR: blue has data. if !greenHasData { - holder.Logger.Errorf("cannot migrate from green '%v' because it has no data in it.", greenSrc) - return fmt.Errorf("error: cannot migrate from green '%v' because it has no data in it.", greenSrc) + holder.Logger.Errorf("cannot migrate from green '%v' because it has no data in it", greenSrc) + return fmt.Errorf("error: cannot migrate from green '%v' because it has no data in it", greenSrc) } nGoro := runtime.NumCPU() From f019cc740923af39eb9d78b6ba42031cfcf48196 Mon Sep 17 00:00:00 2001 From: Seebs Date: Mon, 26 Jul 2021 13:29:55 -0500 Subject: [PATCH 8/9] make import correctly reflect that it needs a single shard always In fact, we have a number of things assuming that values passed to Import always fit within a single known shard, so, drop all the extra complexity around this, drop the computation of fancy view/shard keys, and so on. There's a lot of room left to improve this probably but it's at least better, I think. Unfortunately, there's a handful of things, basically all of which are test cases, which were relying on this, so, we also add functionality for splitting import requests by shards. But this allows us to stop duplicating each shard's inputs one at a time... which turns out to mean that we now care that the import operation can write back to the import request. This only affects test cases, so we adopt a crufty hack involving cloning import requests in those rare cases, and also when reusing the same column IDs to write to the existence field that we'd be using later to write to another field. Note that even if we weren't overwriting the column IDs with positions, we'd be sorting the column/row ID lists by row-then-column, which means we'd still be corrupting the column ID lists. This may want to change at some point. We also reuse a single Tx for all the views, because DB-per-shard means that should work fine, and reduces the cost of doing these updates, probably. --- api.go | 27 +++++---- api_test.go | 4 +- executor_test.go | 11 +++- field.go | 121 ++++++++++++++++++++++++-------------- field_internal_test.go | 11 ++-- fragment_internal_test.go | 4 ++ handler.go | 118 +++++++++++++++++++++++++++++++++++++ http/client.go | 2 +- index.go | 5 -- tx_test.go | 4 +- 10 files changed, 236 insertions(+), 71 deletions(-) diff --git a/api.go b/api.go index 0353b0f42..bbdb8c6b4 100644 --- a/api.go +++ b/api.go @@ -1619,7 +1619,7 @@ func (api *API) ImportWithTx(ctx context.Context, qcx *Qcx, req *ImportRequest, // Note: req.Shard may not be the only shard imported into here, // so don't expect it to be invariant. if !options.Clear { - if err := importExistenceColumns(qcx, idx, req.ColumnIDs); err != nil { + if err := importExistenceColumns(qcx, idx, req.ColumnIDs, req.Shard); err != nil { api.server.logger.Errorf("import existence error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err) return err } @@ -1629,7 +1629,7 @@ func (api *API) ImportWithTx(ctx context.Context, qcx *Qcx, req *ImportRequest, } // Import into fragment. - err = field.Import(qcx, req.RowIDs, req.ColumnIDs, timestamps, opts...) + err = field.Import(qcx, req.RowIDs, req.ColumnIDs, timestamps, req.Shard, opts...) if err != nil { api.server.logger.Errorf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err) return errors.Wrap(err, "importing") @@ -1728,8 +1728,9 @@ func (api *API) ImportValueWithTx(ctx context.Context, qcx *Qcx, req *ImportValu // if we're importing into a specific shard if req.Shard != math.MaxUint64 { // Check that column IDs match the stated shard. - if s1, s2 := req.ColumnIDs[0]/ShardWidth, req.ColumnIDs[len(req.ColumnIDs)-1]/ShardWidth; s1 != s2 && s2 != req.Shard { - return errors.Errorf("shard %d specified, but import spans shards %d to %d", req.Shard, s1, s2) + shard := req.ColumnIDs[0] / ShardWidth + if s2 := req.ColumnIDs[len(req.ColumnIDs)-1] / ShardWidth; (shard != s2) || (shard != req.Shard) { + return errors.Errorf("shard %d specified, but import spans shards %d to %d", req.Shard, shard, s2) } // Validate shard ownership. TODO - we should forward to the // correct node rather than barfing here. @@ -1738,7 +1739,7 @@ func (api *API) ImportValueWithTx(ctx context.Context, qcx *Qcx, req *ImportValu } // Import columnIDs into existence field. if !options.Clear { - if err := importExistenceColumns(qcx, idx, req.ColumnIDs); err != nil { + if err := importExistenceColumns(qcx, idx, req.ColumnIDs, shard); err != nil { api.server.logger.Errorf("import existence error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err) return errors.Wrap(err, "importing existence columns") } @@ -1746,17 +1747,17 @@ func (api *API) ImportValueWithTx(ctx context.Context, qcx *Qcx, req *ImportValu // Import into fragment. if len(req.Values) > 0 { - err = field.importValue(qcx, req.ColumnIDs, req.Values, options) + err = field.importValue(qcx, req.ColumnIDs, req.Values, shard, options) if err != nil { api.server.logger.Errorf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err) } } else if len(req.TimestampValues) > 0 { - err = field.importTimestampValue(qcx, req.ColumnIDs, req.TimestampValues, options) + err = field.importTimestampValue(qcx, req.ColumnIDs, req.TimestampValues, shard, options) if err != nil { api.server.logger.Errorf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err) } } else if len(req.FloatValues) > 0 { - err = field.importFloatValue(qcx, req.ColumnIDs, req.FloatValues, options) + err = field.importFloatValue(qcx, req.ColumnIDs, req.FloatValues, shard, options) if err != nil { api.server.logger.Errorf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err) } @@ -1814,14 +1815,20 @@ func (api *API) ImportValueWithTx(ctx context.Context, qcx *Qcx, req *ImportValu return nil } -func importExistenceColumns(qcx *Qcx, index *Index, columnIDs []uint64) error { +func importExistenceColumns(qcx *Qcx, index *Index, columnIDs []uint64, shard uint64) error { ef := index.existenceField() if ef == nil { return nil } existenceRowIDs := make([]uint64, len(columnIDs)) - return ef.Import(qcx, existenceRowIDs, columnIDs, nil) + // If we don't gratuitously hand-duplicate things in field.Import, + // the fact that fragment.bulkImport rewrites its row and column + // lists can burn us if we don't make a copy before doing the + // existence field write. + columnCopy := make([]uint64, len(columnIDs)) + copy(columnCopy, columnIDs) + return ef.Import(qcx, existenceRowIDs, columnCopy, nil, shard) } // ShardDistribution returns an object representing the distribution of shards diff --git a/api_test.go b/api_test.go index 4a4418ed9..201a7c89a 100644 --- a/api_test.go +++ b/api_test.go @@ -482,10 +482,10 @@ func TestAPI_ClearFlagForImportAndImportValues(t *testing.T) { } qcx := m0api.Txf().NewQcx() - if err := m0api.Import(ctx, qcx, ir0); err != nil { + if err := m0api.Import(ctx, qcx, ir0.Clone()); err != nil { t.Fatal(err) } - if err := m0api.ImportValue(ctx, qcx, ivr0); err != nil { + if err := m0api.ImportValue(ctx, qcx, ivr0.Clone()); err != nil { t.Fatal(err) } PanicOn(qcx.Finish()) diff --git a/executor_test.go b/executor_test.go index 46b3e00ef..8440adb9d 100644 --- a/executor_test.go +++ b/executor_test.go @@ -4131,10 +4131,17 @@ func TestExecutor_Execute_All(t *testing.T) { req.ColumnIDs[bitCount-1] = uint64((3 * ShardWidth) + 2) m0 := c.GetNode(0) + // the request gets altered by the Import operation now... + reqs, err := req.Clone().ShardSplit() + if err != nil { + t.Fatalf("splitting request into shards: %v", err) + } qcx := m0.API.Txf().NewQcx() - if err := m0.API.Import(context.Background(), qcx, req); err != nil { - t.Fatal(err) + for _, r := range reqs { + if err := m0.API.Import(context.Background(), qcx, r); err != nil { + t.Fatal(err) + } } PanicOn(qcx.Finish()) diff --git a/field.go b/field.go index bfffc3509..dee36745d 100644 --- a/field.go +++ b/field.go @@ -1437,7 +1437,7 @@ func (f *Field) Range(qcx *Qcx, name string, op pql.Token, predicate int64) (*Ro } // Import bulk imports data. -func (f *Field) Import(qcx *Qcx, rowIDs, columnIDs []uint64, timestamps []int64, opts ...ImportOption) (err0 error) { +func (f *Field) Import(qcx *Qcx, rowIDs, columnIDs []uint64, timestamps []int64, shard uint64, opts ...ImportOption) (err0 error) { // Set up import options. options := &ImportOptions{} @@ -1456,12 +1456,57 @@ func (f *Field) Import(qcx *Qcx, rowIDs, columnIDs []uint64, timestamps []int64, } else if options.Clear { return errors.New("import clear is not supported with timestamps") } + } else { + // short path: if we don't have any timestamps, we only need + // to write to exactly one view, which is always viewStandard, + // and *every* bit goes into that view, and we already verified that + // everything is in the same shard, so we can skip most of this. + fieldType := f.Type() + if fieldType == FieldTypeBool { + for _, rowID := range rowIDs { + if rowID > 1 { + return errors.New("bool field imports only support values 0 and 1") + } + } + } + tx, finisher, err := qcx.GetTx(Txo{Write: true, Index: f.idx, Shard: shard}) + if err != nil { + return errors.Wrap(err, "qcx.GetTx") + } + var err1 error + defer finisher(&err1) + view, err := f.createViewIfNotExists(viewStandard) + if err != nil { + return errors.Wrapf(err, "creating view %s", viewStandard) + } + + frag, err := view.CreateFragmentIfNotExists(shard) + if err != nil { + return errors.Wrap(err, "creating fragment") + } + + err1 = frag.bulkImport(tx, rowIDs, columnIDs, options) + return err1 } fieldType := f.Type() // Split import data by fragment. - dataByFragment := make(map[importKey]importData) + views := make(map[string]int) + var allData []importData + see := func(name string, columnID uint64, rowID uint64) { + var ok bool + var idx int + if idx, ok = views[name]; !ok { + allData = append(allData, importData{}) + idx = len(allData) + views[name] = idx + } + data := allData[idx] + data.RowIDs = append(data.RowIDs, rowID) + data.ColumnIDs = append(data.ColumnIDs, columnID) + allData[idx] = data + } for i := range rowIDs { rowID, columnID := rowIDs[i], columnIDs[i] @@ -1472,53 +1517,41 @@ func (f *Field) Import(qcx *Qcx, rowIDs, columnIDs []uint64, timestamps []int64, hasTime := len(timestamps) > i && timestamps[i] != 0 - var standard []string - if !hasTime { - standard = []string{viewStandard} - } else { - // Yes, we mean `0, ts`; ts is an int64 in UnixNano units, - // time.Unix takes seconds-and-nanoseconds. - standard = viewsByTime(viewStandard, time.Unix(0, timestamps[i]).UTC(), q) - if !f.options.NoStandardView { - // In order to match the logic of `SetBit()`, we want bits - // with timestamps to write to both time and standard views. - standard = append(standard, viewStandard) + // attach bit to standard view unless we have a timestamp and + // have the NoStandardView option set + if !hasTime || !f.options.NoStandardView { + see(viewStandard, columnID, rowID) + } + if hasTime { + // attach bit to all the views for this timestamp + views := viewsByTime(viewStandard, time.Unix(0, timestamps[i]).UTC(), q) + for _, view := range views { + see(view, columnID, rowID) } } - - // Attach bit to each standard view. - for _, name := range standard { - key := importKey{View: name, Shard: columnID / ShardWidth} - data := dataByFragment[key] - data.RowIDs = append(data.RowIDs, rowID) - data.ColumnIDs = append(data.ColumnIDs, columnID) - dataByFragment[key] = data - } } - - // Import into each fragment. - for key, data := range dataByFragment { - view, err := f.createViewIfNotExists(key.View) + tx, finisher, err := qcx.GetTx(Txo{Write: true, Index: f.idx, Shard: shard}) + if err != nil { + return errors.Wrap(err, "qcx.GetTx") + } + var err1 error + defer finisher(&err1) + for viewName, idx := range views { + data := allData[idx] + view, err := f.createViewIfNotExists(viewName) if err != nil { - return errors.Wrap(err, "creating view") + return errors.Wrapf(err, "creating view %s", viewName) } - frag, err := view.CreateFragmentIfNotExists(key.Shard) + frag, err := view.CreateFragmentIfNotExists(shard) if err != nil { return errors.Wrap(err, "creating fragment") } - tx, finisher, err := qcx.GetTx(Txo{Write: true, Index: frag.idx, Fragment: frag, Shard: frag.shard}) - if err != nil { - return errors.Wrap(err, "qcx.GetTx") - } - - err1 := frag.bulkImport(tx, data.RowIDs, data.ColumnIDs, options) + err1 = frag.bulkImport(tx, data.RowIDs, data.ColumnIDs, options) if err1 != nil { - finisher(&err1) return err1 } - finisher(nil) } return nil } @@ -1526,7 +1559,7 @@ func (f *Field) Import(qcx *Qcx, rowIDs, columnIDs []uint64, timestamps []int64, // importFloatValue imports floating point values. In current usage, this // should only ever be called with data for a single shard; the API calls // around this are splitting it up per shard. -func (f *Field) importFloatValue(qcx *Qcx, columnIDs []uint64, values []float64, options *ImportOptions) error { +func (f *Field) importFloatValue(qcx *Qcx, columnIDs []uint64, values []float64, shard uint64, options *ImportOptions) error { // convert values to int64 values based on scale ivalues := make([]int64, len(values)) bsig := f.bsiGroup(f.name) @@ -1538,13 +1571,13 @@ func (f *Field) importFloatValue(qcx *Qcx, columnIDs []uint64, values []float64, ivalues[i] = int64(fval * mult) } // then call importValue - return f.importValue(qcx, columnIDs, ivalues, options) + return f.importValue(qcx, columnIDs, ivalues, shard, options) } // importFloatValue imports timestamp values. In current usage, this // should only ever be called with data for a single shard; the API calls // around this are splitting it up per shard. -func (f *Field) importTimestampValue(qcx *Qcx, columnIDs []uint64, values []time.Time, options *ImportOptions) error { +func (f *Field) importTimestampValue(qcx *Qcx, columnIDs []uint64, values []time.Time, shard uint64, options *ImportOptions) error { ivalues := make([]int64, len(values)) bsig := f.bsiGroup(f.name) if bsig == nil { @@ -1554,13 +1587,13 @@ func (f *Field) importTimestampValue(qcx *Qcx, columnIDs []uint64, values []time for i, t := range values { ivalues[i] = t.UnixNano() / TimeUnitNanos(f.options.TimeUnit) } - return f.importValue(qcx, columnIDs, ivalues, options) + return f.importValue(qcx, columnIDs, ivalues, shard, options) } // importValue bulk imports range-encoded value data. This function should // only be called with data for a single shard; the API calls that wrap // this handle splitting the data up per-shard. -func (f *Field) importValue(qcx *Qcx, columnIDs []uint64, values []int64, options *ImportOptions) (err0 error) { +func (f *Field) importValue(qcx *Qcx, columnIDs []uint64, values []int64, shard uint64, options *ImportOptions) (err0 error) { // no data to import if len(columnIDs) == 0 { return nil @@ -1613,9 +1646,9 @@ func (f *Field) importValue(qcx *Qcx, columnIDs []uint64, values []int64, option } f.mu.Unlock() - // Since all data should be for the same shard, we can just compute - // this from the first value. - shard := columnIDs[0] / ShardWidth + if columnIDs[0]/ShardWidth != shard { + return fmt.Errorf("requested import for shard %d, got record ID for shard %d", shard, columnIDs[0]/ShardWidth) + } view, err := f.createViewIfNotExists(viewName) if err != nil { diff --git a/field_internal_test.go b/field_internal_test.go index 41f288016..fc2296845 100644 --- a/field_internal_test.go +++ b/field_internal_test.go @@ -494,7 +494,7 @@ func TestBSIGroup_importValue(t *testing.T) { []uint64{100}, }, } { - if err := f.importValue(qcx, tt.columnIDs, tt.values, options); err != nil { + if err := f.importValue(qcx, tt.columnIDs, tt.values, 0, options); err != nil { t.Fatalf("test %d, importing values: %s", i, err.Error()) } PanicOn(qcx.Finish()) @@ -514,7 +514,8 @@ func TestBSIGroup_importValue(t *testing.T) { func benchmarkFieldImportValues(b *testing.B, qcx *Qcx, bitDepth uint64, f *TestField, cfunc func(uint64) uint64) { batches := makeBenchmarkImportValueData(b, bitDepth, cfunc) for _, req := range batches { - err := f.importValue(qcx, req.ColumnIDs, req.Values, &ImportOptions{}) + // NOTE: We assume everything's in Shard 0 for now. + err := f.importValue(qcx, req.ColumnIDs, req.Values, 0, &ImportOptions{}) if err != nil { b.Fatalf("error importing values: %s", err) } @@ -591,7 +592,7 @@ func TestIntField_MinMaxForShard(t *testing.T) { }, } { t.Run(test.name+strconv.Itoa(i), func(t *testing.T) { - if err := f.importValue(qcx, test.columnIDs, test.values, options); err != nil { + if err := f.importValue(qcx, test.columnIDs, test.values, 0, options); err != nil { t.Fatalf("test %d, importing values: %s", i, err.Error()) } PanicOn(qcx.Finish()) @@ -751,7 +752,7 @@ func TestDecimalField_MinMaxForShard(t *testing.T) { }, } { t.Run(test.name+strconv.Itoa(i), func(t *testing.T) { - if err := f.importFloatValue(qcx, test.columnIDs, test.values, options); err != nil { + if err := f.importFloatValue(qcx, test.columnIDs, test.values, 0, options); err != nil { t.Fatalf("test %d, importing values: %s", i, err.Error()) } @@ -811,7 +812,7 @@ func TestBSIGroup_TxReopenDB(t *testing.T) { []uint64{100}, }, } { - if err := f.importValue(qcx, tt.columnIDs, tt.values, options); err != nil { + if err := f.importValue(qcx, tt.columnIDs, tt.values, 0, options); err != nil { t.Fatalf("test %d, importing values: %s", i, err.Error()) } PanicOn(qcx.Finish()) diff --git a/fragment_internal_test.go b/fragment_internal_test.go index 7de7b5c1f..3820d25fa 100644 --- a/fragment_internal_test.go +++ b/fragment_internal_test.go @@ -1020,6 +1020,10 @@ func BenchmarkFragment_SetValue(b *testing.B) { } } +// makeBenchmarkImportValueData produces data that's supposed to be all within +// the same shard; for fragment purposes, implicitly shard 0. This also gets +// used by the field tests, but import requests are supposed to be per-shard, +// so it's important that we generate values only within a given shard. func makeBenchmarkImportValueData(b *testing.B, bitDepth uint64, cfunc func(uint64) uint64) []ImportValueRequest { b.StopTimer() column := uint64(0) diff --git a/handler.go b/handler.go index 49c4b1ff5..a07c3d6c9 100644 --- a/handler.go +++ b/handler.go @@ -18,6 +18,7 @@ import ( "encoding/json" "time" + "github.com/molecula/featurebase/v2/shardwidth" "github.com/molecula/featurebase/v2/tracing" "github.com/pkg/errors" ) @@ -129,6 +130,41 @@ type ImportValueRequest struct { scratch []int // scratch space to allow us to get a stable sort in reasonable time } +func (ivr *ImportValueRequest) Clone() *ImportValueRequest { + newIVR := &ImportValueRequest{} + if ivr == nil { + return newIVR + } + *newIVR = *ivr + // don't copy the internal scratch buffer + newIVR.scratch = nil + if len(ivr.ColumnIDs) > 0 { + newIVR.ColumnIDs = make([]uint64, len(ivr.ColumnIDs)) + copy(newIVR.ColumnIDs, ivr.ColumnIDs) + } + if len(ivr.ColumnKeys) > 0 { + newIVR.ColumnKeys = make([]string, len(ivr.ColumnKeys)) + copy(newIVR.ColumnKeys, ivr.ColumnKeys) + } + if len(ivr.Values) > 0 { + newIVR.Values = make([]int64, len(ivr.Values)) + copy(newIVR.Values, ivr.Values) + } + if len(ivr.FloatValues) > 0 { + newIVR.FloatValues = make([]float64, len(ivr.FloatValues)) + copy(newIVR.FloatValues, ivr.FloatValues) + } + if len(ivr.TimestampValues) > 0 { + newIVR.TimestampValues = make([]time.Time, len(ivr.TimestampValues)) + copy(newIVR.TimestampValues, ivr.TimestampValues) + } + if len(ivr.StringValues) > 0 { + newIVR.StringValues = make([]string, len(ivr.StringValues)) + copy(newIVR.StringValues, ivr.StringValues) + } + return newIVR +} + // AtomicRecord applies all its Ivr and Ivr atomically, in a Tx. // The top level Shard has to agree with Ivr[i].Shard and the Iv[i].Shard // for all i included (in Ivr and Ir). The same goes for the top level Index: all records @@ -142,6 +178,19 @@ type AtomicRecord struct { Ir []*ImportRequest // other field types, e.g. single bit } +func (ar *AtomicRecord) Clone() *AtomicRecord { + newAR := &AtomicRecord{Index: ar.Index, Shard: ar.Shard} + newAR.Ivr = make([]*ImportValueRequest, len(ar.Ivr)) + for i, vr := range ar.Ivr { + newAR.Ivr[i] = vr.Clone() + } + newAR.Ir = make([]*ImportRequest, len(ar.Ir)) + for i, vr := range ar.Ir { + newAR.Ir[i] = vr.Clone() + } + return newAR +} + func (ivr *ImportValueRequest) Len() int { return len(ivr.ColumnIDs) } func (ivr *ImportValueRequest) Less(i, j int) bool { if ivr.ColumnIDs[i] < ivr.ColumnIDs[j] { @@ -225,6 +274,75 @@ type ImportRequest struct { Clear bool } +// Clone allows copying an import request. Normally you wouldn't, but +// some import functions are destructive on their inputs, and if you +// want to *re-use* an import request, you might need this. If you're +// using this outside tx_test, something is probably wrong. +func (ir *ImportRequest) Clone() *ImportRequest { + newIR := &ImportRequest{} + if ir == nil { + return newIR + } + *newIR = *ir + if ir.RowIDs != nil { + newIR.RowIDs = make([]uint64, len(ir.RowIDs)) + copy(newIR.RowIDs, ir.RowIDs) + } + if ir.ColumnIDs != nil { + newIR.ColumnIDs = make([]uint64, len(ir.ColumnIDs)) + copy(newIR.ColumnIDs, ir.ColumnIDs) + } + if ir.RowKeys != nil { + newIR.RowKeys = make([]string, len(ir.RowKeys)) + copy(newIR.RowKeys, ir.RowKeys) + } + if ir.ColumnKeys != nil { + newIR.ColumnKeys = make([]string, len(ir.ColumnKeys)) + copy(newIR.ColumnKeys, ir.ColumnKeys) + } + if ir.Timestamps != nil { + newIR.Timestamps = make([]int64, len(ir.Timestamps)) + copy(newIR.Timestamps, ir.Timestamps) + } + return newIR +} + +// ShardSplit splits the request into a slice of import requests. It requires +// that the original request have all elements sorted, and already have +// column IDs, not column keys. +func (ir *ImportRequest) ShardSplit() ([]*ImportRequest, error) { + if ir == nil { + return nil, nil + } + // fix shard + if len(ir.ColumnIDs) < 2 { + ir.Shard = ir.ColumnIDs[0] >> shardwidth.Exponent + return []*ImportRequest{ir}, nil + } + shards, ends := shardwidth.FindShards(ir.ColumnIDs) + out := make([]*ImportRequest, len(shards)) + prev := 0 + for i, shard := range shards { + next := ends[i] + newIR := &ImportRequest{} + *newIR = *ir + newIR.ColumnIDs = ir.ColumnIDs[prev:next:next] + if ir.RowIDs != nil { + newIR.RowIDs = ir.RowIDs[prev:next:next] + } + if ir.RowKeys != nil { + newIR.RowKeys = ir.RowKeys[prev:next:next] + } + if ir.Timestamps != nil { + newIR.Timestamps = ir.Timestamps[prev:next:next] + } + newIR.Shard = shard + out[i] = newIR + prev = next + } + return out, nil +} + // ValidateWithTimestamp ensures that the payload of the request is valid. func (ir *ImportRequest) ValidateWithTimestamp(indexCreatedAt, fieldCreatedAt int64) error { if (ir.IndexCreatedAt != 0 && ir.IndexCreatedAt != indexCreatedAt) || diff --git a/http/client.go b/http/client.go index a51f54233..ce31012d7 100644 --- a/http/client.go +++ b/http/client.go @@ -30,7 +30,7 @@ import ( "strings" "time" - "github.com/molecula/featurebase/v2" + pilosa "github.com/molecula/featurebase/v2" "github.com/molecula/featurebase/v2/encoding/proto" pnet "github.com/molecula/featurebase/v2/net" "github.com/molecula/featurebase/v2/topology" diff --git a/index.go b/index.go index 370a328fb..25b7bd5ee 100644 --- a/index.go +++ b/index.go @@ -840,11 +840,6 @@ type IndexOptions struct { TrackExistence bool `json:"trackExistence"` } -type importKey struct { - View string - Shard uint64 -} - type importData struct { RowIDs []uint64 ColumnIDs []uint64 diff --git a/tx_test.go b/tx_test.go index 756fae6e5..e383a86b1 100644 --- a/tx_test.go +++ b/tx_test.go @@ -203,7 +203,7 @@ func TestAPI_ImportAtomicRecord(t *testing.T) { qcx = m0api.Txf().NewQcx() //vv("just before the SECOND ImportAtomicRecord, qcx is %p, should NOT BE NIL", qcx) - err = m0api.ImportAtomicRecord(ctx, qcx, air, opt) + err = m0api.ImportAtomicRecord(ctx, qcx, air.Clone(), opt) //err = m0api.ImportAtomicRecord(ctx, nil, air, opt) if err != pilosa.ErrAborted { PanicOn(fmt.Sprintf("expected ErrTxnAborted but got err='%#v'", err)) @@ -229,7 +229,7 @@ func TestAPI_ImportAtomicRecord(t *testing.T) { // happy path with no power failure half-way through. qcx = m0api.Txf().NewQcx() - err = m0api.ImportAtomicRecord(ctx, qcx, air) + err = m0api.ImportAtomicRecord(ctx, qcx, air.Clone()) PanicOn(err) if err := qcx.Finish(); err != nil { t.Fatal(err) From 423cddbdbb7e5f345ac05e56e8eab0ea16ba3f3d Mon Sep 17 00:00:00 2001 From: Seebs Date: Wed, 4 Aug 2021 15:55:51 -0500 Subject: [PATCH 9/9] avoid allocations in viewsByTime This is sort of horrible, but viewsByTime was about 25% of total CPU time in the ingest path, NOT including increased GC overhead. This overoptimized approach to letting us recycle a buffer, and use the same buffer for multiple time views at once, reduces that to about 2.5%. Sorry for the mess. We also streamline the process of building the per-view data sets a bit, and streamline it a lot in the non-time-quantum case. --- field.go | 55 +++++++++++++++++++++++++------------- time.go | 62 ++++++++++++++++++++++++++++++++++++++++--- time_internal_test.go | 32 ++++++++++++++++++++++ 3 files changed, 127 insertions(+), 22 deletions(-) diff --git a/field.go b/field.go index dee36745d..5f147d735 100644 --- a/field.go +++ b/field.go @@ -935,7 +935,7 @@ func (f *Field) RowTime(tx Tx, rowID uint64, time time.Time, quantum string) (*R if !TimeQuantum(quantum).Valid() { return nil, ErrInvalidTimeQuantum } - viewname := viewsByTime(viewStandard, time, TimeQuantum(quantum[len(quantum)-1:]))[0] + viewname := viewByTimeUnit(viewStandard, time, rune(quantum[len(quantum)-1])) view := f.view(viewname) if view == nil { return nil, errors.Errorf("view with quantum %v not found.", quantum) @@ -1492,20 +1492,38 @@ func (f *Field) Import(qcx *Qcx, rowIDs, columnIDs []uint64, timestamps []int64, fieldType := f.Type() // Split import data by fragment. - views := make(map[string]int) - var allData []importData - see := func(name string, columnID uint64, rowID uint64) { + views := make(map[string]*importData) + var timeStringBuf []byte + var timeViews [][]byte + if len(q) > 0 { + // We're supporting time quantums, so we need to store bits in a + // number of views for every entry with a timestamp. We want to compute + // time quantum view names for whatever combination of YMDH views + // we have. But we don't want to allocate four strings per entry, or + // recompute and recreate the entire string. We know that only the + // YYYYMMDDHH part of the string changes over time. + timeStringBuf = make([]byte, len(viewStandard) + 11) + copy(timeStringBuf, []byte(viewStandard)) + copy(timeStringBuf[len(viewStandard):], []byte("_YYYYMMDDHH")) + // Now we have a buffer that contains + // `standard_YYYYMMDDHH`. We also need storage space to hold several + // slice headers, one per entry in q. These will hold the view names + // corresponding to each letter in q. + timeViews = make([][]byte, len(q)) + } + // This helper function records that a given column/row pair is relevant + // to a specific view. We use a map lookup for the strings, but do the + // actual operations using a slice so we're only writing each map entry + // once, not once on every update. + see := func(name []byte, columnID uint64, rowID uint64) { var ok bool - var idx int - if idx, ok = views[name]; !ok { - allData = append(allData, importData{}) - idx = len(allData) - views[name] = idx + var data *importData + if data, ok = views[string(name)]; !ok { + data = &importData{} + views[string(name)] = data } - data := allData[idx] data.RowIDs = append(data.RowIDs, rowID) data.ColumnIDs = append(data.ColumnIDs, columnID) - allData[idx] = data } for i := range rowIDs { rowID, columnID := rowIDs[i], columnIDs[i] @@ -1520,13 +1538,15 @@ func (f *Field) Import(qcx *Qcx, rowIDs, columnIDs []uint64, timestamps []int64, // attach bit to standard view unless we have a timestamp and // have the NoStandardView option set if !hasTime || !f.options.NoStandardView { - see(viewStandard, columnID, rowID) + see([]byte(viewStandard), columnID, rowID) } if hasTime { - // attach bit to all the views for this timestamp - views := viewsByTime(viewStandard, time.Unix(0, timestamps[i]).UTC(), q) - for _, view := range views { - see(view, columnID, rowID) + // attach bit to all the views for this timestamp. note that the + // `timeViews` slice gets resliced and reused by this process, so + // we don't have to allocate millions of tiny slices of slice headers. + timeViews = viewsByTimeInto(timeStringBuf, timeViews, time.Unix(0, timestamps[i]).UTC(), q) + for _, v := range timeViews { + see(v, columnID, rowID) } } } @@ -1536,8 +1556,7 @@ func (f *Field) Import(qcx *Qcx, rowIDs, columnIDs []uint64, timestamps []int64, } var err1 error defer finisher(&err1) - for viewName, idx := range views { - data := allData[idx] + for viewName, data := range views { view, err := f.createViewIfNotExists(viewName) if err != nil { return errors.Wrapf(err, "creating view %s", viewName) diff --git a/time.go b/time.go index a35ba7dd4..a17d37726 100644 --- a/time.go +++ b/time.go @@ -89,15 +89,69 @@ func viewByTimeUnit(name string, t time.Time, unit rune) string { } } +// YYYYMMDDHH lengths. Note that this is a []int, not a map[byte]int, so +// the lookups can be cheaper. +var lengthsByQuantum = []int{ + 'Y': 4, + 'M': 6, + 'D': 8, + 'H': 10, +} + +// viewsByTimeInto computes the list of views for a given time. It expects +// to be given an initial buffer of the form `name_YYYYMMDDHH`, and a slice +// of []bytes. This allows us to reuse the buffer for all the sub-buffers, +// and also to reuse the slice of slices, to eliminate all those allocations. +// This might seem crazy, but even including the JSON parsing and all the +// disk activity, the straightforward viewsByTime implementation was 25% +// of runtime in an ingest test. +func viewsByTimeInto(fullBuf []byte, into [][]byte, t time.Time, q TimeQuantum) [][]byte { + l := len(fullBuf) - 10 + date := fullBuf[l : l+10] + y, m, d := t.Date() + h := t.Hour() + // Did you know that Sprintf, Printf, and other things like that all + // do allocations, and that doing allocations in a tight loop like this + // is stunningly expensive? viewsByTime was 25% of an ingest test's + // total CPU, not counting the garbage collector overhead. This is about + // 3%. No, I'm not totally sure that justifies it. + if y < 1000 { + ys := fmt.Sprintf("%04d", y) + copy(date[0:4], []byte(ys)) + } else if y >= 10000 { + // This is probably a bad answer but there isn't really a + // good answer. + ys := fmt.Sprintf("%04d", y%1000) + copy(date[0:4], []byte(ys)) + } else { + strconv.AppendInt(date[:0], int64(y), 10) + } + date[4] = '0' + byte(m/10) + date[5] = '0' + byte(m%10) + date[6] = '0' + byte(d/10) + date[7] = '0' + byte(d%10) + date[8] = '0' + byte(h/10) + date[9] = '0' + byte(h%10) + into = into[:0] + for _, unit := range q { + if int(unit) < len(lengthsByQuantum) && lengthsByQuantum[unit] != 0 { + into = append(into, fullBuf[:l+lengthsByQuantum[unit]]) + } + } + return into +} + // viewsByTime returns a list of views for a given timestamp. func viewsByTime(name string, t time.Time, q TimeQuantum) []string { // nolint: unparam + y, m, d := t.Date() + h := t.Hour() + full := fmt.Sprintf("%s_%04d%02d%02d%02d", name, y, m, d, h) + l := len(name) + 1 a := make([]string, 0, len(q)) for _, unit := range q { - view := viewByTimeUnit(name, t, unit) - if view == "" { - continue + if int(unit) < len(lengthsByQuantum) && lengthsByQuantum[unit] != 0 { + a = append(a, full[:l+lengthsByQuantum[unit]]) } - a = append(a, view) } return a } diff --git a/time_internal_test.go b/time_internal_test.go index 8356d5dbe..2627ec65d 100644 --- a/time_internal_test.go +++ b/time_internal_test.go @@ -83,6 +83,38 @@ func TestViewsByTime(t *testing.T) { }) } +func TestViewsByTimeInto(t *testing.T) { + ts := time.Date(2000, time.January, 2, 3, 4, 5, 6, time.UTC) + s := []byte("F_YYYYMMDDHH") + var timeViews [][]byte + + t.Run("YMDH", func(t *testing.T) { + a := viewsByTime("F", ts, mustParseTimeQuantum("YMDH")) + b := viewsByTimeInto(s, timeViews, ts, mustParseTimeQuantum("YMDH")) + if len(a) != len(b) { + t.Fatalf("mismatch: viewsByTime: %q, viewsByTimeInto: %q", a, b) + } + for i := range a { + if a[i] != string(b[i]) { + t.Fatalf("mismatch: viewsByTime: %q, viewsByTimeInto: %q", a, b) + } + } + }) + + t.Run("D", func(t *testing.T) { + a := viewsByTime("F", ts, mustParseTimeQuantum("D")) + b := viewsByTimeInto(s, timeViews, ts, mustParseTimeQuantum("D")) + if len(a) != len(b) { + t.Fatalf("mismatch: viewsByTime: %q, viewsByTimeInto: %q", a, b) + } + for i := range a { + if a[i] != string(b[i]) { + t.Fatalf("mismatch: viewsByTime: %q, viewsByTimeInto: %q", a, b) + } + } + }) +} + // Ensure sets of fields can be returned for a given time range. func TestViewsByTimeRange(t *testing.T) { t.Run("Y", func(t *testing.T) {