Compare commits

...
Author SHA1 Message Date
chrislu aaba9ad659 logs on non-zero data 2022-01-15 23:26:33 -08:00
chrislu eb058e931a add file size 2022-01-15 19:18:20 -08:00
chrislu d2cd6a6f6e adjust logs 2022-01-15 19:06:02 -08:00
chrislu 50d9383c3b adjust logging 2022-01-15 18:48:05 -08:00
chrislu cbffbd9025 fix 2022-01-15 15:06:36 -08:00
chrislu ab61386d29 fix checking 2022-01-15 15:01:36 -08:00
chrislu 4d06e18554 fix compilation 2022-01-15 14:51:15 -08:00
chrislu cb25d100b7 add debug zero 2022-01-15 14:49:59 -08:00
chrislu fd83009c05 add debug message when only zeros are copied 2022-01-15 14:39:23 -08:00
chrislu fcf5b6cc93 delete only when not used 2022-01-15 14:30:18 -08:00
chrislu 1dc25218cd delay deleting from memory unless the metadata chunks is updated 2022-01-15 07:40:29 -08:00
chrislu 2c95008a1a Revert "temp"
This reverts commit 98e11895de.
2022-01-15 06:45:51 -08:00
chrislu 98e11895de temp 2022-01-15 06:42:36 -08:00
chrislu 8f9d1c1e3c upload only not flushed chunks 2022-01-15 06:41:42 -08:00
11 changed files with 62 additions and 18 deletions
+2 -2
View File
@@ -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))
+1 -2
View File
@@ -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()
}
+2 -2
View File
@@ -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 -1
View File
@@ -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
}
+6
View File
@@ -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)
+21
View File
@@ -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)
}
}