mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-08 20:07:51 +00:00
add concurrent update benchmark and clean up temp frags
This commit is contained in:
parent
054926fd6c
commit
4f459028d9
1 changed files with 125 additions and 69 deletions
|
|
@ -21,6 +21,7 @@ import (
|
|||
"io/ioutil"
|
||||
"math"
|
||||
"math/rand"
|
||||
"os"
|
||||
"reflect"
|
||||
"sort"
|
||||
"testing"
|
||||
|
|
@ -43,7 +44,7 @@ var (
|
|||
// Ensure a fragment can set a bit and retrieve it.
|
||||
func TestFragment_SetBit(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set bits on the fragment.
|
||||
if _, err := f.setBit(120, 1); err != nil {
|
||||
|
|
@ -74,7 +75,7 @@ func TestFragment_SetBit(t *testing.T) {
|
|||
// Ensure a fragment can clear a set bit.
|
||||
func TestFragment_ClearBit(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set and then clear bits on the fragment.
|
||||
if _, err := f.setBit(1000, 1); err != nil {
|
||||
|
|
@ -101,7 +102,7 @@ func TestFragment_ClearBit(t *testing.T) {
|
|||
// Ensure a fragment can clear a row.
|
||||
func TestFragment_ClearRow(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set and then clear bits on the fragment.
|
||||
if _, err := f.setBit(1000, 1); err != nil {
|
||||
|
|
@ -128,7 +129,7 @@ func TestFragment_ClearRow(t *testing.T) {
|
|||
// Ensure a fragment can set a row.
|
||||
func TestFragment_SetRow(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 7, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
rowID := uint64(1000)
|
||||
|
||||
|
|
@ -177,7 +178,7 @@ func TestFragment_SetRow(t *testing.T) {
|
|||
func TestFragment_SetValue(t *testing.T) {
|
||||
t.Run("OK", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set value.
|
||||
if changed, err := f.setValue(100, 16, 3829); err != nil {
|
||||
|
|
@ -205,7 +206,7 @@ func TestFragment_SetValue(t *testing.T) {
|
|||
|
||||
t.Run("Overwrite", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set value.
|
||||
if changed, err := f.setValue(100, 16, 3829); err != nil {
|
||||
|
|
@ -233,7 +234,7 @@ func TestFragment_SetValue(t *testing.T) {
|
|||
|
||||
t.Run("Clear", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set value.
|
||||
if changed, err := f.setValue(100, 16, 3829); err != nil {
|
||||
|
|
@ -261,7 +262,7 @@ func TestFragment_SetValue(t *testing.T) {
|
|||
|
||||
t.Run("NotExists", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set value.
|
||||
if changed, err := f.setValue(100, 10, 20); err != nil {
|
||||
|
|
@ -291,7 +292,7 @@ func TestFragment_SetValue(t *testing.T) {
|
|||
}
|
||||
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set values.
|
||||
m := make(map[uint64]int64)
|
||||
|
|
@ -329,7 +330,7 @@ func TestFragment_Sum(t *testing.T) {
|
|||
const bitDepth = 16
|
||||
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set values.
|
||||
if _, err := f.setValue(1000, bitDepth, 382); err != nil {
|
||||
|
|
@ -368,7 +369,7 @@ func TestFragment_MinMax(t *testing.T) {
|
|||
const bitDepth = 16
|
||||
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set values.
|
||||
if _, err := f.setValue(1000, bitDepth, 382); err != nil {
|
||||
|
|
@ -442,7 +443,7 @@ func TestFragment_Range(t *testing.T) {
|
|||
|
||||
t.Run("EQ", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set values.
|
||||
if _, err := f.setValue(1000, bitDepth, 382); err != nil {
|
||||
|
|
@ -465,7 +466,7 @@ func TestFragment_Range(t *testing.T) {
|
|||
|
||||
t.Run("NEQ", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set values.
|
||||
if _, err := f.setValue(1000, bitDepth, 382); err != nil {
|
||||
|
|
@ -488,7 +489,7 @@ func TestFragment_Range(t *testing.T) {
|
|||
|
||||
t.Run("LT", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set values.
|
||||
if _, err := f.setValue(1000, bitDepth, 382); err != nil {
|
||||
|
|
@ -536,7 +537,7 @@ func TestFragment_Range(t *testing.T) {
|
|||
|
||||
t.Run("GT", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set values.
|
||||
if _, err := f.setValue(1000, bitDepth, 382); err != nil {
|
||||
|
|
@ -584,7 +585,7 @@ func TestFragment_Range(t *testing.T) {
|
|||
|
||||
t.Run("BETWEEN", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set values.
|
||||
if _, err := f.setValue(1000, bitDepth, 382); err != nil {
|
||||
|
|
@ -634,7 +635,7 @@ func TestFragment_Range(t *testing.T) {
|
|||
// Ensure a fragment can snapshot correctly.
|
||||
func TestFragment_Snapshot(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set and then clear bits on the fragment.
|
||||
if _, err := f.setBit(1000, 1); err != nil {
|
||||
|
|
@ -663,7 +664,7 @@ func TestFragment_Snapshot(t *testing.T) {
|
|||
// Ensure a fragment can iterate over all bits in order.
|
||||
func TestFragment_ForEachBit(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set bits on the fragment.
|
||||
if _, err := f.setBit(100, 20); err != nil {
|
||||
|
|
@ -692,7 +693,7 @@ func TestFragment_ForEachBit(t *testing.T) {
|
|||
// Ensure a fragment can return the top n results.
|
||||
func TestFragment_Top(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, CacheTypeRanked)
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
// Set bits on the rows 100, 101, & 102.
|
||||
f.mustSetBits(100, 1, 3, 200)
|
||||
f.mustSetBits(101, 1)
|
||||
|
|
@ -714,7 +715,7 @@ func TestFragment_Top(t *testing.T) {
|
|||
// Ensure a fragment can filter rows when retrieving the top n rows.
|
||||
func TestFragment_Top_Filter(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, CacheTypeRanked)
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set bits on the rows 100, 101, & 102.
|
||||
f.mustSetBits(100, 1, 3, 200)
|
||||
|
|
@ -744,7 +745,7 @@ func TestFragment_Top_Filter(t *testing.T) {
|
|||
// Ensure a fragment can return top rows that intersect with an input row.
|
||||
func TestFragment_TopN_Intersect(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, CacheTypeRanked)
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Create an intersecting input row.
|
||||
src := NewRow(1, 2, 3)
|
||||
|
|
@ -775,7 +776,7 @@ func TestFragment_TopN_Intersect_Large(t *testing.T) {
|
|||
}
|
||||
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, CacheTypeRanked)
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Create an intersecting input row.
|
||||
src := NewRow(
|
||||
|
|
@ -813,7 +814,7 @@ func TestFragment_TopN_Intersect_Large(t *testing.T) {
|
|||
// Ensure a fragment can return top rows when specified by ID.
|
||||
func TestFragment_TopN_IDs(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, CacheTypeRanked)
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set bits on various rows.
|
||||
f.mustSetBits(100, 1, 2, 3)
|
||||
|
|
@ -834,7 +835,7 @@ func TestFragment_TopN_IDs(t *testing.T) {
|
|||
// Ensure a fragment return none if CacheTypeNone is set
|
||||
func TestFragment_TopN_NopCache(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, CacheTypeNone)
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set bits on various rows.
|
||||
f.mustSetBits(100, 1, 2, 3)
|
||||
|
|
@ -882,7 +883,7 @@ func TestFragment_TopN_CacheSize(t *testing.T) {
|
|||
if err := f.Open(); err != nil {
|
||||
panic(err)
|
||||
}
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set bits on various rows.
|
||||
f.mustSetBits(100, 1, 2, 3)
|
||||
|
|
@ -915,7 +916,7 @@ func TestFragment_TopN_CacheSize(t *testing.T) {
|
|||
// Ensure fragment can return a checksum for its blocks.
|
||||
func TestFragment_Checksum(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Retrieve checksum and set bits.
|
||||
orig := f.Checksum()
|
||||
|
|
@ -934,7 +935,7 @@ func TestFragment_Checksum(t *testing.T) {
|
|||
// Ensure fragment can return a checksum for a given block.
|
||||
func TestFragment_Blocks(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Retrieve initial checksum.
|
||||
var prev []FragmentBlock
|
||||
|
|
@ -972,7 +973,7 @@ func TestFragment_Blocks(t *testing.T) {
|
|||
// Ensure fragment returns an empty checksum if no data exists for a block.
|
||||
func TestFragment_Blocks_Empty(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set bits on a different block.
|
||||
if _, err := f.setBit(100, 1); err != nil {
|
||||
|
|
@ -990,7 +991,7 @@ func TestFragment_Blocks_Empty(t *testing.T) {
|
|||
// Ensure a fragment's cache can be persisted between restarts.
|
||||
func TestFragment_LRUCache_Persistence(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, CacheTypeLRU)
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set bits on the fragment.
|
||||
for i := uint64(0); i < 1000; i++ {
|
||||
|
|
@ -1075,7 +1076,7 @@ func TestFragment_RankCache_Persistence(t *testing.T) {
|
|||
// Ensure a fragment can be copied to another fragment.
|
||||
func TestFragment_WriteTo_ReadFrom(t *testing.T) {
|
||||
f0 := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f0.Close()
|
||||
defer f0.Clean()
|
||||
|
||||
// Set and then clear bits on the fragment.
|
||||
if _, err := f0.setBit(1000, 1); err != nil {
|
||||
|
|
@ -1136,7 +1137,7 @@ func BenchmarkFragment_Blocks(b *testing.B) {
|
|||
if err := f.Open(); err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Reset timer and execute benchmark.
|
||||
b.ResetTimer()
|
||||
|
|
@ -1149,7 +1150,7 @@ func BenchmarkFragment_Blocks(b *testing.B) {
|
|||
|
||||
func BenchmarkFragment_IntersectionCount(b *testing.B) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
f.MaxOpN = math.MaxInt32
|
||||
|
||||
// Generate some intersecting data.
|
||||
|
|
@ -1180,7 +1181,7 @@ func BenchmarkFragment_IntersectionCount(b *testing.B) {
|
|||
|
||||
func TestFragment_Tanimoto(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, CacheTypeRanked)
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
src := NewRow(1, 2, 3)
|
||||
|
||||
|
|
@ -1203,7 +1204,7 @@ func TestFragment_Tanimoto(t *testing.T) {
|
|||
|
||||
func TestFragment_Zero_Tanimoto(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, CacheTypeRanked)
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
src := NewRow(1, 2, 3)
|
||||
|
||||
|
|
@ -1228,7 +1229,7 @@ func TestFragment_Zero_Tanimoto(t *testing.T) {
|
|||
|
||||
func TestFragment_Snapshot_Run(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set bits on the fragment.
|
||||
for i := uint64(1); i < 3; i++ {
|
||||
|
|
@ -1255,7 +1256,7 @@ func TestFragment_Snapshot_Run(t *testing.T) {
|
|||
// Ensure a fragment can set mutually exclusive values.
|
||||
func TestFragment_SetMutex(t *testing.T) {
|
||||
f := mustOpenMutexFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
var cols []uint64
|
||||
|
||||
|
|
@ -1369,7 +1370,7 @@ func TestFragment_ImportSet(t *testing.T) {
|
|||
for i, test := range tests {
|
||||
t.Run(fmt.Sprintf("importset%d", i), func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set import.
|
||||
err := f.bulkImport(test.setRowIDs, test.setColIDs, &ImportOptions{})
|
||||
|
|
@ -1405,7 +1406,7 @@ func TestFragment_ImportSet(t *testing.T) {
|
|||
func TestFragment_ConcurrentImport(t *testing.T) {
|
||||
t.Run("bulkImportStandard", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
eg := errgroup.Group{}
|
||||
eg.Go(func() error { return f.bulkImportStandard([]uint64{1, 2}, []uint64{1, 2}, &ImportOptions{}) })
|
||||
|
|
@ -1502,7 +1503,7 @@ func TestFragment_ImportMutex(t *testing.T) {
|
|||
for i, test := range tests {
|
||||
t.Run(fmt.Sprintf("importmutex%d", i), func(t *testing.T) {
|
||||
f := mustOpenMutexFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set import.
|
||||
err := f.bulkImport(test.setRowIDs, test.setColIDs, &ImportOptions{})
|
||||
|
|
@ -1621,7 +1622,7 @@ func TestFragment_ImportBool(t *testing.T) {
|
|||
for i, test := range tests {
|
||||
t.Run(fmt.Sprintf("importmutex%d", i), func(t *testing.T) {
|
||||
f := mustOpenBoolFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
// Set import.
|
||||
err := f.bulkImport(test.setRowIDs, test.setColIDs, &ImportOptions{})
|
||||
|
|
@ -1665,7 +1666,7 @@ func BenchmarkFragment_Snapshot(b *testing.B) {
|
|||
if err := f.Open(); err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
b.ResetTimer()
|
||||
|
||||
// Reset timer and execute benchmark.
|
||||
|
|
@ -1681,7 +1682,7 @@ func BenchmarkFragment_Snapshot(b *testing.B) {
|
|||
|
||||
func BenchmarkFragment_FullSnapshot(b *testing.B) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
// Generate some intersecting data.
|
||||
maxX := 1048576 / 2
|
||||
sz := maxX
|
||||
|
|
@ -1719,7 +1720,7 @@ func BenchmarkFragment_FullSnapshot(b *testing.B) {
|
|||
|
||||
func BenchmarkFragment_Import(b *testing.B) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
maxX := 1048576 * 5 * 2
|
||||
sz := maxX
|
||||
rows := make([]uint64, sz)
|
||||
|
|
@ -1747,8 +1748,14 @@ func BenchmarkFragment_Import(b *testing.B) {
|
|||
}
|
||||
}
|
||||
|
||||
var (
|
||||
rowCases = []uint64{2, 50, 1000, 100000}
|
||||
colCases = []uint64{20, 1000, 50000, 500000}
|
||||
concurrencyCases = []int{2, 4, 8, 16}
|
||||
)
|
||||
|
||||
func BenchmarkImportRoaring(b *testing.B) {
|
||||
for _, numRows := range []uint64{10, 100, 1000, 10000, 100000} {
|
||||
for _, numRows := range rowCases {
|
||||
data := getZipfRowsSliceRoaring(numRows, 1)
|
||||
b.Logf("%dRows: %.2fMB\n", numRows, float64(len(data))/1024/1024)
|
||||
for _, cacheType := range []string{CacheTypeRanked} { // CacheTypeNone didn't seem to affect the results much
|
||||
|
|
@ -1759,10 +1766,11 @@ func BenchmarkImportRoaring(b *testing.B) {
|
|||
b.StartTimer()
|
||||
err := f.importRoaring(data, false)
|
||||
if err != nil {
|
||||
f.Clean()
|
||||
b.Fatalf("import error: %v", err)
|
||||
}
|
||||
b.StopTimer()
|
||||
f.Close()
|
||||
f.Clean()
|
||||
}
|
||||
})
|
||||
}
|
||||
|
|
@ -1770,10 +1778,10 @@ func BenchmarkImportRoaring(b *testing.B) {
|
|||
}
|
||||
|
||||
func BenchmarkImportRoaringConcurrent(b *testing.B) {
|
||||
for _, numRows := range []uint64{10, 100, 1000, 10000, 100000} {
|
||||
for _, numRows := range rowCases {
|
||||
data := getZipfRowsSliceRoaring(numRows, 1)
|
||||
b.Logf("%dRows: %.2fMB\n", numRows, float64(len(data))/1024/1024)
|
||||
for _, concurrency := range []int{2, 4, 8} {
|
||||
for _, concurrency := range concurrencyCases {
|
||||
b.Run(fmt.Sprintf("%dRows%dConcurrency", numRows, concurrency), func(b *testing.B) {
|
||||
b.StopTimer()
|
||||
frags := make([]*fragment, concurrency)
|
||||
|
|
@ -1791,22 +1799,57 @@ func BenchmarkImportRoaringConcurrent(b *testing.B) {
|
|||
}
|
||||
err := eg.Wait()
|
||||
if err != nil {
|
||||
b.Fatalf("importing fragment: %v", err)
|
||||
b.Errorf("importing fragment: %v", err)
|
||||
}
|
||||
b.StopTimer()
|
||||
for j := 0; j < concurrency; j++ {
|
||||
frags[j].Close()
|
||||
frags[j].Clean()
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
func BenchmarkImportRoaringUpdateConcurrent(b *testing.B) {
|
||||
for _, numRows := range rowCases {
|
||||
for _, numCols := range colCases {
|
||||
data := getZipfRowsSliceRoaring(numRows, 1)
|
||||
updata := getUpdataRoaring(numRows, numCols, 1)
|
||||
for _, concurrency := range concurrencyCases {
|
||||
b.Run(fmt.Sprintf("%dRows%dCols%dConcurrency", numRows, numCols, concurrency), func(b *testing.B) {
|
||||
b.StopTimer()
|
||||
frags := make([]*fragment, concurrency)
|
||||
for i := 0; i < b.N; i++ {
|
||||
for j := 0; j < concurrency; j++ {
|
||||
frags[j] = mustOpenFragment("i", "f", viewStandard, uint64(j), CacheTypeRanked)
|
||||
frags[j].importRoaring(data, false)
|
||||
}
|
||||
eg := errgroup.Group{}
|
||||
b.StartTimer()
|
||||
for j := 0; j < concurrency; j++ {
|
||||
j := j
|
||||
eg.Go(func() error {
|
||||
return frags[j].importRoaring(updata, false)
|
||||
})
|
||||
}
|
||||
err := eg.Wait()
|
||||
if err != nil {
|
||||
b.Errorf("importing fragment: %v", err)
|
||||
}
|
||||
b.StopTimer()
|
||||
for j := 0; j < concurrency; j++ {
|
||||
frags[j].Clean()
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func BenchmarkImportStandard(b *testing.B) {
|
||||
for _, cacheType := range []string{CacheTypeRanked} {
|
||||
for _, numRows := range []uint64{10, 100, 1000, 10000, 100000} {
|
||||
for _, numRows := range rowCases {
|
||||
rowIDs, columnIDs := getZipfRowsSliceStandard(numRows, 1)
|
||||
b.Run(fmt.Sprintf("Rows%dCache_%s", numRows, cacheType), func(b *testing.B) {
|
||||
b.StopTimer()
|
||||
|
|
@ -1815,10 +1858,10 @@ func BenchmarkImportStandard(b *testing.B) {
|
|||
b.StartTimer()
|
||||
err := f.bulkImport(rowIDs, columnIDs, &ImportOptions{})
|
||||
if err != nil {
|
||||
b.Fatalf("import error: %v", err)
|
||||
b.Errorf("import error: %v", err)
|
||||
}
|
||||
b.StopTimer()
|
||||
f.Close()
|
||||
f.Clean()
|
||||
}
|
||||
})
|
||||
}
|
||||
|
|
@ -1829,8 +1872,8 @@ func BenchmarkImportRoaringUpdate(b *testing.B) {
|
|||
fileSize := make(map[string]int64)
|
||||
names := []string{}
|
||||
for _, cacheType := range []string{CacheTypeRanked} {
|
||||
for _, numRows := range []uint64{10, 100, 1000, 10000, 100000} {
|
||||
for _, numCols := range []uint64{20, 1000, 50000, 500000} {
|
||||
for _, numRows := range rowCases {
|
||||
for _, numCols := range colCases {
|
||||
data := getZipfRowsSliceRoaring(numRows, 1)
|
||||
updata := getUpdataRoaring(numRows, numCols, 1)
|
||||
name := fmt.Sprintf("%s%dRows%dCols", cacheType, numRows, numCols)
|
||||
|
|
@ -1841,17 +1884,18 @@ func BenchmarkImportRoaringUpdate(b *testing.B) {
|
|||
f := mustOpenFragment("i", fmt.Sprintf("r%dc%s", numRows, cacheType), viewStandard, 0, cacheType)
|
||||
err := f.importRoaring(data, false)
|
||||
if err != nil {
|
||||
b.Fatalf("import error: %v", err)
|
||||
b.Errorf("import error: %v", err)
|
||||
}
|
||||
b.StartTimer()
|
||||
err = f.importRoaring(updata, false)
|
||||
if err != nil {
|
||||
b.Fatalf("import error: %v", err)
|
||||
f.Clean()
|
||||
b.Errorf("import error: %v", err)
|
||||
}
|
||||
b.StopTimer()
|
||||
stat, _ := f.file.Stat()
|
||||
fileSize[name] = stat.Size()
|
||||
f.Close()
|
||||
f.Clean()
|
||||
}
|
||||
})
|
||||
|
||||
|
|
@ -1936,7 +1980,7 @@ func getZipfRowsSliceStandard(numRows uint64, seed int64) (rowIDs, columnIDs []u
|
|||
}
|
||||
|
||||
func BenchmarkFileWrite(b *testing.B) {
|
||||
for _, numRows := range []uint64{10, 100, 1000, 10000, 100000} {
|
||||
for _, numRows := range rowCases {
|
||||
data := getZipfRowsSliceRoaring(numRows, 1)
|
||||
b.Run(fmt.Sprintf("Rows%d", numRows), func(b *testing.B) {
|
||||
b.StopTimer()
|
||||
|
|
@ -1959,6 +2003,7 @@ func BenchmarkFileWrite(b *testing.B) {
|
|||
b.Fatal(err)
|
||||
}
|
||||
b.StopTimer()
|
||||
os.Remove(f.Name())
|
||||
}
|
||||
})
|
||||
}
|
||||
|
|
@ -1967,6 +2012,18 @@ func BenchmarkFileWrite(b *testing.B) {
|
|||
|
||||
/////////////////////////////////////////////////////////////////////
|
||||
|
||||
func (f *fragment) Clean() error {
|
||||
errc := f.Close()
|
||||
errf := os.Remove(f.path)
|
||||
errp := os.Remove(f.cachePath())
|
||||
if errc != nil {
|
||||
return errc
|
||||
} else if errf != nil {
|
||||
return errf
|
||||
}
|
||||
return errp
|
||||
}
|
||||
|
||||
// mustOpenFragment returns a new instance of Fragment with a temporary path.
|
||||
func mustOpenFragment(index, field, view string, shard uint64, cacheType string) *fragment {
|
||||
file, err := ioutil.TempFile("", "pilosa-fragment-")
|
||||
|
|
@ -2030,7 +2087,7 @@ func (f *fragment) mustSetBits(rowID uint64, columnIDs ...uint64) {
|
|||
func TestFragment_RowsIteration(t *testing.T) {
|
||||
t.Run("firstContainer", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
expectedAll := make([]uint64, 0)
|
||||
expectedOdd := make([]uint64, 0)
|
||||
|
|
@ -2057,7 +2114,7 @@ func TestFragment_RowsIteration(t *testing.T) {
|
|||
|
||||
t.Run("secondRow", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
expected := []uint64{1, 2}
|
||||
if _, err := f.setBit(1, 66000); err != nil {
|
||||
|
|
@ -2081,7 +2138,7 @@ func TestFragment_RowsIteration(t *testing.T) {
|
|||
|
||||
t.Run("combinations", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
expectedRows := make([]uint64, 0)
|
||||
for r := uint64(1); r < uint64(10000); r += 100 {
|
||||
|
|
@ -2128,7 +2185,7 @@ func TestFragment_RoaringImport(t *testing.T) {
|
|||
for i, test := range tests {
|
||||
t.Run(fmt.Sprintf("importroaring%d", i), func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, "")
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
for num, input := range test {
|
||||
buf := &bytes.Buffer{}
|
||||
bm := roaring.NewBitmap(input...)
|
||||
|
|
@ -2171,7 +2228,7 @@ func TestFragment_RoaringImportTopN(t *testing.T) {
|
|||
for i, test := range tests {
|
||||
t.Run(fmt.Sprintf("importroaring%d", i), func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", viewStandard, 0, CacheTypeRanked)
|
||||
defer f.Close()
|
||||
defer f.Clean()
|
||||
|
||||
options := &ImportOptions{}
|
||||
err := f.bulkImport(test.rowIDs, test.colIDs, options)
|
||||
|
|
@ -2305,6 +2362,7 @@ func calcExpected(inputs ...[]uint64) [][]uint64 {
|
|||
func TestFragmentRowIterator(t *testing.T) {
|
||||
t.Run("basic", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", "v", 0, CacheTypeRanked)
|
||||
defer f.Clean()
|
||||
f.mustSetBits(0, 0)
|
||||
f.mustSetBits(1, 0)
|
||||
f.mustSetBits(2, 0)
|
||||
|
|
@ -2333,11 +2391,11 @@ func TestFragmentRowIterator(t *testing.T) {
|
|||
if !wrapped {
|
||||
t.Fatalf("wrapped should be true after iterator is exhausted")
|
||||
}
|
||||
f.Close()
|
||||
})
|
||||
|
||||
t.Run("skipped rows", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", "v", 0, CacheTypeRanked)
|
||||
defer f.Clean()
|
||||
f.mustSetBits(1, 0)
|
||||
f.mustSetBits(3, 0)
|
||||
f.mustSetBits(5, 0)
|
||||
|
|
@ -2366,11 +2424,11 @@ func TestFragmentRowIterator(t *testing.T) {
|
|||
if !wrapped {
|
||||
t.Fatalf("wrapped should be true after iterator is exhausted")
|
||||
}
|
||||
f.Close()
|
||||
})
|
||||
|
||||
t.Run("basic wrapped", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", "v", 0, CacheTypeRanked)
|
||||
defer f.Clean()
|
||||
f.mustSetBits(0, 0)
|
||||
f.mustSetBits(1, 0)
|
||||
f.mustSetBits(2, 0)
|
||||
|
|
@ -2391,11 +2449,11 @@ func TestFragmentRowIterator(t *testing.T) {
|
|||
t.Fatalf("got wrong columns back on iteration %d - should just be 0 but %v", i, row.Columns())
|
||||
}
|
||||
}
|
||||
f.Close()
|
||||
})
|
||||
|
||||
t.Run("skipped rows wrapped", func(t *testing.T) {
|
||||
f := mustOpenFragment("i", "f", "v", 0, CacheTypeRanked)
|
||||
defer f.Clean()
|
||||
f.mustSetBits(1, 0)
|
||||
f.mustSetBits(3, 0)
|
||||
f.mustSetBits(5, 0)
|
||||
|
|
@ -2416,7 +2474,5 @@ func TestFragmentRowIterator(t *testing.T) {
|
|||
t.Fatalf("got wrong columns back on iteration %d - should just be 0 but %v", i, row.Columns())
|
||||
}
|
||||
}
|
||||
f.Close()
|
||||
})
|
||||
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue