2018-12-26 19:22:33 -08:00
|
|
|
|
/*
|
2025-12-10 17:11:18 -05:00
|
|
|
|
* SPDX-FileCopyrightText: © 2017-2025 Istari Digital, Inc.
|
2025-02-05 17:00:32 -05:00
|
|
|
|
* SPDX-License-Identifier: Apache-2.0
|
2018-12-26 19:22:33 -08:00
|
|
|
|
*/
|
|
|
|
|
|
|
|
|
|
|
|
package badger
|
|
|
|
|
|
|
|
|
|
|
|
import (
|
2024-10-14 14:11:59 +02:00
|
|
|
|
stderrors "errors"
|
2018-12-26 19:22:33 -08:00
|
|
|
|
"sync"
|
|
|
|
|
|
"time"
|
|
|
|
|
|
|
2026-04-11 00:01:08 -07:00
|
|
|
|
"github.com/luxfi/zapdb/y"
|
2024-10-25 22:03:34 +05:30
|
|
|
|
"github.com/dgraph-io/ristretto/v2/z"
|
2018-12-26 19:22:33 -08:00
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
// MergeOperator represents a Badger merge operator.
|
|
|
|
|
|
type MergeOperator struct {
|
|
|
|
|
|
sync.RWMutex
|
|
|
|
|
|
f MergeFunc
|
|
|
|
|
|
db *DB
|
|
|
|
|
|
key []byte
|
2020-09-07 23:20:46 +05:30
|
|
|
|
closer *z.Closer
|
2018-12-26 19:22:33 -08:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// MergeFunc accepts two byte slices, one representing an existing value, and
|
|
|
|
|
|
// another representing a new value that needs to be ‘merged’ into it. MergeFunc
|
|
|
|
|
|
// contains the logic to perform the ‘merge’ and return an updated value.
|
|
|
|
|
|
// MergeFunc could perform operations like integer addition, list appends etc.
|
2019-06-04 13:38:32 +05:30
|
|
|
|
// Note that the ordering of the operands is maintained.
|
|
|
|
|
|
type MergeFunc func(existingVal, newVal []byte) []byte
|
2018-12-26 19:22:33 -08:00
|
|
|
|
|
|
|
|
|
|
// GetMergeOperator creates a new MergeOperator for a given key and returns a
|
|
|
|
|
|
// pointer to it. It also fires off a goroutine that performs a compaction using
|
|
|
|
|
|
// the merge function that runs periodically, as specified by dur.
|
|
|
|
|
|
func (db *DB) GetMergeOperator(key []byte,
|
|
|
|
|
|
f MergeFunc, dur time.Duration) *MergeOperator {
|
|
|
|
|
|
op := &MergeOperator{
|
|
|
|
|
|
f: f,
|
|
|
|
|
|
db: db,
|
|
|
|
|
|
key: key,
|
2020-09-07 23:20:46 +05:30
|
|
|
|
closer: z.NewCloser(1),
|
2018-12-26 19:22:33 -08:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
go op.runCompactions(dur)
|
|
|
|
|
|
return op
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2024-10-14 14:11:59 +02:00
|
|
|
|
var errNoMerge = stderrors.New("No need for merge")
|
2018-12-26 19:22:33 -08:00
|
|
|
|
|
2019-06-04 13:38:32 +05:30
|
|
|
|
func (op *MergeOperator) iterateAndMerge() (newVal []byte, latest uint64, err error) {
|
2019-05-31 12:51:43 +05:30
|
|
|
|
txn := op.db.NewTransaction(false)
|
|
|
|
|
|
defer txn.Discard()
|
2018-12-26 19:22:33 -08:00
|
|
|
|
opt := DefaultIteratorOptions
|
|
|
|
|
|
opt.AllVersions = true
|
2019-03-07 23:43:24 +05:30
|
|
|
|
it := txn.NewKeyIterator(op.key, opt)
|
2018-12-26 19:22:33 -08:00
|
|
|
|
defer it.Close()
|
|
|
|
|
|
|
|
|
|
|
|
var numVersions int
|
2019-03-07 23:43:24 +05:30
|
|
|
|
for it.Rewind(); it.Valid(); it.Next() {
|
2018-12-26 19:22:33 -08:00
|
|
|
|
item := it.Item()
|
2021-03-03 13:41:05 +05:30
|
|
|
|
if item.IsDeletedOrExpired() {
|
|
|
|
|
|
break
|
|
|
|
|
|
}
|
2018-12-26 19:22:33 -08:00
|
|
|
|
numVersions++
|
|
|
|
|
|
if numVersions == 1 {
|
2019-06-04 13:38:32 +05:30
|
|
|
|
// This should be the newVal, considering this is the latest version.
|
|
|
|
|
|
newVal, err = item.ValueCopy(newVal)
|
2018-12-26 19:22:33 -08:00
|
|
|
|
if err != nil {
|
2019-05-31 12:51:43 +05:30
|
|
|
|
return nil, 0, err
|
2018-12-26 19:22:33 -08:00
|
|
|
|
}
|
2019-05-31 12:51:43 +05:30
|
|
|
|
latest = item.Version()
|
2018-12-26 19:22:33 -08:00
|
|
|
|
} else {
|
2019-06-04 13:38:32 +05:30
|
|
|
|
if err := item.Value(func(oldVal []byte) error {
|
|
|
|
|
|
// The merge should always be on the newVal considering it has the merge result of
|
|
|
|
|
|
// the latest version. The value read should be the oldVal.
|
|
|
|
|
|
newVal = op.f(oldVal, newVal)
|
2018-12-26 19:22:33 -08:00
|
|
|
|
return nil
|
|
|
|
|
|
}); err != nil {
|
2019-05-31 12:51:43 +05:30
|
|
|
|
return nil, 0, err
|
2018-12-26 19:22:33 -08:00
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
if item.DiscardEarlierVersions() {
|
|
|
|
|
|
break
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
if numVersions == 0 {
|
2019-05-31 12:51:43 +05:30
|
|
|
|
return nil, latest, ErrKeyNotFound
|
2018-12-26 19:22:33 -08:00
|
|
|
|
} else if numVersions == 1 {
|
2019-06-04 13:38:32 +05:30
|
|
|
|
return newVal, latest, errNoMerge
|
2018-12-26 19:22:33 -08:00
|
|
|
|
}
|
2019-06-04 13:38:32 +05:30
|
|
|
|
return newVal, latest, nil
|
2018-12-26 19:22:33 -08:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
func (op *MergeOperator) compact() error {
|
|
|
|
|
|
op.Lock()
|
|
|
|
|
|
defer op.Unlock()
|
2019-05-31 12:51:43 +05:30
|
|
|
|
val, version, err := op.iterateAndMerge()
|
2018-12-26 19:22:33 -08:00
|
|
|
|
if err == ErrKeyNotFound || err == errNoMerge {
|
2019-05-31 12:51:43 +05:30
|
|
|
|
return nil
|
2018-12-26 19:22:33 -08:00
|
|
|
|
} else if err != nil {
|
|
|
|
|
|
return err
|
|
|
|
|
|
}
|
2019-05-31 12:51:43 +05:30
|
|
|
|
entries := []*Entry{
|
|
|
|
|
|
{
|
|
|
|
|
|
Key: y.KeyWithTs(op.key, version),
|
|
|
|
|
|
Value: val,
|
|
|
|
|
|
meta: bitDiscardEarlierVersions,
|
|
|
|
|
|
},
|
|
|
|
|
|
}
|
|
|
|
|
|
// Write value back to the DB. It is important that we do not set the bitMergeEntry bit
|
|
|
|
|
|
// here. When compaction happens, all the older merged entries will be removed.
|
|
|
|
|
|
return op.db.batchSetAsync(entries, func(err error) {
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
op.db.opt.Errorf("failed to insert the result of merge compaction: %s", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
})
|
2018-12-26 19:22:33 -08:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
func (op *MergeOperator) runCompactions(dur time.Duration) {
|
|
|
|
|
|
ticker := time.NewTicker(dur)
|
|
|
|
|
|
defer op.closer.Done()
|
|
|
|
|
|
var stop bool
|
|
|
|
|
|
for {
|
|
|
|
|
|
select {
|
|
|
|
|
|
case <-op.closer.HasBeenClosed():
|
|
|
|
|
|
stop = true
|
|
|
|
|
|
case <-ticker.C: // wait for tick
|
|
|
|
|
|
}
|
|
|
|
|
|
if err := op.compact(); err != nil {
|
2019-01-02 16:07:41 -08:00
|
|
|
|
op.db.opt.Errorf("failure while running merge operation: %s", err)
|
2018-12-26 19:22:33 -08:00
|
|
|
|
}
|
|
|
|
|
|
if stop {
|
|
|
|
|
|
ticker.Stop()
|
|
|
|
|
|
break
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// Add records a value in Badger which will eventually be merged by a background
|
|
|
|
|
|
// routine into the values that were recorded by previous invocations to Add().
|
|
|
|
|
|
func (op *MergeOperator) Add(val []byte) error {
|
|
|
|
|
|
return op.db.Update(func(txn *Txn) error {
|
2019-05-28 10:58:37 +05:30
|
|
|
|
return txn.SetEntry(NewEntry(op.key, val).withMergeBit())
|
2018-12-26 19:22:33 -08:00
|
|
|
|
})
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// Get returns the latest value for the merge operator, which is derived by
|
|
|
|
|
|
// applying the merge function to all the values added so far.
|
|
|
|
|
|
//
|
|
|
|
|
|
// If Add has not been called even once, Get will return ErrKeyNotFound.
|
|
|
|
|
|
func (op *MergeOperator) Get() ([]byte, error) {
|
|
|
|
|
|
op.RLock()
|
|
|
|
|
|
defer op.RUnlock()
|
|
|
|
|
|
var existing []byte
|
|
|
|
|
|
err := op.db.View(func(txn *Txn) (err error) {
|
2019-05-31 12:51:43 +05:30
|
|
|
|
existing, _, err = op.iterateAndMerge()
|
2018-12-26 19:22:33 -08:00
|
|
|
|
return err
|
|
|
|
|
|
})
|
|
|
|
|
|
if err == errNoMerge {
|
|
|
|
|
|
return existing, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
return existing, err
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// Stop waits for any pending merge to complete and then stops the background
|
|
|
|
|
|
// goroutine.
|
|
|
|
|
|
func (op *MergeOperator) Stop() {
|
|
|
|
|
|
op.closer.SignalAndWait()
|
|
|
|
|
|
}
|