From 6b2e3ece698e84d8dde915bd85d298bb0e245c27 Mon Sep 17 00:00:00 2001 From: cnderrauber Date: Wed, 23 Sep 2026 13:54:40 +0800 Subject: [PATCH] Increase receiver loadbalance threshold when batch io enabled (#4899) * Increase receiver loadbalance threshold when batch io enabled With batch IO the send syscall moves to the BatchConn flush goroutine, so WriteRTP only translates the packet and enqueues it. Increash the threshold to reduce * lint --- pkg/rtc/mediatrack.go | 15 ++++++++++++++- pkg/rtc/participant.go | 2 ++ pkg/service/roommanager.go | 3 +++ 3 files changed, 19 insertions(+), 1 deletion(-) diff --git a/pkg/rtc/mediatrack.go b/pkg/rtc/mediatrack.go index 3c167e0f2..be1436fce 100644 --- a/pkg/rtc/mediatrack.go +++ b/pkg/rtc/mediatrack.go @@ -95,6 +95,7 @@ type MediaTrackParams struct { EnableRTPStreamRestartDetection bool UpdateTrackInfoByVideoSizeChange bool ForceBackupCodecPolicySimulcast bool + MediaBatchIOEnabled bool OnSubscribedMaxQualityChange func( trackID livekit.TrackID, trackInfo *livekit.TrackInfo, @@ -272,6 +273,18 @@ func (t *MediaTrack) ToProto() *livekit.TrackInfo { return t.MediaTrackReceiver.TrackInfoClone() } +func (t *MediaTrack) getReceiverLBThreshold() int { + // With batch IO the send syscall moves to the BatchConn flush goroutine, so WriteRTP only + // translates the packet and enqueues it: ~3.3us per downtrack, against ~26us when the send runs + // inline. The last subscriber of 300 subscribers at a serial forwarder is enqueued ~1ms + // after the first. + if t.params.MediaBatchIOEnabled { + return 300 + } + + return 20 +} + // AddReceiver adds a new RTP receiver to the track, returns true when receiver represents a new codec // and if a receiver was added successfully func (t *MediaTrack) AddReceiver(receiver *webrtc.RTPReceiver, track sfu.TrackRemote, mid string) (bool, bool) { @@ -401,7 +414,7 @@ func (t *MediaTrack) AddReceiver(receiver *webrtc.RTPReceiver, track sfu.TrackRe t.params.VideoConfig.StreamTrackerManager, sfu.WithPliThrottleConfig(t.params.PLIThrottleConfig), sfu.WithAudioConfig(t.params.AudioConfig), - sfu.WithLoadBalanceThreshold(20), + sfu.WithLoadBalanceThreshold(t.getReceiverLBThreshold()), sfu.WithForwardStats(t.params.ForwardStats), sfu.WithEnableRTPStreamRestartDetection(t.params.EnableRTPStreamRestartDetection), ) diff --git a/pkg/rtc/participant.go b/pkg/rtc/participant.go index b0a663f27..a63529fb5 100644 --- a/pkg/rtc/participant.go +++ b/pkg/rtc/participant.go @@ -241,6 +241,7 @@ type ParticipantParams struct { MigrationWaitDuration time.Duration ExcludeIPv6LocalCandidates bool EnableWarp bool + MediaBatchIOEnabled bool } type ParticipantImpl struct { @@ -3573,6 +3574,7 @@ func (p *ParticipantImpl) addMediaTrack(signalCid string, ti *livekit.TrackInfo) EnableRTPStreamRestartDetection: p.params.EnableRTPStreamRestartDetection, UpdateTrackInfoByVideoSizeChange: p.params.UseOneShotSignallingMode, ForceBackupCodecPolicySimulcast: p.params.ForceBackupCodecPolicySimulcast, + MediaBatchIOEnabled: p.params.MediaBatchIOEnabled, OnSubscribedMaxQualityChange: p.onSubscribedMaxQualityChange, OnSubscribedAudioCodecChange: p.onSubscribedAudioCodecChange, }, ti) diff --git a/pkg/service/roommanager.go b/pkg/service/roommanager.go index 5aace9845..b667344f3 100644 --- a/pkg/service/roommanager.go +++ b/pkg/service/roommanager.go @@ -543,6 +543,9 @@ func (r *RoomManager) StartSession( EnableParticipantDataBlob: r.config.EnableParticipantDataBlob, EnableRTPStreamRestartDetection: r.config.RTC.EnableRTPStreamRestartDetection, EnableWarp: enableWarp, + MediaBatchIOEnabled: r.config.RTC.BatchIO.BatchSize > 0 && + (r.config.RTC.ICEPortRangeStart == 0 || r.config.RTC.ICEPortRangeEnd == 0) && + r.config.RTC.UDPPort.Valid(), }) if err != nil { prometheus.IncrementParticipantRtcCanceled(1, enableWarp, pi.Client.GetSdk())