This commit is contained in:
2026-06-12 14:21:14 +00:00
parent 58db58a0a0
commit bf22e9a85e
7 changed files with 168 additions and 147 deletions

View File

@@ -61,9 +61,8 @@ func (s *_metric) DeleteMeasures() {
} }
type CapturedState struct { type CapturedState struct {
Pages []storage.DataPayload LastTimestamp uint32
TimeState qb.TimeDeltaCapturedState LastValue float64
ValueState qb.ValueDeltaCapturedState
} }
// func (s *_metric) startAppendMeasures(req tryAppendMeasuresReq, sendToStorage func(storage.AppendedMeasures)) { // func (s *_metric) startAppendMeasures(req tryAppendMeasuresReq, sendToStorage func(storage.AppendedMeasures)) {
@@ -76,19 +75,24 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, send
values = s.values values = s.values
// офсети на head сторінці // офсети на head сторінці
timestampsOffset int timestampsOffset int
valuesOffset int timestampsRewindOffset int
valuesOffset int
valuesRewindOffset int
headTimestamps []byte headTimestamps []byte
headValues []byte headValues []byte
pages []storage.DataPayload
written int written int
resultCode byte resultCode byte
) )
s.capturedState = &CapturedState{ s.capturedState = &CapturedState{
// LastTimestamp: timestamps.LastTimestamp(),
LastValue: s.lastValue,
//LastValue: values.LastValue(),
} }
s.timestamps.CaptureState() s.timestamps.CaptureState()
@@ -111,20 +115,25 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, send
if totalRequiredSpace <= len(s.buffer) { if totalRequiredSpace <= len(s.buffer) {
// якщо на сторінці є місце // якщо на сторінці є місце
timestamps.Append(tReport.RewindOffset, tmp[:tReport.ChangeSize], measure.Timestamp)
values.Append(vReport.RewindOffset, tmp[7:7+vReport.ChangeSize], measure.Value, vReport.Delta)
if idx == 0 { if idx == 0 {
timestampsOffset = tReport.Offset timestampsOffset = tReport.Offset
timestampsRewindOffset = tReport.RewindOffset
valuesOffset = vReport.Offset valuesOffset = vReport.Offset
valuesRewindOffset = vReport.RewindOffset
} }
} else { } else {
// сторінка заповнена // сторінка заповнена
since := s.timestamps.ReplaceSinceWithUntil() since := s.timestamps.ReplaceSinceWithUntil()
if len(s.capturedState.Pages) == 0 { if len(pages) == 0 {
headTimestamps = timestamps.Tail(timestampsOffset) headTimestamps = timestamps.Tail(timestampsOffset)
headValues = values.Tail(valuesOffset) headValues = values.Tail(valuesOffset)
} }
s.capturedState.Pages = append(s.capturedState.Pages, storage.DataPayload{ pages = append(pages, storage.DataPayload{
Since: since, Since: since,
Content: s.buffer, Content: s.buffer,
TimestampsSize: timestamps.Size(), TimestampsSize: timestamps.Size(),
@@ -139,10 +148,11 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, send
// renew // renew
s.buffer = buf s.buffer = buf
timestamps.Append(tReport.RewindOffset, tmp[:tReport.ChangeSize], measure.Timestamp)
values.Append(vReport.RewindOffset, tmp[7:7+vReport.ChangeSize], measure.Value, vReport.Delta)
} }
timestamps.Append(tReport.RewindOffset, tmp[:tReport.ChangeSize], measure.Timestamp)
values.Append(vReport.RewindOffset, tmp[7:7+vReport.ChangeSize], measure.Value, vReport.Delta)
// //
s.lastValue = measure.Value s.lastValue = measure.Value
written++ written++
@@ -159,26 +169,26 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, send
// виділити змінені байти. // виділити змінені байти.
// скопіювати. Причому можна скопіювати зрізи chunks // скопіювати. Причому можна скопіювати зрізи chunks
if len(s.capturedState.Pages) > 0 { if len(pages) > 0 {
// пишу в storage довгим шляхом через redo файл і запис в data файл // пишу в storage довгим шляхом через redo файл і запис в data файл
sendToStorage(storage.AppendedMeasures{ sendToStorage(storage.AppendedMeasuresWithGrow{
MetricID: req.MetricID, MetricID: req.MetricID,
LastPageNo: s.lastPageNo, LastPageNo: s.lastPageNo,
TimestampsOffset: timestampsOffset, TimestampsOffset: timestampsRewindOffset,
Timestamps: headTimestamps, Timestamps: headTimestamps,
ValuesOffset: valuesOffset, ValuesOffset: valuesRewindOffset,
Values: headValues, Values: headValues,
IndexLevelTails: s.indexLevelTails, IndexLevelTails: s.indexLevelTails,
DataPages: s.capturedState.Pages, DataPages: pages,
TailTimestamps: timestamps.Tail(0), // payload TailTimestamps: timestamps.Tail(0), // payload
TailValues: values.Tail(0), TailValues: values.Tail(0),
ResultCode: resultCode, ResultCode: resultCode,
WrittenCount: written, WrittenCount: written,
//ResultCh: req.ResultCh, ResultCh: req.ResultCh,
}) })
} else { } else {
// короткий шлях - запис лише в storage // короткий шлях - запис лише в storage
sendToStorage(storage.AppendedMeasures{ sendToStorage(storage.AppendedMeasuresWithGrow{
MetricID: req.MetricID, MetricID: req.MetricID,
TimestampsOffset: timestampsOffset, TimestampsOffset: timestampsOffset,
ValuesOffset: valuesOffset, ValuesOffset: valuesOffset,
@@ -186,7 +196,7 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, send
Values: values.Tail(valuesOffset), Values: values.Tail(valuesOffset),
ResultCode: resultCode, ResultCode: resultCode,
WrittenCount: written, WrittenCount: written,
//ResultCh: req.ResultCh, ResultCh: req.ResultCh,
}) })
} }
@@ -207,74 +217,71 @@ func (s *_metric) FinAppendMeasures(rec storage.AppendMeasuresSummary) {
// READ // READ
func (s *_metric) StartRangeScan(req tryRangeScanReq) { func (s *_metric) StartRangeScan(req tryRangeScanReq) {
// if s.Since == 0 { if s.Since == 0 {
// req.ResultCh <- rangeScanResult{ req.ResultCh <- rangeScanResult{
// ResultCode: QueryDone, ResultCode: QueryDone,
// } }
// return return
// } }
// if req.Since > s.Until { if req.Since > s.Until {
// req.ResultCh <- rangeScanResult{ req.ResultCh <- rangeScanResult{
// ResultCode: QueryDone, ResultCode: QueryDone,
// } }
// return return
// } }
// if req.Until < s.Since { if req.Until < s.Since {
// if s.RootPageNo > 0 { if s.RootPageNo > 0 {
// req.ResultCh <- rangeScanResult{ req.ResultCh <- rangeScanResult{
// ResultCode: UntilNotFound, ResultCode: UntilNotFound,
// RootPageNo: s.RootPageNo, RootPageNo: s.RootPageNo,
// FracDigits: s.FracDigits, FracDigits: s.FracDigits,
// } }
// s.RLocks++ s.RLocks++
// return return
// } else { } else {
// req.ResultCh <- rangeScanResult{ req.ResultCh <- rangeScanResult{
// ResultCode: QueryDone, ResultCode: QueryDone,
// } }
// return return
// } }
// } }
// timestampDecompressor := s.Timestamps.CreateDecompressor() timestampDecompressor := s.timestamps.CreateDecompressor()
// valueDecompressor := s.Values.CreateDecompressor(s.FracDigits) valueDecompressor := s.values.CreateDecompressor(s.metricType, s.fracDigits)
// for { for {
// timestamp, done := timestampDecompressor.NextValue() timestamp, done := timestampDecompressor.NextValue()
// if done { if done {
// break break
// } }
value, done := valueDecompressor.NextValue()
// value, done := valueDecompressor.NextValue() if done {
// if done { qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug)
// qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug) }
// } if timestamp <= req.Until {
req.ResponseWriter.FeedNoSend(timestamp, value)
// if timestamp <= req.Until { if timestamp < req.Since {
// req.ResponseWriter.FeedNoSend(timestamp, value) req.ResultCh <- rangeScanResult{
// if timestamp < req.Since { ResultCode: QueryDone,
// req.ResultCh <- rangeScanResult{ }
// ResultCode: QueryDone, return
// } }
// return }
// } }
// } if s.lastPageNo > 0 {
// } req.ResultCh <- rangeScanResult{
ResultCode: UntilFound,
// if s.LastPageNo > 0 { LastPageNo: s.lastPageNo,
// req.ResultCh <- rangeScanResult{ FracDigits: s.fracDigits,
// ResultCode: UntilFound, }
// LastPageNo: s.LastPageNo, s.RLocks++
// FracDigits: s.FracDigits, } else {
// } req.ResultCh <- rangeScanResult{
// s.RLocks++ ResultCode: QueryDone,
// } else { }
// req.ResultCh <- rangeScanResult{ }
// ResultCode: QueryDone,
// }
// }
} }
func (s *_metric) StartFullScan(req tryFullScanReq) { func (s *_metric) StartFullScan(req tryFullScanReq) {
@@ -285,33 +292,33 @@ func (s *_metric) StartFullScan(req tryFullScanReq) {
// return // return
// } // }
// timestampDecompressor := s.Timestamps.CreateDecompressor() timestampDecompressor := s.timestamps.CreateDecompressor()
// valueDecompressor := s.Values.CreateDecompressor() valueDecompressor := s.values.CreateDecompressor(s.metricType, s.fracDigits)
// for { for {
// timestamp, done := timestampDecompressor.NextValue() timestamp, done := timestampDecompressor.NextValue()
// if done { if done {
// break break
// } }
// value, done := valueDecompressor.NextValue() value, done := valueDecompressor.NextValue()
// if done { if done {
// qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug) qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug)
// } }
// req.ResponseWriter.FeedNoSend(timestamp, value) req.ResponseWriter.FeedNoSend(timestamp, value)
// } }
// if s.LastPageNo > 0 { if s.lastPageNo > 0 {
// req.ResultCh <- fullScanResult{ req.ResultCh <- fullScanResult{
// ResultCode: UntilFound, ResultCode: UntilFound,
// LastPageNo: s.LastPageNo, LastPageNo: s.lastPageNo,
// FracDigits: s.FracDigits, FracDigits: s.fracDigits,
// } }
// s.RLocks++ s.RLocks++
// } else { } else {
// req.ResultCh <- fullScanResult{ req.ResultCh <- fullScanResult{
// ResultCode: QueryDone, ResultCode: QueryDone,
// } }
// } }
} }
// індекси // індекси
@@ -342,29 +349,21 @@ func (s *_metric) WriteTo(w io.Writer) (err error) {
if err != nil { if err != nil {
return return
} }
err = bin.WriteUint16(w, uint16(s.timestamps.Size())) err = bin.WriteUint16(w, uint16(s.timestamps.StoredSize()))
if err != nil { if err != nil {
return return
} }
err = bin.WriteUint16(w, uint16(s.values.Size())) err = bin.WriteUint16(w, uint16(s.values.StoredSize()))
if err != nil { if err != nil {
return return
} }
var (
timestampsCapturedState *qb.TimeDeltaCapturedState
valuesCapturedState *qb.ValueDeltaCapturedState
)
if s.capturedState != nil {
timestampsCapturedState = &s.capturedState.TimeState
valuesCapturedState = &s.capturedState.ValueState
}
// timestamps payload // timestamps payload
err = s.timestamps.WritePayloadTo(w, timestampsCapturedState) err = s.timestamps.WriteStoredTo(w)
if err != nil { if err != nil {
return return
} }
// values payload // values payload
err = s.values.WritePayloadTo(w, valuesCapturedState) err = s.values.WriteStoredTo(w)
if err != nil { if err != nil {
return return
} }

View File

@@ -87,13 +87,13 @@ func (s *ReplayMetric) ReadFrom(r io.Reader) (err error) {
// } // }
func (s *ReplayMetric) MeasuresAppend(rec storage.MeasuresAppendRecord) { func (s *ReplayMetric) MeasuresAppend(rec storage.MeasuresAppendRecord) {
timestampsSize := rec.TimestampsOffset - len(rec.Timestamps) timestampsSize := rec.TimestampsRewindOffset - len(rec.Timestamps)
// копіюю нові дані із WAL з урахуванням offset-ів // копіюю нові дані із WAL з урахуванням offset-ів
copy(s.databuf[rec.ValuesOffset:], rec.Values) copy(s.databuf[rec.ValuesRewindOffset:], rec.Values)
copy(s.databuf[len(s.databuf)-timestampsSize:], rec.Timestamps) copy(s.databuf[len(s.databuf)-timestampsSize:], rec.Timestamps)
// //
s.tSize = timestampsSize s.tSize = timestampsSize
s.vSize = rec.ValuesOffset - len(rec.Values) s.vSize = rec.ValuesRewindOffset - len(rec.Values)
} }
type AppendMeasuresResult struct { type AppendMeasuresResult struct {

View File

@@ -236,21 +236,21 @@ func (s *TimeDeltaCompressor) StoredSize() int {
// } // }
// } // }
func (s *TimeDeltaCompressor) WritePayloadTo(w io.Writer, state *qb.TimeDeltaCapturedState) (err error) { func (s *TimeDeltaCompressor) WriteStoredTo(w io.Writer) (err error) {
if state == nil { if s.state == nil {
_, err = w.Write(s.buf[:s.pos]) _, err = w.Write(s.buf[:s.pos])
return return
} else { } else {
_, err = w.Write(state.Payload) _, err = w.Write(s.state.Payload)
if err != nil { if err != nil {
return return
} }
_, err = bin.WriteVarUint64(w, uint64(state.LastDelta)) _, err = bin.WriteVarUint64(w, uint64(s.state.LastDelta))
if err != nil { if err != nil {
return return
} }
_, err = w.Write([]byte{ _, err = w.Write([]byte{
state.H, s.state.H,
}) })
return return
} }

View File

@@ -198,21 +198,21 @@ func (s *ValueDeltaCompressor) StoredSize() int {
// } // }
// } // }
func (s *ValueDeltaCompressor) WritePayloadTo(w io.Writer, state *qb.ValueDeltaCapturedState) (err error) { func (s *ValueDeltaCompressor) WriteStoredTo(w io.Writer) (err error) {
if state == nil { if s.state == nil {
_, err = w.Write(s.buf[:s.pos]) _, err = w.Write(s.buf[:s.pos])
return return
} else { } else {
_, err = w.Write(state.Payload) _, err = w.Write(s.state.Payload)
if err != nil { if err != nil {
return return
} }
_, err = bin.WriteVarUint64(w, state.LastDelta) _, err = bin.WriteVarUint64(w, s.state.LastDelta)
if err != nil { if err != nil {
return return
} }
_, err = w.Write([]byte{ _, err = w.Write([]byte{
state.H, s.state.H,
}) })
return return
} }

7
qb.go
View File

@@ -29,7 +29,7 @@ type TimestampCompressor interface {
// (rewindOffset, change, timestamp) // (rewindOffset, change, timestamp)
Append(int, []byte, uint32) Append(int, []byte, uint32)
Size() int Size() int
//Chunks() [][]byte StoredSize() int
//DeleteLast() //DeleteLast()
CaptureState() CaptureState()
// (offset) => payload // (offset) => payload
@@ -38,7 +38,7 @@ type TimestampCompressor interface {
//Payload() []byte // для снапшота //Payload() []byte // для снапшота
// Offset() int // Offset() int
ReplaceBuffer([]byte) ReplaceBuffer([]byte)
WritePayloadTo(io.Writer, *TimeDeltaCapturedState) error WriteStoredTo(io.Writer) error
LastTimestamp() uint32 LastTimestamp() uint32
ReplaceSinceWithUntil() uint32 ReplaceSinceWithUntil() uint32
} }
@@ -63,6 +63,7 @@ type ValueCompressor interface {
// (rewindOffset, change, value, delta) // (rewindOffset, change, value, delta)
Append(int, []byte, float64, uint64) Append(int, []byte, float64, uint64)
Size() int Size() int
StoredSize() int
//Chunks() [][]byte //Chunks() [][]byte
//DeleteLast() //DeleteLast()
CaptureState() CaptureState()
@@ -73,7 +74,7 @@ type ValueCompressor interface {
//Payload() []byte // для снапшота //Payload() []byte // для снапшота
//Offset() int //Offset() int
ReplaceBuffer([]byte) ReplaceBuffer([]byte)
WritePayloadTo(io.Writer, *ValueDeltaCapturedState) error WriteStoredTo(io.Writer) error
LastValue() float64 LastValue() float64
} }

View File

@@ -57,7 +57,28 @@ func (s *DataPreparer) Prepare(input []any) PreparedData {
// if len(x.FreePageNumbers) > 0 { // if len(x.FreePageNumbers) > 0 {
// s.freeList.AddPageNumbers(x.FreePageNumbers) // s.freeList.AddPageNumbers(x.FreePageNumbers)
// } // }
case AppendedMeasures: case MeasuresAppendRecord:
rec := MeasuresAppendRecord{
MetricID: x.MetricID,
TimestampsRewindOffset: x.TimestampsRewindOffset,
Timestamps: x.Timestamps,
ValuesRewindOffset: x.ValuesRewindOffset,
Values: x.Values,
}
rec.WriteTo(s.w)
// Дані для Worker
s.writeResults = append(s.writeResults, AppendMeasuresSummary{
// MetricID: x.MetricID,
// LastPageNo: sealResult.DataPages[len(sealResult.DataPages)-1].PageNo,
// Index: indexLevelTails,
// ResultCode: x.ResultCode,
// WrittenCount: x.WrittenCount,
// ResultCh: x.ResultCh,
})
case AppendedMeasuresWithGrow:
sealResult := s.pagePreparer.SealPages(SealPagesIn{ sealResult := s.pagePreparer.SealPages(SealPagesIn{
LastPageNo: x.LastPageNo, LastPageNo: x.LastPageNo,
IndexLevelTails: x.IndexLevelTails, IndexLevelTails: x.IndexLevelTails,
@@ -178,7 +199,7 @@ func (s *DataPreparer) Reset() {
} }
// MeasuresToWrite // MeasuresToWrite
type AppendedMeasures struct { type AppendedMeasuresWithGrow struct {
MetricID uint32 MetricID uint32
LastPageNo uint32 LastPageNo uint32
TimestampsOffset int // offset on 1st page TimestampsOffset int // offset on 1st page

View File

@@ -306,20 +306,20 @@ func (s *MeasuresAppendWithGrowRecord) Parse(r *bytes.Buffer) (err error) {
} }
type MeasuresAppendRecord struct { type MeasuresAppendRecord struct {
MetricID uint32 MetricID uint32
TimestampsOffset int TimestampsRewindOffset int
Timestamps []byte Timestamps []byte
ValuesOffset int ValuesRewindOffset int
Values []byte Values []byte
} }
func (s MeasuresAppendRecord) WriteTo(w *bytes.Buffer) { func (s MeasuresAppendRecord) WriteTo(w *bytes.Buffer) {
w.WriteByte(CodeMeasuresAppendWithGrow) w.WriteByte(CodeMeasuresAppendWithGrow)
bin.WriteUint32(w, s.MetricID) bin.WriteUint32(w, s.MetricID)
bin.WriteVarSize(w, s.TimestampsOffset) bin.WriteVarSize(w, s.TimestampsRewindOffset)
bin.WriteVarSize(w, len(s.Timestamps)) bin.WriteVarSize(w, len(s.Timestamps))
w.Write(s.Timestamps) w.Write(s.Timestamps)
bin.WriteVarSize(w, s.ValuesOffset) bin.WriteVarSize(w, s.ValuesRewindOffset)
bin.WriteVarSize(w, len(s.Values)) bin.WriteVarSize(w, len(s.Values))
w.Write(s.Values) w.Write(s.Values)
} }
@@ -330,7 +330,7 @@ func (s *MeasuresAppendRecord) Parse(r *bytes.Buffer) (err error) {
if err != nil { if err != nil {
return return
} }
s.TimestampsOffset, err = bin.ReadVarSize(r) s.TimestampsRewindOffset, err = bin.ReadVarSize(r)
if err != nil { if err != nil {
return return
} }
@@ -342,7 +342,7 @@ func (s *MeasuresAppendRecord) Parse(r *bytes.Buffer) (err error) {
if err != nil { if err != nil {
return return
} }
s.ValuesOffset, err = bin.ReadVarSize(r) s.ValuesRewindOffset, err = bin.ReadVarSize(r)
if err != nil { if err != nil {
return return
} }