본문으로 건너뛰기
김신건의 로그

[Concurrency] Backpressure: 흐름 제어로 시스템 보호

· 수정 · 📖 약 2분 · 825자/단어 #backpressure #flow-control #reactive #concurrency #backend
Backpressure, flow control, reactive streams, rate shaping, buffer overflow, shedding

정의

Backpressure = 생산자가 소비자 속도에 맞춰 자기 속도 조절. 분산 시스템의 cascade failure 방지.

IMPORTANT

Backpressure 의 핵심: “받는 쪽이 처리 못 함” 을 주는 쪽이 알아야 한다. 모르면 버퍼 폭발 + OOM + 시스템 다운.

문제: 비대칭 속도

flowchart LR
    Fast[Fast Producer<br/>10000 msg/s] -->|push| Buf[Unbounded Buffer]
    Buf -->|consume| Slow[Slow Consumer<br/>100 msg/s]
    Buf -.OOM.-> Crash[💥]

Producer-Consumer 에서 Producer 가 항상 빠르면 buffer 가 무한 으로. backpressure 의 본질적 시나리오.

4가지 전략

flowchart TD
    Q[Producer > Consumer]
    Q --> Buf[1. Bounded buffer + Block]
    Q --> Drop["2. Drop (lossy)"]
    Q --> Sample[3. Sampling]
    Q --> Pull["4. Pull-based (demand)"]

1. Bounded Buffer + Block

queue = BoundedQueue(maxsize=1000)

# Producer
queue.put(msg)   # 가득 차면 block 또는 timeout
  • 단순.
  • Producer 가 blocking → 상위 backpressure 자동 전파.

2. Drop (Lossy)

try:
    queue.put_nowait(msg)
except QueueFull:
    pass  # 드롭
  • metrics / log 같은 손실 OK 데이터.
  • Drop policy: oldest, newest, sample.

3. Sampling

if random.random() < 0.1:   # 10% 만 처리
    queue.put(msg)
  • 대규모 telemetry.
  • 통계적으로 유의미.

4. Pull-based (Reactive Streams)

sequenceDiagram
    autonumber
    participant P as Producer
    participant C as Consumer

    C->>P: request(10) (demand)
    P->>C: next(msg) × 10
    C->>P: request(5) (more demand)
    P->>C: next(msg) × 5

Consumer 가 요청한 만큼만 전송. Producer 가 자기 속도 자동 조절.

라이브러리: Project Reactor (Java), RxJava, ReactiveX, Akka Streams.

시스템 레벨 Backpressure

flowchart LR
    LB --> Sv1[Service 1]
    Sv1 --> Sv2[Service 2]
    Sv2 --> DB[(DB)]
    DB -.포화.-> Sv2
    Sv2 -.503.-> Sv1
    Sv1 -.429.-> LB
    LB -.queue.-> Client
레이어신호
LB429 Too Many Requests
Service503 Service Unavailable
DBconnection pool 고갈
NetworkTCP window 축소

HTTP/2, QUIC 의 flow control

flowchart LR
    Sender -->|DATA| Receiver
    Receiver -->|WINDOW_UPDATE +N| Sender
    Note[Sender 는 window 안에서만]

자세한 건 HTTP/2 의 flow control 절.

TCP 의 flow control

TCP 의 sliding window전송 레이어 backpressure. 자세한 건 tcp.

Bulkhead

flowchart TB
    Pool1["Pool A (10 threads)<br/>→ service A"]
    Pool2["Pool B (10 threads)<br/>→ service B"]
    Pool3["Pool C (10 threads)<br/>→ service C"]
    Fail["service A 다운"]
    Fail --> Pool1
    Pool1 -.포화.-> Pool1
    Note[Pool A 만 영향. B/C 는 정상]

Thread pool 을 service 별로 분리 → 한 service 다운이 전체 다운 안 되게.

Shedding (Load Shedding)

flowchart TD
    Q[현재 부하 > 한도?]
    Q -->|예| Shed["일부 요청 거절 (503)"]
    Q -->|아니오| Process[정상 처리]

모든 요청을 느리게 처리 하는 것보다 일부 거절하고 나머지는 빠르게. Google SRE 의 graceful degradation.

Adaptive Concurrency

