Implementation:Apache Flink CongestionControlRateLimitingStrategy Builder
| Knowledge Sources | |
|---|---|
| Domains | Stream_Processing, Flow_Control |
| Last Updated | 2026-02-09 00:00 GMT |
Overview
Concrete tool for configuring TCP-inspired congestion control with AIMD scaling for async sink rate limiting provided by the Apache Flink connector-base module.
Description
The CongestionControlRateLimitingStrategy implements RateLimitingStrategy with three key methods: shouldBlock (checks if new requests should wait), registerInFlightRequest (tracks new submissions), and registerCompletedRequest (adjusts rate based on success/failure). It uses a pluggable ScalingStrategy<Integer> — the default AIMDScalingStrategy increases the rate linearly on success and decreases multiplicatively on failure.
Usage
This is the default rate limiting strategy for async sinks. Configure it through the builder to tune the initial rate, maximum concurrent requests, and scaling parameters.
Code Reference
Source Location
- Repository: Apache Flink
- File: flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/sink/writer/strategy/CongestionControlRateLimitingStrategy.java
- Lines: L35-124
- File: flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/sink/writer/strategy/AIMDScalingStrategy.java
- Lines: L28-85
Signature
@PublicEvolving
public class CongestionControlRateLimitingStrategy implements RateLimitingStrategy {
public static CongestionControlRateLimitingStrategyBuilder builder();
public boolean shouldBlock(RequestInfo requestInfo);
public void registerInFlightRequest(RequestInfo requestInfo);
public void registerCompletedRequest(ResultInfo resultInfo);
public int getMaxBatchSize();
public static class CongestionControlRateLimitingStrategyBuilder {
public CongestionControlRateLimitingStrategyBuilder setMaxInFlightRequests(int v);
public CongestionControlRateLimitingStrategyBuilder setInitialMaxInFlightMessages(int v);
public CongestionControlRateLimitingStrategyBuilder setScalingStrategy(ScalingStrategy<Integer> v);
public CongestionControlRateLimitingStrategy build();
}
}
@PublicEvolving
public class AIMDScalingStrategy implements ScalingStrategy<Integer> {
public AIMDScalingStrategy(int increaseRate, double decreaseFactor, int rateThreshold);
public Integer scaleUp(Integer currentRate);
public Integer scaleDown(Integer currentRate);
public static AIMDScalingStrategyBuilder builder(int rateThreshold);
}
Import
import org.apache.flink.connector.base.sink.writer.strategy.CongestionControlRateLimitingStrategy;
import org.apache.flink.connector.base.sink.writer.strategy.AIMDScalingStrategy;
I/O Contract
Inputs
| Name | Type | Required | Description |
|---|---|---|---|
| maxInFlightRequests | int | Yes | Absolute cap on concurrent requests |
| initialMaxInFlightMessages | int | Yes | Starting concurrency level |
| scalingStrategy | ScalingStrategy<Integer> | Yes | AIMD or custom scaling |
Outputs
| Name | Type | Description |
|---|---|---|
| strategy | CongestionControlRateLimitingStrategy | Configured rate limiter |
Usage Examples
Configuring Rate Limiting
CongestionControlRateLimitingStrategy rateLimiter =
CongestionControlRateLimitingStrategy.builder()
.setMaxInFlightRequests(10)
.setInitialMaxInFlightMessages(100)
.setScalingStrategy(
AIMDScalingStrategy.builder(10)
.setIncreaseRate(10)
.setDecreaseFactor(0.5)
.build())
.build();