mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-05 14:09:24 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
aaba9ad659 | ||
|
|
eb058e931a | ||
|
|
d2cd6a6f6e | ||
|
|
50d9383c3b | ||
|
|
cbffbd9025 | ||
|
|
ab61386d29 | ||
|
|
4d06e18554 | ||
|
|
cb25d100b7 | ||
|
|
fd83009c05 | ||
|
|
fcf5b6cc93 | ||
|
|
1dc25218cd | ||
|
|
2c95008a1a | ||
|
|
98e11895de | ||
|
|
8f9d1c1e3c |
@@ -141,7 +141,7 @@ func (c *ChunkReadAt) doReadAt(p []byte, offset int64) (n int, err error) {
|
||||
if chunkStart >= chunkStop {
|
||||
continue
|
||||
}
|
||||
// glog.V(4).Infof("read [%d,%d), %d/%d chunk %s [%d,%d)", chunkStart, chunkStop, i, len(c.chunkViews), chunk.FileId, chunk.LogicOffset-chunk.Offset, chunk.LogicOffset-chunk.Offset+int64(chunk.Size))
|
||||
glog.V(4).Infof("read [%d,%d), %d/%d chunk %s [%d,%d)", chunkStart, chunkStop, i, len(c.chunkViews), chunk.FileId, chunk.LogicOffset-chunk.Offset, chunk.LogicOffset-chunk.Offset+int64(chunk.Size))
|
||||
var buffer []byte
|
||||
bufferOffset := chunkStart - chunk.LogicOffset + chunk.Offset
|
||||
bufferLength := chunkStop - chunkStart
|
||||
@@ -156,7 +156,7 @@ func (c *ChunkReadAt) doReadAt(p []byte, offset int64) (n int, err error) {
|
||||
startOffset, remaining = startOffset+int64(copied), remaining-int64(copied)
|
||||
}
|
||||
|
||||
// glog.V(4).Infof("doReadAt [%d,%d), n:%v, err:%v", offset, offset+int64(len(p)), n, err)
|
||||
glog.V(4).Infof("doReadAt [%d,%d), n:%v, err:%v", offset, offset+int64(len(p)), n, err)
|
||||
|
||||
if err == nil && remaining > 0 && c.fileSize > startOffset {
|
||||
delta := int(min(remaining, c.fileSize-startOffset))
|
||||
|
||||
@@ -52,7 +52,6 @@ func (pages *StreamDirtyPages) FlushData() error {
|
||||
if pages.lastErr != nil {
|
||||
return fmt.Errorf("flush data: %v", pages.lastErr)
|
||||
}
|
||||
pages.chunkedStream.Reset()
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -102,5 +101,5 @@ func (pages *StreamDirtyPages) saveChunkedFileIntevalToStorage(reader io.Reader,
|
||||
}
|
||||
|
||||
func (pages StreamDirtyPages) Destroy() {
|
||||
pages.chunkedStream.Reset()
|
||||
pages.chunkedStream.Destroy()
|
||||
}
|
||||
|
||||
@@ -51,7 +51,6 @@ func (pages *TempFileDirtyPages) FlushData() error {
|
||||
if pages.lastErr != nil {
|
||||
return fmt.Errorf("flush data: %v", pages.lastErr)
|
||||
}
|
||||
pages.chunkedFile.Reset()
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -68,6 +67,7 @@ func (pages *TempFileDirtyPages) saveChunkedFileToStorage() {
|
||||
pages.chunkedFile.ProcessEachInterval(func(file *os.File, logicChunkIndex page_writer.LogicChunkIndex, interval *page_writer.ChunkWrittenInterval) {
|
||||
reader := page_writer.NewFileIntervalReader(pages.chunkedFile, logicChunkIndex, interval)
|
||||
pages.saveChunkedFileIntevalToStorage(reader, int64(logicChunkIndex)*pages.chunkedFile.ChunkSize+interval.StartOffset, interval.Size())
|
||||
interval.MarkFlushed()
|
||||
})
|
||||
|
||||
}
|
||||
@@ -102,5 +102,5 @@ func (pages *TempFileDirtyPages) saveChunkedFileIntevalToStorage(reader io.Reade
|
||||
}
|
||||
|
||||
func (pages TempFileDirtyPages) Destroy() {
|
||||
pages.chunkedFile.Reset()
|
||||
pages.chunkedFile.Destroy()
|
||||
}
|
||||
|
||||
@@ -3,6 +3,7 @@ package filesys
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"github.com/chrislusf/seaweedfs/weed/filesys/page_writer"
|
||||
"io"
|
||||
"math"
|
||||
"net/http"
|
||||
@@ -161,7 +162,8 @@ func (fh *FileHandle) readFromChunks(buff []byte, offset int64) (int64, error) {
|
||||
glog.Errorf("file handle read %s: %v", fileFullPath, err)
|
||||
}
|
||||
|
||||
// glog.V(4).Infof("file handle read %s [%d,%d] %d : %v", fileFullPath, offset, offset+int64(totalRead), totalRead, err)
|
||||
glog.V(4).Infof("file handle read %s [%d,%d) %d : %v", fileFullPath, offset, offset+int64(totalRead), totalRead, err)
|
||||
page_writer.CheckByteZero(fmt.Sprintf("read %d chunks", len(entry.Chunks)), buff, 0, int64(totalRead))
|
||||
|
||||
return int64(totalRead), err
|
||||
}
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package filesys
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"github.com/chrislusf/seaweedfs/weed/filesys/page_writer"
|
||||
"github.com/chrislusf/seaweedfs/weed/glog"
|
||||
)
|
||||
@@ -69,6 +70,9 @@ func (pw *PageWriter) FlushData() error {
|
||||
func (pw *PageWriter) ReadDirtyDataAt(data []byte, offset int64) (maxStop int64) {
|
||||
glog.V(4).Infof("ReadDirtyDataAt %v [%d, %d)", pw.f.fullpath(), offset, offset+int64(len(data)))
|
||||
|
||||
originalData := data
|
||||
originalOffset := offset
|
||||
|
||||
chunkIndex := offset / pw.chunkSize
|
||||
for i := chunkIndex; len(data) > 0; i++ {
|
||||
readSize := min(int64(len(data)), (i+1)*pw.chunkSize-offset)
|
||||
@@ -84,6 +88,8 @@ func (pw *PageWriter) ReadDirtyDataAt(data []byte, offset int64) (maxStop int64)
|
||||
data = data[readSize:]
|
||||
}
|
||||
|
||||
page_writer.CheckByteZero(fmt.Sprintf("page writer read [%d,%d) of size %d", originalOffset, originalOffset+int64(len(originalData)), pw.f.entry.Attributes.FileSize), originalData, 0, maxStop-originalOffset)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -6,6 +6,7 @@ import "math"
|
||||
type ChunkWrittenInterval struct {
|
||||
StartOffset int64
|
||||
stopOffset int64
|
||||
flushed bool
|
||||
prev *ChunkWrittenInterval
|
||||
next *ChunkWrittenInterval
|
||||
}
|
||||
@@ -18,6 +19,10 @@ func (interval *ChunkWrittenInterval) isComplete(chunkSize int64) bool {
|
||||
return interval.stopOffset-interval.StartOffset == chunkSize
|
||||
}
|
||||
|
||||
func (interval *ChunkWrittenInterval) MarkFlushed() {
|
||||
interval.flushed = true
|
||||
}
|
||||
|
||||
// ChunkWrittenIntervalList mark written intervals within one page chunk
|
||||
type ChunkWrittenIntervalList struct {
|
||||
head *ChunkWrittenInterval
|
||||
@@ -64,18 +69,21 @@ func (list *ChunkWrittenIntervalList) addInterval(interval *ChunkWrittenInterval
|
||||
if interval.StartOffset <= p.stopOffset && q.StartOffset <= interval.stopOffset {
|
||||
// merge p and q together
|
||||
p.stopOffset = q.stopOffset
|
||||
p.flushed = false
|
||||
unlinkNodesBetween(p, q.next)
|
||||
return
|
||||
}
|
||||
if interval.StartOffset <= p.stopOffset {
|
||||
// merge new interval into p
|
||||
p.stopOffset = interval.stopOffset
|
||||
p.flushed = false
|
||||
unlinkNodesBetween(p, q)
|
||||
return
|
||||
}
|
||||
if q.StartOffset <= interval.stopOffset {
|
||||
// merge new interval into q
|
||||
q.StartOffset = interval.StartOffset
|
||||
q.flushed = false
|
||||
unlinkNodesBetween(p, q)
|
||||
return
|
||||
}
|
||||
|
||||
@@ -73,6 +73,7 @@ func (cw *ChunkedFileWriter) ReadDataAt(p []byte, off int64) (maxStop int64) {
|
||||
glog.Errorf("reading temp file: %v", err)
|
||||
break
|
||||
}
|
||||
CheckByteZero("temp file writer read", p, logicStart-off, logicStop-off)
|
||||
maxStop = max(maxStop, logicStop)
|
||||
}
|
||||
}
|
||||
@@ -106,13 +107,15 @@ func (cw *ChunkedFileWriter) ProcessEachInterval(process func(file *os.File, log
|
||||
for logicChunkIndex, actualChunkIndex := range cw.logicToActualChunkIndex {
|
||||
chunkUsage := cw.chunkUsages[actualChunkIndex]
|
||||
for t := chunkUsage.head.next; t != chunkUsage.tail; t = t.next {
|
||||
process(cw.file, logicChunkIndex, t)
|
||||
if !t.flushed {
|
||||
process(cw.file, logicChunkIndex, t)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Reset releases used resources
|
||||
func (cw *ChunkedFileWriter) Reset() {
|
||||
// Destroy releases used resources
|
||||
func (cw *ChunkedFileWriter) Destroy() {
|
||||
if cw.file != nil {
|
||||
cw.file.Close()
|
||||
os.Remove(cw.file.Name())
|
||||
|
||||
@@ -35,9 +35,9 @@ func writeToFile(cw *ChunkedFileWriter, startOffset int64, stopOffset int64) {
|
||||
|
||||
func TestWriteChunkedFile(t *testing.T) {
|
||||
x := NewChunkedFileWriter(os.TempDir(), 20)
|
||||
defer x.Reset()
|
||||
defer x.Destroy()
|
||||
y := NewChunkedFileWriter(os.TempDir(), 12)
|
||||
defer y.Reset()
|
||||
defer y.Destroy()
|
||||
|
||||
batchSize := 4
|
||||
buf := make([]byte, batchSize)
|
||||
|
||||
@@ -57,7 +57,6 @@ func (cw *ChunkedStreamWriter) WriteAt(p []byte, off int64) (n int, err error) {
|
||||
if memChunk.usage.IsComplete(cw.ChunkSize) {
|
||||
if cw.saveToStorageFn != nil {
|
||||
cw.saveOneChunk(memChunk, logicChunkIndex)
|
||||
delete(cw.activeChunks, logicChunkIndex)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -81,6 +80,9 @@ func (cw *ChunkedStreamWriter) ReadDataAt(p []byte, off int64) (maxStop int64) {
|
||||
if logicStart < logicStop {
|
||||
copy(p[logicStart-off:logicStop-off], memChunk.buf[logicStart-memChunkBaseOffset:logicStop-memChunkBaseOffset])
|
||||
maxStop = max(maxStop, logicStop)
|
||||
|
||||
CheckByteZero("stream writer read", p, logicStart-off, logicStop-off)
|
||||
|
||||
}
|
||||
}
|
||||
return
|
||||
@@ -92,7 +94,6 @@ func (cw *ChunkedStreamWriter) FlushAll() {
|
||||
for logicChunkIndex, memChunk := range cw.activeChunks {
|
||||
if cw.saveToStorageFn != nil {
|
||||
cw.saveOneChunk(memChunk, logicChunkIndex)
|
||||
delete(cw.activeChunks, logicChunkIndex)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -105,13 +106,17 @@ func (cw *ChunkedStreamWriter) saveOneChunk(memChunk *MemChunk, logicChunkIndex
|
||||
atomic.AddInt32(&referenceCounter, -1)
|
||||
if atomic.LoadInt32(&referenceCounter) == 0 {
|
||||
mem.Free(memChunk.buf)
|
||||
delete(cw.activeChunks, logicChunkIndex)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// Reset releases used resources
|
||||
func (cw *ChunkedStreamWriter) Reset() {
|
||||
// Destroy releases used resources
|
||||
func (cw *ChunkedStreamWriter) Destroy() {
|
||||
cw.Lock()
|
||||
defer cw.Unlock()
|
||||
|
||||
for t, memChunk := range cw.activeChunks {
|
||||
mem.Free(memChunk.buf)
|
||||
delete(cw.activeChunks, t)
|
||||
|
||||
@@ -8,9 +8,9 @@ import (
|
||||
|
||||
func TestWriteChunkedStream(t *testing.T) {
|
||||
x := NewChunkedStreamWriter(20)
|
||||
defer x.Reset()
|
||||
defer x.Destroy()
|
||||
y := NewChunkedFileWriter(os.TempDir(), 12)
|
||||
defer y.Reset()
|
||||
defer y.Destroy()
|
||||
|
||||
batchSize := 4
|
||||
buf := make([]byte, batchSize)
|
||||
|
||||
@@ -0,0 +1,21 @@
|
||||
package page_writer
|
||||
|
||||
import "github.com/chrislusf/seaweedfs/weed/glog"
|
||||
|
||||
func CheckByteZero(message string, p []byte, start, stop int64) {
|
||||
isAllZero := true
|
||||
for i := start; i < stop; i++ {
|
||||
if p[i] != 0 {
|
||||
isAllZero = false
|
||||
break
|
||||
}
|
||||
}
|
||||
if isAllZero {
|
||||
if start != stop {
|
||||
glog.Errorf("%s is all zeros [%d,%d)", message, start, stop)
|
||||
}
|
||||
} else {
|
||||
glog.V(4).Infof("%s read some non-zero data [%d,%d)", message, start, stop)
|
||||
}
|
||||
|
||||
}
|
||||
Reference in New Issue
Block a user