Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -166,7 +166,7 @@ This will compile and link the C++ characterization executable.
## Build Instructions (Go)

### Dependencies
* Go 1.24.9 or later is required to compile the Go code.
* Go 1.25.0 or later is required to compile the Go code.

### Build
* The project uses Go modules, so you can build the project by running the following command:
Expand Down
126 changes: 126 additions & 0 deletions go/bloom_filter_accuracy_profile.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,126 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package main

import (
"fmt"
"math"
"strings"

"github.com/apache/datasketches-go/filters"
)

type BloomFilterAccuracyProfile struct {
config filterJobConfig

sketch filters.BloomFilter
filterLengthBits uint64
numItemsInserted uint64
vIn uint64
}

func MustNewBloomFilterAccuracyProfile(cfg filterJobConfig) *BloomFilterAccuracyProfile {
return &BloomFilterAccuracyProfile{
config: cfg,
numItemsInserted: cfg.numItemsInserted(),
vIn: 1,
}
}

func (p *BloomFilterAccuracyProfile) run() {
fmt.Println(p.getHeader())

numQueries := uint64(1) << (p.config.minNumHashes + 1)

sb := &strings.Builder{}
for nh := p.config.minNumHashes; nh <= p.config.maxNumHashes; nh++ {
fpr := 0.0
filterNumBits := uint64(0)

numTrials := p.config.getNumTrials(nh)
for t := 0; t < numTrials; t++ {
fpr += p.doTrial(nh, numQueries)
filterNumBits += p.getFilterLengthBits()
}
fpr /= float64(numTrials)
filterNumBits /= uint64(numTrials)

p.process(nh, fpr, filterNumBits, numQueries, numTrials, sb)
fmt.Println(sb.String())

numQueries = pwr2SeriesNext(p.config.tppo, uint64(1)<<(nh+1))
}
}

func (p *BloomFilterAccuracyProfile) doTrial(numHashes int, numQueries uint64) float64 {
p.filterLengthBits = uint64(float64(uint64(numHashes)*p.numItemsInserted) / math.Ln2)

sketch, err := filters.NewBloomFilterBySize(p.filterLengthBits, uint16(numHashes))
if err != nil {
panic(err)
}
p.sketch = sketch

for i := uint64(0); i < p.numItemsInserted; i++ {
p.vIn++
if err := p.sketch.UpdateUInt64(p.vIn); err != nil {
panic(err)
}
}

numFalsePositive := uint64(0)
for i := uint64(0); i < numQueries; i++ {
p.vIn++
if p.sketch.QueryUInt64(p.vIn) {
numFalsePositive++
}
}
return float64(numFalsePositive) / float64(numQueries)
}

func (p *BloomFilterAccuracyProfile) getFilterLengthBits() uint64 {
return p.sketch.Capacity()
}

func (p *BloomFilterAccuracyProfile) getBitsPerEntry(numHashes int) int {
return int(float64(numHashes) / math.Ln2)
}

func (p *BloomFilterAccuracyProfile) getHeader() string {
return strings.Join([]string{
"numHashes",
"FPR",
"filterSizeBits",
"numQueryPoints",
"numTrials",
}, "\t")
}

func (p *BloomFilterAccuracyProfile) process(numHashes int, falsePositiveRate float64,
filterSizeBits, numQueryPoints uint64, numTrials int, sb *strings.Builder) {
sb.Reset()
sb.WriteString(fmt.Sprintf("%d", numHashes))
sb.WriteString("\t")
sb.WriteString(fmt.Sprintf("%.5e", falsePositiveRate))
sb.WriteString("\t")
sb.WriteString(fmt.Sprintf("%d", filterSizeBits))
sb.WriteString("\t")
sb.WriteString(fmt.Sprintf("%d", numQueryPoints))
sb.WriteString("\t")
sb.WriteString(fmt.Sprintf("%d", numTrials))
}
125 changes: 125 additions & 0 deletions go/bloom_filter_space_profile.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,125 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package main

import (
"fmt"
"strconv"
"strings"

"github.com/apache/datasketches-go/filters"
)

type BloomFilterSpaceProfile struct {
config filterSpaceJobConfig

vIn uint64
}

func MustNewBloomFilterSpaceProfile(cfg filterSpaceJobConfig) *BloomFilterSpaceProfile {
return &BloomFilterSpaceProfile{config: cfg}
}

type spaceTrialResults struct {
filterSizeBits uint64
measuredFPR float64
numHashes uint16
}

func (p *BloomFilterSpaceProfile) run() {
fmt.Println(p.getHeader())

maxU := uint64(1) << p.config.lgMaxU
inputCardinality := pwr2SeriesNext(p.config.uppo, uint64(1)<<p.config.lgMinU)

sb := &strings.Builder{}
for inputCardinality < maxU {
numTrials := p.config.getNumTrials(inputCardinality)

inputCardinality = pwr2SeriesNext(p.config.uppo, inputCardinality)

res := p.doTrial(inputCardinality)

p.process(inputCardinality, res.filterSizeBits, numTrials, res.measuredFPR, res.numHashes, sb)
fmt.Println(sb.String())
}
}

