Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
52 commits
Select commit Hold shift + click to select a range
84dced1
improve memory control v1
asddongmen Jan 20, 2026
1d6d73c
fix bug
asddongmen Jan 20, 2026
381e57d
add debug log
asddongmen Jan 21, 2026
d2852f6
add debug log 2
asddongmen Jan 21, 2026
f1f9fe1
add debug log 3
asddongmen Jan 21, 2026
62e1546
add debug log 4
asddongmen Jan 21, 2026
f714b8d
add debug log 5
asddongmen Jan 21, 2026
2f04643
adjust = 0.7
asddongmen Jan 21, 2026
75b2d16
adjust 5
asddongmen Jan 21, 2026
8dcb0a2
for debug
asddongmen Jan 21, 2026
994ef9b
adjust 7
asddongmen Jan 22, 2026
4073610
fix bug
asddongmen Jan 22, 2026
87c5e6f
fix bug 2
asddongmen Jan 22, 2026
fcf76f2
adjust 10
asddongmen Jan 22, 2026
0853113
adjust 11
asddongmen Jan 22, 2026
9139d93
allow scan interval grow larger
asddongmen Jan 23, 2026
274878f
allow scan interval grow larger 2
asddongmen Jan 23, 2026
36be880
remove verbose log
asddongmen Jan 23, 2026
96b9ae8
fix bug
asddongmen Jan 23, 2026
fc271cd
Merge remote-tracking branch 'upstream/master' into 0119-refine-dispa…
asddongmen Jan 26, 2026
c9df6b0
comment hardcode scan interval
asddongmen Jan 26, 2026
2ebeef6
add cool down
asddongmen Jan 26, 2026
abc0015
add metrics
asddongmen Jan 26, 2026
3d506ab
adjust scan window algorithm
asddongmen Jan 26, 2026
a0febab
adjust scan window algorithm 2
asddongmen Jan 26, 2026
460b74d
adjust scan window algorithm 3
asddongmen Jan 27, 2026
f8c34a1
add comment
asddongmen Feb 2, 2026
1a86e53
remove useless codes
asddongmen Feb 2, 2026
6c6babf
add ddl workload
asddongmen Feb 2, 2026
6464520
remove useless code
asddongmen Feb 2, 2026
8ff0d67
remove useless code
asddongmen Feb 4, 2026
7902d26
remove useless code 3
asddongmen Feb 4, 2026
8602eb6
add debug log
asddongmen Feb 5, 2026
ee4f5ff
add debug sleep for ddl
asddongmen Feb 6, 2026
606b70c
Merge remote-tracking branch 'upstream/master' into 0119-refine-dispa…
asddongmen Feb 10, 2026
d33e238
add debug log
asddongmen Feb 13, 2026
4f82641
fix
asddongmen Feb 13, 2026
d4ec73c
remove debug code
asddongmen Feb 24, 2026
47f7b12
Merge remote-tracking branch 'upstream/master' into 0119-refine-dispa…
asddongmen Feb 25, 2026
f6ddfdc
improve code
asddongmen Feb 25, 2026
c6f9f05
improve code 2
asddongmen Feb 25, 2026
3ebfaa3
fix make check
asddongmen Feb 25, 2026
02d7550
Merge remote-tracking branch 'upstream/master' into 0119-refine-dispa…
asddongmen Feb 26, 2026
2557a99
metrics: add scan window related panel
asddongmen Feb 26, 2026
272746f
metrics: add scan window related panel 2
asddongmen Feb 26, 2026
5e9c211
resolve comment
asddongmen Feb 27, 2026
9842dce
eventservice: remove syncpoint enable
asddongmen Feb 28, 2026
aabfa67
adjust for test
asddongmen Feb 28, 2026
7726f51
resolve comment 3
asddongmen Feb 28, 2026
5171bd6
resolve comment 4
asddongmen Mar 2, 2026
715175f
resolve comment 5
asddongmen Mar 2, 2026
1b161d3
fix panic
asddongmen Mar 2, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
81 changes: 67 additions & 14 deletions downstreamadapter/eventcollector/event_collector.go
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,7 @@ type changefeedStat struct {
metricMemoryUsageMaxRedo prometheus.Gauge
metricMemoryUsageUsedRedo prometheus.Gauge
dispatcherCount atomic.Int32
memoryReleaseCount atomic.Uint32
}

