Jump to content

Connect SuperML | Leeroopedia MCP: Equip your AI agents with best practices, code verification, and debugging knowledge. Powered by Leeroo — building Organizational Superintelligence. Contact us at founders@leeroo.com.

Implementation:Apache Flink CongestionControlRateLimitingStrategy Builder

From Leeroopedia


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();

Related Pages

Implements Principle

Page Connections

Double-click a node to navigate. Hold to expand connections.
Principle
Implementation
Heuristic
Environment