diff --git a/fragment.go b/fragment.go index 2b2925e16..f864fdc1c 100644 --- a/fragment.go +++ b/fragment.go @@ -1785,7 +1785,9 @@ func (f *fragment) importValue(columnIDs, values []uint64, bitDepth uint, clear for i := uint(0); i < bitDepth+1; i++ { rowSet[uint64(i)] = struct{}{} } + f.mu.Lock() err := f.importPositions(toSet, toClear, rowSet) + f.mu.Unlock() return errors.Wrap(err, "importing positions") } err := f.snapshot() diff --git a/fragment_internal_test.go b/fragment_internal_test.go index b726d4fe5..71d3b5ebe 100644 --- a/fragment_internal_test.go +++ b/fragment_internal_test.go @@ -3024,3 +3024,24 @@ func TestSmallImportRestart(t *testing.T) { t.Errorf("row 1 should be [1], but got %v", r1) } } + +func TestImportValueConcurrent(t *testing.T) { + f := mustOpenFragment("i", "f", viewBSIGroupPrefix+"foo", 0, "none") + eg := &errgroup.Group{} + for i := 0; i < 4; i++ { + i := i + eg.Go(func() error { + for j := uint64(0); j < 10; j++ { + err := f.importValue([]uint64{j}, []uint64{uint64(rand.Int63n(1000))}, 10, i%2 == 0) + if err != nil { + return err + } + } + return nil + }) + } + err := eg.Wait() + if err != nil { + t.Fatalf("concurrently importing values: %v", err) + } +}