seaweedfs/weed/filesys/dirty_pages_mem_chunk.go

104 lines
2.8 KiB
Go
Raw Normal View History

2022-01-17 09:53:56 +00:00
package filesys
import (
"fmt"
"github.com/chrislusf/seaweedfs/weed/filesys/page_writer"
"github.com/chrislusf/seaweedfs/weed/glog"
"github.com/chrislusf/seaweedfs/weed/pb/filer_pb"
"io"
"sync"
"time"
)
type MemoryChunkPages struct {
2022-01-17 21:53:30 +00:00
fh *FileHandle
2022-01-17 09:53:56 +00:00
writeWaitGroup sync.WaitGroup
chunkAddLock sync.Mutex
lastErr error
collection string
replication string
uploadPipeline *page_writer.UploadPipeline
2022-01-17 11:19:00 +00:00
hasWrites bool
2022-01-17 09:53:56 +00:00
}
var (
_ = page_writer.DirtyPages(&MemoryChunkPages{})
)
2022-01-17 21:53:30 +00:00
func newMemoryChunkPages(fh *FileHandle, chunkSize int64) *MemoryChunkPages {
2022-01-17 09:53:56 +00:00
dirtyPages := &MemoryChunkPages{
2022-01-17 21:53:30 +00:00
fh: fh,
2022-01-17 09:53:56 +00:00
}
swapFileDir := fh.f.wfs.option.getTempFilePageDir()
2022-01-17 22:15:10 +00:00
dirtyPages.uploadPipeline = page_writer.NewUploadPipeline(fh.f.fullpath(),
fh.f.wfs.concurrentWriters, chunkSize, dirtyPages.saveChunkedFileIntevalToStorage, swapFileDir)
2022-01-17 09:53:56 +00:00
return dirtyPages
}
func (pages *MemoryChunkPages) AddPage(offset int64, data []byte) {
2022-01-17 11:19:00 +00:00
pages.hasWrites = true
2022-01-17 09:53:56 +00:00
2022-01-17 21:53:30 +00:00
glog.V(4).Infof("%v memory AddPage [%d, %d)", pages.fh.f.fullpath(), offset, offset+int64(len(data)))
2022-01-17 09:53:56 +00:00
pages.uploadPipeline.SaveDataAt(data, offset)
return
}
func (pages *MemoryChunkPages) FlushData() error {
2022-01-17 11:19:00 +00:00
if !pages.hasWrites {
return nil
}
2022-01-17 23:50:11 +00:00
pages.uploadPipeline.FlushAll()
2022-01-17 09:53:56 +00:00
if pages.lastErr != nil {
return fmt.Errorf("flush data: %v", pages.lastErr)
}
return nil
}
func (pages *MemoryChunkPages) ReadDirtyDataAt(data []byte, startOffset int64) (maxStop int64) {
2022-01-17 11:19:00 +00:00
if !pages.hasWrites {
return
}
2022-01-17 09:53:56 +00:00
return pages.uploadPipeline.MaybeReadDataAt(data, startOffset)
}
func (pages *MemoryChunkPages) GetStorageOptions() (collection, replication string) {
return pages.collection, pages.replication
}
func (pages *MemoryChunkPages) saveChunkedFileIntevalToStorage(reader io.Reader, offset int64, size int64, cleanupFn func()) {
mtime := time.Now().UnixNano()
2022-01-17 23:50:11 +00:00
defer cleanupFn()
2022-01-17 09:53:56 +00:00
2022-01-17 23:50:11 +00:00
chunk, collection, replication, err := pages.fh.f.wfs.saveDataAsChunk(pages.fh.f.fullpath())(reader, pages.fh.f.Name, offset)
if err != nil {
glog.V(0).Infof("%s saveToStorage [%d,%d): %v", pages.fh.f.fullpath(), offset, offset+size, err)
pages.lastErr = err
return
2022-01-17 09:53:56 +00:00
}
2022-01-17 23:50:11 +00:00
chunk.Mtime = mtime
pages.collection, pages.replication = collection, replication
pages.chunkAddLock.Lock()
pages.fh.f.addChunks([]*filer_pb.FileChunk{chunk})
pages.fh.entryViewCache = nil
glog.V(3).Infof("%s saveToStorage %s [%d,%d)", pages.fh.f.fullpath(), chunk.FileId, offset, offset+size)
pages.chunkAddLock.Unlock()
2022-01-17 09:53:56 +00:00
}
func (pages MemoryChunkPages) Destroy() {
pages.uploadPipeline.Shutdown()
}
func (pages *MemoryChunkPages) LockForRead(startOffset, stopOffset int64) {
pages.uploadPipeline.LockForRead(startOffset, stopOffset)
}
func (pages *MemoryChunkPages) UnlockForRead(startOffset, stopOffset int64) {
pages.uploadPipeline.UnlockForRead(startOffset, stopOffset)
}