func (p *BloomFilterSpaceProfile) doTrial(inputCardinality uint64) spaceTrialResults {
numBits := filters.SuggestNumFilterBits(inputCardinality, p.config.targetFpp)
suggestedNumHashes := filters.SuggestNumHashesFromSize(inputCardinality, numBits)

numHashes := int(suggestedNumHashes) + p.config.numHashesDelta
if numHashes < 1 {
panic(fmt.Sprintf("numHashesDelta %d drives the hash count to %d at cardinality %d; "+
"must stay >= 1", p.config.numHashesDelta, numHashes, inputCardinality))
}

sketch, err := filters.NewBloomFilterBySize(numBits, uint16(numHashes),
filters.WithSeed(p.config.seed))
if err != nil {
panic(err)
}

numQueries := pwr2SeriesNext(p.config.tppo, uint64(1)<<suggestedNumHashes)

for i := uint64(0); i < inputCardinality; i++ {
p.vIn++
if err := sketch.UpdateUInt64(p.vIn); err != nil {
panic(err)
}
}

numFalsePositive := uint64(0)
for i := uint64(0); i < numQueries; i++ {
p.vIn++
if sketch.QueryUInt64(p.vIn) {
numFalsePositive++
}
}

return spaceTrialResults{
filterSizeBits: sketch.Capacity(),
measuredFPR: float64(numFalsePositive) / float64(numQueries),
numHashes: sketch.NumHashes(),
}
}

func (p *BloomFilterSpaceProfile) getHeader() string {
return strings.Join([]string{
"TrueU",
"Size",
"NumTrials",
"FalsePositiveRate",
"NumHashBits",
}, "\t")
}

func (p *BloomFilterSpaceProfile) process(inputCardinality, sizeInBits uint64, numTrials int,
falsePositiveRate float64, numHashes uint16, sb *strings.Builder) {
sb.Reset()
sb.WriteString(fmt.Sprintf("%d", inputCardinality))
sb.WriteString("\t")
sb.WriteString(fmt.Sprintf("%d", sizeInBits))
sb.WriteString("\t")
sb.WriteString(fmt.Sprintf("%d", numTrials))
sb.WriteString("\t")
sb.WriteString(strconv.FormatFloat(falsePositiveRate, 'g', -1, 64))
sb.WriteString("\t")
sb.WriteString(fmt.Sprintf("%d", numHashes))
}
120 changes: 120 additions & 0 deletions go/bloom_filter_update_speed_profile.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,120 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package main

import (
"fmt"
"math"
"runtime"
"runtime/debug"
"strings"
"time"

"github.com/apache/datasketches-go/filters"
)

type BloomFilterUpdateSpeedProfile struct {
config filterSpeedJobConfig

sketch filters.BloomFilter
vIn uint64
}

func MustNewBloomFilterUpdateSpeedProfile(cfg filterSpeedJobConfig) *BloomFilterUpdateSpeedProfile {
sketch, err := filters.NewBloomFilterBySize(cfg.numBits, cfg.numHashes)
if err != nil {
panic(err)
}

return &BloomFilterUpdateSpeedProfile{
config: cfg,
sketch: sketch,
vIn: 1,
}
}

func (p *BloomFilterUpdateSpeedProfile) run() {
debug.SetMemoryLimit(math.MaxInt64)

fmt.Println(p.getHeader())

maxU := uint64(1) << p.config.lgMaxU
minU := uint64(1) << p.config.lgMinU

sb := &strings.Builder{}

limit := 0.9 * float64(maxU)
lastU := uint64(0)
for float64(lastU) < limit {
nextU := minU
if lastU != 0 {
nextU = pwr2SeriesNext(p.config.uppo, lastU)
}
lastU = nextU

trials := p.config.getNumTrials(nextU)

runtime.GC()

sumUpdateTimePerUNanoSec := 0.0
for t := 0; t < trials; t++ {
sumUpdateTimePerUNanoSec += p.doTrial(nextU)
}
meanUpdateTimePerUNanoSec := sumUpdateTimePerUNanoSec / float64(trials)

p.process(meanUpdateTimePerUNanoSec, trials, nextU, sb)
fmt.Println(sb.String())
}
}

func (p *BloomFilterUpdateSpeedProfile) doTrial(uPerTrial uint64) float64 {
if err := p.sketch.Reset(); err != nil {
panic(err)
}

start := time.Now()
for u := uPerTrial; u > 0; u-- {
p.vIn++
p.sketch.UpdateUInt64(p.vIn)
}
elapsed := time.Since(start)

return float64(elapsed.Nanoseconds()) / float64(uPerTrial)
}

func (p *BloomFilterUpdateSpeedProfile) getHeader() string {
cols := []string{"InU", "Trials", "nS/Set"}
if p.config.numSketches > 1 {
cols = append(cols, "nS/Sketch")
}
return strings.Join(cols, "\t")
}

func (p *BloomFilterUpdateSpeedProfile) process(meanUpdateTimePerSetNanoSec float64,
trials int, uPerTrial uint64, sb *strings.Builder) {
sb.Reset()
sb.WriteString(fmt.Sprintf("%d", uPerTrial))
sb.WriteString("\t")
sb.WriteString(fmt.Sprintf("%d", trials))
sb.WriteString("\t")
sb.WriteString(fmt.Sprintf("%e", meanUpdateTimePerSetNanoSec))
if p.config.numSketches > 1 {
sb.WriteString("\t")
sb.WriteString(fmt.Sprintf("%e", meanUpdateTimePerSetNanoSec/float64(p.config.numSketches)))
}
}
Loading