Merge pull request #27 from chrislusf/master

sync
This commit is contained in:
hilimd 2020-10-15 15:15:01 +08:00 committed by GitHub
commit 5c2e409ffe
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
4 changed files with 66 additions and 116 deletions

View file

@ -156,7 +156,7 @@ func (c *ChunkReadAt) doReadAt(p []byte, offset int64) (n int, err error) {
n += delta n += delta
} }
if err == nil && offset+int64(len(p)) > c.fileSize { if err == nil && offset+int64(len(p)) >= c.fileSize {
err = io.EOF err = io.EOF
} }
// fmt.Printf("~~~ filled %d, err: %v\n\n", n, err) // fmt.Printf("~~~ filled %d, err: %v\n\n", n, err)

View file

@ -66,9 +66,9 @@ func TestReaderAt(t *testing.T) {
chunkCache: &mockChunkCache{}, chunkCache: &mockChunkCache{},
} }
testReadAt(t, readerAt, 0, 10, 10, nil) testReadAt(t, readerAt, 0, 10, 10, io.EOF)
testReadAt(t, readerAt, 0, 12, 10, io.EOF) testReadAt(t, readerAt, 0, 12, 10, io.EOF)
testReadAt(t, readerAt, 2, 8, 8, nil) testReadAt(t, readerAt, 2, 8, 8, io.EOF)
testReadAt(t, readerAt, 3, 6, 6, nil) testReadAt(t, readerAt, 3, 6, 6, nil)
} }
@ -116,7 +116,7 @@ func TestReaderAt0(t *testing.T) {
chunkCache: &mockChunkCache{}, chunkCache: &mockChunkCache{},
} }
testReadAt(t, readerAt, 0, 10, 10, nil) testReadAt(t, readerAt, 0, 10, 10, io.EOF)
testReadAt(t, readerAt, 3, 16, 7, io.EOF) testReadAt(t, readerAt, 3, 16, 7, io.EOF)
testReadAt(t, readerAt, 3, 5, 5, nil) testReadAt(t, readerAt, 3, 5, 5, nil)
@ -144,7 +144,7 @@ func TestReaderAt1(t *testing.T) {
chunkCache: &mockChunkCache{}, chunkCache: &mockChunkCache{},
} }
testReadAt(t, readerAt, 0, 20, 20, nil) testReadAt(t, readerAt, 0, 20, 20, io.EOF)
testReadAt(t, readerAt, 1, 7, 7, nil) testReadAt(t, readerAt, 1, 7, 7, nil)
testReadAt(t, readerAt, 0, 1, 1, nil) testReadAt(t, readerAt, 0, 1, 1, nil)
testReadAt(t, readerAt, 18, 4, 2, io.EOF) testReadAt(t, readerAt, 18, 4, 2, io.EOF)

View file

