vendor: sync the lazy compaction PR

This commit is contained in:
rjl493456442 2019-03-12 18:24:17 +08:00
parent 228f080ebd
commit 2ac01b4238
4 changed files with 210 additions and 35 deletions

View file

@ -53,9 +53,18 @@ type session struct {
manifestWriter storage.Writer
manifestFd storage.FileDesc
stCompPtrs []internalKey // compaction pointers; need external synchronization
stVersion *version // current version
vmu sync.Mutex
stCompPtrs []internalKey // compaction pointers; need external synchronization
stVersion *version // current version
ntVersionId int64 // next version id to assign
refCh chan *vTask
relCh chan *vTask
deltaCh chan *vDelta
closeC chan struct{}
closeW sync.WaitGroup
vmu sync.Mutex
// Testing fields
fileRefCh chan chan map[int64]int // channel used to pass current reference stat
}
// Creates new initialized session instance.
@ -68,13 +77,21 @@ func newSession(stor storage.Storage, o *opt.Options) (s *session, err error) {
return
}
s = &session{
stor: newIStorage(stor),
storLock: storLock,
fileRef: make(map[int64]int),
stor: newIStorage(stor),
storLock: storLock,
fileRef: make(map[int64]int),
refCh: make(chan *vTask),
relCh: make(chan *vTask),
deltaCh: make(chan *vDelta),
fileRefCh: make(chan chan map[int64]int),
closeC: make(chan struct{}),
}
s.setOptions(o)
s.tops = newTableOps(s)
s.setVersion(newVersion(s))
s.closeW.Add(1)
go s.refLoop()
s.setVersion(nil, newVersion(s))
s.log("log@legend F·NumFile S·FileSize N·Entry C·BadEntry B·BadBlock Ke·KeyError D·DroppedEntry L·Level Q·SeqNum T·TimeElapsed")
return
}
@ -90,7 +107,11 @@ func (s *session) close() {
}
s.manifest = nil
s.manifestWriter = nil
s.setVersion(&version{s: s, closing: true})
s.setVersion(nil, &version{s: s, closing: true, id: s.ntVersionId})
// Close all background goroutines
close(s.closeC)
s.closeW.Wait()
}
// Release session lock.
@ -180,7 +201,7 @@ func (s *session) recover() (err error) {
}
s.manifestFd = fd
s.setVersion(staging.finish(false))
s.setVersion(rec, staging.finish(false))
s.setNextFileNum(rec.nextFileNum)
s.recordCommited(rec)
return nil
@ -203,7 +224,7 @@ func (s *session) commit(r *sessionRecord, trivial bool) (err error) {
// finally, apply new version if no error rise
if err == nil {
s.setVersion(nv)
s.setVersion(r, nv)
}
return

View file

@ -39,20 +39,155 @@ func (s *session) newTemp() storage.FileDesc {
return storage.FileDesc{Type: storage.TypeTemp, Num: num}
}
func (s *session) addFileRef(fd storage.FileDesc, ref int) int {
ref += s.fileRef[fd.Num]
// addFileRef adds file reference counter with specified file number and
// reference value; need external synchronization.
func (s *session) addFileRef(fnum int64, ref int) int {
ref += s.fileRef[fnum]
if ref > 0 {
s.fileRef[fd.Num] = ref
s.fileRef[fnum] = ref
} else if ref == 0 {
delete(s.fileRef, fd.Num)
delete(s.fileRef, fnum)
} else {
panic(fmt.Sprintf("negative ref: %v", fd))
panic(fmt.Sprintf("negative ref: %v", fnum))
}
return ref
}
// Session state.
const cachedTaskBound = 256
// vDelta indicates the change information between the next version
// and the currently specified version
type vDelta struct {
vid int64
added []int64
deleted []int64
}
// vTask defines a version task for either reference or release.
type vTask struct {
vid int64
files []tFiles
}
func (s *session) refLoop() {
var (
ref = make(map[int64][]tFiles)
deltas = make(map[int64]*vDelta)
referenced = make(map[int64]struct{})
released = make(map[int64]*vDelta)
next, last int64
)
// processTasks processes version tasks in strict order.
//
// If we want to use delta to reduce the cost of file references and dereferences,
// we must strictly follow the id of the version, otherwise some files that are
// being referenced will be deleted.
//
// In addition, some db operations (such as iterators) may cause a version to be
// referenced for a long time. In order to prevent such operations from blocking
// the entire processing queue, we will properly convert some of the version tasks
// into full file references and releases.
processTasks := func() {
// Make sure we don't cache too many version tasks.
for {
if last-next < cachedTaskBound {
break
}
if _, exist := released[next]; exist {
break
}
if _, exist := ref[next]; !exist {
break
}
for _, tt := range ref[next] {
for _, t := range tt {
s.addFileRef(t.fd.Num, 1)
}
}
referenced[next] = struct{}{}
delete(ref, next)
delete(deltas, next)
next += 1
}
// Use delta information to process all released versions.
for {
if d, exist := released[next]; exist {
if d != nil {
for _, t := range d.added {
s.addFileRef(t, 1)
}
for _, t := range d.deleted {
if s.addFileRef(t, -1) == 0 {
s.tops.remove(storage.FileDesc{Type: storage.TypeTable, Num: t})
}
}
}
delete(released, next)
next += 1
continue
}
return
}
}
for {
processTasks()
select {
case t := <-s.refCh:
if _, exist := ref[t.vid]; exist {
panic("duplicate reference request")
}
ref[t.vid] = t.files
if t.vid > last {
last = t.vid
}
case d := <-s.deltaCh:
if _, exist := ref[d.vid]; !exist {
if _, exist2 := referenced[d.vid]; !exist2 {
panic("invalid release request")
}
continue
}
deltas[d.vid] = d
case t := <-s.relCh:
if _, exist := referenced[t.vid]; exist {
for _, tt := range t.files {
for _, t := range tt {
if s.addFileRef(t.fd.Num, -1) == 0 {
s.tops.remove(t.fd)
}
}
}
delete(referenced, t.vid)
continue
}
if _, exist := ref[t.vid]; !exist {
panic("invalid release request")
}
released[t.vid] = deltas[t.vid]
delete(deltas, t.vid)
delete(ref, t.vid)
case r := <-s.fileRefCh:
ref := make(map[int64]int)
for f, c := range s.fileRef {
ref[f] = c
}
r <- ref
case <-s.closeC:
s.closeW.Done()
return
}
}
}
// Get current version. This will incr version ref, must call
// version.release (exactly once) after use.
func (s *session) version() *version {
@ -69,13 +204,30 @@ func (s *session) tLen(level int) int {
}
// Set current version to v.
func (s *session) setVersion(v *version) {
func (s *session) setVersion(r *sessionRecord, v *version) {
s.vmu.Lock()
defer s.vmu.Unlock()
// Hold by session. It is important to call this first before releasing
// current version, otherwise the still used files might get released.
v.incref()
if s.stVersion != nil {
if r != nil {
var (
added = make([]int64, 0, len(r.addedTables))
deleted = make([]int64, 0, len(r.deletedTables))
)
for _, t := range r.addedTables {
added = append(added, t.num)
}
for _, t := range r.deletedTables {
deleted = append(deleted, t.num)
}
select {
case s.deltaCh <- &vDelta{vid: s.stVersion.id, added: added, deleted: deleted}:
case <-v.s.closeC:
s.log("reference loop already exist")
}
}
// Release current version.
s.stVersion.releaseNB()
}

View file

@ -483,15 +483,15 @@ func (t *tOps) newIterator(f *tFile, slice *util.Range, ro *opt.ReadOptions) ite
// Removes table from persistent storage. It waits until
// no one use the the table.
func (t *tOps) remove(f *tFile) {
t.cache.Delete(0, uint64(f.fd.Num), func() {
if err := t.s.stor.Remove(f.fd); err != nil {
t.s.logf("table@remove removing @%d %q", f.fd.Num, err)
func (t *tOps) remove(fd storage.FileDesc) {
t.cache.Delete(0, uint64(fd.Num), func() {
if err := t.s.stor.Remove(fd); err != nil {
t.s.logf("table@remove removing @%d %q", fd.Num, err)
} else {
t.s.logf("table@remove removed @%d", f.fd.Num)
t.s.logf("table@remove removed @%d", fd.Num)
}
if t.evictRemoved && t.bcache != nil {
t.bcache.EvictNS(uint64(f.fd.Num))
t.bcache.EvictNS(uint64(fd.Num))
}
})
}

View file

@ -22,7 +22,8 @@ type tSet struct {
}
type version struct {
s *session
id int64 // unique monotonous increasing version id
s *session
levels []tFiles
@ -40,7 +41,9 @@ type version struct {
}
func newVersion(s *session) *version {
return &version{s: s}
nv := &version{s: s, id: s.ntVersionId}
s.ntVersionId += 1
return nv
}
func (v *version) incref() {
@ -50,11 +53,11 @@ func (v *version) incref() {
v.ref++
if v.ref == 1 {
// Incr file ref.
for _, tt := range v.levels {
for _, t := range tt {
v.s.addFileRef(t.fd, 1)
}
select {
case v.s.refCh <- &vTask{vid: v.id, files: v.levels}:
// We can use v.levels directly here since it is immutable.
case <-v.s.closeC:
v.s.log("reference loop already exist")
}
}
}
@ -67,12 +70,11 @@ func (v *version) releaseNB() {
panic("negative version ref")
}
for _, tt := range v.levels {
for _, t := range tt {
if v.s.addFileRef(t.fd, -1) == 0 {
v.s.tops.remove(t)
}
}
select {
case v.s.relCh <- &vTask{vid: v.id, files: v.levels}:
// We can use v.levels directly here since it is immutable.
case <-v.s.closeC:
v.s.log("reference loop already exist")
}
v.released = true