flowchart LR
    Lat[현재 latency 측정] --> Calc[Little's Law 로<br/>적정 동시성 계산]
    Calc --> Cap[concurrency 한도 조정]
    Cap --> Limit[새 요청 제한]
  • Netflix Concurrency Limits.
  • Envoy adaptive concurrency.
  • latency 가 늘면 동시성 줄임.

흔한 함정

WARNING

  1. Unbounded queue = OOM 보장. 모든 buffer 는 bounded.
  2. Drop 정책 명확하지 않음 = 운영자가 어떤 데이터 사라지는지 모름. 메트릭 필수.
  3. Backpressure 전파 안 됨 = 한 계층만 잡고 에서 압박 지속. cascade fail.
  4. Pull-based 의 너무 큰 request(N) = 사실상 push 와 같음. 적절한 batch size.

Java: Reactive Streams API 구현

Project Reactor (Spring WebFlux 기반) 를 사용한 pull-based backpressure 예시:

Flux<String> upstream = Flux.generate(sink -> sink.next(fetchNextItem()));

upstream
    .onBackpressureBuffer(1000)           // bounded buffer
    .publishOn(Schedulers.boundedElastic())
    .subscribe(new BaseSubscriber<String>() {
        @Override
        protected void hookOnSubscribe(Subscription subscription) {
            request(50);  // 최초 50개 요청
        }

        @Override
        protected void hookOnNext(String value) {
            process(value);
            request(1);   // 처리 완료 후 1개씩 추가 요청
        }
    });

onBackpressureBuffer / onBackpressureDrop / onBackpressureLatest 로 정책 선택:

연산자전략사용 시나리오
onBackpressureBuffer(n)Bounded buffer순간 burst 허용, OOM 방지
onBackpressureDrop()Drop newest실시간 센서 데이터, 손실 OK
onBackpressureLatest()Keep latest onlyUI 화면 갱신, 최신값만 유의미
onBackpressureError()즉시 오류손실 불허 파이프라인

자세히는 spring-webflux 참고.

Kafka Producer 설정

Kafka 에서 producer 의 backpressure 는 max.block.msbuffer.memory 로 제어한다.

# Producer 설정
buffer.memory=33554432        # 총 버퍼 크기 (32 MB)
max.block.ms=60000            # 버퍼가 가득 찼을 때 block 최대 시간
batch.size=16384              # batch 크기
linger.ms=5                   # 배치 대기 시간

버퍼가 가득 차면 KafkaProducer.send()max.block.ms 동안 blocking. 그래도 공간이 안 나면 TimeoutException. 이것이 Kafka producer 레벨의 backpressure. Consumer lag 기반 스케일링은 kafka-consumer-group 참고.

Resilience4j 와의 연계

Backpressure 는 단독으로 사용하기보다 circuit breaker, rate limiter 와 조합한다.

flowchart LR
    Client["클라이언트 요청"] --> RL["Rate Limiter\n초당 N건 제한"]
    RL --> CB["Circuit Breaker\n실패율 임계 초과 시 차단"]
    CB --> BP["Bulkhead\n동시 실행 수 제한"]
    BP --> Svc["Service"]
    Svc -.실패.-> CB
    CB -.open.-> FBK["Fallback 응답"]

Resilience4j 에서 셋을 조합한 예시:

Supplier<String> decorated = Decorators.ofSupplier(() -> remoteService.call())
    .withRateLimiter(rateLimiter)
    .withCircuitBreaker(circuitBreaker)
    .withBulkhead(bulkhead)
    .decorate();

circuit-breakerrate-limiting 함께 설계해야 실질적인 cascade failure 방지가 된다.

관찰 가능성 (Observability)

Backpressure 가 발동하고 있는지 모르면 운영 중에 발견하지 못한다. 필수 메트릭:

메트릭설명경보 기준
queue.size / buffer.used현재 버퍼 점유율> 80%
drop.rate초당 드롭된 메시지 수> 0 (손실 불허 시)
block.timeproducer block 시간> SLO 임계
consumer.lagKafka consumer lag지속 증가
http.server.requests + 429/503 비율HTTP 레이어 shed 비율급증 시 알람

Micrometer + Prometheus + Grafana 스택으로 위 메트릭을 시각화하면 backpressure 발동 시점과 원인을 즉시 파악할 수 있다.

관련 위키

이 글의 용어 (8개)
[Concurrency] Circuit Breaker: cascade failure 방어concurrency
정의 Circuit Breaker = 전기 회로의 차단기처럼, 백엔드 다운 / 느림 시 호출 자체를 차단 → 빠른 실패 + 백엔드 회복 시간. 와 함께 cascade failur…
[Concurrency] Connection Pool: DB, HTTP, gRPC 의 공통 패턴concurrency
정의 Connection Pool = 연결을 미리 만들어 두고 재사용. 매 요청마다 TCP + TLS + auth 비용을 회피. Pool 의 동작 = producer (reque…
[Concurrency] Rate Limiting: token bucket, sliding windowconcurrency
정의 Rate Limiting = 시간당 요청 수 제한. API 보호, fair use, 비용 제어, DDoS 완화. 5가지 알고리즘 1. Fixed Window - 단순. - …
[Concurrency] Retry + Exponential Backoff + Jitterconcurrency
정의 Retry with Exponential Backoff + Jitter = 실패 시 점점 긴 간격으로 재시도, 랜덤 분산. 분산 시스템의 thundering herd / r…
[Distributed] Kafka Consumer Group: rebalancing, offset, lagdistributed-systems
정의 Consumer Group = 같은 를 가진 consumer들이 함께 한 topic을 분담 소비. partition 단위로 분배. 핵심 특성: - 한 partition → …
[Network] HTTP/2: multiplexing, HPACK, server push의 종말network
정의 HTTP/2 (2015, RFC 9113) 는 HTTP/1.1 의 문법은 유지하되 전송을 바이너리 + 멀티플렉싱 으로 바꾼다. Head-of-Line Blocking (HT…
[Spring] WebFlux: Reactive, Mono, Fluxspring
정의 Spring WebFlux는 Project Reactor 기반의 비동기·non-blocking 웹 프레임워크. Spring MVC와 별개. Netty 기본 (Tomcat도 …
TCPnetwork
정의 TCP (Transmission Control Protocol) 는 신뢰성 있는 연결 지향 전송 계층 프로토콜이다. RFC 9293 으로 정의되어 있다. HTTP/1.1·2…

💬 댓글

사이트 검색 / 명령어

검색

스크롤 = 확대/축소 · 드래그 = 이동 · 0 = 원래 크기 · ESC = 닫기