@ -2,17 +2,17 @@ package filesys
import ( import (
"bytes" "bytes"
"io"
"sync"
"time"
"github.com/chrislusf/seaweedfs/weed/glog" "github.com/chrislusf/seaweedfs/weed/glog"
"github.com/chrislusf/seaweedfs/weed/pb/filer_pb" "github.com/chrislusf/seaweedfs/weed/pb/filer_pb"
"io"
"sync"
) )
type ContinuousDirtyPages struct { type ContinuousDirtyPages struct {
intervals *ContinuousIntervals intervals *ContinuousIntervals
f *File f *File
writeWaitGroup sync.WaitGroup
chunkSaveErrChan chan error
lock sync.Mutex lock sync.Mutex
collection string collection string
replication string replication string
@ -22,127 +22,82 @@ func newDirtyPages(file *File) *ContinuousDirtyPages {
return &ContinuousDirtyPages{ return &ContinuousDirtyPages{
intervals: &ContinuousIntervals{}, intervals: &ContinuousIntervals{},
f: file, f: file,
chunkSaveErrChan: make(chan error, 8),
} }
} }
var counter = int32(0) func (pages *ContinuousDirtyPages) AddPage(offset int64, data []byte) {
func (pages *ContinuousDirtyPages) AddPage(offset int64, data []byte) (chunks []*filer_pb.FileChunk, err error) {
glog.V(4).Infof("%s AddPage [%d,%d) of %d bytes", pages.f.fullpath(), offset, offset+int64(len(data)), pages.f.entry.Attributes.FileSize) glog.V(4).Infof("%s AddPage [%d,%d) of %d bytes", pages.f.fullpath(), offset, offset+int64(len(data)), pages.f.entry.Attributes.FileSize)
if len(data) > int(pages.f.wfs.option.ChunkSizeLimit) { if len(data) > int(pages.f.wfs.option.ChunkSizeLimit) {
// this is more than what buffer can hold. // this is more than what buffer can hold.
return pages.flushAndSave(offset, data) pages.flushAndSave(offset, data)
} }
pages.intervals.AddInterval(data, offset) pages.intervals.AddInterval(data, offset)
var chunk *filer_pb.FileChunk
var hasSavedData bool
if pages.intervals.TotalSize() > pages.f.wfs.option.ChunkSizeLimit { if pages.intervals.TotalSize() > pages.f.wfs.option.ChunkSizeLimit {
chunk, hasSavedData, err = pages.saveExistingLargestPageToStorage() pages.saveExistingLargestPageToStorage()
if hasSavedData {
chunks = append(chunks, chunk)
}
} }
return return
} }
func (pages *ContinuousDirtyPages) flushAndSave(offset int64, data []byte) (chunks []*filer_pb.FileChunk, err error) { func (pages *ContinuousDirtyPages) flushAndSave(offset int64, data []byte) {
var chunk *filer_pb.FileChunk
var newChunks []*filer_pb.FileChunk
// flush existing // flush existing
if newChunks, err = pages.saveExistingPagesToStorage(); err == nil { pages.saveExistingPagesToStorage()
if newChunks != nil {
chunks = append(chunks, newChunks...)
}
} else {
return
}
// flush the new page // flush the new page
if chunk, err = pages.saveToStorage(bytes.NewReader(data), offset, int64(len(data))); err == nil { pages.saveToStorage(bytes.NewReader(data), offset, int64(len(data)))
if chunk != nil {
glog.V(4).Infof("%s/%s flush big request [%d,%d) to %s", pages.f.dir.FullPath(), pages.f.Name, chunk.Offset, chunk.Offset+int64(chunk.Size), chunk.FileId)
chunks = append(chunks, chunk)
}
} else {
glog.V(0).Infof("%s/%s failed to flush2 [%d,%d): %v", pages.f.dir.FullPath(), pages.f.Name, chunk.Offset, chunk.Offset+int64(chunk.Size), err)
return
}
return return
} }
func (pages *ContinuousDirtyPages) saveExistingPagesToStorage() (chunks []*filer_pb.FileChunk, err error) { func (pages *ContinuousDirtyPages) saveExistingPagesToStorage() {
for pages.saveExistingLargestPageToStorage() {
var hasSavedData bool
var chunk *filer_pb.FileChunk
for {
chunk, hasSavedData, err = pages.saveExistingLargestPageToStorage()
if !hasSavedData {
return chunks, err
}
if err == nil {
if chunk != nil {
chunks = append(chunks, chunk)
}
} else {
return
} }
} }
} func (pages *ContinuousDirtyPages) saveExistingLargestPageToStorage() (hasSavedData bool) {
func (pages *ContinuousDirtyPages) saveExistingLargestPageToStorage() (chunk *filer_pb.FileChunk, hasSavedData bool, err error) {
maxList := pages.intervals.RemoveLargestIntervalLinkedList() maxList := pages.intervals.RemoveLargestIntervalLinkedList()
if maxList == nil { if maxList == nil {
return nil, false, nil return false
} }
fileSize := int64(pages.f.entry.Attributes.FileSize) fileSize := int64(pages.f.entry.Attributes.FileSize)
for {
chunkSize := min(maxList.Size(), fileSize-maxList.Offset()) chunkSize := min(maxList.Size(), fileSize-maxList.Offset())
if chunkSize == 0 { if chunkSize == 0 {
return return false
}
chunk, err = pages.saveToStorage(maxList.ToReader(), maxList.Offset(), chunkSize)
if err == nil {
if chunk != nil {
hasSavedData = true
}
glog.V(4).Infof("saveToStorage %s %s [%d,%d) of %d bytes", pages.f.fullpath(), chunk.GetFileIdString(), maxList.Offset(), maxList.Offset()+chunkSize, fileSize)
return
} else {
glog.V(0).Infof("%s saveToStorage [%d,%d): %v", pages.f.fullpath(), maxList.Offset(), maxList.Offset()+chunkSize, err)
time.Sleep(5 * time.Second)
}
} }
pages.saveToStorage(maxList.ToReader(), maxList.Offset(), chunkSize)
return true
} }
func (pages *ContinuousDirtyPages) saveToStorage(reader io.Reader, offset int64, size int64) (*filer_pb.FileChunk, error) { func (pages *ContinuousDirtyPages) saveToStorage(reader io.Reader, offset int64, size int64) {
pages.writeWaitGroup.Add(1)
go func() {
defer pages.writeWaitGroup.Done()
dir, _ := pages.f.fullpath().DirAndName() dir, _ := pages.f.fullpath().DirAndName()
reader = io.LimitReader(reader, size) reader = io.LimitReader(reader, size)
chunk, collection, replication, err := pages.f.wfs.saveDataAsChunk(dir)(reader, pages.f.Name, offset) chunk, collection, replication, err := pages.f.wfs.saveDataAsChunk(dir)(reader, pages.f.Name, offset)
if err != nil { if err != nil {
return nil, err glog.V(0).Infof("%s saveToStorage [%d,%d): %v", pages.f.fullpath(), offset, offset+size, err)
pages.chunkSaveErrChan <- err
return
} }
pages.collection, pages.replication = collection, replication pages.collection, pages.replication = collection, replication
pages.f.addChunks([]*filer_pb.FileChunk{chunk})
return chunk, nil pages.chunkSaveErrChan <- nil
}()
} }
func maxUint64(x, y uint64) uint64 { func maxUint64(x, y uint64) uint64 {

View file

@ -148,11 +148,7 @@ func (fh *FileHandle) Write(ctx context.Context, req *fuse.WriteRequest, resp *f
fh.f.entry.Attributes.FileSize = uint64(max(req.Offset+int64(len(data)), int64(fh.f.entry.Attributes.FileSize))) fh.f.entry.Attributes.FileSize = uint64(max(req.Offset+int64(len(data)), int64(fh.f.entry.Attributes.FileSize)))
glog.V(4).Infof("%v write [%d,%d) %d", fh.f.fullpath(), req.Offset, req.Offset+int64(len(req.Data)), len(req.Data)) glog.V(4).Infof("%v write [%d,%d) %d", fh.f.fullpath(), req.Offset, req.Offset+int64(len(req.Data)), len(req.Data))
chunks, err := fh.dirtyPages.AddPage(req.Offset, data) fh.dirtyPages.AddPage(req.Offset, data)
if err != nil {
glog.Errorf("%v write fh %d: [%d,%d): %v", fh.f.fullpath(), fh.handle, req.Offset, req.Offset+int64(len(data)), err)
return fuse.EIO
}
resp.Size = len(data) resp.Size = len(data)
@ -162,12 +158,7 @@ func (fh *FileHandle) Write(ctx context.Context, req *fuse.WriteRequest, resp *f
fh.f.dirtyMetadata = true fh.f.dirtyMetadata = true
} }
if len(chunks) > 0 {
fh.f.addChunks(chunks)
fh.f.dirtyMetadata = true fh.f.dirtyMetadata = true
}
return nil return nil
} }
@ -204,20 +195,24 @@ func (fh *FileHandle) Flush(ctx context.Context, req *fuse.FlushRequest) error {
} }
func (fh *FileHandle) doFlush(ctx context.Context, header fuse.Header) error { func (fh *FileHandle) doFlush(ctx context.Context, header fuse.Header) error {
// fflush works at fh level // flush works at fh level
// send the data to the OS // send the data to the OS
glog.V(4).Infof("doFlush %s fh %d", fh.f.fullpath(), fh.handle) glog.V(4).Infof("doFlush %s fh %d", fh.f.fullpath(), fh.handle)
chunks, err := fh.dirtyPages.saveExistingPagesToStorage() fh.dirtyPages.saveExistingPagesToStorage()
if err != nil {
glog.Errorf("flush %s: %v", fh.f.fullpath(), err) var err error
return fuse.EIO go func() {
for t := range fh.dirtyPages.chunkSaveErrChan {
if t != nil {
err = t
} }
}
}()
fh.dirtyPages.writeWaitGroup.Wait()
if len(chunks) > 0 { if err != nil {
return err
fh.f.addChunks(chunks)
fh.f.dirtyMetadata = true
} }
if !fh.f.dirtyMetadata { if !fh.f.dirtyMetadata {