From d2b925d2964892e8a21a0cb0c6da495d7d3402b9 Mon Sep 17 00:00:00 2001 From: Seebs Date: Tue, 18 May 2021 15:17:29 -0500 Subject: [PATCH] make importValueSmallWrite faster and also the only path Since we don't always have "snapshots" anymore, the arguable benefit of avoiding the snapshot is reduced, and the primary expense of importPositions has been dramatically reduced as well, so let's just use that all the time, and simplify life. We also want to make it faster. We don't know how many bits there are to set or clear in the input set, but we do know exactly how many bits there are to set AND clear. We can subdivide these into batches by rows, then process each batch by storing sets at the bottom and clears at the top. We can also do batches by columns, reducing the memory overhead of unpacking all the bits at once. (For extra credit, we could alternate set/clear settings, and thus do batches of "the clears from row 0, followed by the clears from row 1" and "the sets from row 1, followed by the sets from row 2", and so on, but this is too fancy.) Every caller of importValue is in fact already providing values with column IDs sorted. As such, we don't need a map for checking the previously-set columns; we just need to check against the previous value. --- fragment.go | 231 ++++++++++++++++------------------------------------ 1 file changed, 70 insertions(+), 161 deletions(-) diff --git a/fragment.go b/fragment.go index d81caa460..c458a4cb3 100644 --- a/fragment.go +++ b/fragment.go @@ -1136,74 +1136,6 @@ func (f *fragment) setValueBase(txOrig Tx, columnID uint64, bitDepth uint64, val return changed, err } -// importSetValue is a more efficient SetValue just for imports. -func (f *fragment) importSetValue(txb *TxBitmap, columnID uint64, bitDepth uint64, value int64, clear bool) (changed int, err error) { // nolint: unparam - // Convert value to an unsigned representation. - uvalue := uint64(value) - if value < 0 { - uvalue = uint64(-value) - } - - for i := uint64(0); i < bitDepth; i++ { - bit, err := f.pos(uint64(bsiOffsetBit+i), columnID) - if err != nil { - return changed, errors.Wrap(err, "getting pos") - } - - if uvalue&(1<= 0 || clear { - if c, err := txb.Remove(p); err != nil { - return changed, errors.Wrap(err, "removing sign from storage") - } else if c { - changed++ - } - } else { - if c, err := txb.Add(p); err != nil { - return changed, errors.Wrap(err, "adding sign to storage") - } else if c { - changed++ - } - } - - return changed, nil -} - // 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(tx Tx, filter *Row, bitDepth uint64) (sum int64, count uint64, err error) { @@ -2613,57 +2545,6 @@ func (f *fragment) bulkImportMutex(tx Tx, rowIDs, columnIDs []uint64) error { return errors.Wrap(f.importPositions(tx, toSet, toClear, rowSet), "importing positions") } -func (f *fragment) importValueSmallWrite(tx Tx, columnIDs []uint64, values []int64, bitDepth uint64, clear bool) error { - // TODO figure out how to avoid re-allocating these each time. Probably - // possible to store them on the fragment with a capacity based on - // MaxOpN. For now, we know that the total number of bits to be - // set+cleared is len(values)*(bitDepth+1), so we make each slice - // slightly more than half of that to try to avoid reallocation. - toSet := make([]uint64, 0, len(columnIDs)*int(bitDepth+1)*(5/8)) - toClear := make([]uint64, 0, len(columnIDs)*int(bitDepth+1)*(5/8)) - colSet := make(map[uint64]struct{}, len(columnIDs)) - - if err := func() (err error) { - for i := len(columnIDs) - 1; i >= 0; i-- { - columnID, value := columnIDs[i], values[i] - if _, ok := colSet[columnID]; ok { - continue - } - - colSet[columnID] = struct{}{} - toSet, toClear, err = f.positionsForValue(columnID, bitDepth, value, clear, toSet, toClear) - if err != nil { - return errors.Wrap(err, "getting positions for value") - } - } - return nil - }(); err != nil { - errOpenStorage := f.openStorage(true) - if errOpenStorage != nil { - f.Logger.Errorf("failed to import data into fragment: %v", err) - f.Logger.Errorf("recovery with openStorage failed for fragment: %v", errOpenStorage) - f.Logger.Debugf("%s", debug.Stack()) - os.Exit(1) - } - return err - } - rowSet := make(map[uint64]struct{}, bitDepth+1) - for i := uint64(0); i < bitDepth+1; i++ { - rowSet[uint64(i)] = struct{}{} - } - err := f.importPositions(tx, toSet, toClear, rowSet) - if err != nil { - return errors.Wrap(err, "importing positions") - } - - if tx.UseRowCache() { - // Reset the rowCache. - f.rowCache = newSimpleCache() - } - - return nil -} - // importValue bulk imports a set of range-encoded values. func (f *fragment) importValue(tx Tx, columnIDs []uint64, values []int64, bitDepth uint64, clear bool) error { f.mu.Lock() @@ -2673,56 +2554,84 @@ func (f *fragment) importValue(tx Tx, columnIDs []uint64, values []int64, bitDep if len(columnIDs) != len(values) { return fmt.Errorf("mismatch of column/value len: %d != %d", len(columnIDs), len(values)) } - - if len(columnIDs)*int(bitDepth+1)+f.opN < f.MaxOpN { - return errors.Wrap(f.importValueSmallWrite(tx, columnIDs, values, bitDepth, clear), "import small write") + positionsByDepth := make([][]uint64, bitDepth+2) + toSetByDepth := make([]int, bitDepth+2) + toClearByDepth := make([]int, bitDepth+2) + batchSize := len(columnIDs) + if batchSize > 65536 { + batchSize = 65536 + } + for i := 0; i < int(bitDepth)+2; i++ { + positionsByDepth[i] = make([]uint64, batchSize) + toClearByDepth[i] = batchSize } - // Process every value. - // If an error occurs then reopen the storage. - if f.storage != nil { - f.storage.OpWriter = nil - } - - var totalChanges int - if err := func() (err error) { - // Build changes into temporary bitmap. - txb := NewTxBitmap(tx, f.index(), f.field(), f.view(), f.shard) - for i := range columnIDs { - columnID, value := columnIDs[i], values[i] - if _, err := f.importSetValue(txb, columnID, bitDepth, value, clear); err != nil { - return errors.Wrapf(err, "importSetValue") + row := 0 + columnID := uint64(0) + value := int64(0) + // arbitrarily set prev to be not equal to the first column ID + // we will encounter. + prev := columnIDs[len(columnIDs)-1] + 1 + for len(columnIDs) > 0 { + downTo := len(columnIDs) - batchSize + if downTo < 0 { + downTo = 0 + } + for i := range positionsByDepth { + toSetByDepth[i] = 0 + toClearByDepth[i] = batchSize + } + for i := len(columnIDs) - 1; i >= downTo; i-- { + columnID, value = columnIDs[i], values[i] + columnID = columnID % ShardWidth + if columnID == prev { + continue + } + prev = columnID + row = 0 + if clear { + toClearByDepth[row]-- + positionsByDepth[row][toClearByDepth[row]] = columnID + } else { + positionsByDepth[row][toSetByDepth[row]] = columnID + toSetByDepth[row]++ + } + row++ + columnID += ShardWidth + if value < 0 { + positionsByDepth[row][toSetByDepth[row]] = columnID + toSetByDepth[row]++ + value *= -1 + } else { + toClearByDepth[row]-- + positionsByDepth[row][toClearByDepth[row]] = columnID + } + row++ + columnID += ShardWidth + for j := 0; j < int(bitDepth); j++ { + if value&1 != 0 { + positionsByDepth[row][toSetByDepth[row]] = columnID + toSetByDepth[row]++ + } else { + toClearByDepth[row]-- + positionsByDepth[row][toClearByDepth[row]] = columnID + } + row++ + columnID += ShardWidth + value >>= 1 } } - // Flush changes in bulk back to the transaction. - return txb.Flush() - }(); err != nil { - errOpenStorage := f.openStorage(true) - if errOpenStorage != nil { - f.Logger.Errorf("failed to import data into fragment: %v", err) - f.Logger.Errorf("recovery with openStorage failed for fragment: %v", errOpenStorage) - f.Logger.Debugf("%s", debug.Stack()) - os.Exit(1) + for i := range positionsByDepth { + err := f.importPositions(tx, positionsByDepth[i][:toSetByDepth[i]], positionsByDepth[i][toClearByDepth[i]:], nil) + if err != nil { + return errors.Wrap(err, "importing positions") + } } - return err - } - // Keep stats accurate. We don't call incrementOpN here because it may - // or may not enqueue a request, which would then be in the queue - // taking up space and otherwise being a possible nuisance, when we're - // about to force a snapshot anyway. - f.opN += totalChanges - f.ops++ - - if tx.UseRowCache() { - // Reset the rowCache. - f.rowCache = newSimpleCache() + columnIDs = columnIDs[:downTo] } - // in theory, this should probably have been queued anyway, but if enough - // of the bits matched existing bits, we'll be under our opN estimate, and - // we want to ensure that the snapshot happens. - return f.holder.SnapshotQueue.Immediate(f) + return nil } // importRoaring imports from the official roaring data format defined at