@@ -83,6 +83,7 @@ type changefeedStat struct {
8383 metricMemoryUsageMaxRedo prometheus.Gauge
8484 metricMemoryUsageUsedRedo prometheus.Gauge
8585 dispatcherCount atomic.Int32
86+ memoryReleaseCount atomic.Uint32
8687}
8788
8889func newChangefeedStat (changefeedID common.ChangeFeedID ) * changefeedStat {
@@ -421,11 +422,17 @@ func (c *EventCollector) processDSFeedback(ctx context.Context) error {
421422 return context .Cause (ctx )
422423 case feedback := <- c .ds .Feedback ():
423424 if feedback .FeedbackType == dynstream .ReleasePath {
425+ if v , ok := c .changefeedMap .Load (feedback .Area ); ok {
426+ v .(* changefeedStat ).memoryReleaseCount .Add (1 )
427+ }
424428 log .Info ("release dispatcher memory in DS" , zap .Any ("dispatcherID" , feedback .Path ))
425429 c .ds .Release (feedback .Path )
426430 }
427431 case feedback := <- c .redoDs .Feedback ():
428432 if feedback .FeedbackType == dynstream .ReleasePath {
433+ if v , ok := c .changefeedMap .Load (feedback .Area ); ok {
434+ v .(* changefeedStat ).memoryReleaseCount .Add (1 )
435+ }
429436 log .Info ("release dispatcher memory in redo DS" , zap .Any ("dispatcherID" , feedback .Path ))
430437 c .redoDs .Release (feedback .Path )
431438 }
@@ -597,9 +604,24 @@ func (c *EventCollector) controlCongestion(ctx context.Context) error {
597604}
598605
599606func (c * EventCollector ) newCongestionControlMessages () map [node.ID ]* event.CongestionControl {
607+ changefeedMemoryReleaseCount := make (map [common.ChangeFeedID ]uint32 )
608+ getAndResetMemoryReleaseCount := func (changefeedID common.ChangeFeedID ) uint32 {
609+ if count , ok := changefeedMemoryReleaseCount [changefeedID ]; ok {
610+ return count
611+ }
612+ v , ok := c .changefeedMap .Load (changefeedID .ID ())
613+ if ! ok {
614+ return 0
615+ }
616+ count := v .(* changefeedStat ).memoryReleaseCount .Swap (0 )
617+ changefeedMemoryReleaseCount [changefeedID ] = count
618+ return count
619+ }
620+
600621 // collect path-level available memory and total available memory for each changefeed
601622 changefeedPathMemory := make (map [common.ChangeFeedID ]map [common.DispatcherID ]uint64 )
602623 changefeedTotalMemory := make (map [common.ChangeFeedID ]uint64 )
624+ changefeedUsageRatio := make (map [common.ChangeFeedID ]float64 )
603625
604626 // collect from main dynamic stream
605627 for _ , quota := range c .ds .GetMetrics ().MemoryControl .AreaMemoryMetrics {
@@ -617,6 +639,7 @@ func (c *EventCollector) newCongestionControlMessages() map[node.ID]*event.Conge
617639 }
618640 // store total available memory from AreaMemoryMetric
619641 changefeedTotalMemory [cfID ] = uint64 (quota .AvailableMemory ())
642+ changefeedUsageRatio [cfID ] = calcUsageRatio (quota .MemoryUsage (), quota .MaxMemory ())
620643 }
621644
622645 // collect from redo dynamic stream and take minimum
@@ -638,11 +661,9 @@ func (c *EventCollector) newCongestionControlMessages() map[node.ID]*event.Conge
638661 }
639662 }
640663 // take minimum total available memory between main and redo streams
641- if existing , exists := changefeedTotalMemory [cfID ]; exists {
642- changefeedTotalMemory [cfID ] = min (existing , uint64 (quota .AvailableMemory ()))
643- } else {
644- changefeedTotalMemory [cfID ] = uint64 (quota .AvailableMemory ())
645- }
664+ updateMinUint64MapValue (changefeedTotalMemory , cfID , uint64 (quota .AvailableMemory ()))
665+ // take maximum usage ratio between main and redo streams
666+ changefeedUsageRatio [cfID ] = max (changefeedUsageRatio [cfID ], calcUsageRatio (quota .MemoryUsage (), quota .MaxMemory ()))
646667 }
647668
648669 if len (changefeedPathMemory ) == 0 {
@@ -679,32 +700,64 @@ func (c *EventCollector) newCongestionControlMessages() map[node.ID]*event.Conge
679700 // build congestion control messages for each node
680701 result := make (map [node.ID ]* event.CongestionControl )
681702 for nodeID , changefeedDispatchers := range nodeDispatcherMemory {
682- congestionControl := event .NewCongestionControl ( )
703+ congestionControl := event .NewCongestionControlWithVersion ( event . CongestionControlVersion2 )
683704
684705 for changefeedID , dispatcherMemory := range changefeedDispatchers {
685706 if len (dispatcherMemory ) == 0 {
686707 continue
687708 }
688709
689710 // get total available memory directly from AreaMemoryMetric
690- totalAvailable := uint64 (changefeedTotalMemory [changefeedID ])
691- if totalAvailable > 0 {
692- congestionControl .AddAvailableMemoryWithDispatchers (
693- changefeedID .ID (),
694- totalAvailable ,
695- dispatcherMemory ,
696- )
711+ totalAvailable , ok := changefeedTotalMemory [changefeedID ]
712+ if ! ok {
713+ continue
697714 }
715+ congestionControl .AddAvailableMemoryWithDispatchersAndUsageAndReleaseCount (
716+ changefeedID .ID (),
717+ totalAvailable ,
718+ changefeedUsageRatio [changefeedID ],
719+ dispatcherMemory ,
720+ getAndResetMemoryReleaseCount (changefeedID ),
721+ )
698722 }
699723
700724 if len (congestionControl .GetAvailables ()) > 0 {
701725 result [nodeID ] = congestionControl
702726 }
703727 }
704-
705728 return result
706729}
707730
731+ func updateMinUint64MapValue (m map [common.ChangeFeedID ]uint64 , key common.ChangeFeedID , value uint64 ) {
732+ if existing , exists := m [key ]; exists {
733+ m [key ] = min (existing , value )
734+ } else {
735+ m [key ] = value
736+ }
737+ }
738+
739+ func updateMaxUint64MapValue (m map [common.ChangeFeedID ]uint64 , key common.ChangeFeedID , value uint64 ) {
740+ if existing , exists := m [key ]; exists {
741+ m [key ] = max (existing , value )
742+ } else {
743+ m [key ] = value
744+ }
745+ }
746+
747+ func calcUsageRatio (usedMemory int64 , maxMemory int64 ) float64 {
748+ if maxMemory <= 0 {
749+ return 0
750+ }
751+ ratio := float64 (usedMemory ) / float64 (maxMemory )
752+ if ratio < 0 {
753+ return 0
754+ }
755+ if ratio > 1 {
756+ return 1
757+ }
758+ return ratio
759+ }
760+
708761func (c * EventCollector ) updateMetrics (ctx context.Context ) error {
709762 ticker := time .NewTicker (5 * time .Second )
710763 defer ticker .Stop ()
0 commit comments