add span around fragment lock, bytes written metadata

This commit is contained in:
Matt Jaffee 2019-04-30 16:55:46 -05:00
parent 61bf3d929d
commit 00911d024b
No known key found for this signature in database
GPG key ID: 08A3DFFF987B11BF
2 changed files with 16 additions and 12 deletions

View file

@ -1814,8 +1814,10 @@ func (f *fragment) importValue(columnIDs, values []uint64, bitDepth uint, clear
func (f *fragment) importRoaring(ctx context.Context, data []byte, clear bool) error {
span, ctx := tracing.StartSpanFromContext(ctx, "fragment.importRoaring")
defer span.Finish()
span, ctx = tracing.StartSpanFromContext(ctx, "importRoaring.AcquireFragmentLock")
f.mu.Lock()
defer f.mu.Unlock()
span.Finish()
bm := roaring.NewBTreeBitmap()
span, ctx = tracing.StartSpanFromContext(ctx, "importRoaring.UnmarshalBinary")
err := bm.UnmarshalBinary(data)
@ -1883,7 +1885,8 @@ func (f *fragment) importRoaring(ctx context.Context, data []byte, clear bool) e
f.cache.Recalculate()
span, _ = tracing.StartSpanFromContext(ctx, "importRoaring.WriteToFragment")
err = unprotectedWriteToFragment(f, bm)
n, err := unprotectedWriteToFragment(f, bm)
span.LogKV("bytesWritten", n)
span.Finish()
return err
}
@ -1915,12 +1918,13 @@ func track(start time.Time, message string, stats stats.StatsClient, logger logg
}
func (f *fragment) snapshot() error {
return unprotectedWriteToFragment(f, f.storage)
_, err := unprotectedWriteToFragment(f, f.storage)
return err
}
// unprotectedWriteToFragment writes the fragment f with bm as the data. It is unprotected, and
// f.mu must be locked when calling it.
func unprotectedWriteToFragment(f *fragment, bm *roaring.Bitmap) error { // nolint: interfacer
func unprotectedWriteToFragment(f *fragment, bm *roaring.Bitmap) (n int64, err error) { // nolint: interfacer
completeMessage := fmt.Sprintf("fragment: snapshot complete %s/%s/%s/%d", f.index, f.field, f.view, f.shard)
start := time.Now()
@ -1930,39 +1934,39 @@ func unprotectedWriteToFragment(f *fragment, bm *roaring.Bitmap) error { // noli
snapshotPath := f.path + snapshotExt
file, err := os.Create(snapshotPath)
if err != nil {
return fmt.Errorf("create snapshot file: %s", err)
return n, fmt.Errorf("create snapshot file: %s", err)
}
defer file.Close()
// Write storage to snapshot.
bw := bufio.NewWriter(file)
if _, err := bm.WriteTo(bw); err != nil {
return fmt.Errorf("snapshot write to: %s", err)
if n, err = bm.WriteTo(bw); err != nil {
return n, fmt.Errorf("snapshot write to: %s", err)
}
if err := bw.Flush(); err != nil {
return fmt.Errorf("flush: %s", err)
return n, fmt.Errorf("flush: %s", err)
}
// Close current storage.
if err := f.closeStorage(); err != nil {
return fmt.Errorf("close storage: %s", err)
return n, fmt.Errorf("close storage: %s", err)
}
// Move snapshot to data file location.
if err := os.Rename(snapshotPath, f.path); err != nil {
return fmt.Errorf("rename snapshot: %s", err)
return n, fmt.Errorf("rename snapshot: %s", err)
}
// Reopen storage.
if err := f.openStorage(); err != nil {
return fmt.Errorf("open storage: %s", err)
return n, fmt.Errorf("open storage: %s", err)
}
// Reset operation count.
f.opN = 0
return nil
return n, nil
}
// RecalculateCache rebuilds the cache regardless of invalidate time delay.

View file

@ -2947,7 +2947,7 @@ func TestUnionInPlaceMapped(t *testing.T) {
count1 := setBM1.Count()
// now we write setBM0 into f.storage.
err = unprotectedWriteToFragment(f, setBM0)
_, err = unprotectedWriteToFragment(f, setBM0)
if err != nil {
t.Fatalf("trying to flush fragment to disk: %v", err)
}