func newChangefeedStat(changefeedID common.ChangeFeedID) *changefeedStat {
Expand Down Expand Up @@ -421,11 +422,17 @@ func (c *EventCollector) processDSFeedback(ctx context.Context) error {
return context.Cause(ctx)
case feedback := <-c.ds.Feedback():
if feedback.FeedbackType == dynstream.ReleasePath {
if v, ok := c.changefeedMap.Load(feedback.Area); ok {
v.(*changefeedStat).memoryReleaseCount.Add(1)
}
log.Info("release dispatcher memory in DS", zap.Any("dispatcherID", feedback.Path))
c.ds.Release(feedback.Path)
}
case feedback := <-c.redoDs.Feedback():
if feedback.FeedbackType == dynstream.ReleasePath {
if v, ok := c.changefeedMap.Load(feedback.Area); ok {
v.(*changefeedStat).memoryReleaseCount.Add(1)
}
log.Info("release dispatcher memory in redo DS", zap.Any("dispatcherID", feedback.Path))
c.redoDs.Release(feedback.Path)
}
Expand Down Expand Up @@ -597,9 +604,24 @@ func (c *EventCollector) controlCongestion(ctx context.Context) error {
}

func (c *EventCollector) newCongestionControlMessages() map[node.ID]*event.CongestionControl {
changefeedMemoryReleaseCount := make(map[common.ChangeFeedID]uint32)
getAndResetMemoryReleaseCount := func(changefeedID common.ChangeFeedID) uint32 {
if count, ok := changefeedMemoryReleaseCount[changefeedID]; ok {
return count
}
v, ok := c.changefeedMap.Load(changefeedID.ID())
if !ok {
return 0
}
count := v.(*changefeedStat).memoryReleaseCount.Swap(0)
changefeedMemoryReleaseCount[changefeedID] = count
return count
}

// collect path-level available memory and total available memory for each changefeed
changefeedPathMemory := make(map[common.ChangeFeedID]map[common.DispatcherID]uint64)
changefeedTotalMemory := make(map[common.ChangeFeedID]uint64)
changefeedUsageRatio := make(map[common.ChangeFeedID]float64)

// collect from main dynamic stream
for _, quota := range c.ds.GetMetrics().MemoryControl.AreaMemoryMetrics {
Expand All @@ -617,6 +639,7 @@ func (c *EventCollector) newCongestionControlMessages() map[node.ID]*event.Conge
}
// store total available memory from AreaMemoryMetric
changefeedTotalMemory[cfID] = uint64(quota.AvailableMemory())
changefeedUsageRatio[cfID] = calcUsageRatio(quota.MemoryUsage(), quota.MaxMemory())
}

// collect from redo dynamic stream and take minimum
Expand All @@ -638,11 +661,9 @@ func (c *EventCollector) newCongestionControlMessages() map[node.ID]*event.Conge
}
}
// take minimum total available memory between main and redo streams
if existing, exists := changefeedTotalMemory[cfID]; exists {
changefeedTotalMemory[cfID] = min(existing, uint64(quota.AvailableMemory()))
} else {
changefeedTotalMemory[cfID] = uint64(quota.AvailableMemory())
}
updateMinUint64MapValue(changefeedTotalMemory, cfID, uint64(quota.AvailableMemory()))
// take maximum usage ratio between main and redo streams
changefeedUsageRatio[cfID] = max(changefeedUsageRatio[cfID], calcUsageRatio(quota.MemoryUsage(), quota.MaxMemory()))
}

if len(changefeedPathMemory) == 0 {
Expand Down Expand Up @@ -679,32 +700,64 @@ func (c *EventCollector) newCongestionControlMessages() map[node.ID]*event.Conge
// build congestion control messages for each node
result := make(map[node.ID]*event.CongestionControl)
for nodeID, changefeedDispatchers := range nodeDispatcherMemory {
congestionControl := event.NewCongestionControl()
congestionControl := event.NewCongestionControlWithVersion(event.CongestionControlVersion2)

for changefeedID, dispatcherMemory := range changefeedDispatchers {
if len(dispatcherMemory) == 0 {
continue
}

// get total available memory directly from AreaMemoryMetric
totalAvailable := uint64(changefeedTotalMemory[changefeedID])
if totalAvailable > 0 {
congestionControl.AddAvailableMemoryWithDispatchers(
changefeedID.ID(),
totalAvailable,
dispatcherMemory,
)
totalAvailable, ok := changefeedTotalMemory[changefeedID]
if !ok {
continue
}
congestionControl.AddAvailableMemoryWithDispatchersAndUsageAndReleaseCount(
changefeedID.ID(),
totalAvailable,
changefeedUsageRatio[changefeedID],
dispatcherMemory,
getAndResetMemoryReleaseCount(changefeedID),
)
}

if len(congestionControl.GetAvailables()) > 0 {
result[nodeID] = congestionControl
}
}

return result
}

func updateMinUint64MapValue(m map[common.ChangeFeedID]uint64, key common.ChangeFeedID, value uint64) {
if existing, exists := m[key]; exists {
m[key] = min(existing, value)
} else {
m[key] = value
}
}

func updateMaxUint64MapValue(m map[common.ChangeFeedID]uint64, key common.ChangeFeedID, value uint64) {
if existing, exists := m[key]; exists {
m[key] = max(existing, value)
} else {
m[key] = value
}
}

func calcUsageRatio(usedMemory int64, maxMemory int64) float64 {
if maxMemory <= 0 {
return 0
}
ratio := float64(usedMemory) / float64(maxMemory)
if ratio < 0 {
return 0
}
if ratio > 1 {
return 1
}
return ratio
}

func (c *EventCollector) updateMetrics(ctx context.Context) error {
ticker := time.NewTicker(5 * time.Second)
defer ticker.Stop()
Expand Down
Loading