mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-08 03:47:51 +00:00
first pass at un-exporting Fragment methods
This commit is contained in:
parent
8d4cf1abf3
commit
15cb391570
8 changed files with 189 additions and 215 deletions
4
api.go
4
api.go
|
|
@ -306,7 +306,7 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, fieldName strin
|
|||
cw := csv.NewWriter(w)
|
||||
|
||||
// Iterate over each column.
|
||||
if err := f.ForEachBit(func(rowID, columnID uint64) error {
|
||||
if err := f.forEachBit(func(rowID, columnID uint64) error {
|
||||
return cw.Write([]string{
|
||||
strconv.FormatUint(rowID, 10),
|
||||
strconv.FormatUint(columnID, 10),
|
||||
|
|
@ -403,7 +403,7 @@ func (api *API) FragmentBlockData(ctx context.Context, body io.Reader) ([]byte,
|
|||
}
|
||||
|
||||
var resp = internal.BlockDataResponse{}
|
||||
resp.RowIDs, resp.ColumnIDs = f.BlockData(int(req.Block))
|
||||
resp.RowIDs, resp.ColumnIDs = f.blockData(int(req.Block))
|
||||
|
||||
// Encode response.
|
||||
buf, err := proto.Marshal(&resp)
|
||||
|
|
|
|||
|
|
@ -1017,7 +1017,7 @@ type BitsByPos []Bit
|
|||
func (p BitsByPos) Swap(i, j int) { p[i], p[j] = p[j], p[i] }
|
||||
func (p BitsByPos) Len() int { return len(p) }
|
||||
func (p BitsByPos) Less(i, j int) bool {
|
||||
p0, p1 := Pos(p[i].RowID, p[i].ColumnID), Pos(p[j].RowID, p[j].ColumnID)
|
||||
p0, p1 := pos(p[i].RowID, p[i].ColumnID), pos(p[j].RowID, p[j].ColumnID)
|
||||
if p0 == p1 {
|
||||
return p[i].Timestamp < p[j].Timestamp
|
||||
}
|
||||
|
|
|
|||
20
executor.go
20
executor.go
|
|
@ -384,7 +384,7 @@ func (e *Executor) executeSumCountSlice(ctx context.Context, index string, c *pq
|
|||
return ValCount{}, nil
|
||||
}
|
||||
|
||||
vsum, vcount, err := fragment.Sum(filter, bsig.BitDepth())
|
||||
vsum, vcount, err := fragment.sum(filter, bsig.BitDepth())
|
||||
if err != nil {
|
||||
return ValCount{}, errors.Wrap(err, "computing sum")
|
||||
}
|
||||
|
|
@ -422,7 +422,7 @@ func (e *Executor) executeMinSlice(ctx context.Context, index string, c *pql.Cal
|
|||
return ValCount{}, nil
|
||||
}
|
||||
|
||||
fmin, fcount, err := fragment.Min(filter, bsig.BitDepth())
|
||||
fmin, fcount, err := fragment.min(filter, bsig.BitDepth())
|
||||
if err != nil {
|
||||
return ValCount{}, err
|
||||
}
|
||||
|
|
@ -460,7 +460,7 @@ func (e *Executor) executeMaxSlice(ctx context.Context, index string, c *pql.Cal
|
|||
return ValCount{}, nil
|
||||
}
|
||||
|
||||
fmax, fcount, err := fragment.Max(filter, bsig.BitDepth())
|
||||
fmax, fcount, err := fragment.max(filter, bsig.BitDepth())
|
||||
if err != nil {
|
||||
return ValCount{}, err
|
||||
}
|
||||
|
|
@ -587,7 +587,7 @@ func (e *Executor) executeTopNSlice(ctx context.Context, index string, c *pql.Ca
|
|||
if tanimotoThreshold > 100 {
|
||||
return nil, errors.New("Tanimoto Threshold is from 1 to 100 only")
|
||||
}
|
||||
return f.Top(TopOptions{
|
||||
return f.top(TopOptions{
|
||||
N: int(n),
|
||||
Src: src,
|
||||
RowIDs: rowIDs,
|
||||
|
|
@ -793,7 +793,7 @@ func (e *Executor) executeBSIGroupRangeSlice(ctx context.Context, index string,
|
|||
return NewRow(), nil
|
||||
}
|
||||
|
||||
return frag.NotNull(bsig.BitDepth())
|
||||
return frag.notNull(bsig.BitDepth())
|
||||
|
||||
} else if cond.Op == pql.BETWEEN {
|
||||
|
||||
|
|
@ -831,10 +831,10 @@ func (e *Executor) executeBSIGroupRangeSlice(ctx context.Context, index string,
|
|||
// If the query is asking for the entire valid range, just return
|
||||
// the not-null bitmap for the bsiGroup.
|
||||
if predicates[0] <= bsig.Min && predicates[1] >= bsig.Max {
|
||||
return frag.NotNull(bsig.BitDepth())
|
||||
return frag.notNull(bsig.BitDepth())
|
||||
}
|
||||
|
||||
return frag.RangeBetween(bsig.BitDepth(), baseValueMin, baseValueMax)
|
||||
return frag.rangeBetween(bsig.BitDepth(), baseValueMin, baseValueMax)
|
||||
|
||||
} else {
|
||||
|
||||
|
|
@ -864,16 +864,16 @@ func (e *Executor) executeBSIGroupRangeSlice(ctx context.Context, index string,
|
|||
// LT[E] and GT[E] should return all not-null if selected range fully encompasses valid bsiGroup range.
|
||||
if (cond.Op == pql.LT && value > bsig.Max) || (cond.Op == pql.LTE && value >= bsig.Max) ||
|
||||
(cond.Op == pql.GT && value < bsig.Min) || (cond.Op == pql.GTE && value <= bsig.Min) {
|
||||
return frag.NotNull(bsig.BitDepth())
|
||||
return frag.notNull(bsig.BitDepth())
|
||||
}
|
||||
|
||||
// outOfRange for NEQ should return all not-null.
|
||||
if outOfRange && cond.Op == pql.NEQ {
|
||||
return frag.NotNull(bsig.BitDepth())
|
||||
return frag.notNull(bsig.BitDepth())
|
||||
}
|
||||
|
||||
f.Stats.Count("range:bsigroup", 1, 1.0)
|
||||
return frag.RangeOp(cond.Op, bsig.BitDepth(), baseValue)
|
||||
return frag.rangeOp(cond.Op, bsig.BitDepth(), baseValue)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
4
field.go
4
field.go
|
|
@ -905,7 +905,7 @@ func (f *Field) Import(rowIDs, columnIDs []uint64, timestamps []*time.Time) erro
|
|||
return errors.Wrap(err, "creating view")
|
||||
}
|
||||
|
||||
if err := frag.Import(data.RowIDs, data.ColumnIDs); err != nil {
|
||||
if err := frag.bulkImport(data.RowIDs, data.ColumnIDs); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
|
@ -962,7 +962,7 @@ func (f *Field) ImportValue(columnIDs []uint64, values []int64) error {
|
|||
baseValues[i] = uint64(value - bsig.Min)
|
||||
}
|
||||
|
||||
if err := frag.ImportValue(data.ColumnIDs, baseValues, bsig.BitDepth()); err != nil {
|
||||
if err := frag.importValue(data.ColumnIDs, baseValues, bsig.BitDepth()); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
|
|
|||
170
fragment.go
170
fragment.go
|
|
@ -130,27 +130,8 @@ func NewFragment(path, index, field, view string, slice uint64) *Fragment {
|
|||
}
|
||||
}
|
||||
|
||||
// Path returns the path the fragment was initialized with.
|
||||
func (f *Fragment) Path() string { return f.path }
|
||||
|
||||
// CachePath returns the path to the fragment's cache data.
|
||||
func (f *Fragment) CachePath() string { return f.path + CacheExt }
|
||||
|
||||
// Index returns the index that the fragment was initialized with.
|
||||
func (f *Fragment) Index() string { return f.index }
|
||||
|
||||
// Field returns the field the fragment was initialized with.
|
||||
func (f *Fragment) Field() string { return f.field }
|
||||
|
||||
// View returns the view the fragment was initialized with.
|
||||
func (f *Fragment) View() string { return f.view }
|
||||
|
||||
// Slice returns the slice the fragment was initialized with.
|
||||
func (f *Fragment) Slice() uint64 { return f.slice }
|
||||
|
||||
// Cache returns the fragment's cache.
|
||||
// This is not safe for concurrent use.
|
||||
func (f *Fragment) Cache() Cache { return f.cache }
|
||||
// cachePath returns the path to the fragment's cache data.
|
||||
func (f *Fragment) cachePath() string { return f.path + CacheExt }
|
||||
|
||||
// Open opens the underlying storage.
|
||||
func (f *Fragment) Open() error {
|
||||
|
|
@ -261,7 +242,7 @@ func (f *Fragment) openCache() error {
|
|||
}
|
||||
|
||||
// Read cache data from disk.
|
||||
path := f.CachePath()
|
||||
path := f.cachePath()
|
||||
buf, err := ioutil.ReadFile(path)
|
||||
if os.IsNotExist(err) {
|
||||
return nil
|
||||
|
|
@ -486,8 +467,8 @@ func (f *Fragment) bit(rowID, columnID uint64) (bool, error) {
|
|||
return f.storage.Contains(pos), nil
|
||||
}
|
||||
|
||||
// Value uses a column of bits to read a multi-bit value.
|
||||
func (f *Fragment) Value(columnID uint64, bitDepth uint) (value uint64, exists bool, err error) {
|
||||
// value uses a column of bits to read a multi-bit value.
|
||||
func (f *Fragment) value(columnID uint64, bitDepth uint) (value uint64, exists bool, err error) {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
|
||||
|
|
@ -510,8 +491,8 @@ func (f *Fragment) Value(columnID uint64, bitDepth uint) (value uint64, exists b
|
|||
return value, true, nil
|
||||
}
|
||||
|
||||
// SetValue uses a column of bits to set a multi-bit value.
|
||||
func (f *Fragment) SetValue(columnID uint64, bitDepth uint, value uint64) (changed bool, err error) {
|
||||
// setValue uses a column of bits to set a multi-bit value.
|
||||
func (f *Fragment) setValue(columnID uint64, bitDepth uint, value uint64) (changed bool, err error) {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
|
||||
|
|
@ -582,9 +563,9 @@ func (f *Fragment) importSetValue(columnID uint64, bitDepth uint, value uint64)
|
|||
return changed, nil
|
||||
}
|
||||
|
||||
// Sum returns the sum of a given bsiGroup as well as the number of columns involved.
|
||||
// sum returns the sum of a given bsiGroup as well as the number of columns involved.
|
||||
// A bitmap can be passed in to optionally filter the computed columns.
|
||||
func (f *Fragment) Sum(filter *Row, bitDepth uint) (sum, count uint64, err error) {
|
||||
func (f *Fragment) sum(filter *Row, bitDepth uint) (sum, count uint64, err error) {
|
||||
// Compute count based on the existence row.
|
||||
row := f.Row(uint64(bitDepth))
|
||||
if filter != nil {
|
||||
|
|
@ -614,9 +595,9 @@ func (f *Fragment) Sum(filter *Row, bitDepth uint) (sum, count uint64, err error
|
|||
return sum, count, nil
|
||||
}
|
||||
|
||||
// Min returns the min of a given bsiGroup as well as the number of columns involved.
|
||||
// min returns the min of a given bsiGroup as well as the number of columns involved.
|
||||
// A bitmap can be passed in to optionally filter the computed columns.
|
||||
func (f *Fragment) Min(filter *Row, bitDepth uint) (min, count uint64, err error) {
|
||||
func (f *Fragment) min(filter *Row, bitDepth uint) (min, count uint64, err error) {
|
||||
|
||||
consider := f.Row(uint64(bitDepth))
|
||||
if filter != nil {
|
||||
|
|
@ -647,9 +628,9 @@ func (f *Fragment) Min(filter *Row, bitDepth uint) (min, count uint64, err error
|
|||
return min, count, nil
|
||||
}
|
||||
|
||||
// Max returns the max of a given bsiGroup as well as the number of columns involved.
|
||||
// max returns the max of a given bsiGroup as well as the number of columns involved.
|
||||
// A bitmap can be passed in to optionally filter the computed columns.
|
||||
func (f *Fragment) Max(filter *Row, bitDepth uint) (max, count uint64, err error) {
|
||||
func (f *Fragment) max(filter *Row, bitDepth uint) (max, count uint64, err error) {
|
||||
|
||||
consider := f.Row(uint64(bitDepth))
|
||||
if filter != nil {
|
||||
|
|
@ -678,8 +659,8 @@ func (f *Fragment) Max(filter *Row, bitDepth uint) (max, count uint64, err error
|
|||
return max, count, nil
|
||||
}
|
||||
|
||||
// RangeOp returns bitmaps with a bsiGroup value encoding matching the predicate.
|
||||
func (f *Fragment) RangeOp(op pql.Token, bitDepth uint, predicate uint64) (*Row, error) {
|
||||
// rangeOp returns bitmaps with a bsiGroup value encoding matching the predicate.
|
||||
func (f *Fragment) rangeOp(op pql.Token, bitDepth uint, predicate uint64) (*Row, error) {
|
||||
switch op {
|
||||
case pql.EQ:
|
||||
return f.rangeEQ(bitDepth, predicate)
|
||||
|
|
@ -812,13 +793,13 @@ func (f *Fragment) rangeGT(bitDepth uint, predicate uint64, allowEquality bool)
|
|||
return b, nil
|
||||
}
|
||||
|
||||
// NotNull returns the not-null row (stored at bitDepth).
|
||||
func (f *Fragment) NotNull(bitDepth uint) (*Row, error) {
|
||||
// notNull returns the not-null row (stored at bitDepth).
|
||||
func (f *Fragment) notNull(bitDepth uint) (*Row, error) {
|
||||
return f.Row(uint64(bitDepth)), nil
|
||||
}
|
||||
|
||||
// RangeBetween returns bitmaps with a bsiGroup value encoding matching any value between predicateMin and predicateMax.
|
||||
func (f *Fragment) RangeBetween(bitDepth uint, predicateMin, predicateMax uint64) (*Row, error) {
|
||||
// rangeBetween returns bitmaps with a bsiGroup value encoding matching any value between predicateMin and predicateMax.
|
||||
func (f *Fragment) rangeBetween(bitDepth uint, predicateMin, predicateMax uint64) (*Row, error) {
|
||||
b := f.Row(uint64(bitDepth))
|
||||
keep1 := NewRow() // GTE
|
||||
keep2 := NewRow() // LTE
|
||||
|
|
@ -864,12 +845,12 @@ func (f *Fragment) pos(rowID, columnID uint64) (uint64, error) {
|
|||
if columnID < minColumnID || columnID >= minColumnID+SliceWidth {
|
||||
return 0, errors.New("column out of bounds")
|
||||
}
|
||||
return Pos(rowID, columnID), nil
|
||||
return pos(rowID, columnID), nil
|
||||
}
|
||||
|
||||
// ForEachBit executes fn for every bit set in the fragment.
|
||||
// forEachBit executes fn for every bit set in the fragment.
|
||||
// Errors returned from fn are passed through.
|
||||
func (f *Fragment) ForEachBit(fn func(rowID, columnID uint64) error) error {
|
||||
func (f *Fragment) forEachBit(fn func(rowID, columnID uint64) error) error {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
|
||||
|
|
@ -886,10 +867,10 @@ func (f *Fragment) ForEachBit(fn func(rowID, columnID uint64) error) error {
|
|||
return err
|
||||
}
|
||||
|
||||
// Top returns the top rows from the fragment.
|
||||
// top returns the top rows from the fragment.
|
||||
// If opt.Src is specified then only rows which intersect src are returned.
|
||||
// If opt.FilterValues exist then the row attribute specified by field is matched.
|
||||
func (f *Fragment) Top(opt TopOptions) ([]Pair, error) {
|
||||
func (f *Fragment) top(opt TopOptions) ([]Pair, error) {
|
||||
// Retrieve pairs. If no row ids specified then return from cache.
|
||||
pairs := f.topBitmapPairs(opt.RowIDs)
|
||||
|
||||
|
|
@ -1089,13 +1070,6 @@ func (f *Fragment) Checksum() []byte {
|
|||
return h.Sum(nil)
|
||||
}
|
||||
|
||||
// BlockN returns the number of blocks in the fragment.
|
||||
func (f *Fragment) BlockN() int {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
return int(f.storage.Max() / (HashBlockSize * SliceWidth))
|
||||
}
|
||||
|
||||
// InvalidateChecksums clears all cached block checksums.
|
||||
func (f *Fragment) InvalidateChecksums() {
|
||||
f.mu.Lock()
|
||||
|
|
@ -1184,8 +1158,8 @@ func (f *Fragment) readContiguousChecksums(a *[]FragmentBlock, blockID int) (n i
|
|||
}
|
||||
}
|
||||
|
||||
// BlockData returns bits in a block as row & column ID pairs.
|
||||
func (f *Fragment) BlockData(id int) (rowIDs, columnIDs []uint64) {
|
||||
// blockData returns bits in a block as row & column ID pairs.
|
||||
func (f *Fragment) blockData(id int) (rowIDs, columnIDs []uint64) {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
|
||||
|
|
@ -1196,17 +1170,17 @@ func (f *Fragment) BlockData(id int) (rowIDs, columnIDs []uint64) {
|
|||
return
|
||||
}
|
||||
|
||||
// MergeBlock compares the block's bits and computes a diff with another set of block bits.
|
||||
// mergeBlock compares the block's bits and computes a diff with another set of block bits.
|
||||
// The state of a bit is determined by consensus from all blocks being considered.
|
||||
//
|
||||
// For example, if 3 blocks are compared and two have a set bit and one has a
|
||||
// cleared bit then the bit is considered cleared. The function returns the
|
||||
// diff per incoming block so that all can be in sync.
|
||||
func (f *Fragment) MergeBlock(id int, data []PairSet) (sets, clears []PairSet, err error) {
|
||||
func (f *Fragment) mergeBlock(id int, data []pairSet) (sets, clears []pairSet, err error) {
|
||||
// Ensure that all pair sets are of equal length.
|
||||
for i := range data {
|
||||
if len(data[i].RowIDs) != len(data[i].ColumnIDs) {
|
||||
return nil, nil, fmt.Errorf("pair set mismatch(idx=%d): %d != %d", i, len(data[i].RowIDs), len(data[i].ColumnIDs))
|
||||
if len(data[i].rowIDs) != len(data[i].columnIDs) {
|
||||
return nil, nil, fmt.Errorf("pair set mismatch(idx=%d): %d != %d", i, len(data[i].rowIDs), len(data[i].columnIDs))
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1214,8 +1188,8 @@ func (f *Fragment) MergeBlock(id int, data []PairSet) (sets, clears []PairSet, e
|
|||
defer f.mu.Unlock()
|
||||
|
||||
// Track sets and clears for all blocks (including local).
|
||||
sets = make([]PairSet, len(data)+1)
|
||||
clears = make([]PairSet, len(data)+1)
|
||||
sets = make([]pairSet, len(data)+1)
|
||||
clears = make([]pairSet, len(data)+1)
|
||||
|
||||
// Limit upper row/column pair.
|
||||
maxRowID := uint64(id+1) * HashBlockSize
|
||||
|
|
@ -1231,7 +1205,7 @@ func (f *Fragment) MergeBlock(id int, data []PairSet) (sets, clears []PairSet, e
|
|||
|
||||
// Append buffered iterators for each incoming block.
|
||||
for i := range data {
|
||||
var itr Iterator = NewSliceIterator(data[i].RowIDs, data[i].ColumnIDs)
|
||||
var itr Iterator = NewSliceIterator(data[i].rowIDs, data[i].columnIDs)
|
||||
itr = NewLimitIterator(itr, maxRowID, maxColumnID)
|
||||
itrs = append(itrs, NewBufIterator(itr))
|
||||
}
|
||||
|
|
@ -1296,25 +1270,25 @@ func (f *Fragment) MergeBlock(id int, data []PairSet) (sets, clears []PairSet, e
|
|||
|
||||
// Append to either the set or clear diff.
|
||||
if newValue {
|
||||
sets[i].RowIDs = append(sets[i].RowIDs, min.rowID)
|
||||
sets[i].ColumnIDs = append(sets[i].ColumnIDs, min.columnID)
|
||||
sets[i].rowIDs = append(sets[i].rowIDs, min.rowID)
|
||||
sets[i].columnIDs = append(sets[i].columnIDs, min.columnID)
|
||||
} else {
|
||||
clears[i].RowIDs = append(sets[i].RowIDs, min.rowID)
|
||||
clears[i].ColumnIDs = append(sets[i].ColumnIDs, min.columnID)
|
||||
clears[i].rowIDs = append(sets[i].rowIDs, min.rowID)
|
||||
clears[i].columnIDs = append(sets[i].columnIDs, min.columnID)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Set local bits.
|
||||
for i := range sets[0].ColumnIDs {
|
||||
if _, err := f.setBit(sets[0].RowIDs[i], (f.Slice()*SliceWidth)+sets[0].ColumnIDs[i]); err != nil {
|
||||
for i := range sets[0].columnIDs {
|
||||
if _, err := f.setBit(sets[0].rowIDs[i], (f.slice*SliceWidth)+sets[0].columnIDs[i]); err != nil {
|
||||
return nil, nil, errors.Wrap(err, "setting")
|
||||
}
|
||||
}
|
||||
|
||||
// Clear local bits.
|
||||
for i := range clears[0].ColumnIDs {
|
||||
if _, err := f.clearBit(clears[0].RowIDs[i], (f.Slice()*SliceWidth)+clears[0].ColumnIDs[i]); err != nil {
|
||||
for i := range clears[0].columnIDs {
|
||||
if _, err := f.clearBit(clears[0].rowIDs[i], (f.slice*SliceWidth)+clears[0].columnIDs[i]); err != nil {
|
||||
return nil, nil, errors.Wrap(err, "clearing")
|
||||
}
|
||||
}
|
||||
|
|
@ -1322,9 +1296,9 @@ func (f *Fragment) MergeBlock(id int, data []PairSet) (sets, clears []PairSet, e
|
|||
return sets[1:], clears[1:], nil
|
||||
}
|
||||
|
||||
// Import bulk imports a set of bits and then snapshots the storage.
|
||||
// bulkImport bulk imports a set of bits and then snapshots the storage.
|
||||
// This does not affect the fragment's cache.
|
||||
func (f *Fragment) Import(rowIDs, columnIDs []uint64) error {
|
||||
func (f *Fragment) bulkImport(rowIDs, columnIDs []uint64) error {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
// Verify that there are an equal number of row ids and column ids.
|
||||
|
|
@ -1392,8 +1366,8 @@ func (f *Fragment) Import(rowIDs, columnIDs []uint64) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// ImportValue bulk imports a set of range-encoded values.
|
||||
func (f *Fragment) ImportValue(columnIDs, values []uint64, bitDepth uint) error {
|
||||
// importValue bulk imports a set of range-encoded values.
|
||||
func (f *Fragment) importValue(columnIDs, values []uint64, bitDepth uint) error {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
// Verify that there are an equal number of column ids and values.
|
||||
|
|
@ -1529,7 +1503,7 @@ func (f *Fragment) flushCache() error {
|
|||
}
|
||||
|
||||
// Write to disk.
|
||||
if err := ioutil.WriteFile(f.CachePath(), buf, 0666); err != nil {
|
||||
if err := ioutil.WriteFile(f.cachePath(), buf, 0666); err != nil {
|
||||
return errors.Wrap(err, "writing")
|
||||
}
|
||||
|
||||
|
|
@ -1603,7 +1577,7 @@ func (f *Fragment) writeCacheToArchive(tw *tar.Writer) error {
|
|||
defer f.mu.Unlock()
|
||||
|
||||
// Read cache into buffer.
|
||||
buf, err := ioutil.ReadFile(f.CachePath())
|
||||
buf, err := ioutil.ReadFile(f.cachePath())
|
||||
if os.IsNotExist(err) {
|
||||
return nil
|
||||
} else if err != nil {
|
||||
|
|
@ -1697,7 +1671,7 @@ func (f *Fragment) readCacheFromArchive(r io.Reader) error {
|
|||
buf, err := ioutil.ReadAll(r)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "reading")
|
||||
} else if err := ioutil.WriteFile(f.CachePath(), buf, 0666); err != nil {
|
||||
} else if err := ioutil.WriteFile(f.cachePath(), buf, 0666); err != nil {
|
||||
return errors.Wrap(err, "writing")
|
||||
}
|
||||
|
||||
|
|
@ -1762,11 +1736,11 @@ func (s *FragmentSyncer) isClosing() bool {
|
|||
}
|
||||
}
|
||||
|
||||
// SyncFragment compares checksums for the local and remote fragments and
|
||||
// syncFragment compares checksums for the local and remote fragments and
|
||||
// then merges any blocks which have differences.
|
||||
func (s *FragmentSyncer) SyncFragment() error {
|
||||
func (s *FragmentSyncer) syncFragment() error {
|
||||
// Determine replica set.
|
||||
nodes := s.Cluster.SliceNodes(s.Fragment.Index(), s.Fragment.Slice())
|
||||
nodes := s.Cluster.SliceNodes(s.Fragment.index, s.Fragment.slice)
|
||||
if len(nodes) == 1 {
|
||||
return nil
|
||||
}
|
||||
|
|
@ -1783,7 +1757,7 @@ func (s *FragmentSyncer) SyncFragment() error {
|
|||
|
||||
// Retrieve remote blocks.
|
||||
client := NewInternalHTTPClientFromURI(&node.URI, s.RemoteClient)
|
||||
blocks, err := client.FragmentBlocks(context.Background(), s.Fragment.Index(), s.Fragment.Field(), s.Fragment.Slice())
|
||||
blocks, err := client.FragmentBlocks(context.Background(), s.Fragment.index, s.Fragment.field, s.Fragment.slice)
|
||||
if err != nil && err != ErrFragmentNotFound {
|
||||
return errors.Wrap(err, "getting blocks")
|
||||
}
|
||||
|
|
@ -1846,9 +1820,9 @@ func (s *FragmentSyncer) syncBlock(id int) error {
|
|||
f := s.Fragment
|
||||
|
||||
// Read pairs from each remote block.
|
||||
var pairSets []PairSet
|
||||
var pairSets []pairSet
|
||||
var clients []InternalClient
|
||||
for _, node := range s.Cluster.SliceNodes(f.Index(), f.Slice()) {
|
||||
for _, node := range s.Cluster.SliceNodes(f.index, f.slice) {
|
||||
if s.Node.ID == node.ID {
|
||||
continue
|
||||
}
|
||||
|
|
@ -1862,14 +1836,14 @@ func (s *FragmentSyncer) syncBlock(id int) error {
|
|||
clients = append(clients, client)
|
||||
|
||||
// Only sync the standard block.
|
||||
rowIDs, columnIDs, err := client.BlockData(context.Background(), f.Index(), f.Field(), f.Slice(), id)
|
||||
rowIDs, columnIDs, err := client.BlockData(context.Background(), f.index, f.field, f.slice, id)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "getting block")
|
||||
}
|
||||
|
||||
pairSets = append(pairSets, PairSet{
|
||||
ColumnIDs: columnIDs,
|
||||
RowIDs: rowIDs,
|
||||
pairSets = append(pairSets, pairSet{
|
||||
columnIDs: columnIDs,
|
||||
rowIDs: rowIDs,
|
||||
})
|
||||
}
|
||||
|
||||
|
|
@ -1879,7 +1853,7 @@ func (s *FragmentSyncer) syncBlock(id int) error {
|
|||
}
|
||||
|
||||
// Merge blocks together.
|
||||
sets, clears, err := f.MergeBlock(id, pairSets)
|
||||
sets, clears, err := f.mergeBlock(id, pairSets)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "merging")
|
||||
}
|
||||
|
|
@ -1890,12 +1864,12 @@ func (s *FragmentSyncer) syncBlock(id int) error {
|
|||
count := 0
|
||||
|
||||
// Ignore if there are no differences.
|
||||
if len(set.ColumnIDs) == 0 && len(clear.ColumnIDs) == 0 {
|
||||
if len(set.columnIDs) == 0 && len(clear.columnIDs) == 0 {
|
||||
continue
|
||||
}
|
||||
|
||||
// Generate query with sets & clears, and group the requests to not exceed MaxWritesPerRequest.
|
||||
total := len(set.ColumnIDs) + len(clear.ColumnIDs)
|
||||
total := len(set.columnIDs) + len(clear.columnIDs)
|
||||
maxWrites := s.Cluster.MaxWritesPerRequest
|
||||
if maxWrites <= 0 {
|
||||
maxWrites = 5000
|
||||
|
|
@ -1903,12 +1877,12 @@ func (s *FragmentSyncer) syncBlock(id int) error {
|
|||
buffers := make([]bytes.Buffer, int(math.Ceil(float64(total)/float64(maxWrites))))
|
||||
|
||||
// Only sync the standard block.
|
||||
for j := 0; j < len(set.ColumnIDs); j++ {
|
||||
fmt.Fprintf(&(buffers[count/maxWrites]), "SetBit(field=%q, row=%d, col=%d)\n", f.Field(), set.RowIDs[j], (f.Slice()*SliceWidth)+set.ColumnIDs[j])
|
||||
for j := 0; j < len(set.columnIDs); j++ {
|
||||
fmt.Fprintf(&(buffers[count/maxWrites]), "SetBit(field=%q, row=%d, col=%d)\n", f.field, set.rowIDs[j], (f.slice*SliceWidth)+set.columnIDs[j])
|
||||
count++
|
||||
}
|
||||
for j := 0; j < len(clear.ColumnIDs); j++ {
|
||||
fmt.Fprintf(&(buffers[count/maxWrites]), "ClearBit(field=%q, row=%d, col=%d)\n", f.Field(), clear.RowIDs[j], (f.Slice()*SliceWidth)+clear.ColumnIDs[j])
|
||||
for j := 0; j < len(clear.columnIDs); j++ {
|
||||
fmt.Fprintf(&(buffers[count/maxWrites]), "ClearBit(field=%q, row=%d, col=%d)\n", f.field, clear.rowIDs[j], (f.slice*SliceWidth)+clear.columnIDs[j])
|
||||
count++
|
||||
}
|
||||
|
||||
|
|
@ -1924,7 +1898,7 @@ func (s *FragmentSyncer) syncBlock(id int) error {
|
|||
Query: buffers[k].String(),
|
||||
Remote: true,
|
||||
}
|
||||
_, err := clients[i].Query(context.Background(), f.Index(), queryRequest)
|
||||
_, err := clients[i].Query(context.Background(), f.index, queryRequest)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "executing")
|
||||
}
|
||||
|
|
@ -1942,10 +1916,10 @@ func madvise(b []byte, advice int) (err error) {
|
|||
return
|
||||
}
|
||||
|
||||
// PairSet is a list of equal length row and column id lists.
|
||||
type PairSet struct {
|
||||
RowIDs []uint64
|
||||
ColumnIDs []uint64
|
||||
// pairSet is a list of equal length row and column id lists.
|
||||
type pairSet struct {
|
||||
rowIDs []uint64
|
||||
columnIDs []uint64
|
||||
}
|
||||
|
||||
// byteSlicesEqual returns true if all slices are equal.
|
||||
|
|
@ -1962,7 +1936,7 @@ func byteSlicesEqual(a [][]byte) bool {
|
|||
return true
|
||||
}
|
||||
|
||||
// Pos returns the row position of a row/column pair.
|
||||
func Pos(rowID, columnID uint64) uint64 {
|
||||
// pos returns the row position of a row/column pair.
|
||||
func pos(rowID, columnID uint64) uint64 {
|
||||
return (rowID * SliceWidth) + (columnID % SliceWidth)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -99,14 +99,14 @@ func TestFragment_SetValue(t *testing.T) {
|
|||
defer f.Close()
|
||||
|
||||
// Set value.
|
||||
if changed, err := f.SetValue(100, 16, 3829); err != nil {
|
||||
if changed, err := f.setValue(100, 16, 3829); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !changed {
|
||||
t.Fatal("expected change")
|
||||
}
|
||||
|
||||
// Read value.
|
||||
if value, exists, err := f.Value(100, 16); err != nil {
|
||||
if value, exists, err := f.value(100, 16); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if value != 3829 {
|
||||
t.Fatalf("unexpected value: %d", value)
|
||||
|
|
@ -115,7 +115,7 @@ func TestFragment_SetValue(t *testing.T) {
|
|||
}
|
||||
|
||||
// Setting value should return no change.
|
||||
if changed, err := f.SetValue(100, 16, 3829); err != nil {
|
||||
if changed, err := f.setValue(100, 16, 3829); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if changed {
|
||||
t.Fatal("expected no change")
|
||||
|
|
@ -127,21 +127,21 @@ func TestFragment_SetValue(t *testing.T) {
|
|||
defer f.Close()
|
||||
|
||||
// Set value.
|
||||
if changed, err := f.SetValue(100, 16, 3829); err != nil {
|
||||
if changed, err := f.setValue(100, 16, 3829); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !changed {
|
||||
t.Fatal("expected change")
|
||||
}
|
||||
|
||||
// Overwriting value should overwrite all bits.
|
||||
if changed, err := f.SetValue(100, 16, 2028); err != nil {
|
||||
if changed, err := f.setValue(100, 16, 2028); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !changed {
|
||||
t.Fatal("expected change")
|
||||
}
|
||||
|
||||
// Read value.
|
||||
if value, exists, err := f.Value(100, 16); err != nil {
|
||||
if value, exists, err := f.value(100, 16); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if value != 2028 {
|
||||
t.Fatalf("unexpected value: %d", value)
|
||||
|
|
@ -155,14 +155,14 @@ func TestFragment_SetValue(t *testing.T) {
|
|||
defer f.Close()
|
||||
|
||||
// Set value.
|
||||
if changed, err := f.SetValue(100, 10, 20); err != nil {
|
||||
if changed, err := f.setValue(100, 10, 20); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !changed {
|
||||
t.Fatal("expected change")
|
||||
}
|
||||
|
||||
// Non-existent value.
|
||||
if value, exists, err := f.Value(100, 11); err != nil {
|
||||
if value, exists, err := f.value(100, 11); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if value != 0 {
|
||||
t.Fatalf("unexpected value: %d", value)
|
||||
|
|
@ -191,14 +191,14 @@ func TestFragment_SetValue(t *testing.T) {
|
|||
|
||||
m[columnID] = int64(value)
|
||||
|
||||
if _, err := f.SetValue(columnID, bitDepth, value); err != nil {
|
||||
if _, err := f.setValue(columnID, bitDepth, value); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
// Ensure values are set.
|
||||
for columnID, value := range m {
|
||||
v, exists, err := f.Value(columnID, bitDepth)
|
||||
v, exists, err := f.value(columnID, bitDepth)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
} else if value != int64(v) {
|
||||
|
|
@ -223,18 +223,18 @@ func TestFragment_Sum(t *testing.T) {
|
|||
defer f.Close()
|
||||
|
||||
// Set values.
|
||||
if _, err := f.SetValue(1000, bitDepth, 382); err != nil {
|
||||
if _, err := f.setValue(1000, bitDepth, 382); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(2000, bitDepth, 300); err != nil {
|
||||
} else if _, err := f.setValue(2000, bitDepth, 300); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(3000, bitDepth, 2818); err != nil {
|
||||
} else if _, err := f.setValue(3000, bitDepth, 2818); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(4000, bitDepth, 300); err != nil {
|
||||
} else if _, err := f.setValue(4000, bitDepth, 300); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
t.Run("NoFilter", func(t *testing.T) {
|
||||
if sum, n, err := f.Sum(nil, bitDepth); err != nil {
|
||||
if sum, n, err := f.sum(nil, bitDepth); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if n != 4 {
|
||||
t.Fatalf("unexpected count: %d", n)
|
||||
|
|
@ -244,7 +244,7 @@ func TestFragment_Sum(t *testing.T) {
|
|||
})
|
||||
|
||||
t.Run("WithFilter", func(t *testing.T) {
|
||||
if sum, n, err := f.Sum(NewRow(2000, 4000, 5000), bitDepth); err != nil {
|
||||
if sum, n, err := f.sum(NewRow(2000, 4000, 5000), bitDepth); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if n != 2 {
|
||||
t.Fatalf("unexpected count: %d", n)
|
||||
|
|
@ -262,19 +262,19 @@ func TestFragment_MinMax(t *testing.T) {
|
|||
defer f.Close()
|
||||
|
||||
// Set values.
|
||||
if _, err := f.SetValue(1000, bitDepth, 382); err != nil {
|
||||
if _, err := f.setValue(1000, bitDepth, 382); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(2000, bitDepth, 300); err != nil {
|
||||
} else if _, err := f.setValue(2000, bitDepth, 300); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(3000, bitDepth, 2818); err != nil {
|
||||
} else if _, err := f.setValue(3000, bitDepth, 2818); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(4000, bitDepth, 300); err != nil {
|
||||
} else if _, err := f.setValue(4000, bitDepth, 300); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(5000, bitDepth, 2818); err != nil {
|
||||
} else if _, err := f.setValue(5000, bitDepth, 2818); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(6000, bitDepth, 2817); err != nil {
|
||||
} else if _, err := f.setValue(6000, bitDepth, 2817); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(7000, bitDepth, 0); err != nil {
|
||||
} else if _, err := f.setValue(7000, bitDepth, 0); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
|
|
@ -292,7 +292,7 @@ func TestFragment_MinMax(t *testing.T) {
|
|||
{filter: NewRow(7000), exp: 0, cnt: 1},
|
||||
}
|
||||
for i, test := range tests {
|
||||
if min, cnt, err := f.Min(test.filter, bitDepth); err != nil {
|
||||
if min, cnt, err := f.min(test.filter, bitDepth); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if min != test.exp {
|
||||
t.Errorf("test %d expected min: %v, but got: %v", i, test.exp, min)
|
||||
|
|
@ -316,7 +316,7 @@ func TestFragment_MinMax(t *testing.T) {
|
|||
{filter: NewRow(7000), exp: 0, cnt: 1},
|
||||
}
|
||||
for i, test := range tests {
|
||||
if max, cnt, err := f.Max(test.filter, bitDepth); err != nil {
|
||||
if max, cnt, err := f.max(test.filter, bitDepth); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if max != test.exp {
|
||||
t.Errorf("test %d expected max: %v, but got: %v", i, test.exp, max)
|
||||
|
|
@ -336,18 +336,18 @@ func TestFragment_Range(t *testing.T) {
|
|||
defer f.Close()
|
||||
|
||||
// Set values.
|
||||
if _, err := f.SetValue(1000, bitDepth, 382); err != nil {
|
||||
if _, err := f.setValue(1000, bitDepth, 382); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(2000, bitDepth, 300); err != nil {
|
||||
} else if _, err := f.setValue(2000, bitDepth, 300); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(3000, bitDepth, 2818); err != nil {
|
||||
} else if _, err := f.setValue(3000, bitDepth, 2818); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(4000, bitDepth, 300); err != nil {
|
||||
} else if _, err := f.setValue(4000, bitDepth, 300); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// Query for equality.
|
||||
if b, err := f.RangeOp(pql.EQ, bitDepth, 300); err != nil {
|
||||
if b, err := f.rangeOp(pql.EQ, bitDepth, 300); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(b.Columns(), []uint64{2000, 4000}) {
|
||||
t.Fatalf("unexpected columns: %+v", b.Columns())
|
||||
|
|
@ -359,18 +359,18 @@ func TestFragment_Range(t *testing.T) {
|
|||
defer f.Close()
|
||||
|
||||
// Set values.
|
||||
if _, err := f.SetValue(1000, bitDepth, 382); err != nil {
|
||||
if _, err := f.setValue(1000, bitDepth, 382); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(2000, bitDepth, 300); err != nil {
|
||||
} else if _, err := f.setValue(2000, bitDepth, 300); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(3000, bitDepth, 2818); err != nil {
|
||||
} else if _, err := f.setValue(3000, bitDepth, 2818); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(4000, bitDepth, 300); err != nil {
|
||||
} else if _, err := f.setValue(4000, bitDepth, 300); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// Query for inequality.
|
||||
if b, err := f.RangeOp(pql.NEQ, bitDepth, 300); err != nil {
|
||||
if b, err := f.rangeOp(pql.NEQ, bitDepth, 300); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(b.Columns(), []uint64{1000, 3000}) {
|
||||
t.Fatalf("unexpected columns: %+v", b.Columns())
|
||||
|
|
@ -382,43 +382,43 @@ func TestFragment_Range(t *testing.T) {
|
|||
defer f.Close()
|
||||
|
||||
// Set values.
|
||||
if _, err := f.SetValue(1000, bitDepth, 382); err != nil {
|
||||
if _, err := f.setValue(1000, bitDepth, 382); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(2000, bitDepth, 300); err != nil {
|
||||
} else if _, err := f.setValue(2000, bitDepth, 300); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(3000, bitDepth, 2817); err != nil {
|
||||
} else if _, err := f.setValue(3000, bitDepth, 2817); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(4000, bitDepth, 301); err != nil {
|
||||
} else if _, err := f.setValue(4000, bitDepth, 301); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(5000, bitDepth, 1); err != nil {
|
||||
} else if _, err := f.setValue(5000, bitDepth, 1); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(6000, bitDepth, 0); err != nil {
|
||||
} else if _, err := f.setValue(6000, bitDepth, 0); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// Query for values less than (ending with set column).
|
||||
if b, err := f.RangeOp(pql.LT, bitDepth, 301); err != nil {
|
||||
if b, err := f.rangeOp(pql.LT, bitDepth, 301); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(b.Columns(), []uint64{2000, 5000, 6000}) {
|
||||
t.Fatalf("unexpected columns: %+v", b.Columns())
|
||||
}
|
||||
|
||||
// Query for values less than (ending with unset column).
|
||||
if b, err := f.RangeOp(pql.LT, bitDepth, 300); err != nil {
|
||||
if b, err := f.rangeOp(pql.LT, bitDepth, 300); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(b.Columns(), []uint64{5000, 6000}) {
|
||||
t.Fatalf("unexpected columns: %+v", b.Columns())
|
||||
}
|
||||
|
||||
// Query for values less than or equal to (ending with set column).
|
||||
if b, err := f.RangeOp(pql.LTE, bitDepth, 301); err != nil {
|
||||
if b, err := f.rangeOp(pql.LTE, bitDepth, 301); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(b.Columns(), []uint64{2000, 4000, 5000, 6000}) {
|
||||
t.Fatalf("unexpected columns: %+v", b.Columns())
|
||||
}
|
||||
|
||||
// Query for values less than or equal to (ending with unset column).
|
||||
if b, err := f.RangeOp(pql.LTE, bitDepth, 300); err != nil {
|
||||
if b, err := f.rangeOp(pql.LTE, bitDepth, 300); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(b.Columns(), []uint64{2000, 5000, 6000}) {
|
||||
t.Fatalf("unexpected columns: %+v", b.Columns())
|
||||
|
|
@ -430,43 +430,43 @@ func TestFragment_Range(t *testing.T) {
|
|||
defer f.Close()
|
||||
|
||||
// Set values.
|
||||
if _, err := f.SetValue(1000, bitDepth, 382); err != nil {
|
||||
if _, err := f.setValue(1000, bitDepth, 382); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(2000, bitDepth, 300); err != nil {
|
||||
} else if _, err := f.setValue(2000, bitDepth, 300); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(3000, bitDepth, 2817); err != nil {
|
||||
} else if _, err := f.setValue(3000, bitDepth, 2817); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(4000, bitDepth, 301); err != nil {
|
||||
} else if _, err := f.setValue(4000, bitDepth, 301); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(5000, bitDepth, 1); err != nil {
|
||||
} else if _, err := f.setValue(5000, bitDepth, 1); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(6000, bitDepth, 0); err != nil {
|
||||
} else if _, err := f.setValue(6000, bitDepth, 0); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// Query for values greater than (ending with unset bit).
|
||||
if b, err := f.RangeOp(pql.GT, bitDepth, 300); err != nil {
|
||||
if b, err := f.rangeOp(pql.GT, bitDepth, 300); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(b.Columns(), []uint64{1000, 3000, 4000}) {
|
||||
t.Fatalf("unexpected columns: %+v", b.Columns())
|
||||
}
|
||||
|
||||
// Query for values greater than (ending with set bit).
|
||||
if b, err := f.RangeOp(pql.GT, bitDepth, 301); err != nil {
|
||||
if b, err := f.rangeOp(pql.GT, bitDepth, 301); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(b.Columns(), []uint64{1000, 3000}) {
|
||||
t.Fatalf("unexpected columns: %+v", b.Columns())
|
||||
}
|
||||
|
||||
// Query for values greater than or equal to (ending with unset bit).
|
||||
if b, err := f.RangeOp(pql.GTE, bitDepth, 300); err != nil {
|
||||
if b, err := f.rangeOp(pql.GTE, bitDepth, 300); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(b.Columns(), []uint64{1000, 2000, 3000, 4000}) {
|
||||
t.Fatalf("unexpected columns: %+v", b.Columns())
|
||||
}
|
||||
|
||||
// Query for values greater than or equal to (ending with set bit).
|
||||
if b, err := f.RangeOp(pql.GTE, bitDepth, 301); err != nil {
|
||||
if b, err := f.rangeOp(pql.GTE, bitDepth, 301); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(b.Columns(), []uint64{1000, 3000, 4000}) {
|
||||
t.Fatalf("unexpected columns: %+v", b.Columns())
|
||||
|
|
@ -478,43 +478,43 @@ func TestFragment_Range(t *testing.T) {
|
|||
defer f.Close()
|
||||
|
||||
// Set values.
|
||||
if _, err := f.SetValue(1000, bitDepth, 382); err != nil {
|
||||
if _, err := f.setValue(1000, bitDepth, 382); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(2000, bitDepth, 300); err != nil {
|
||||
} else if _, err := f.setValue(2000, bitDepth, 300); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(3000, bitDepth, 2817); err != nil {
|
||||
} else if _, err := f.setValue(3000, bitDepth, 2817); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(4000, bitDepth, 301); err != nil {
|
||||
} else if _, err := f.setValue(4000, bitDepth, 301); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(5000, bitDepth, 1); err != nil {
|
||||
} else if _, err := f.setValue(5000, bitDepth, 1); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if _, err := f.SetValue(6000, bitDepth, 0); err != nil {
|
||||
} else if _, err := f.setValue(6000, bitDepth, 0); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// Query for values greater than (ending with unset column).
|
||||
if b, err := f.RangeBetween(bitDepth, 300, 2817); err != nil {
|
||||
if b, err := f.rangeBetween(bitDepth, 300, 2817); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(b.Columns(), []uint64{1000, 2000, 3000, 4000}) {
|
||||
t.Fatalf("unexpected columns: %+v", b.Columns())
|
||||
}
|
||||
|
||||
// Query for values greater than (ending with set column).
|
||||
if b, err := f.RangeBetween(bitDepth, 301, 2817); err != nil {
|
||||
if b, err := f.rangeBetween(bitDepth, 301, 2817); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(b.Columns(), []uint64{1000, 3000, 4000}) {
|
||||
t.Fatalf("unexpected columns: %+v", b.Columns())
|
||||
}
|
||||
|
||||
// Query for values greater than or equal to (ending with unset column).
|
||||
if b, err := f.RangeBetween(bitDepth, 301, 2816); err != nil {
|
||||
if b, err := f.rangeBetween(bitDepth, 301, 2816); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(b.Columns(), []uint64{1000, 4000}) {
|
||||
t.Fatalf("unexpected columns: %+v", b.Columns())
|
||||
}
|
||||
|
||||
// Query for values greater than or equal to (ending with set column).
|
||||
if b, err := f.RangeBetween(bitDepth, 300, 2816); err != nil {
|
||||
if b, err := f.rangeBetween(bitDepth, 300, 2816); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(b.Columns(), []uint64{1000, 2000, 4000}) {
|
||||
t.Fatalf("unexpected columns: %+v", b.Columns())
|
||||
|
|
@ -567,7 +567,7 @@ func TestFragment_ForEachBit(t *testing.T) {
|
|||
|
||||
// Iterate over bits.
|
||||
var result [][2]uint64
|
||||
if err := f.ForEachBit(func(rowID, columnID uint64) error {
|
||||
if err := f.forEachBit(func(rowID, columnID uint64) error {
|
||||
result = append(result, [2]uint64{rowID, columnID})
|
||||
return nil
|
||||
}); err != nil {
|
||||
|
|
@ -591,7 +591,7 @@ func TestFragment_Top(t *testing.T) {
|
|||
f.RecalculateCache()
|
||||
|
||||
// Retrieve top rows.
|
||||
if pairs, err := f.Top(TopOptions{N: 2}); err != nil {
|
||||
if pairs, err := f.top(TopOptions{N: 2}); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if len(pairs) != 2 {
|
||||
t.Fatalf("unexpected count: %d", len(pairs))
|
||||
|
|
@ -617,7 +617,7 @@ func TestFragment_Top_Filter(t *testing.T) {
|
|||
f.RowAttrStore.SetAttrs(102, map[string]interface{}{"x": int64(20)})
|
||||
|
||||
// Retrieve top rows.
|
||||
if pairs, err := f.Top(TopOptions{
|
||||
if pairs, err := f.top(TopOptions{
|
||||
N: 2,
|
||||
FilterName: "x",
|
||||
FilterValues: []interface{}{int64(10), int64(15), int64(20)},
|
||||
|
|
@ -648,7 +648,7 @@ func TestFragment_TopN_Intersect(t *testing.T) {
|
|||
f.RecalculateCache()
|
||||
|
||||
// Retrieve top rows.
|
||||
if pairs, err := f.Top(TopOptions{N: 3, Src: src}); err != nil {
|
||||
if pairs, err := f.top(TopOptions{N: 3, Src: src}); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(pairs, []Pair{
|
||||
{ID: 101, Count: 3},
|
||||
|
|
@ -683,7 +683,7 @@ func TestFragment_TopN_Intersect_Large(t *testing.T) {
|
|||
f.RecalculateCache()
|
||||
|
||||
// Retrieve top rows.
|
||||
if pairs, err := f.Top(TopOptions{N: 10, Src: src}); err != nil {
|
||||
if pairs, err := f.top(TopOptions{N: 10, Src: src}); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(pairs, []Pair{
|
||||
{ID: 999, Count: 19},
|
||||
|
|
@ -712,7 +712,7 @@ func TestFragment_TopN_IDs(t *testing.T) {
|
|||
f.mustSetBits(102, 8, 9, 10, 11, 12)
|
||||
|
||||
// Retrieve top rows.
|
||||
if pairs, err := f.Top(TopOptions{RowIDs: []uint64{100, 101, 200}}); err != nil {
|
||||
if pairs, err := f.top(TopOptions{RowIDs: []uint64{100, 101, 200}}); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(pairs, []Pair{
|
||||
{ID: 101, Count: 4},
|
||||
|
|
@ -733,7 +733,7 @@ func TestFragment_TopN_NopCache(t *testing.T) {
|
|||
f.mustSetBits(102, 8, 9, 10, 11, 12)
|
||||
|
||||
// Retrieve top rows.
|
||||
if pairs, err := f.Top(TopOptions{RowIDs: []uint64{100, 101, 200}}); err != nil {
|
||||
if pairs, err := f.top(TopOptions{RowIDs: []uint64{100, 101, 200}}); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual(pairs, []Pair{}) {
|
||||
t.Fatalf("unexpected pairs: %s", spew.Sdump(pairs))
|
||||
|
|
@ -792,7 +792,7 @@ func TestFragment_TopN_CacheSize(t *testing.T) {
|
|||
}
|
||||
|
||||
// Retrieve top rows.
|
||||
if pairs, err := f.Top(TopOptions{N: 5}); err != nil {
|
||||
if pairs, err := f.top(TopOptions{N: 5}); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if len(pairs) > int(cacheSize) {
|
||||
t.Fatalf("TopN count cannot exceed cache size: %d", cacheSize)
|
||||
|
|
@ -891,8 +891,8 @@ func TestFragment_LRUCache_Persistence(t *testing.T) {
|
|||
}
|
||||
|
||||
// Verify correct cache type and size.
|
||||
if cache, ok := f.Cache().(*LRUCache); !ok {
|
||||
t.Fatalf("unexpected cache: %T", f.Cache())
|
||||
if cache, ok := f.cache.(*LRUCache); !ok {
|
||||
t.Fatalf("unexpected cache: %T", f.cache)
|
||||
} else if cache.Len() != 1000 {
|
||||
t.Fatalf("unexpected cache len: %d", cache.Len())
|
||||
}
|
||||
|
|
@ -903,8 +903,8 @@ func TestFragment_LRUCache_Persistence(t *testing.T) {
|
|||
}
|
||||
|
||||
// Re-verify correct cache type and size.
|
||||
if cache, ok := f.Cache().(*LRUCache); !ok {
|
||||
t.Fatalf("unexpected cache: %T", f.Cache())
|
||||
if cache, ok := f.cache.(*LRUCache); !ok {
|
||||
t.Fatalf("unexpected cache: %T", f.cache)
|
||||
} else if cache.Len() != 1000 {
|
||||
t.Fatalf("unexpected cache len: %d", cache.Len())
|
||||
}
|
||||
|
|
@ -941,8 +941,8 @@ func TestFragment_RankCache_Persistence(t *testing.T) {
|
|||
}
|
||||
|
||||
// Verify correct cache type and size.
|
||||
if cache, ok := f.Cache().(*RankCache); !ok {
|
||||
t.Fatalf("unexpected cache: %T", f.Cache())
|
||||
if cache, ok := f.cache.(*RankCache); !ok {
|
||||
t.Fatalf("unexpected cache: %T", f.cache)
|
||||
} else if cache.Len() != 1000 {
|
||||
t.Fatalf("unexpected cache len: %d", cache.Len())
|
||||
}
|
||||
|
|
@ -956,8 +956,8 @@ func TestFragment_RankCache_Persistence(t *testing.T) {
|
|||
f = index.Field("f").View(ViewStandard).Fragment(0)
|
||||
|
||||
// Re-verify correct cache type and size.
|
||||
if cache, ok := f.Cache().(*RankCache); !ok {
|
||||
t.Fatalf("unexpected cache: %T", f.Cache())
|
||||
if cache, ok := f.cache.(*RankCache); !ok {
|
||||
t.Fatalf("unexpected cache: %T", f.cache)
|
||||
} else if cache.Len() != 1000 {
|
||||
t.Fatalf("unexpected cache len: %d", cache.Len())
|
||||
}
|
||||
|
|
@ -978,7 +978,7 @@ func TestFragment_WriteTo_ReadFrom(t *testing.T) {
|
|||
}
|
||||
|
||||
// Verify cache is populated.
|
||||
if n := f0.Cache().Len(); n != 1 {
|
||||
if n := f0.cache.Len(); n != 1 {
|
||||
t.Fatalf("unexpected cache size: %d", n)
|
||||
}
|
||||
|
||||
|
|
@ -998,7 +998,7 @@ func TestFragment_WriteTo_ReadFrom(t *testing.T) {
|
|||
}
|
||||
|
||||
// Verify cache is in other fragment.
|
||||
if n := f1.Cache().Len(); n != 1 {
|
||||
if n := f1.cache.Len(); n != 1 {
|
||||
t.Fatalf("unexpected cache size: %d", n)
|
||||
}
|
||||
|
||||
|
|
@ -1010,7 +1010,7 @@ func TestFragment_WriteTo_ReadFrom(t *testing.T) {
|
|||
// Close and reopen the fragment & verify the data.
|
||||
if err := f1.reopen(); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if n := f1.Cache().Len(); n != 1 {
|
||||
} else if n := f1.cache.Len(); n != 1 {
|
||||
t.Fatalf("unexpected cache size (reopen): %d", n)
|
||||
} else if a := f1.Row(1000).Columns(); !reflect.DeepEqual(a, []uint64{2}) {
|
||||
t.Fatalf("unexpected columns (reopen): %+v", a)
|
||||
|
|
@ -1081,7 +1081,7 @@ func TestFragment_Tanimoto(t *testing.T) {
|
|||
f.mustSetBits(102, 1, 2, 10, 12)
|
||||
f.RecalculateCache()
|
||||
|
||||
if pairs, err := f.Top(TopOptions{TanimotoThreshold: 50, Src: src}); err != nil {
|
||||
if pairs, err := f.top(TopOptions{TanimotoThreshold: 50, Src: src}); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if len(pairs) != 2 {
|
||||
t.Fatalf("unexpected count: %d", len(pairs))
|
||||
|
|
@ -1104,7 +1104,7 @@ func TestFragment_Zero_Tanimoto(t *testing.T) {
|
|||
f.mustSetBits(102, 1, 2, 10, 12)
|
||||
f.RecalculateCache()
|
||||
|
||||
if pairs, err := f.Top(TopOptions{TanimotoThreshold: 0, Src: src}); err != nil {
|
||||
if pairs, err := f.top(TopOptions{TanimotoThreshold: 0, Src: src}); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if len(pairs) != 3 {
|
||||
t.Fatalf("unexpected count: %d", len(pairs))
|
||||
|
|
@ -1187,7 +1187,7 @@ func BenchmarkFragment_FullSnapshot(b *testing.B) {
|
|||
val += 2
|
||||
i++
|
||||
}
|
||||
if err := f.Import(rows, cols); err != nil {
|
||||
if err := f.bulkImport(rows, cols); err != nil {
|
||||
b.Fatalf("Error Building Sample: %s", err)
|
||||
}
|
||||
if row > max {
|
||||
|
|
@ -1228,7 +1228,7 @@ func BenchmarkFragment_Import(b *testing.B) {
|
|||
b.ResetTimer()
|
||||
b.ReportAllocs()
|
||||
for i := 0; i < b.N; i++ {
|
||||
if err := f.Import(rows, cols); err != nil {
|
||||
if err := f.bulkImport(rows, cols); err != nil {
|
||||
b.Fatalf("Error Building Sample: %s", err)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -445,7 +445,7 @@ func (h *Holder) flushCaches() {
|
|||
}
|
||||
|
||||
if err := fragment.FlushCache(); err != nil {
|
||||
h.Logger.Printf("error flushing cache: err=%s, path=%s", err, fragment.CachePath())
|
||||
h.Logger.Printf("error flushing cache: err=%s, path=%s", err, fragment.cachePath())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -765,7 +765,7 @@ func (s *HolderSyncer) syncFragment(index, field, view string, slice uint64) err
|
|||
Closing: s.Closing,
|
||||
RemoteClient: s.RemoteClient,
|
||||
}
|
||||
if err := fs.SyncFragment(); err != nil {
|
||||
if err := fs.syncFragment(); err != nil {
|
||||
return errors.Wrap(err, "syncing fragment")
|
||||
}
|
||||
|
||||
|
|
@ -809,7 +809,7 @@ func (c *HolderCleaner) CleanHolder() error {
|
|||
for _, field := range index.Fields() {
|
||||
for _, view := range field.Views() {
|
||||
for _, fragment := range view.Fragments() {
|
||||
fragSlice := fragment.Slice()
|
||||
fragSlice := fragment.slice
|
||||
// Ignore fragments that should be present.
|
||||
if uint64InSlice(fragSlice, containedSlices) {
|
||||
continue
|
||||
|
|
|
|||
22
view.go
22
view.go
|
|
@ -151,10 +151,10 @@ func (v *View) openFragments() error {
|
|||
|
||||
frag := v.newFragment(v.FragmentPath(slice), slice)
|
||||
if err := frag.Open(); err != nil {
|
||||
return fmt.Errorf("open fragment: slice=%d, err=%s", frag.Slice(), err)
|
||||
return fmt.Errorf("open fragment: slice=%d, err=%s", frag.slice, err)
|
||||
}
|
||||
frag.RowAttrStore = v.RowAttrStore
|
||||
v.fragments[frag.Slice()] = frag
|
||||
v.fragments[frag.slice] = frag
|
||||
}
|
||||
|
||||
return nil
|
||||
|
|
@ -289,12 +289,12 @@ func (v *View) DeleteFragment(slice uint64) error {
|
|||
}
|
||||
|
||||
// Delete fragment file.
|
||||
if err := os.Remove(fragment.Path()); err != nil {
|
||||
if err := os.Remove(fragment.path); err != nil {
|
||||
return errors.Wrap(err, "deleting fragment file")
|
||||
}
|
||||
|
||||
// Delete fragment cache file.
|
||||
if err := os.Remove(fragment.CachePath()); err != nil {
|
||||
if err := os.Remove(fragment.cachePath()); err != nil {
|
||||
v.Logger.Printf("no cache file to delete for slice %d", slice)
|
||||
}
|
||||
|
||||
|
|
@ -330,7 +330,7 @@ func (v *View) value(columnID uint64, bitDepth uint) (value uint64, exists bool,
|
|||
if err != nil {
|
||||
return value, exists, err
|
||||
}
|
||||
return frag.Value(columnID, bitDepth)
|
||||
return frag.value(columnID, bitDepth)
|
||||
}
|
||||
|
||||
// setValue uses a column of bits to set a multi-bit value.
|
||||
|
|
@ -340,13 +340,13 @@ func (v *View) setValue(columnID uint64, bitDepth uint, value uint64) (changed b
|
|||
if err != nil {
|
||||
return changed, err
|
||||
}
|
||||
return frag.SetValue(columnID, bitDepth, value)
|
||||
return frag.setValue(columnID, bitDepth, value)
|
||||
}
|
||||
|
||||
// sum returns the sum & count of a field.
|
||||
func (v *View) sum(filter *Row, bitDepth uint) (sum, count uint64, err error) {
|
||||
for _, f := range v.Fragments() {
|
||||
fsum, fcount, err := f.Sum(filter, bitDepth)
|
||||
fsum, fcount, err := f.sum(filter, bitDepth)
|
||||
if err != nil {
|
||||
return sum, count, err
|
||||
}
|
||||
|
|
@ -360,7 +360,7 @@ func (v *View) sum(filter *Row, bitDepth uint) (sum, count uint64, err error) {
|
|||
func (v *View) min(filter *Row, bitDepth uint) (min, count uint64, err error) {
|
||||
var minHasValue bool
|
||||
for _, f := range v.Fragments() {
|
||||
fmin, fcount, err := f.Min(filter, bitDepth)
|
||||
fmin, fcount, err := f.min(filter, bitDepth)
|
||||
if err != nil {
|
||||
return min, count, err
|
||||
}
|
||||
|
|
@ -387,7 +387,7 @@ func (v *View) min(filter *Row, bitDepth uint) (min, count uint64, err error) {
|
|||
// max returns the max and count of a field.
|
||||
func (v *View) max(filter *Row, bitDepth uint) (max, count uint64, err error) {
|
||||
for _, f := range v.Fragments() {
|
||||
fmax, fcount, err := f.Max(filter, bitDepth)
|
||||
fmax, fcount, err := f.max(filter, bitDepth)
|
||||
if err != nil {
|
||||
return max, count, err
|
||||
}
|
||||
|
|
@ -403,7 +403,7 @@ func (v *View) max(filter *Row, bitDepth uint) (max, count uint64, err error) {
|
|||
func (v *View) rangeOp(op pql.Token, bitDepth uint, predicate uint64) (*Row, error) {
|
||||
r := NewRow()
|
||||
for _, frag := range v.Fragments() {
|
||||
other, err := frag.RangeOp(op, bitDepth, predicate)
|
||||
other, err := frag.rangeOp(op, bitDepth, predicate)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -417,7 +417,7 @@ func (v *View) rangeOp(op pql.Token, bitDepth uint, predicate uint64) (*Row, err
|
|||
func (v *View) rangeBetween(bitDepth uint, predicateMin, predicateMax uint64) (*Row, error) {
|
||||
r := NewRow()
|
||||
for _, frag := range v.Fragments() {
|
||||
other, err := frag.RangeBetween(bitDepth, predicateMin, predicateMax)
|
||||
other, err := frag.rangeBetween(bitDepth, predicateMin, predicateMax)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue