How does consumer-side retry with backoff work in Spring Cloud Stream, and which properties control it?
answer
- RetryTemplate wraps handler, maxAttempts default 3
- exponential: initial 1000ms × 2.0, cap 10000ms
- maxAttempts=1 disables
- blocking + in-memory = head-of-line block, lost on crash
- exhausted → DLQ or discarded; handler must be idempotent
basics
~10 sIf your consumer throws, Spring Cloud Stream retries the same message a few times with growing delays before giving up. You tune it with maxAttempts, backOffInitialInterval, backOffMultiplier, and backOffMaxInterval.
solid answer
~50 sBy default each consumer binding wraps handling in a `RetryTemplate`. On an exception it re-invokes your function up to `maxAttempts` (default 3, meaning 1 original + 2 retries) with **exponential backoff**: first wait `backOffInitialInterval` (1000ms), multiplied by `backOffMultiplier` (2.0) each attempt, capped at `backOffMaxInterval` (10000ms). Set `maxAttempts=1` to disable retry. This retry is **in-memory and blocking** — the same thread sleeps between attempts, so on Kafka it holds up that partition; you don't advance to newer records until this one succeeds or exhausts. You can restrict which exceptions retry via `retryableExceptions` and `defaultRetryable`. When attempts are exhausted, the framework hands off to error handling: if `enableDlq` is on the message goes to the dead-letter topic, otherwise it's logged and (Kafka) the offset is committed, effectively dropping it. For non-blocking retry on Kafka you supply a custom error handler.
code
java · 24 lines// Consumer retry tuning (application.yml shown as reference):
//
// spring.cloud.stream.bindings.process-in-0.group: order-workers
// spring.cloud.stream.bindings.process-in-0.consumer:
// max-attempts: 4 # 1 original + 3 retries; set 1 to DISABLE retry
// back-off-initial-interval: 500
// back-off-multiplier: 2.0 # 500ms, 1000ms, 2000ms ...
// back-off-max-interval: 5000
// default-retryable: false # only retry the listed exceptions
// retryable-exceptions:
// org.springframework.dao.TransientDataAccessException: true
// java.lang.IllegalArgumentException: false # a poison message, don't retry
import java.util.function.Consumer;
import org.springframework.context.annotation.Bean;
public class Handlers {
@Bean
public Consumer<Order> process(OrderService svc) {
// Throwing a TransientDataAccessException -> retried with backoff.
// Throwing IllegalArgumentException -> not retried -> straight to DLQ/discard.
return svc::handle;
}
}go deeper
Name the four backoff properties and that maxAttempts=1 disables retry.
Explain exponential backoff math and that maxAttempts counts the original attempt.
Discuss blocking/in-memory nature, head-of-line blocking, and rebalance risk from long backoff.
Weigh blocking RetryTemplate vs non-blocking DefaultErrorHandler / retry-topic patterns for throughput and ordering, and mandate idempotency.
**The problem.** A consumer may fail *transiently* — a downstream DB is briefly unavailable, a network call times out. Reprocessing the same message a moment later often succeeds. Spring Cloud Stream builds this in as **consumer retry**. **Mechanism.** For each consumer binding, the framework wraps your message handler in a Spring Retry **`RetryTemplate`**. When your function throws, the template catches it and re-invokes the *same* message according to a **backoff policy**. The relevant *consumer* properties (`spring.cloud.stream.bindings.<name>.consumer.*`): - **`maxAttempts`** — total attempts including the first. Default **3** (so 2 retries after the original). `maxAttempts=1` **disables** retry entirely. - **`backOffInitialInterval`** — delay before the first retry, default **1000** ms. - **`backOffMultiplier`** — factor applied each subsequent retry, default **2.0** (exponential). So delays go 1000ms, 2000ms, 4000ms… - **`backOffMaxInterval`** — ceiling on the delay, default **10000** ms. - **`retryableExceptions`** — a map of exception types → retry (true) / don't retry (false). - **`defaultRetryable`** — whether exceptions *not* listed are retried (default true). Set false to retry only the whitelisted types. **Blocking, in-memory nature — the key gotcha.** This retry happens **on the consumer thread, in the same JVM**, using `Thread.sleep`-style waits. It is *not* persisted. Consequences: - On **Kafka**, the consumer does not commit the record's offset until it succeeds or exhausts, and it does not poll newer records for that partition meanwhile → **head-of-line blocking**: one poison message with long backoff stalls the whole partition. - If the app **crashes mid-retry**, the in-memory attempt count is lost; after restart Kafka redelivers from the last committed offset and retries start over from zero. - Long backoffs can exceed Kafka's `max.poll.interval.ms`, making the broker think the consumer is dead and triggering a **rebalance**. Keep total blocking time bounded. **After exhaustion.** When `maxAttempts` is reached the framework routes the failure: - With **`enableDlq=true`** (Kafka) / **`autoBindDlq=true`** (Rabbit), the record is published to the **dead-letter** destination. - Otherwise it's logged; on Kafka the offset is committed so the app *moves on*, meaning the message is effectively **discarded** (data loss unless you have a DLQ). **Non-blocking alternative.** For Kafka, you can disable the binding-level RetryTemplate (`maxAttempts=1`) and instead register a `ListenerContainerCustomizer` that installs a Spring Kafka `DefaultErrorHandler` with a `BackOff` — or use dead-letter-and-retry topics — so retries don't block the partition. This is the pattern for high-throughput or long-backoff scenarios. **Idempotency caveat.** Because retry (and Kafka's at-least-once redelivery) can reprocess a message, your handler must be **idempotent** — re-running it must not double-charge, double-insert, etc. **When to use defaults.** Built-in blocking retry is fine for short, transient faults with small backoff. Move to non-blocking / DLQ-based retry when backoffs are long, throughput is high, or ordering per partition must not stall.
- Why can a long backoff on Kafka trigger a consumer group rebalance?The retry sleeps on the poll thread. If the total blocking time exceeds `max.poll.interval.ms`, the broker assumes the consumer died and reassigns its partitions — causing a rebalance and duplicate processing. Keep cumulative backoff under that limit or use non-blocking retry.
- How would you avoid retrying a message you know can never succeed (a poison/validation error)?Mark that exception non-retryable via `retryableExceptions` (or set `defaultRetryable=false` and whitelist only transient types), so it skips retries and goes straight to the DLQ. Throwing it once then routing to DLQ avoids wasting attempts.
saying these in an interview costs you the question
- Thinking retry state is persisted and survives an app restart (it's in-memory)
- Saying maxAttempts=3 means 3 retries (it's 3 total = 2 retries)
- Assuming retries happen on a background thread without blocking the partition
- Forgetting handlers must be idempotent because messages get reprocessed