From 721a9e80a5accfbb5cb8e2afa093397b1297aed9 Mon Sep 17 00:00:00 2001 From: wk989898 Date: Wed, 12 Aug 2026 04:04:27 +0000 Subject: [PATCH 1/2] init Signed-off-by: wk989898 --- .../dispatcher_manager_redo.go | 11 +- .../sink/cloudstorage/buffer_manager.go | 2 +- .../sink/cloudstorage/buffer_manager_test.go | 2 +- .../sink/cloudstorage/dml_writers.go | 3 +- .../sink/cloudstorage/spool/budget.go | 100 ------------- .../sink/cloudstorage/spool/budget_test.go | 81 ----------- .../sink/cloudstorage/spool_metrics.go | 45 ++++++ downstreamadapter/sink/cloudstorage/writer.go | 2 +- .../sink/cloudstorage/writer_test.go | 2 +- downstreamadapter/sink/helper/row_callback.go | 8 +- downstreamadapter/sink/redo/sink.go | 28 ++-- downstreamadapter/sink/redo/sink_test.go | 86 +++++++---- pkg/common/event/redo.go | 8 ++ pkg/redo/config.go | 3 - pkg/redo/writer/blackhole/writer.go | 1 + pkg/redo/writer/factory/factory.go | 9 +- pkg/redo/writer/factory/factory_test.go | 2 +- pkg/redo/writer/file/file.go | 15 +- pkg/redo/writer/file/file_test.go | 32 ++--- pkg/redo/writer/memory/ddl_writer.go | 6 +- pkg/redo/writer/memory/dml_writer.go | 134 +++++++++++++++++- pkg/redo/writer/memory/dml_writer_test.go | 121 +++++++++++++++- pkg/redo/writer/memory/encoding_worker.go | 19 +-- pkg/redo/writer/memory/file_worker.go | 14 +- pkg/sink/spool/budget.go | 107 ++++++++++++++ pkg/sink/spool/budget_test.go | 67 +++++++++ .../cloudstorage => pkg/sink}/spool/codec.go | 0 .../sink}/spool/codec_test.go | 0 .../cloudstorage => pkg/sink}/spool/quota.go | 59 ++++---- .../cloudstorage => pkg/sink}/spool/spool.go | 83 ++++++++--- .../sink}/spool/spool_test.go | 21 +-- 31 files changed, 705 insertions(+), 366 deletions(-) delete mode 100644 downstreamadapter/sink/cloudstorage/spool/budget.go delete mode 100644 downstreamadapter/sink/cloudstorage/spool/budget_test.go create mode 100644 downstreamadapter/sink/cloudstorage/spool_metrics.go create mode 100644 pkg/sink/spool/budget.go create mode 100644 pkg/sink/spool/budget_test.go rename {downstreamadapter/sink/cloudstorage => pkg/sink}/spool/codec.go (100%) rename {downstreamadapter/sink/cloudstorage => pkg/sink}/spool/codec_test.go (100%) rename {downstreamadapter/sink/cloudstorage => pkg/sink}/spool/quota.go (69%) rename {downstreamadapter/sink/cloudstorage => pkg/sink}/spool/spool.go (90%) rename {downstreamadapter/sink/cloudstorage => pkg/sink}/spool/spool_test.go (97%) diff --git a/downstreamadapter/dispatchermanager/dispatcher_manager_redo.go b/downstreamadapter/dispatchermanager/dispatcher_manager_redo.go index 7ae077afb2..c6bbba6cc3 100644 --- a/downstreamadapter/dispatchermanager/dispatcher_manager_redo.go +++ b/downstreamadapter/dispatchermanager/dispatcher_manager_redo.go @@ -57,18 +57,17 @@ func initRedoComponet( } var err error redoDispatcherMap := newDispatcherMap[*dispatcher.RedoDispatcher]() - redoSink, err := redo.New(ctx, changefeedID, manager.config.Consistent) - if err != nil { - return err - } - redoSchemaIDToDispatchers := dispatcher.NewSchemaIDToDispatchers() - totalQuota := manager.sinkQuota consistentMemoryUsage := manager.config.Consistent.MemoryUsage if consistentMemoryUsage == nil { consistentMemoryUsage = config.GetDefaultReplicaConfig().Consistent.MemoryUsage } redoQuota := totalQuota * consistentMemoryUsage.MemoryQuotaPercentage / 100 + redoSink, err := redo.New(ctx, changefeedID, manager.config.Consistent, redoQuota) + if err != nil { + return err + } + redoSchemaIDToDispatchers := dispatcher.NewSchemaIDToDispatchers() manager.writePathMu.Lock() if manager.writePathClosed.Load() { diff --git a/downstreamadapter/sink/cloudstorage/buffer_manager.go b/downstreamadapter/sink/cloudstorage/buffer_manager.go index fa49c404e1..2e4610c3c6 100644 --- a/downstreamadapter/sink/cloudstorage/buffer_manager.go +++ b/downstreamadapter/sink/cloudstorage/buffer_manager.go @@ -18,11 +18,11 @@ import ( "context" "time" - "github.com/pingcap/ticdc/downstreamadapter/sink/cloudstorage/spool" "github.com/pingcap/ticdc/downstreamadapter/sink/metrics" "github.com/pingcap/ticdc/pkg/cloudstorage" "github.com/pingcap/ticdc/pkg/common" "github.com/pingcap/ticdc/pkg/errors" + "github.com/pingcap/ticdc/pkg/sink/spool" ) const ( diff --git a/downstreamadapter/sink/cloudstorage/buffer_manager_test.go b/downstreamadapter/sink/cloudstorage/buffer_manager_test.go index 370ae3a986..2f64dc9a27 100644 --- a/downstreamadapter/sink/cloudstorage/buffer_manager_test.go +++ b/downstreamadapter/sink/cloudstorage/buffer_manager_test.go @@ -18,12 +18,12 @@ import ( "testing" "time" - "github.com/pingcap/ticdc/downstreamadapter/sink/cloudstorage/spool" "github.com/pingcap/ticdc/pkg/cloudstorage" commonType "github.com/pingcap/ticdc/pkg/common" commonEvent "github.com/pingcap/ticdc/pkg/common/event" "github.com/pingcap/ticdc/pkg/errors" "github.com/pingcap/ticdc/pkg/sink/codec/common" + "github.com/pingcap/ticdc/pkg/sink/spool" "github.com/stretchr/testify/require" ) diff --git a/downstreamadapter/sink/cloudstorage/dml_writers.go b/downstreamadapter/sink/cloudstorage/dml_writers.go index 6f4d8e9eb2..640ae3c00e 100644 --- a/downstreamadapter/sink/cloudstorage/dml_writers.go +++ b/downstreamadapter/sink/cloudstorage/dml_writers.go @@ -17,7 +17,6 @@ import ( "context" "time" - "github.com/pingcap/ticdc/downstreamadapter/sink/cloudstorage/spool" "github.com/pingcap/ticdc/downstreamadapter/sink/columnselector" sinkmetrics "github.com/pingcap/ticdc/downstreamadapter/sink/metrics" "github.com/pingcap/ticdc/pkg/cloudstorage" @@ -25,6 +24,7 @@ import ( commonEvent "github.com/pingcap/ticdc/pkg/common/event" "github.com/pingcap/ticdc/pkg/metrics" "github.com/pingcap/ticdc/pkg/sink/codec/common" + "github.com/pingcap/ticdc/pkg/sink/spool" "github.com/pingcap/ticdc/utils/chann" "github.com/pingcap/tidb/pkg/objstore/storeapi" "go.uber.org/atomic" @@ -68,6 +68,7 @@ func newDMLWriters( changefeedID, spool.WithRootDir(config.SpoolBaseDir), spool.WithDiskQuotaBytes(config.SpoolDiskQuota), + spool.WithMetrics(newSpoolMetrics(changefeedID)), ) if err != nil { return nil, err diff --git a/downstreamadapter/sink/cloudstorage/spool/budget.go b/downstreamadapter/sink/cloudstorage/spool/budget.go deleted file mode 100644 index 8fbc4a240d..0000000000 --- a/downstreamadapter/sink/cloudstorage/spool/budget.go +++ /dev/null @@ -1,100 +0,0 @@ -// Copyright 2026 PingCAP, Inc. -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// See the License for the specific language governing permissions and -// limitations under the License. - -package spool - -// budget stores the current queued byte counts and the byte limits derived from spool config. -type budget struct { - // diskQuotaBytes is the largest byte count allowed in local spool files. - diskQuotaBytes int64 - - // memoryQuotaBytes is the largest byte count we still keep in memory. - // If adding a new entry would cross this value, spool writes that entry - // to local spool files instead of keeping it in memory. - memoryQuotaBytes int64 - - // highWatermarkBytes is the byte count that makes spool stop running the - // new PostEnqueue callback immediately. The callback is saved in memory and - // will be run later. - highWatermarkBytes int64 - - // lowWatermarkBytes is the byte count that lets spool run the saved - // PostEnqueue callbacks again after some queued data has been flushed to the - // downstream storage or discarded locally. - lowWatermarkBytes int64 - - // memoryBytes is the number of queued bytes that are still kept in memory. - memoryBytes int64 - // diskBytes is the number of queued bytes that have already been written to local spool files. - diskBytes int64 -} - -func newBudget(options *options) *budget { - return &budget{ - diskQuotaBytes: options.diskQuotaBytes, - memoryQuotaBytes: int64(float64(options.diskQuotaBytes) * options.memoryRatio), - highWatermarkBytes: int64(float64(options.diskQuotaBytes) * options.highWatermarkRatio), - lowWatermarkBytes: int64(float64(options.diskQuotaBytes) * options.lowWatermarkRatio), - } -} - -// shouldSpill decides whether a new entry should stay in memory or be written to local spool files. -func (b *budget) shouldSpill(entryBytes int64) bool { - return b.memoryBytes+entryBytes > b.memoryQuotaBytes -} - -// entryExceedsDiskQuota returns true when a single spilled entry is larger -// than the configured disk quota by itself. -func (b *budget) entryExceedsDiskQuota(entryBytes int64) bool { - return entryBytes > b.diskQuotaBytes -} - -// spillWouldExceedDiskQuota returns true when adding one more spilled entry to -// the current on-disk usage would exceed the configured disk quota. -func (b *budget) spillWouldExceedDiskQuota(entryBytes int64) bool { - return b.diskBytes+entryBytes > b.diskQuotaBytes -} - -// acquire adds a newly accepted entry to the current byte counters and returns -// whether total queued bytes are now above the high watermark. -func (b *budget) acquire(entryBytes int64, spilled bool) bool { - if spilled { - b.diskBytes += entryBytes - return b.totalBytes() > b.highWatermarkBytes - } - b.memoryBytes += entryBytes - return b.totalBytes() > b.highWatermarkBytes -} - -// release removes an entry from the current byte counters after the entry has -// been flushed or discarded, and returns whether total queued bytes are now at -// or below the low watermark. -func (b *budget) release(entryBytes int64, spilled bool) bool { - if spilled { - b.diskBytes -= entryBytes - } - if !spilled { - b.memoryBytes -= entryBytes - } - if b.memoryBytes < 0 { - b.memoryBytes = 0 - } - if b.diskBytes < 0 { - b.diskBytes = 0 - } - return b.totalBytes() <= b.lowWatermarkBytes -} - -func (b *budget) totalBytes() int64 { - return b.memoryBytes + b.diskBytes -} diff --git a/downstreamadapter/sink/cloudstorage/spool/budget_test.go b/downstreamadapter/sink/cloudstorage/spool/budget_test.go deleted file mode 100644 index 700df31145..0000000000 --- a/downstreamadapter/sink/cloudstorage/spool/budget_test.go +++ /dev/null @@ -1,81 +0,0 @@ -// Copyright 2026 PingCAP, Inc. -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// See the License for the specific language governing permissions and -// limitations under the License. - -package spool - -import ( - "testing" - - "github.com/stretchr/testify/require" -) - -func TestBudgetTracksMemoryAndDiskBytes(t *testing.T) { - t.Parallel() - - core := newBudget(&options{ - diskQuotaBytes: 100, - memoryRatio: 0.2, - highWatermarkRatio: 0.8, - lowWatermarkRatio: 0.6, - }) - - require.False(t, core.shouldSpill(10)) - - overHighWatermark := core.acquire(10, false) - require.False(t, overHighWatermark) - require.Equal(t, int64(10), core.memoryBytes) - require.Equal(t, int64(0), core.diskBytes) - require.Equal(t, int64(10), core.totalBytes()) - - require.True(t, core.shouldSpill(11)) - - overHighWatermark = core.acquire(11, true) - require.False(t, overHighWatermark) - require.Equal(t, int64(10), core.memoryBytes) - require.Equal(t, int64(11), core.diskBytes) - require.Equal(t, int64(21), core.totalBytes()) - - atOrBelowLowWatermark := core.release(50, false) - require.True(t, atOrBelowLowWatermark) - require.Equal(t, int64(0), core.memoryBytes) - require.Equal(t, int64(11), core.diskBytes) - require.Equal(t, int64(11), core.totalBytes()) - - atOrBelowLowWatermark = core.release(50, true) - require.True(t, atOrBelowLowWatermark) - require.Equal(t, int64(0), core.memoryBytes) - require.Equal(t, int64(0), core.diskBytes) - require.Equal(t, int64(0), core.totalBytes()) -} - -func TestBudgetTracksWatermarkState(t *testing.T) { - t.Parallel() - - core := newBudget(&options{ - diskQuotaBytes: 100, - memoryRatio: 0.2, - highWatermarkRatio: 0.8, - lowWatermarkRatio: 0.6, - }) - - require.LessOrEqual(t, core.totalBytes(), core.lowWatermarkBytes) - require.LessOrEqual(t, core.totalBytes(), core.highWatermarkBytes) - - overHighWatermark := core.acquire(81, true) - require.True(t, overHighWatermark) - require.Greater(t, core.totalBytes(), core.lowWatermarkBytes) - - atOrBelowLowWatermark := core.release(21, true) - require.True(t, atOrBelowLowWatermark) - require.LessOrEqual(t, core.totalBytes(), core.highWatermarkBytes) -} diff --git a/downstreamadapter/sink/cloudstorage/spool_metrics.go b/downstreamadapter/sink/cloudstorage/spool_metrics.go new file mode 100644 index 0000000000..54e569e518 --- /dev/null +++ b/downstreamadapter/sink/cloudstorage/spool_metrics.go @@ -0,0 +1,45 @@ +// Copyright 2026 PingCAP, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// See the License for the specific language governing permissions and +// limitations under the License. + +package cloudstorage + +import ( + sinkmetrics "github.com/pingcap/ticdc/downstreamadapter/sink/metrics" + "github.com/pingcap/ticdc/pkg/common" + "github.com/pingcap/ticdc/pkg/sink/spool" +) + +func newSpoolMetrics(changefeedID common.ChangeFeedID) *spool.Metrics { + keyspace := changefeedID.Keyspace() + changefeed := changefeedID.Name() + return &spool.Metrics{ + MemoryBytes: sinkmetrics.CloudStorageSpoolMemoryBytesGauge.WithLabelValues(keyspace, changefeed), + DiskBytes: sinkmetrics.CloudStorageSpoolDiskBytesGauge.WithLabelValues(keyspace, changefeed), + PendingPostEnqueue: sinkmetrics.CloudStoragePendingPostEnqueueGauge.WithLabelValues(keyspace, changefeed), + DiskQuotaWaiters: sinkmetrics.CloudStorageSpoolDiskQuotaWaitersGauge.WithLabelValues(keyspace, changefeed), + DiskQuotaWait: sinkmetrics.CloudStorageSpoolDiskQuotaWaitDurationHistogram.WithLabelValues(keyspace, changefeed), + LoadedBytes: sinkmetrics.CloudStorageLoadBytesHistogram.WithLabelValues(keyspace, changefeed), + RotatedCount: sinkmetrics.CloudStorageRotateCountCounter.WithLabelValues(keyspace, changefeed), + SegmentCount: sinkmetrics.CloudStorageSpoolSegmentCountGauge.WithLabelValues(keyspace, changefeed), + Close: func() { + sinkmetrics.CloudStorageSpoolMemoryBytesGauge.DeleteLabelValues(keyspace, changefeed) + sinkmetrics.CloudStorageSpoolDiskBytesGauge.DeleteLabelValues(keyspace, changefeed) + sinkmetrics.CloudStoragePendingPostEnqueueGauge.DeleteLabelValues(keyspace, changefeed) + sinkmetrics.CloudStorageSpoolDiskQuotaWaitersGauge.DeleteLabelValues(keyspace, changefeed) + sinkmetrics.CloudStorageSpoolDiskQuotaWaitDurationHistogram.DeleteLabelValues(keyspace, changefeed) + sinkmetrics.CloudStorageLoadBytesHistogram.DeleteLabelValues(keyspace, changefeed) + sinkmetrics.CloudStorageRotateCountCounter.DeleteLabelValues(keyspace, changefeed) + sinkmetrics.CloudStorageSpoolSegmentCountGauge.DeleteLabelValues(keyspace, changefeed) + }, + } +} diff --git a/downstreamadapter/sink/cloudstorage/writer.go b/downstreamadapter/sink/cloudstorage/writer.go index e31a1387c6..635de4e78b 100644 --- a/downstreamadapter/sink/cloudstorage/writer.go +++ b/downstreamadapter/sink/cloudstorage/writer.go @@ -20,12 +20,12 @@ import ( "time" "github.com/pingcap/log" - "github.com/pingcap/ticdc/downstreamadapter/sink/cloudstorage/spool" "github.com/pingcap/ticdc/downstreamadapter/sink/metrics" "github.com/pingcap/ticdc/pkg/cloudstorage" "github.com/pingcap/ticdc/pkg/common" "github.com/pingcap/ticdc/pkg/errors" pmetrics "github.com/pingcap/ticdc/pkg/metrics" + "github.com/pingcap/ticdc/pkg/sink/spool" "github.com/pingcap/tidb/pkg/objstore/storeapi" "github.com/prometheus/client_golang/prometheus" "go.uber.org/zap" diff --git a/downstreamadapter/sink/cloudstorage/writer_test.go b/downstreamadapter/sink/cloudstorage/writer_test.go index 7aaf02dd82..7499167519 100644 --- a/downstreamadapter/sink/cloudstorage/writer_test.go +++ b/downstreamadapter/sink/cloudstorage/writer_test.go @@ -26,7 +26,6 @@ import ( "testing" "time" - "github.com/pingcap/ticdc/downstreamadapter/sink/cloudstorage/spool" "github.com/pingcap/ticdc/pkg/cloudstorage" commonType "github.com/pingcap/ticdc/pkg/common" commonEvent "github.com/pingcap/ticdc/pkg/common/event" @@ -34,6 +33,7 @@ import ( "github.com/pingcap/ticdc/pkg/metrics" "github.com/pingcap/ticdc/pkg/pdutil" "github.com/pingcap/ticdc/pkg/sink/codec/common" + "github.com/pingcap/ticdc/pkg/sink/spool" "github.com/pingcap/ticdc/pkg/util" "github.com/pingcap/tidb/pkg/meta/model" "github.com/pingcap/tidb/pkg/objstore/objectio" diff --git a/downstreamadapter/sink/helper/row_callback.go b/downstreamadapter/sink/helper/row_callback.go index 9acbc6e671..c9fe9ab480 100644 --- a/downstreamadapter/sink/helper/row_callback.go +++ b/downstreamadapter/sink/helper/row_callback.go @@ -21,10 +21,16 @@ import ( // NewPostFlushRowCallback returns a row-level callback that triggers txn-level // PostFlush exactly once when the callback has been invoked totalCount times. func NewPostFlushRowCallback(event *event.DMLEvent, totalCount uint64) func() { + return NewRowCallback(totalCount, event.PostFlush) +} + +// NewRowCallback returns a row-level callback that triggers callback exactly +// once after it has been invoked totalCount times. +func NewRowCallback(totalCount uint64, callback func()) func() { var calledCount atomic.Uint64 return func() { if calledCount.Inc() == totalCount { - event.PostFlush() + callback() } } } diff --git a/downstreamadapter/sink/redo/sink.go b/downstreamadapter/sink/redo/sink.go index d19fbf9c1e..a1b66ee719 100644 --- a/downstreamadapter/sink/redo/sink.go +++ b/downstreamadapter/sink/redo/sink.go @@ -61,6 +61,7 @@ func Verify(ctx context.Context, changefeedID common.ChangeFeedID, cfg *config.C // New creates a new redo sink. func New(ctx context.Context, changefeedID common.ChangeFeedID, cfg *config.ConsistentConfig, + spoolQuotaBytes uint64, ) (*Sink, error) { var err error config, err := writer.NewConfig(changefeedID, cfg) @@ -112,7 +113,7 @@ func New(ctx context.Context, changefeedID common.ChangeFeedID, zap.Error(err)) return nil, err } - dmlWriter, err = factory.NewRedoDMLWriter(ctx, config) + dmlWriter, err = factory.NewRedoDMLWriter(ctx, config, spoolQuotaBytes) if err != nil { log.Error("redo: failed to create redo log writer", zap.String("keyspace", changefeedID.Keyspace()), @@ -168,7 +169,9 @@ func (s *Sink) WriteBlockEvent(event commonEvent.BlockEvent) error { func (s *Sink) AddDMLEvent(event *commonEvent.DMLEvent) { rowsCount := event.Len() events := make([]*commonEvent.RedoRowEvent, 0, rowsCount) - rowCallback := helper.NewPostFlushRowCallback(event, uint64(rowsCount)) + postEnqueue, postFlush := event.DetachPostCallbacks() + rowPostEnqueue := helper.NewRowCallback(uint64(rowsCount), postEnqueue) + rowPostFlush := helper.NewRowCallback(uint64(rowsCount), postFlush) var ( startTs = event.GetStartTs() @@ -187,7 +190,8 @@ func (s *Sink) AddDMLEvent(event *commonEvent.DMLEvent) { Event: row, PhysicalTableID: physicalTableID, TableInfo: event.TableInfo, - Callback: rowCallback, + Callback: rowPostFlush, + EnqueueCallback: rowPostEnqueue, }) } s.logBuffer.Push(events...) @@ -237,29 +241,21 @@ func (s *Sink) Close() { } func (s *Sink) sendMessages(ctx context.Context) error { - buffer := make([]*commonEvent.RedoRowEvent, 0, redo.DefaultFlushBatchSize) for { - select { - case <-ctx.Done(): - return errors.Trace(context.Cause(ctx)) - default: + event, ok, err := s.logBuffer.GetWithContext(ctx) + if err != nil { + return errors.Trace(err) } - events, ok := s.logBuffer.GetMultipleNoGroup(buffer) if !ok { return nil } - if len(events) == 0 { - continue - } - buffer = events[:0] start := time.Now() - err := s.dmlWriter.AddDMLEvents(ctx, events...) - if err != nil { + if err := s.dmlWriter.AddDMLEvents(ctx, event); err != nil { return err } if s.metricCollector != nil { - s.metricCollector.observeRowWrite(len(events), time.Since(start)) + s.metricCollector.observeRowWrite(1, time.Since(start)) } } } diff --git a/downstreamadapter/sink/redo/sink_test.go b/downstreamadapter/sink/redo/sink_test.go index 1207572f4d..73f083ad35 100644 --- a/downstreamadapter/sink/redo/sink_test.go +++ b/downstreamadapter/sink/redo/sink_test.go @@ -111,6 +111,7 @@ func TestRedoSinkBatchConfig(t *testing.T) { context.Background(), common.NewChangeFeedIDWithName("test", common.DefaultKeyspaceName), cfg, + config.DefaultChangefeedMemoryQuota, ) require.NoError(t, err) defer sink.Close() @@ -119,6 +120,52 @@ func TestRedoSinkBatchConfig(t *testing.T) { require.Equal(t, int(32*redo.Megabyte), sink.BatchBytes()) } +func TestRedoSinkTwoStageAck(t *testing.T) { + helper := commonEvent.NewEventTestHelper(t) + defer helper.Close() + + helper.Tk().MustExec("use test") + job := helper.DDL2Job("create table t (id int primary key)") + require.NotNil(t, job) + event := helper.DML2Event("test", "t", "insert into t values (1), (2), (3)") + + callbacks := make([]string, 0, 2) + event.AddPostEnqueueFunc(func() { + callbacks = append(callbacks, "enqueue") + }) + event.AddPostFlushFunc(func() { + callbacks = append(callbacks, "flush") + }) + + sink := &Sink{ + ctx: context.Background(), + logBuffer: chann.NewUnlimitedChannelDefault[*commonEvent.RedoRowEvent](), + } + sink.AddDMLEvent(event) + require.Empty(t, callbacks) + + sink.logBuffer.Close() + rowEvents, ok := sink.logBuffer.GetMultipleNoGroup( + make([]*commonEvent.RedoRowEvent, 0, event.Len())) + require.True(t, ok) + require.Len(t, rowEvents, int(event.Len())) + + for _, rowEvent := range rowEvents[:len(rowEvents)-1] { + rowEvent.PostEnqueue() + } + require.Empty(t, callbacks) + rowEvents[len(rowEvents)-1].PostEnqueue() + require.Equal(t, []string{"enqueue"}, callbacks) + + for _, rowEvent := range rowEvents[:len(rowEvents)-1] { + rowEvent.PostFlush() + } + require.Equal(t, []string{"enqueue"}, callbacks) + + rowEvents[len(rowEvents)-1].PostFlush() + require.Equal(t, []string{"enqueue", "flush"}, callbacks) +} + // TestRedoSinkInProcessor tests how redo log manager is used in processor. func TestRedoSinkInProcessor(t *testing.T) { helper := commonEvent.NewEventTestHelper(t) @@ -144,7 +191,7 @@ func TestRedoSinkInProcessor(t *testing.T) { ctx, cancel := context.WithCancel(ctx) cfg := newTestConsistentConfig(storage) cfg.UseFileBackend = util.AddressOf(useFileBackend) - dmlMgr, err := New(ctx, common.NewChangeFeedIDWithName("test", common.DefaultKeyspaceName), cfg) + dmlMgr, err := New(ctx, common.NewChangeFeedIDWithName("test", common.DefaultKeyspaceName), cfg, config.DefaultChangefeedMemoryQuota) require.NoError(t, err) defer dmlMgr.Close() @@ -227,7 +274,7 @@ func TestRedoSinkError(t *testing.T) { defer cancel() cfg := newTestConsistentConfig("blackhole-invalid://") - logMgr, err := New(ctx, common.NewChangeFeedIDWithName("test", common.DefaultKeyspaceName), cfg) + logMgr, err := New(ctx, common.NewChangeFeedIDWithName("test", common.DefaultKeyspaceName), cfg, config.DefaultChangefeedMemoryQuota) require.NoError(t, err) defer logMgr.Close() @@ -281,7 +328,7 @@ func runBenchTest(b *testing.B, storage string, useFileBackend bool) { cfg.EncodingWorkerNum = util.AddressOf(redo.DefaultEncodingWorkerNum) cfg.FlushWorkerNum = util.AddressOf(redo.DefaultFlushWorkerNum) cfg.UseFileBackend = util.AddressOf(useFileBackend) - dmlMgr, err := New(ctx, common.NewChangeFeedIDWithName("test", common.DefaultKeyspaceName), cfg) + dmlMgr, err := New(ctx, common.NewChangeFeedIDWithName("test", common.DefaultKeyspaceName), cfg, config.DefaultChangefeedMemoryQuota) require.NoError(b, err) defer dmlMgr.Close() @@ -340,7 +387,7 @@ func runBenchTest(b *testing.B, storage string, useFileBackend bool) { require.ErrorIs(b, eg.Wait(), context.Canceled) } -func TestRedoSinkSendMessagesInBatch(t *testing.T) { +func TestRedoSinkSendMessages(t *testing.T) { t.Parallel() ctx, cancel := context.WithCancel(context.Background()) @@ -350,25 +397,13 @@ func TestRedoSinkSendMessagesInBatch(t *testing.T) { defer ctrl.Finish() mockWriter := writer.NewMockRedoDMLWriter(ctrl) - expectWriteBatch := func(batchSize int) *gomock.Call { - args := make([]interface{}, 0, batchSize+1) - args = append(args, gomock.Any()) // context - for range batchSize { - args = append(args, gomock.Any()) - } - return mockWriter.EXPECT(). - AddDMLEvents(args[0], args[1:]...). - DoAndReturn(func(_ context.Context, events ...*commonEvent.RedoRowEvent) error { - require.Len(t, events, batchSize) - return nil - }) - } - - gomock.InOrder( - expectWriteBatch(redo.DefaultFlushBatchSize), - expectWriteBatch(redo.DefaultFlushBatchSize), - expectWriteBatch(17), - ) + mockWriter.EXPECT(). + AddDMLEvents(gomock.Any(), gomock.Any()). + DoAndReturn(func(_ context.Context, events ...*commonEvent.RedoRowEvent) error { + require.Len(t, events, 1) + return nil + }). + Times(3) s := &Sink{ dmlWriter: mockWriter, @@ -380,9 +415,8 @@ func TestRedoSinkSendMessagesInBatch(t *testing.T) { doneCh <- s.sendMessages(ctx) }() - totalEvents := redo.DefaultFlushBatchSize*2 + 17 - events := make([]*commonEvent.RedoRowEvent, 0, totalEvents) - for range totalEvents { + events := make([]*commonEvent.RedoRowEvent, 0, 3) + for range 3 { events = append(events, &commonEvent.RedoRowEvent{}) } s.logBuffer.Push(events...) diff --git a/pkg/common/event/redo.go b/pkg/common/event/redo.go index 36b4906588..2ef562b28f 100644 --- a/pkg/common/event/redo.go +++ b/pkg/common/event/redo.go @@ -117,6 +117,7 @@ type RedoRowEvent struct { TableInfo *common.TableInfo Event RowChange Callback func() + EnqueueCallback func() } const ( @@ -132,6 +133,13 @@ func (r *RedoRowEvent) PostFlush() { } } +// PostEnqueue marks this encoded row as accepted by the redo spool. +func (r *RedoRowEvent) PostEnqueue() { + if r.EnqueueCallback != nil { + r.EnqueueCallback() + } +} + func (r *RedoRowEvent) ToRedoLog() *RedoLog { redoRow := &RedoDMLEvent{ Row: &DMLEventInRedoLog{ diff --git a/pkg/redo/config.go b/pkg/redo/config.go index 6cce071a54..16b9ae5d0f 100644 --- a/pkg/redo/config.go +++ b/pkg/redo/config.go @@ -49,9 +49,6 @@ const ( DefaultMetaFlushIntervalInMs = 200 // MinFlushIntervalInMs is the minimum flush interval for redo log. MinFlushIntervalInMs = 50 - // DefaultFlushBatchSize is the default flush batch size for redo log. - DefaultFlushBatchSize = 1024 - // DefaultEncodingWorkerNum is the default number of encoding workers. DefaultEncodingWorkerNum = 16 // DefaultEncodingInputChanSize is the default size of input channel for encoding worker. diff --git a/pkg/redo/writer/blackhole/writer.go b/pkg/redo/writer/blackhole/writer.go index ec276ac9f2..ec44de3b66 100644 --- a/pkg/redo/writer/blackhole/writer.go +++ b/pkg/redo/writer/blackhole/writer.go @@ -76,6 +76,7 @@ func (bs *blackHoleDMLWriter) AddDMLEvents(_ context.Context, events ...*event.R log.Debug("write redo events", fields...) for _, e := range events { if e != nil { + e.PostEnqueue() e.PostFlush() } } diff --git a/pkg/redo/writer/factory/factory.go b/pkg/redo/writer/factory/factory.go index beba889ff7..b0656bf785 100644 --- a/pkg/redo/writer/factory/factory.go +++ b/pkg/redo/writer/factory/factory.go @@ -26,7 +26,7 @@ import ( // NewRedoDMLWriter creates a new RedoDMLWriter. func NewRedoDMLWriter( - ctx context.Context, cfg *writer.Config, + ctx context.Context, cfg *writer.Config, spoolQuotaBytes uint64, ) (writer.RedoDMLWriter, error) { uri := cfg.URI() if redo.IsBlackholeStorage(uri.Scheme) { @@ -34,10 +34,9 @@ func NewRedoDMLWriter( return blackhole.NewDMLWriter(invalid), nil } - if cfg.UseFileBackend() { - return file.NewDMLWriter(ctx, cfg) - } - return memory.NewDMLWriter(ctx, cfg) + // DML always uses the spooled writer, regardless of use-file-backend, so + // both backend configurations have the same acknowledgement semantics. + return memory.NewDMLWriter(ctx, cfg, spoolQuotaBytes) } // NewRedoDDLWriter creates a new RedoDDLWriter. diff --git a/pkg/redo/writer/factory/factory_test.go b/pkg/redo/writer/factory/factory_test.go index c78021286d..417d55a6d8 100644 --- a/pkg/redo/writer/factory/factory_test.go +++ b/pkg/redo/writer/factory/factory_test.go @@ -32,7 +32,7 @@ func TestNewRedoWriters(t *testing.T) { ) require.NoError(t, err) - dmlWriter, err := NewRedoDMLWriter(context.Background(), cfg) + dmlWriter, err := NewRedoDMLWriter(context.Background(), cfg, 1024) require.NoError(t, err) require.Implements(t, (*writer.RedoDMLWriter)(nil), dmlWriter) diff --git a/pkg/redo/writer/file/file.go b/pkg/redo/writer/file/file.go index 759d429337..11d68ec5da 100644 --- a/pkg/redo/writer/file/file.go +++ b/pkg/redo/writer/file/file.go @@ -355,8 +355,7 @@ func (w *Writer) encode(ctx context.Context) error { d := time.Duration(w.cfg.FlushIntervalInMs()) * time.Millisecond ticker := time.NewTicker(d) defer ticker.Stop() - num := 0 - cacheEventPostFlush := make([]func(), 0, redo.DefaultFlushBatchSize) + var cacheEventPostFlush []func() flush := func() error { err := w.Flush() if err != nil { @@ -365,7 +364,6 @@ func (w *Writer) encode(ctx context.Context) error { for _, fn := range cacheEventPostFlush { fn() } - num = 0 cacheEventPostFlush = cacheEventPostFlush[:0] return nil } @@ -383,16 +381,7 @@ func (w *Writer) encode(ctx context.Context) error { if err != nil { return err } - num++ - if num >= redo.DefaultFlushBatchSize { - err := flush() - if err != nil { - return errors.Trace(err) - } - e.PostFlush() - } else { - cacheEventPostFlush = append(cacheEventPostFlush, e.PostFlush) - } + cacheEventPostFlush = append(cacheEventPostFlush, e.PostFlush) } } } diff --git a/pkg/redo/writer/file/file_test.go b/pkg/redo/writer/file/file_test.go index 02f890ec47..37d2dc2aa4 100644 --- a/pkg/redo/writer/file/file_test.go +++ b/pkg/redo/writer/file/file_test.go @@ -436,22 +436,22 @@ func TestRotateFileWithoutFileAllocator(t *testing.T) { w.Close() } -func TestRunFlushesOnBatchBoundaryAndExecutesPostFlush(t *testing.T) { +func TestRunFlushesOnIntervalAndExecutesPostFlush(t *testing.T) { t.Parallel() dir := t.TempDir() - flushIntervalInMs := int64(60 * 1000) + flushIntervalInMs := int64(100) flushWorkerNum := 9 - batchWriterCfg := newTestWriterConfig( + writerCfg := newTestWriterConfig( t, - common.NewChangeFeedIDWithName("test-run-batch", common.DefaultKeyspaceName), + common.NewChangeFeedIDWithName("test-run-interval", common.DefaultKeyspaceName), &config.ConsistentConfig{ FlushIntervalInMs: &flushIntervalInMs, FlushWorkerNum: &flushWorkerNum, Storage: util.AddressOf("file://" + dir), }, ) - w, err := NewFileWriter(context.Background(), batchWriterCfg, redo.RedoRowLogFileType) + w, err := NewFileWriter(context.Background(), writerCfg, redo.RedoRowLogFileType) require.NoError(t, err) ctx, cancel := context.WithCancel(context.Background()) @@ -461,7 +461,8 @@ func TestRunFlushesOnBatchBoundaryAndExecutesPostFlush(t *testing.T) { }() postFlushCnt := atomic.NewInt64(0) - for i := 0; i < redo.DefaultFlushBatchSize-1; i++ { + const eventCount = 3 + for i := 0; i < eventCount; i++ { ts := uint64(i + 1) w.GetInputCh() <- &pevent.RedoRowEvent{ StartTs: ts, @@ -472,25 +473,8 @@ func TestRunFlushesOnBatchBoundaryAndExecutesPostFlush(t *testing.T) { } } - // The callback should not be executed before the batch reaches the boundary. - require.Equal(t, int64(0), postFlushCnt.Load()) - select { - case err := <-runErrCh: - require.Failf(t, "run exited unexpectedly", "run returned before cancel: %v", err) - default: - } - - ts := uint64(redo.DefaultFlushBatchSize) - w.GetInputCh() <- &pevent.RedoRowEvent{ - StartTs: ts, - CommitTs: ts, - Callback: func() { - postFlushCnt.Inc() - }, - } - require.Eventually(t, func() bool { - return postFlushCnt.Load() == int64(redo.DefaultFlushBatchSize) + return postFlushCnt.Load() == eventCount }, 10*time.Second, 20*time.Millisecond) cancel() diff --git a/pkg/redo/writer/memory/ddl_writer.go b/pkg/redo/writer/memory/ddl_writer.go index e21c0ebbfe..8c541a9de7 100644 --- a/pkg/redo/writer/memory/ddl_writer.go +++ b/pkg/redo/writer/memory/ddl_writer.go @@ -223,8 +223,8 @@ func toPolymorphicDDLEvent( copy(data[8:], rawData) return &polymorphicRedoEvent{ - commitTs: rl.GetCommitTs(), - callback: event.PostFlush, - data: data, + commitTs: rl.GetCommitTs(), + postFlush: event.PostFlush, + data: data, }, nil } diff --git a/pkg/redo/writer/memory/dml_writer.go b/pkg/redo/writer/memory/dml_writer.go index 649126f7b5..a1af450994 100644 --- a/pkg/redo/writer/memory/dml_writer.go +++ b/pkg/redo/writer/memory/dml_writer.go @@ -15,11 +15,20 @@ package memory import ( "context" + "encoding/binary" + "math" + "os" + "path/filepath" "github.com/pingcap/log" commonEvent "github.com/pingcap/ticdc/pkg/common/event" + "github.com/pingcap/ticdc/pkg/config" + "github.com/pingcap/ticdc/pkg/errors" "github.com/pingcap/ticdc/pkg/redo" "github.com/pingcap/ticdc/pkg/redo/writer" + "github.com/pingcap/ticdc/pkg/sink/codec/common" + "github.com/pingcap/ticdc/pkg/sink/spool" + "github.com/pingcap/ticdc/utils/chann" "github.com/pingcap/tidb/pkg/objstore/storeapi" "go.uber.org/zap" "golang.org/x/sync/errgroup" @@ -31,13 +40,22 @@ type dmlWriter struct { cfg *writer.Config encodeWorkers *encodingWorkerGroup fileWorkers *fileWorkerGroup + spool *spool.Spool + spoolEntries *chann.UnlimitedChannel[*redoSpoolEntry, any] extStorage storeapi.Storage cancel context.CancelFunc } +type redoSpoolEntry struct { + entry *spool.Entry + flushImmediately bool +} + +const redoSpoolDirectory = "redo-sink-spool" + // NewDMLWriter creates a new memory DML writer. func NewDMLWriter( - ctx context.Context, cfg *writer.Config, opts ...writer.Option, + ctx context.Context, cfg *writer.Config, spoolQuotaBytes uint64, opts ...writer.Option, ) (writer.RedoDMLWriter, error) { extStorage, err := redo.InitExternalStorage(ctx, *cfg.URI()) if err != nil { @@ -45,13 +63,31 @@ func NewDMLWriter( } encodeWorkers := newEncodingWorkerGroup(cfg) + fileWorkerInput := make(chan *polymorphicRedoEvent, redo.DefaultEncodingOutputChanSize) fileWorkers := newFileWorkerGroup( - cfg, encodeWorkers.outputCh, extStorage, opts...) + cfg, fileWorkerInput, extStorage, opts...) + rootDir := config.GetGlobalServerConfig().DataDir + if rootDir == "" { + rootDir = os.TempDir() + } + rootDir = filepath.Join(rootDir, redoSpoolDirectory) + quotaBytes := int64(min(spoolQuotaBytes, uint64(math.MaxInt64))) + spoolBuffer, err := spool.New( + cfg.ChangeFeedID(), + spool.WithRootDir(rootDir), + spool.WithDiskQuotaBytes(quotaBytes), + ) + if err != nil { + extStorage.Close() + return nil, err + } return &dmlWriter{ cfg: cfg, encodeWorkers: encodeWorkers, fileWorkers: fileWorkers, + spool: spoolBuffer, + spoolEntries: chann.NewUnlimitedChannelDefault[*redoSpoolEntry](), extStorage: extStorage, }, nil } @@ -64,12 +100,102 @@ func (l *dmlWriter) Run(ctx context.Context) error { eg.Go(func() error { return l.encodeWorkers.Run(egCtx) }) + eg.Go(func() error { + return l.writeEncodedEventsToSpool(egCtx) + }) + eg.Go(func() error { + return l.readEncodedEventsFromSpool(egCtx) + }) eg.Go(func() error { return l.fileWorkers.Run(egCtx) }) return eg.Wait() } +func (l *dmlWriter) writeEncodedEventsToSpool(ctx context.Context) error { + for { + select { + case <-ctx.Done(): + return errors.Trace(context.Cause(ctx)) + case event := <-l.encodeWorkers.outputCh: + if event == nil { + return errors.ErrUnexpected.FastGenByArgs("encoded redo event is nil") + } + key := make([]byte, 8) + binary.LittleEndian.PutUint64(key, event.commitTs) + msg := common.NewMsg(key, event.data) + msg.Callback = event.postFlush + + for { + action, entry, err := l.spool.TryEnqueue( + []*common.Message{msg}, event.postEnqueue) + if err != nil { + return err + } + if action == spool.EnqueueActionWaitDiskQuota { + if err := l.spool.WaitForDiskQuota(ctx, []*common.Message{msg}); err != nil { + return err + } + continue + } + l.spoolEntries.Push(&redoSpoolEntry{ + entry: entry, + flushImmediately: action == spool.EnqueueActionAcceptedOversized, + }) + break + } + } + } +} + +func (l *dmlWriter) readEncodedEventsFromSpool(ctx context.Context) error { + for { + spooled, ok, err := l.spoolEntries.GetWithContext(ctx) + if err != nil { + return err + } + if !ok { + return nil + } + entry := spooled.entry + reader, err := l.spool.NewMessageReader(entry) + if err != nil { + return err + } + key, data, _, ok, err := reader.Next() + if err != nil { + return err + } + if !ok || len(key) != 8 || len(data) == 0 { + return errors.ErrUnexpected.FastGenByArgs("invalid encoded redo spool entry") + } + _, _, _, hasMore, err := reader.Next() + if err != nil { + return err + } + if hasMore { + return errors.ErrUnexpected.FastGenByArgs("encoded redo spool entry contains multiple messages") + } + postFlushCallbacks := reader.PostFlushCallbacks() + encodedEvent := &polymorphicRedoEvent{ + commitTs: binary.LittleEndian.Uint64(key), + data: data, + flushImmediately: spooled.flushImmediately, + postFlush: func() { + for _, callback := range postFlushCallbacks { + callback() + } + l.spool.Release(entry) + }, + } + select { + case <-ctx.Done(): + return errors.Trace(context.Cause(ctx)) + case l.fileWorkers.inputCh <- encodedEvent: + } + } +} + func (l *dmlWriter) AddDMLEvents(ctx context.Context, events ...*commonEvent.RedoRowEvent) error { for _, event := range events { if event == nil { @@ -94,5 +220,9 @@ func (l *dmlWriter) Close() error { l.extStorage.Close() l.extStorage = nil } + if l.spool != nil { + l.spool.Close() + l.spool = nil + } return nil } diff --git a/pkg/redo/writer/memory/dml_writer_test.go b/pkg/redo/writer/memory/dml_writer_test.go index f7d79241cc..a27d3864f8 100644 --- a/pkg/redo/writer/memory/dml_writer_test.go +++ b/pkg/redo/writer/memory/dml_writer_test.go @@ -15,12 +15,17 @@ package memory import ( "context" + "strings" + "sync/atomic" "testing" + "time" "github.com/pingcap/ticdc/pkg/common" "github.com/pingcap/ticdc/pkg/redo/testutil" "github.com/pingcap/ticdc/pkg/redo/writer" + "github.com/pingcap/ticdc/pkg/sink/spool" "github.com/pingcap/ticdc/pkg/util" + "github.com/pingcap/ticdc/utils/chann" "github.com/stretchr/testify/require" ) @@ -38,7 +43,121 @@ func TestNewDMLWriter(t *testing.T) { ) require.NoError(t, err) - lw, err := NewDMLWriter(ctx, cfg) + lw, err := NewDMLWriter(ctx, cfg, 1024) require.NoError(t, err) require.NoError(t, lw.Close()) } + +func TestDMLWriterSpoolsEncodedBytesBeforePostEnqueue(t *testing.T) { + changefeedID := common.NewChangeFeedIDWithName(t.Name(), common.DefaultKeyspaceName) + spoolBuffer, err := spool.New( + changefeedID, + spool.WithRootDir(t.TempDir()), + spool.WithDiskQuotaBytes(1000), + spool.WithSegmentBytes(1<<20), + spool.WithMemoryRatio(0.2), + spool.WithHighWatermarkRatio(0.6), + spool.WithLowWatermarkRatio(0.3), + ) + require.NoError(t, err) + defer spoolBuffer.Close() + + encodedCh := make(chan *polymorphicRedoEvent, 2) + dmlWriter := &dmlWriter{ + encodeWorkers: &encodingWorkerGroup{outputCh: encodedCh}, + spool: spoolBuffer, + spoolEntries: chann.NewUnlimitedChannelDefault[*redoSpoolEntry](), + } + + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error, 1) + go func() { + done <- dmlWriter.writeEncodedEventsToSpool(ctx) + }() + + var firstEnqueued atomic.Int64 + var secondEnqueued atomic.Int64 + firstData := []byte(strings.Repeat("a", 350)) + secondData := []byte(strings.Repeat("b", 350)) + encodedCh <- &polymorphicRedoEvent{ + commitTs: 1, + data: firstData, + postEnqueue: func() { firstEnqueued.Add(1) }, + } + encodedCh <- &polymorphicRedoEvent{ + commitTs: 2, + data: secondData, + postEnqueue: func() { secondEnqueued.Add(1) }, + } + + readCtx, readCancel := context.WithTimeout(context.Background(), 5*time.Second) + defer readCancel() + firstEntry, ok, err := dmlWriter.spoolEntries.GetWithContext(readCtx) + require.NoError(t, err) + require.True(t, ok) + secondEntry, ok, err := dmlWriter.spoolEntries.GetWithContext(readCtx) + require.NoError(t, err) + require.True(t, ok) + + require.True(t, firstEntry.entry.IsSpilled()) + require.True(t, secondEntry.entry.IsSpilled()) + require.False(t, firstEntry.flushImmediately) + require.False(t, secondEntry.flushImmediately) + require.Equal(t, int64(1), firstEnqueued.Load()) + require.Equal(t, int64(0), secondEnqueued.Load()) + + reader, err := spoolBuffer.NewMessageReader(firstEntry.entry) + require.NoError(t, err) + _, encodedData, _, ok, err := reader.Next() + require.NoError(t, err) + require.True(t, ok) + require.Equal(t, firstData, encodedData) + + spoolBuffer.Release(firstEntry.entry) + require.Equal(t, int64(0), secondEnqueued.Load()) + spoolBuffer.Release(secondEntry.entry) + require.Equal(t, int64(1), secondEnqueued.Load()) + + cancel() + require.ErrorIs(t, <-done, context.Canceled) +} + +func TestDMLWriterMarksOversizedEncodedBytesForImmediateFlush(t *testing.T) { + changefeedID := common.NewChangeFeedIDWithName(t.Name(), common.DefaultKeyspaceName) + spoolBuffer, err := spool.New( + changefeedID, + spool.WithRootDir(t.TempDir()), + spool.WithDiskQuotaBytes(100), + ) + require.NoError(t, err) + defer spoolBuffer.Close() + + encodedCh := make(chan *polymorphicRedoEvent, 1) + dmlWriter := &dmlWriter{ + encodeWorkers: &encodingWorkerGroup{outputCh: encodedCh}, + spool: spoolBuffer, + spoolEntries: chann.NewUnlimitedChannelDefault[*redoSpoolEntry](), + } + + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error, 1) + go func() { + done <- dmlWriter.writeEncodedEventsToSpool(ctx) + }() + + encodedCh <- &polymorphicRedoEvent{ + commitTs: 1, + data: []byte(strings.Repeat("a", 200)), + } + readCtx, readCancel := context.WithTimeout(context.Background(), 5*time.Second) + defer readCancel() + entry, ok, err := dmlWriter.spoolEntries.GetWithContext(readCtx) + require.NoError(t, err) + require.True(t, ok) + require.True(t, entry.entry.InMemory()) + require.True(t, entry.flushImmediately) + spoolBuffer.Release(entry.entry) + + cancel() + require.ErrorIs(t, <-done, context.Canceled) +} diff --git a/pkg/redo/writer/memory/encoding_worker.go b/pkg/redo/writer/memory/encoding_worker.go index 8787dea0c3..e39c0feae2 100644 --- a/pkg/redo/writer/memory/encoding_worker.go +++ b/pkg/redo/writer/memory/encoding_worker.go @@ -31,14 +31,16 @@ import ( // polymorphicRedoEvent wraps RedoLog and callback for file worker. type polymorphicRedoEvent struct { - commitTs common.Ts - data []byte - callback func() + commitTs common.Ts + data []byte + postEnqueue func() + postFlush func() + flushImmediately bool } func (e *polymorphicRedoEvent) PostFlush() { - if e.callback != nil { - e.callback() + if e.postFlush != nil { + e.postFlush() } } @@ -56,9 +58,10 @@ func toPolymorphicDMLEvent( binary.LittleEndian.PutUint64(data[:8], lenField) copy(data[8:], rawData) return &polymorphicRedoEvent{ - commitTs: rl.GetCommitTs(), - callback: event.PostFlush, - data: data, + commitTs: rl.GetCommitTs(), + postEnqueue: event.PostEnqueue, + postFlush: event.PostFlush, + data: data, }, nil } diff --git a/pkg/redo/writer/memory/file_worker.go b/pkg/redo/writer/memory/file_worker.go index ba67651544..13b269283c 100644 --- a/pkg/redo/writer/memory/file_worker.go +++ b/pkg/redo/writer/memory/file_worker.go @@ -212,8 +212,7 @@ func (f *fileWorkerGroup) bgWriteLogs( d := time.Duration(f.cfg.FlushIntervalInMs()) * time.Millisecond ticker := time.NewTicker(d) defer ticker.Stop() - num := 0 - cacheEventPostFlush := make([]func(), 0, redo.DefaultFlushBatchSize) + var cacheEventPostFlush []func() flush := func() error { err := f.flushAll(egCtx) if err != nil { @@ -222,7 +221,6 @@ func (f *fileWorkerGroup) bgWriteLogs( for _, fn := range cacheEventPostFlush { fn() } - num = 0 cacheEventPostFlush = cacheEventPostFlush[:0] return nil } @@ -244,15 +242,11 @@ func (f *fileWorkerGroup) bgWriteLogs( if err != nil { return errors.Trace(err) } - num++ - if num > redo.DefaultFlushBatchSize { - err := flush() - if err != nil { + cacheEventPostFlush = append(cacheEventPostFlush, event.PostFlush) + if event.flushImmediately { + if err := flush(); err != nil { return errors.Trace(err) } - event.PostFlush() - } else { - cacheEventPostFlush = append(cacheEventPostFlush, event.PostFlush) } } } diff --git a/pkg/sink/spool/budget.go b/pkg/sink/spool/budget.go new file mode 100644 index 0000000000..5e3f1cc4e6 --- /dev/null +++ b/pkg/sink/spool/budget.go @@ -0,0 +1,107 @@ +// Copyright 2026 PingCAP, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// See the License for the specific language governing permissions and +// limitations under the License. + +package spool + +// Limits defines the byte limits used by a spool budget. +type Limits struct { + DiskQuotaBytes int64 + MemoryQuotaBytes int64 + HighWatermarkBytes int64 + LowWatermarkBytes int64 +} + +// Budget tracks the memory and disk bytes owned by a spool. It is not +// thread-safe; callers should synchronize compound admission decisions. +type Budget struct { + limits Limits + + memoryBytes int64 + diskBytes int64 +} + +// NewBudget creates an empty spool budget with the supplied limits. +func NewBudget(limits Limits) *Budget { + return &Budget{limits: limits} +} + +// CanFitMemory reports whether an entry can be admitted into memory. +func (b *Budget) CanFitMemory(entryBytes int64) bool { + return b.memoryBytes+entryBytes <= b.limits.MemoryQuotaBytes +} + +// ShouldSpill reports whether a new entry should be written to disk instead +// of being retained in memory. +func (b *Budget) ShouldSpill(entryBytes int64) bool { + return !b.CanFitMemory(entryBytes) +} + +// EntryExceedsDiskQuota reports whether one entry is larger than the entire +// disk quota. +func (b *Budget) EntryExceedsDiskQuota(entryBytes int64) bool { + return entryBytes > b.limits.DiskQuotaBytes +} + +// SpillWouldExceedDiskQuota reports whether admitting one more spilled entry +// would exceed the disk quota. +func (b *Budget) SpillWouldExceedDiskQuota(entryBytes int64) bool { + return b.diskBytes+entryBytes > b.limits.DiskQuotaBytes +} + +// Acquire records an admitted entry and reports whether total staged bytes +// are above the high watermark. +func (b *Budget) Acquire(entryBytes int64, spilled bool) bool { + if spilled { + b.diskBytes += entryBytes + } else { + b.memoryBytes += entryBytes + } + return b.TotalBytes() > b.limits.HighWatermarkBytes +} + +// Release removes a flushed or discarded entry and reports whether total +// staged bytes are at or below the low watermark. +func (b *Budget) Release(entryBytes int64, spilled bool) bool { + if spilled { + b.diskBytes -= entryBytes + } else { + b.memoryBytes -= entryBytes + } + if b.memoryBytes < 0 { + b.memoryBytes = 0 + } + if b.diskBytes < 0 { + b.diskBytes = 0 + } + return b.TotalBytes() <= b.limits.LowWatermarkBytes +} + +// MemoryBytes returns currently staged in-memory bytes. +func (b *Budget) MemoryBytes() int64 { + return b.memoryBytes +} + +// DiskBytes returns currently staged on-disk bytes. +func (b *Budget) DiskBytes() int64 { + return b.diskBytes +} + +// TotalBytes returns all currently staged bytes. +func (b *Budget) TotalBytes() int64 { + return b.memoryBytes + b.diskBytes +} + +// Limits returns the immutable limits of this budget. +func (b *Budget) Limits() Limits { + return b.limits +} diff --git a/pkg/sink/spool/budget_test.go b/pkg/sink/spool/budget_test.go new file mode 100644 index 0000000000..53a7aca6a4 --- /dev/null +++ b/pkg/sink/spool/budget_test.go @@ -0,0 +1,67 @@ +// Copyright 2026 PingCAP, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// See the License for the specific language governing permissions and +// limitations under the License. + +package spool + +import ( + "testing" + + "github.com/stretchr/testify/require" +) + +func TestBudgetTracksMemoryAndDiskBytes(t *testing.T) { + t.Parallel() + + budget := NewBudget(Limits{ + DiskQuotaBytes: 100, + MemoryQuotaBytes: 20, + HighWatermarkBytes: 80, + LowWatermarkBytes: 60, + }) + + require.False(t, budget.ShouldSpill(10)) + require.False(t, budget.Acquire(10, false)) + require.Equal(t, int64(10), budget.MemoryBytes()) + require.Equal(t, int64(0), budget.DiskBytes()) + require.Equal(t, int64(10), budget.TotalBytes()) + + require.True(t, budget.ShouldSpill(11)) + require.False(t, budget.Acquire(11, true)) + require.Equal(t, int64(10), budget.MemoryBytes()) + require.Equal(t, int64(11), budget.DiskBytes()) + require.Equal(t, int64(21), budget.TotalBytes()) + + require.True(t, budget.Release(50, false)) + require.Equal(t, int64(0), budget.MemoryBytes()) + require.Equal(t, int64(11), budget.DiskBytes()) + + require.True(t, budget.Release(50, true)) + require.Equal(t, int64(0), budget.TotalBytes()) +} + +func TestBudgetTracksWatermarkAndDiskQuota(t *testing.T) { + t.Parallel() + + budget := NewBudget(Limits{ + DiskQuotaBytes: 100, + MemoryQuotaBytes: 20, + HighWatermarkBytes: 80, + LowWatermarkBytes: 60, + }) + + require.True(t, budget.EntryExceedsDiskQuota(101)) + require.False(t, budget.SpillWouldExceedDiskQuota(81)) + require.True(t, budget.Acquire(81, true)) + require.True(t, budget.SpillWouldExceedDiskQuota(20)) + require.True(t, budget.Release(21, true)) +} diff --git a/downstreamadapter/sink/cloudstorage/spool/codec.go b/pkg/sink/spool/codec.go similarity index 100% rename from downstreamadapter/sink/cloudstorage/spool/codec.go rename to pkg/sink/spool/codec.go diff --git a/downstreamadapter/sink/cloudstorage/spool/codec_test.go b/pkg/sink/spool/codec_test.go similarity index 100% rename from downstreamadapter/sink/cloudstorage/spool/codec_test.go rename to pkg/sink/spool/codec_test.go diff --git a/downstreamadapter/sink/cloudstorage/spool/quota.go b/pkg/sink/spool/quota.go similarity index 69% rename from downstreamadapter/sink/cloudstorage/spool/quota.go rename to pkg/sink/spool/quota.go index 1178bcfac8..a4851e4ddb 100644 --- a/downstreamadapter/sink/cloudstorage/spool/quota.go +++ b/pkg/sink/spool/quota.go @@ -16,8 +16,6 @@ package spool import ( "sync" - "github.com/pingcap/ticdc/downstreamadapter/sink/metrics" - "github.com/pingcap/ticdc/pkg/common" "github.com/prometheus/client_golang/prometheus" ) @@ -26,7 +24,7 @@ import ( // state are we in"; this adapter decides how spool reacts to that state. type quotaController struct { // budget owns threshold math and byte accounting. - budget *budget + budget *Budget // postEnqueuePaused is true will hold PostEnqueue callbacks in memory postEnqueuePaused bool @@ -40,31 +38,30 @@ type quotaController struct { metricDiskQuotaWaiters prometheus.Gauge metricDiskQuotaWait prometheus.Observer - keyspace string - changefeed string + closeMetrics func() waitersMu sync.Mutex nextWaiterID uint64 waiters map[uint64]chan struct{} } -func newQuotaController( - changefeedID common.ChangeFeedID, - options *options, -) *quotaController { - keyspace := changefeedID.Keyspace() - changefeed := changefeedID.Name() +func newQuotaController(options *options) *quotaController { + spoolMetrics := normalizeMetrics(options.metrics) controller := "aController{ - keyspace: keyspace, - changefeed: changefeed, - - budget: newBudget(options), - - metricMemoryBytes: metrics.CloudStorageSpoolMemoryBytesGauge.WithLabelValues(keyspace, changefeed), - metricDiskBytes: metrics.CloudStorageSpoolDiskBytesGauge.WithLabelValues(keyspace, changefeed), - metricPendingPostEnqueue: metrics.CloudStoragePendingPostEnqueueGauge.WithLabelValues(keyspace, changefeed), - metricDiskQuotaWaiters: metrics.CloudStorageSpoolDiskQuotaWaitersGauge.WithLabelValues(keyspace, changefeed), - metricDiskQuotaWait: metrics.CloudStorageSpoolDiskQuotaWaitDurationHistogram.WithLabelValues(keyspace, changefeed), + closeMetrics: spoolMetrics.Close, + + budget: NewBudget(Limits{ + DiskQuotaBytes: options.diskQuotaBytes, + MemoryQuotaBytes: int64(float64(options.diskQuotaBytes) * options.memoryRatio), + HighWatermarkBytes: int64(float64(options.diskQuotaBytes) * options.highWatermarkRatio), + LowWatermarkBytes: int64(float64(options.diskQuotaBytes) * options.lowWatermarkRatio), + }), + + metricMemoryBytes: spoolMetrics.MemoryBytes, + metricDiskBytes: spoolMetrics.DiskBytes, + metricPendingPostEnqueue: spoolMetrics.PendingPostEnqueue, + metricDiskQuotaWaiters: spoolMetrics.DiskQuotaWaiters, + metricDiskQuotaWait: spoolMetrics.DiskQuotaWait, waiters: make(map[uint64]chan struct{}), } controller.metricDiskQuotaWaiters.Set(0) @@ -72,15 +69,15 @@ func newQuotaController( } func (q *quotaController) shouldSpill(entryBytes int64) bool { - return q.budget.shouldSpill(entryBytes) + return q.budget.ShouldSpill(entryBytes) } func (q *quotaController) entryExceedsDiskQuota(entryBytes int64) bool { - return q.budget.entryExceedsDiskQuota(entryBytes) + return q.budget.EntryExceedsDiskQuota(entryBytes) } func (q *quotaController) spillWouldExceedDiskQuota(entryBytes int64) bool { - return q.budget.spillWouldExceedDiskQuota(entryBytes) + return q.budget.SpillWouldExceedDiskQuota(entryBytes) } func (q *quotaController) addDiskQuotaWaiter() (uint64, <-chan struct{}) { @@ -111,7 +108,7 @@ func (q *quotaController) acquire( spilled bool, postEnqueue func(), ) func() { - if q.budget.acquire(entryBytes, spilled) { + if q.budget.Acquire(entryBytes, spilled) { q.postEnqueuePaused = true } @@ -130,7 +127,7 @@ func (q *quotaController) acquire( // discarded. It returns all pending PostEnqueue callbacks once local usage has // dropped back to the low watermark. func (q *quotaController) release(entryBytes int64, spilled bool) []func() { - atOrBelowLowWatermark := q.budget.release(entryBytes, spilled) + atOrBelowLowWatermark := q.budget.Release(entryBytes, spilled) if spilled { q.wakeDiskQuotaWaiters() } @@ -151,16 +148,12 @@ func (q *quotaController) release(entryBytes int64, spilled bool) []func() { // deleteMetrics removes per-changefeed label values owned by this adapter. func (q *quotaController) deleteMetrics() { - metrics.CloudStorageSpoolMemoryBytesGauge.DeleteLabelValues(q.keyspace, q.changefeed) - metrics.CloudStorageSpoolDiskBytesGauge.DeleteLabelValues(q.keyspace, q.changefeed) - metrics.CloudStoragePendingPostEnqueueGauge.DeleteLabelValues(q.keyspace, q.changefeed) - metrics.CloudStorageSpoolDiskQuotaWaitersGauge.DeleteLabelValues(q.keyspace, q.changefeed) - metrics.CloudStorageSpoolDiskQuotaWaitDurationHistogram.DeleteLabelValues(q.keyspace, q.changefeed) + q.closeMetrics() } func (q *quotaController) updateMetrics() { - q.metricMemoryBytes.Set(float64(q.budget.memoryBytes)) - q.metricDiskBytes.Set(float64(q.budget.diskBytes)) + q.metricMemoryBytes.Set(float64(q.budget.MemoryBytes())) + q.metricDiskBytes.Set(float64(q.budget.DiskBytes())) q.metricPendingPostEnqueue.Set(float64(len(q.pendingPostEnqueue))) } diff --git a/downstreamadapter/sink/cloudstorage/spool/spool.go b/pkg/sink/spool/spool.go similarity index 90% rename from downstreamadapter/sink/cloudstorage/spool/spool.go rename to pkg/sink/spool/spool.go index ffe5dbe7ca..3731a165bc 100644 --- a/downstreamadapter/sink/cloudstorage/spool/spool.go +++ b/pkg/sink/spool/spool.go @@ -22,7 +22,6 @@ import ( "time" "github.com/pingcap/log" - "github.com/pingcap/ticdc/downstreamadapter/sink/metrics" commonType "github.com/pingcap/ticdc/pkg/common" "github.com/pingcap/ticdc/pkg/config" "github.com/pingcap/ticdc/pkg/errors" @@ -79,6 +78,8 @@ type options struct { highWatermarkRatio float64 // lowWatermarkRatio is the ratio that resumes pending PostEnqueue callbacks. lowWatermarkRatio float64 + + metrics *Metrics } type option func(*options) @@ -179,19 +180,31 @@ func WithLowWatermarkRatio(lowWatermarkRatio float64) option { } } +// Metrics contains component-owned metric handles updated by a spool. +type Metrics struct { + MemoryBytes prometheus.Gauge + DiskBytes prometheus.Gauge + PendingPostEnqueue prometheus.Gauge + DiskQuotaWaiters prometheus.Gauge + DiskQuotaWait prometheus.Observer + LoadedBytes prometheus.Observer + RotatedCount prometheus.Counter + SegmentCount prometheus.Gauge + Close func() +} + +// WithMetrics supplies component-owned metrics to the shared spool. +func WithMetrics(metrics *Metrics) option { + return func(options *options) { + options.metrics = metrics + } +} + type segmentID uint64 -// Spool keeps encoded DML messages after a writer shard has accepted them and -// before that writer shard has flushed them to external storage. -// -// The producer is the cloud storage writer path: after encoderGroup has -// produced encoded messages for a task, writer.Enqueue calls Spool.Enqueue to -// hand those messages to local spool storage. -// -// The consumer is also the cloud storage writer path: when the writer flushes a -// batch, it calls Spool.Load to read the queued messages back, then calls -// Spool.Release after a successful flush or Spool.Discard when the batch is -// ignored. +// Spool keeps encoded sink messages after the sink has accepted them and before +// it has flushed them to external storage. A sink releases an entry only after +// a successful flush, or discards it when the corresponding data is ignored. type Spool struct { keyspace string changefeed string @@ -326,20 +339,55 @@ func New( keyspace = changefeedID.Keyspace() changefeed = changefeedID.Name() ) + spoolMetrics := normalizeMetrics(cfg.metrics) spool := &Spool{ keyspace: keyspace, changefeed: changefeed, workDir: workDir, - quota: newQuotaController(changefeedID, cfg), + quota: newQuotaController(cfg), segmentCapacity: cfg.segmentCapacity, - metricLoadedBytes: metrics.CloudStorageLoadBytesHistogram.WithLabelValues(keyspace, changefeed), - metricRotatedCount: metrics.CloudStorageRotateCountCounter.WithLabelValues(keyspace, changefeed), - metricSegmentCount: metrics.CloudStorageSpoolSegmentCountGauge.WithLabelValues(keyspace, changefeed), + metricLoadedBytes: spoolMetrics.LoadedBytes, + metricRotatedCount: spoolMetrics.RotatedCount, + metricSegmentCount: spoolMetrics.SegmentCount, segments: make(map[segmentID]*segment), } return spool, nil } +func normalizeMetrics(spoolMetrics *Metrics) *Metrics { + if spoolMetrics == nil { + spoolMetrics = &Metrics{} + } + if spoolMetrics.MemoryBytes == nil { + spoolMetrics.MemoryBytes = prometheus.NewGauge(prometheus.GaugeOpts{}) + } + if spoolMetrics.DiskBytes == nil { + spoolMetrics.DiskBytes = prometheus.NewGauge(prometheus.GaugeOpts{}) + } + if spoolMetrics.PendingPostEnqueue == nil { + spoolMetrics.PendingPostEnqueue = prometheus.NewGauge(prometheus.GaugeOpts{}) + } + if spoolMetrics.DiskQuotaWaiters == nil { + spoolMetrics.DiskQuotaWaiters = prometheus.NewGauge(prometheus.GaugeOpts{}) + } + if spoolMetrics.DiskQuotaWait == nil { + spoolMetrics.DiskQuotaWait = prometheus.NewHistogram(prometheus.HistogramOpts{}) + } + if spoolMetrics.LoadedBytes == nil { + spoolMetrics.LoadedBytes = prometheus.NewHistogram(prometheus.HistogramOpts{}) + } + if spoolMetrics.RotatedCount == nil { + spoolMetrics.RotatedCount = prometheus.NewCounter(prometheus.CounterOpts{}) + } + if spoolMetrics.SegmentCount == nil { + spoolMetrics.SegmentCount = prometheus.NewGauge(prometheus.GaugeOpts{}) + } + if spoolMetrics.Close == nil { + spoolMetrics.Close = func() {} + } + return spoolMetrics +} + func defaultOptions() *options { return &options{ diskQuotaBytes: defaultDiskQuotaBytes, @@ -678,9 +726,6 @@ func (s *Spool) Close() { zap.String("keyspace", s.keyspace), zap.String("changefeed", s.changefeed), zap.String("path", s.workDir), zap.Error(err)) } - metrics.CloudStorageLoadBytesHistogram.DeleteLabelValues(s.keyspace, s.changefeed) - metrics.CloudStorageRotateCountCounter.DeleteLabelValues(s.keyspace, s.changefeed) - metrics.CloudStorageSpoolSegmentCountGauge.DeleteLabelValues(s.keyspace, s.changefeed) s.quota.deleteMetrics() } diff --git a/downstreamadapter/sink/cloudstorage/spool/spool_test.go b/pkg/sink/spool/spool_test.go similarity index 97% rename from downstreamadapter/sink/cloudstorage/spool/spool_test.go rename to pkg/sink/spool/spool_test.go index 4c230305fc..62aa8e727f 100644 --- a/downstreamadapter/sink/cloudstorage/spool/spool_test.go +++ b/pkg/sink/spool/spool_test.go @@ -201,9 +201,10 @@ func TestNewUsesDefaultOptionsWhenValuesAreMissing(t *testing.T) { require.NoError(t, err) require.NotNil(t, manager) require.Equal(t, defaultSegmentCapacity, manager.segmentCapacity) - require.Equal(t, int64(float64(expectedQuotaBytes)*defaultMemoryRatio), manager.quota.budget.memoryQuotaBytes) - require.Equal(t, int64(float64(expectedQuotaBytes)*defaultHighWatermarkRatio), manager.quota.budget.highWatermarkBytes) - require.Equal(t, int64(float64(expectedQuotaBytes)*defaultLowWatermarkRatio), manager.quota.budget.lowWatermarkBytes) + limits := manager.quota.budget.Limits() + require.Equal(t, int64(float64(expectedQuotaBytes)*defaultMemoryRatio), limits.MemoryQuotaBytes) + require.Equal(t, int64(float64(expectedQuotaBytes)*defaultHighWatermarkRatio), limits.HighWatermarkBytes) + require.Equal(t, int64(float64(expectedQuotaBytes)*defaultLowWatermarkRatio), limits.LowWatermarkBytes) manager.Close() } @@ -272,9 +273,10 @@ func TestNewSanitizesInvalidOptions(t *testing.T) { require.NoError(t, err) require.NotNil(t, manager) require.Equal(t, defaultSegmentCapacity, manager.segmentCapacity) - require.Equal(t, int64(float64(expectedQuotaBytes)*defaultMemoryRatio), manager.quota.budget.memoryQuotaBytes) - require.Equal(t, int64(float64(expectedQuotaBytes)*defaultHighWatermarkRatio), manager.quota.budget.highWatermarkBytes) - require.Equal(t, int64(float64(expectedQuotaBytes)*defaultLowWatermarkRatio), manager.quota.budget.lowWatermarkBytes) + limits := manager.quota.budget.Limits() + require.Equal(t, int64(float64(expectedQuotaBytes)*defaultMemoryRatio), limits.MemoryQuotaBytes) + require.Equal(t, int64(float64(expectedQuotaBytes)*defaultHighWatermarkRatio), limits.HighWatermarkBytes) + require.Equal(t, int64(float64(expectedQuotaBytes)*defaultLowWatermarkRatio), limits.LowWatermarkBytes) manager.Close() } @@ -295,9 +297,10 @@ func TestNewAppliesFunctionalOptions(t *testing.T) { require.NoError(t, err) require.Equal(t, filepath.Join(baseDir, changefeedID.Keyspace(), changefeedID.Name()), manager.workDir) require.Equal(t, int64(4096), manager.segmentCapacity) - require.Equal(t, int64(512), manager.quota.budget.memoryQuotaBytes) - require.Equal(t, int64(1536), manager.quota.budget.highWatermarkBytes) - require.Equal(t, int64(1024), manager.quota.budget.lowWatermarkBytes) + limits := manager.quota.budget.Limits() + require.Equal(t, int64(512), limits.MemoryQuotaBytes) + require.Equal(t, int64(1536), limits.HighWatermarkBytes) + require.Equal(t, int64(1024), limits.LowWatermarkBytes) manager.Close() } From ec06d23e7dd938b57e898e27053f3ae35d50abe6 Mon Sep 17 00:00:00 2001 From: wk989898 Date: Wed, 12 Aug 2026 09:16:55 +0000 Subject: [PATCH 2/2] fmt Signed-off-by: wk989898 --- pkg/redo/writer/file/file_test.go | 2 +- pkg/sink/spool/spool.go | 18 ++++++++---------- 2 files changed, 9 insertions(+), 11 deletions(-) diff --git a/pkg/redo/writer/file/file_test.go b/pkg/redo/writer/file/file_test.go index 37d2dc2aa4..b48f6cb328 100644 --- a/pkg/redo/writer/file/file_test.go +++ b/pkg/redo/writer/file/file_test.go @@ -462,7 +462,7 @@ func TestRunFlushesOnIntervalAndExecutesPostFlush(t *testing.T) { postFlushCnt := atomic.NewInt64(0) const eventCount = 3 - for i := 0; i < eventCount; i++ { + for i := range eventCount { ts := uint64(i + 1) w.GetInputCh() <- &pevent.RedoRowEvent{ StartTs: ts, diff --git a/pkg/sink/spool/spool.go b/pkg/sink/spool/spool.go index 3731a165bc..10c0a9f58a 100644 --- a/pkg/sink/spool/spool.go +++ b/pkg/sink/spool/spool.go @@ -82,15 +82,13 @@ type options struct { metrics *Metrics } -type option func(*options) - -func WithRootDir(rootDir string) option { +func WithRootDir(rootDir string) func(*options) { return func(options *options) { options.rootDir = rootDir } } -func WithDiskQuotaBytes(quotaBytes int64) option { +func WithDiskQuotaBytes(quotaBytes int64) func(*options) { return func(options *options) { if quotaBytes == 0 { return @@ -108,7 +106,7 @@ func WithDiskQuotaBytes(quotaBytes int64) option { } } -func WithSegmentBytes(segmentBytes int64) option { +func WithSegmentBytes(segmentBytes int64) func(*options) { return func(options *options) { if segmentBytes == 0 { return @@ -126,7 +124,7 @@ func WithSegmentBytes(segmentBytes int64) option { } } -func WithMemoryRatio(memoryRatio float64) option { +func WithMemoryRatio(memoryRatio float64) func(*options) { return func(options *options) { if memoryRatio == 0 { return @@ -144,7 +142,7 @@ func WithMemoryRatio(memoryRatio float64) option { } } -func WithHighWatermarkRatio(highWatermarkRatio float64) option { +func WithHighWatermarkRatio(highWatermarkRatio float64) func(*options) { return func(options *options) { if highWatermarkRatio == 0 { return @@ -162,7 +160,7 @@ func WithHighWatermarkRatio(highWatermarkRatio float64) option { } } -func WithLowWatermarkRatio(lowWatermarkRatio float64) option { +func WithLowWatermarkRatio(lowWatermarkRatio float64) func(*options) { return func(options *options) { if lowWatermarkRatio == 0 { return @@ -194,7 +192,7 @@ type Metrics struct { } // WithMetrics supplies component-owned metrics to the shared spool. -func WithMetrics(metrics *Metrics) option { +func WithMetrics(metrics *Metrics) func(*options) { return func(options *options) { options.metrics = metrics } @@ -321,7 +319,7 @@ func (e *Entry) InMemory() bool { // New return a spool that manages unflushed data. func New( changefeedID commonType.ChangeFeedID, - opts ...option, + opts ...func(*options), ) (*Spool, error) { cfg := defaultOptions() for _, opt := range opts {