wip on adding mutex support to random import perf

This commit is contained in:
Matt Jaffee 2019-02-27 20:14:06 -06:00
parent f06a9f0e6e
commit 4082ce655a
No known key found for this signature in database
GPG key ID: 08A3DFFF987B11BF

View file

@ -1639,88 +1639,43 @@ func (f *fragment) bulkImportMutex(rowIDs, columnIDs []uint64) error {
f.mu.Lock()
defer f.mu.Unlock()
// Disconnect op writer so we don't append updates.
f.storage.OpWriter = nil
// If an error occurs then reopen the storage.
if err := func() error {
// rowSet maintains the set of rowIDs present in this import.
// It allows the cache to be updated once per row, instead of once
// per bit.
rowSet := make(map[uint64]struct{})
lastRowID := uint64(0)
// Process every bit.
for i := range rowIDs {
rowID, columnID := rowIDs[i], columnIDs[i]
// Handle mutex vector (i.e. clear an existing row).
if existingRowID, found, err := f.mutexVector.Get(columnID); err != nil {
return errors.Wrap(err, "getting mutex vector data")
} else if found && existingRowID != rowID {
// Determine the position of the bit in the storage.
pos, err := f.pos(existingRowID, columnID)
if err != nil {
return err
}
// Clear storage.
_, err = f.storage.Remove(pos)
if err != nil {
return err
}
rowSet[existingRowID] = struct{}{}
}
rowSet := make(map[uint64]struct{})
// Since each imported bit will at most set one bit and clear one bit, we
// can reuse the rowIDs and columnIDs slices as the set and clear slice
// arguments to importPositions. We maintain indexes into them as we loop
// through them since the indexes of set and cleared bits may lag behind the
// current index.
setIdx := 0
clearIdx := 0
for i := range rowIDs {
rowID, columnID := rowIDs[i], columnIDs[i]
if existingRowID, found, err := f.mutexVector.Get(columnID); err != nil {
return errors.Wrap(err, "getting mutex vector data")
} else if found && existingRowID != rowID {
// Determine the position of the bit in the storage.
pos, err := f.pos(rowID, columnID)
clearPos, err := f.pos(existingRowID, columnID)
if err != nil {
return err
}
columnIDs[clearIdx] = clearPos
clearIdx++
// Write to storage.
_, err = f.storage.Add(pos)
if err != nil {
return err
}
// Reduce the StatsD rate for high volume stats
f.stats.Count("ImportBit", 1, 0.0001)
// Add row to rowSet.
if i == 0 || rowID != lastRowID {
lastRowID = rowID
rowSet[rowID] = struct{}{}
}
// Invalidate block checksum.
delete(f.checksums, int(rowID/HashBlockSize))
rowSet[existingRowID] = struct{}{}
} else if found && existingRowID == rowID {
continue
}
// Update cache counts for all rows.
for rowID := range rowSet {
// Import should ALWAYS have row() load a new bm from fragment.storage
// because the row that's in rowCache hasn't been updated with
// this import's data.
f.cache.BulkAdd(rowID, f.unprotectedRow(rowID).Count())
pos, err := f.pos(rowID, columnID)
if err != nil {
return err
}
f.cache.Invalidate()
return nil
}(); err != nil {
_ = f.closeStorage()
_ = f.openStorage()
return err
rowIDs[setIdx] = pos
setIdx++
rowSet[rowID] = struct{}{}
}
toSet := rowIDs[:setIdx]
toClear := columnIDs[:clearIdx]
// Write the storage to disk and reload.
if err := f.snapshot(); err != nil {
return err
}
return nil
return errors.Wrap(f.importPositions(toSet, toClear, rowSet), "importing positions")
}
// importValue bulk imports a set of range-encoded values.