In a Spring Integration flow, how do threading and error propagation differ between an event-driven service activator and a polling one, and how would you design for reliability?
answer
- sender's thread => error bubbles to caller (DirectChannel)
- pool/scheduler thread => error channel / poller errorHandler
- ErrorMessage wraps MessagingException + failedMessage
- transactional poller = rollback-to-source = at-least-once
- retry advice + dead-letter channel + idempotency
basics
~20 sWith a DirectChannel the handler runs on the sender's thread, so exceptions propagate straight back to the caller in the caller's transaction. With a polling consumer (or ExecutorChannel) the handler runs on a scheduler/pool thread, so errors go to an error channel or the poller's errorHandler instead. Design reliability with transactional pollers, error channels, and retry advice.
solid answer
~50 sThreading follows the channel and consumer type. A DirectChannel + EventDrivenConsumer runs the handler synchronously on the sender's thread and transaction, so an exception propagates back to the producer — the sender owns retry/rollback. An ExecutorChannel or a PollingConsumer decouples threads: the handler runs on a pool/scheduler thread, the sender has already returned, and an unhandled exception can't reach it — it's routed to an error channel (a MessagingException wrapping the failed message) via MessagePublishingErrorHandler, or handled by the poller's errorHandler. For reliability I make the poll transactional (transactionManager/advice on the poller) so a failed message rolls back to the source for at-least-once delivery, add a RequestHandlerRetryAdvice or ExpressionEvaluatingRequestHandlerAdvice for retry/recovery, wire a dedicated errorChannel with a service activator for dead-lettering, and size taskExecutors and maxMessagesPerPoll for throughput and back-pressure. I also mind idempotency because at-least-once means possible reprocessing.
code
java · 32 lines@Configuration
public class ReliableFlow {
// Transactional poll (rollback-to-source) + retry advice on the handler.
@Bean(name = PollerMetadata.DEFAULT_POLLER_META_DATA_BEAN_NAME)
public PollerMetadata poller(PlatformTransactionManager txm) {
PollerMetadata pm = new PollerMetadata();
pm.setTrigger(new PeriodicTrigger(Duration.ofMillis(200)));
pm.setMaxMessagesPerPoll(10);
pm.setAdviceChain(List.of(new TransactionInterceptor(txm,
new MatchAlwaysTransactionAttributeSource())));
return pm;
}
@Bean
public Advice retryAdvice() {
RequestHandlerRetryAdvice advice = new RequestHandlerRetryAdvice();
advice.setRecoveryCallback(ctx -> null); // route/ log poison messages
return advice;
}
// Async processing: failures go to errorChannel, not back to the caller.
@ServiceActivator(inputChannel = "work", adviceChain = "retryAdvice")
public void process(Task t) { /* idempotent, at-least-once */ }
// Dead-letter / observe failures that happened on background threads.
@ServiceActivator(inputChannel = "errorChannel")
public void onError(ErrorMessage em) {
MessagingException ex = (MessagingException) em.getPayload();
// ex.getFailedMessage() -> persist / alert
}
}go deeper
Know sync vs async changes where errors go.
Explain error channels and that async handler failures don't reach the sender.
Detail transactional pollers, retry advice, and ErrorMessage/MessagePublishingErrorHandler.
Architect delivery semantics: at-least-once + idempotency, dead-lettering, back-pressure, ordering vs concurrency trade-offs, and observability of background failures.
## Why threading and error handling are coupled In Spring Integration, **the thread that carries a message determines where an error can go**. There is no magic: if the sender's thread is running the handler, exceptions bubble back to the sender; if a different thread runs it, they can't. ### Event-driven on a DirectChannel (synchronous) - Handler runs **on the sender's thread**, inside the **sender's transaction**. - A thrown exception **propagates straight back to the producer** — so retry/rollback are the *sender's* responsibility, and the whole hop is atomic if the sender is transactional. - Simple and low-latency, but the sender is blocked and coupled to downstream failures. ### Thread hand-off: ExecutorChannel or PollingConsumer (asynchronous) - **`ExecutorChannel`** dispatches to a `TaskExecutor`; **`PollingConsumer`** runs on the `TaskScheduler`/its own `taskExecutor`. - The sender has already returned, so **the exception cannot reach it**. Instead: - It goes to an **error channel** — Spring Integration wraps the failure in an **`ErrorMessage`** (payload = `MessagingException` that carries the original `failedMessage`) and publishes it, by default to the global `errorChannel`, via **`MessagePublishingErrorHandler`**. - For a poller, the poller's **`errorHandler`** (or `errorChannel` on the poller) handles it. - You can set a per-flow error channel header or configure the endpoint's error channel to segregate handling. ## Designing for reliability 1. **Transactional polling for at-least-once.** Put a `transactionManager`/advice on the **poller** so `receive()` + handler run in one transaction; a failure **rolls the message back to the source** (queue/DB), so it's redelivered. This gives at-least-once (hence **idempotency** is required downstream). 2. **Retry + recovery advice.** Attach a **`RequestHandlerRetryAdvice`** (backed by Spring Retry's `RetryTemplate`) to the handler's `adviceChain` for transient-failure retries with backoff, and a `RecoveryCallback` to route exhausted retries to a dead-letter channel. 3. **Explicit error channels / dead-lettering.** Wire a `@ServiceActivator(inputChannel = "errorChannel")` (or a custom one) to log, alert, and persist poison messages instead of losing them. 4. **Back-pressure & throughput.** Tune `maxMessagesPerPoll`, `receiveTimeout`, and a bounded `taskExecutor`; a `QueueChannel` with capacity provides buffering; unbounded `maxMessagesPerPoll` on a hot source can starve the scheduler. 5. **Ordering vs concurrency trade-off.** Concurrency (executor / multiple poll threads) breaks strict ordering; if order matters, keep single-threaded or partition by key. ## Gotchas - Putting `@Transactional` on the handler method makes the *handler* transactional but **not the receive** — for rollback-to-source you need the advice on the **poller**. - With an `ExecutorChannel`, callers see success even though processing may later fail asynchronously — monitoring the error channel is mandatory. - The default global `errorChannel` is a `PublishSubscribeChannel`; without a subscriber, errors are just logged — easy to miss in production. - At-least-once + retries means **duplicate processing**; design idempotent handlers (dedupe keys, upserts). ## When this matters This is the core of designing a robust asynchronous integration: choosing sync vs async hops deliberately, guaranteeing delivery semantics, and ensuring failures are observable and recoverable rather than silently swallowed on a background thread.
- With an ExecutorChannel, why can't the original sender catch a handler exception?ExecutorChannel hands the message to a TaskExecutor thread and the send() returns immediately, so the sender's stack has already unwound. The handler now runs on a different thread; its exception has nowhere to propagate back to, so Spring Integration routes it to an error channel via MessagePublishingErrorHandler instead.
- Why does at-least-once delivery force you to make handlers idempotent?A transactional poller rolls a failed message back to the source, which then redelivers it; a message that partially succeeded before failing can be processed again. Idempotent handlers (dedup keys, upserts, natural-key checks) ensure reprocessing a duplicate has no harmful side effect.
saying these in an interview costs you the question
- Claiming exceptions from an async handler propagate back to the sender
- Putting @Transactional on the handler and expecting rollback-to-source
- Ignoring that the default errorChannel silently logs if unsubscribed
- Adding concurrency without acknowledging loss of message ordering
- Assuming at-least-once delivery is exactly-once (no idempotency needed)