Logging & Error Handling
Setting up SeriLog
To set up SeriLog to capture the desired logging information you will need to have the following NuGet packages installed:
- Serilog AspNetCore
dotnet add package Serilog.AspNetCore
- Serilog Loki
dotnet add package Serilog.Sinks.Grafana.Loki
- Serilog Exceptions
dotnet add package Serilog.Exceptions
- Serilog Encrichers Span
dotnet add package Serilog.Enrichers.Span
- Serilog Expressions
dotnet add package Serilog.Expressions
To configure SeriLog, in the RegisterServices method in the service startup (Program.cs) add the following:
// Logging using Serilog
builder.Logging.AddSerilog();
var loggerOptions = new ConfigurationReaderOptions { SectionName = AuditConstants.AppSettingsSectionNames.Serilog };
Log.Logger = new LoggerConfiguration()
.ReadFrom.Configuration(builder.Configuration, loggerOptions)
.Filter.ByExcluding("RequestPath like '/health%'")
.Enrich.WithExceptionDetails()
.Enrich.FromLogContext()
.Enrich.WithSpan()
.Enrich.With<ActivityEnricher>()
.CreateLogger();
//Serilog.Debugging.SelfLog.Enable(Console.Error);
You can add additional filters to exclude certain request paths you don’t want captured in logging (swagger for example).
Next add the following to appsettings.json
"Link:Audit:Logging:Serilog": {
"Using": [ "Serilog.Sinks.Console", "Serilog.Sinks.Grafana.Loki" ],
"MinimumLevel": {
"Default": "Information",
"Override": {
"Microsoft": "Warning",
"System": "Warning"
}
},
"WriteTo": [
{ "Name": "Console" },
{
"Name": "GrafanaLoki",
"Args": {
"uri": "http://localhost:3100",
"labels": [
{
"key": "app",
"value": "Link-BoTW"
},
{
"key": "component",
"value": "Audit"
}
],
"propertiesAsLabels": [ "app", "component" ]
}
}
]
}
Additionally you can add the following to enhance the format of the console:
"Logging": {
"LogLevel": {
"Default": "Information",
"Microsoft": "Warning",
"System": "Warning"
},
"Console": {
"FormatterName": "json",
"FormatterOptions": {
"SingleLine": true,
"IncludeScopes": true,
"TimestampFormat": "HH:mm:ss ",
"UseUtcTimestamp": true,
"JsonWriterOptions": {
"Indented": true
}
}
}
}
Adding Problem Details
Problem details is an attempt at an industry standard to carry machine readable details of errors in a HTTP response to avoid the need to define new error response formats for HTTP APIs.
This can be accomplished by using the .NET problem details service.
To enable to sending of Problem Details add the following in the RegisterServices method in the service startup (Program.cs):
//Add problem details
builder.Services.AddProblemDetails(options => {
options.CustomizeProblemDetails = ctx =>
{
ctx.ProblemDetails.Detail = "An error occured in our API. Please use the trace id when requesting assistence.";
if (!ctx.ProblemDetails.Extensions.ContainsKey("traceId"))
{
string? traceId = Activity.Current?.Id ?? ctx.HttpContext.TraceIdentifier;
ctx.ProblemDetails.Extensions.Add(new KeyValuePair<string, object?>("traceId", traceId));
}
if (builder.Environment.IsDevelopment())
{
ctx.ProblemDetails.Extensions.Add("service", "Audit");
}
else
{
ctx.ProblemDetails.Extensions.Remove("exception");
}
};
});
Add the following to the SetupMiddleware method in the service startup (Program.cs):
if (app.Environment.IsDevelopment() || app.Environment.EnvironmentName.Equals("Local", StringComparison.InvariantCultureIgnoreCase))
{
app.UseDeveloperExceptionPage();
}
else
{
app.UseExceptionHandler();
}
This will consume any 500 error returned by the API and turn it into a ProblemDetails model that is returned to the user. It will appear as follows:
{
"type": "https://httpstatuses.io/500",
"title": "Internal Server Error",
"status": 500,
"detail": "An error occured in our API. Please use the trace id when requesting assistence.",
"traceId": "00-09f4ab47f959288745f3ef337ac773b7-fac718cf545d0a9d-01"
}
We can use the Trace Id to find relevant logs pertaining to the issue.
For this to work correctly you will need to return a 500 response type along with the exception in your API.
catch (Exception ex)
{
_logger.LogError(new EventId(LoggingIds.GetItem, "Get Audit Event"), ex, "An exception occurred while attempting to retrieve an audit event with an id of {id}", id);
throw;
}
Logging Ids
The AuditLoggingIds.GetItem in the above _logger.LogError method is a new option in .NET that allows you to assign a Logging Id to log types. In the example above it will give the log an id of 1002 from the following static class.
public static class AuditLoggingIds
{
public const int GenerateItems = 1000;
public const int ListItems = 1001;
public const int GetItem = 1002;
public const int InsertItem = 1003;
public const int UpdateItem = 1004;
public const int DeleteItem = 1005;
public const int GetItemNotFound = 1006;
public const int UpdateItemNotFound = 1007;
public const int EventConsumer = 2000;
public const int EventProducer = 2001;
public const int HealthCheck = 9000;
}
This can allow for finding specific log types by the id and will show up in logs as follows:

UserScope Middleware
The LantanaGroup.Link.Shared.Application.Middleware middleware has been added to the shared project. This middleware will capture the username as well as the id of the user if an authenticated user is present in the .NET HttpContext.User.Identity object. This object will be populated based on if/how you have authentication enabled. For more information about enabled authentication see TBD.
This scoped user information will automatically be added to any logs created. To add the middleware simply add the following in the SetupMiddleware method in the service startup (Program.cs):
app.UseRouting();
app.UseCors("CorsPolicy");
app.UseAuthentication();
app.UseMiddleware<UserScopeMiddleware>();
app.UseAuthorization();
app.UseEndpoints(endpoints => endpoints.MapControllers());
This middlware should come after app.UseAuthentication(). When this is added the following will be included in a scope within your logs.

Adding data to exceptions
For additional context, you may add some additional data about the issue that caused the exception to be thrown as follows:
catch (NullReferenceException ex)
{
_logger.LogDebug(new EventId(LoggingIds.GetItem, "Get Audit Event"), ex, "Failed to get audit event with an id of {id}.", id);
var queryEx = new ApplicationException("Failed to get audit event", ex);
queryEx.Data.Add("Id", id);
throw queryEx;
}
The exception.Data.Add() method will add this information in the following locations:

Additional information about .NET logging can be found here.
LoggerMessage Attribute
Introduced in .NET 6, the LoggerMessage allows for performant logging.
public static partial class Logging
{
[LoggerMessage(
AuditLoggingIds.GenerateItems,
LogLevel.Information,
"New audit event created")]
public static partial void LogAuditEventCreation(this ILogger logger,
[LogProperties]AuditEntity auditEvent);
}
The LogProperties attribute, requires .NET 8, will include all of the properties of the object in your log.

Kafka Consumer Error Handling — Dotnet Services
When consuming an event that cannot be deserialized, the consumer should catch the error through a ConsumeException. In the catch, the service should do the following:
- Create an audit event to trace that an error occurred.
- Produce an error event to the Error topic. Storing improperly structured events in an separate topic will allow for further investigate
- Commit the consumer result to acknowledge back to the Kafka broker that the consumed event was processed.
flowchart LR SourceTopic["SOURCE TOPIC"] --> KafkaApp["KAFKA<br>APP"] KafkaApp -->|1| TargetTopic["TARGET TOPIC"] KafkaApp -->|2| ErrorTopic["ERROR TOPIC"] classDef source fill:#e6f4fa,stroke:#aaa; classDef target fill:#e6f4fa,stroke:#aaa; classDef error fill:#f8d7da,stroke:#c00,color:#000; class SourceTopic source; class TargetTopic target; class ErrorTopic error;
Example:
ConsumeResult<ReportScheduledKey, ReportScheduledValue> consumeResult;
try
{
consumeResult = _reportScheduledConsumer.Consume(cancellationToken);
}
catch (ConsumeException e)
{
_logger.LogError($"Consume failure, potentially schema related: {e.Error.Reason}");
var potentialFacilityId = Encoding.UTF8.GetString(e.ConsumerRecord.Message.Key);
var potentialValue = Encoding.UTF8.GetString(e.ConsumerRecord.Message.Value);
ProduceErrorEvent(potentialFacilityId, potentialValue);
var auditValue = new AuditEventMessage
{
FacilityId = potentialFacilityId,
Action = AuditEventType.Query,
ServiceName = "QueryDispatch",
EventDate = DateTime.UtcNow,
Notes = $"Kafka ReportScheduled consume failure, potentially schema related \nException Message: {e.Error}",
};
ProduceAuditEvent(auditValue, e.ConsumerRecord.Message.Headers);
_reportScheduledConsumer.Commit();
continue;
}
private void ProduceErrorEvent(string key, string value) {
var config = new ProducerConfig()
{
ClientId = "Error-QueryDispatch-ReportScheduled"
};
using (var producer = _errorProducerFactory.CreateProducer(config))
{
var headers = new Headers
{
new Header("X-Consumer", Encoding.UTF8.GetBytes("QueryDispatch")),
new Header("X-Topic", Encoding.UTF8.GetBytes("ReportScheduledEvent"))
};
producer.Produce("Error", new Message<string, string>
{
Key = key,
Value = value,
Headers = headers
});
producer.Flush();
}
}
Kafka Consumer Retry Handling
When a problem occurs while processing a consumed event (Facility not properly configured, API connection could not be established, could not connect to database, etc), we will want to attempt to perform retries to ensure that BotW processes each event successfully. When an exception is caught during the processing of an event, the service should do the following:
- Create an audit event to trace that an error occurred.
- Produce a retry event to the services Retry topic.
- Commit the consumer result to acknowledge back to the Kafka broker that the consumed event was processed.
flowchart LR ReportScheduled --> QueryDispatch QueryDispatch --> Error QueryDispatch --> Retry QueryDispatch --> ConsumeSuccess Retry --> QueryDispatchRetry QueryDispatchRetry --> Retry QueryDispatchRetry --> ConsumeSuccess ConsumeSuccess:::final -->|Success or Retry| Retry classDef final fill:#d4af37,stroke:#333,color:white;
Kafka Consumer Error Handling — Java Services
How the Link Java services — MeasureEval and Validation — handle failures while consuming Kafka messages: what gets retried, what gets dead-lettered, on what schedule, and how to observe and test all of it. The mechanism is shared; each service plugs in its own topics.
This page is self-contained — it assumes familiarity with Kafka basics (topics, consumers, offsets) but nothing about the Link codebase.
| Service | Consumed topics (each with -Retry/-Error companions) |
Downstream event on success | Per-message side resources |
|---|---|---|---|
| MeasureEval | ResourcesNormalized, EvaluationRequested |
MeasureReportGenerated |
Resource cache (Redis/blob) — released only on success or terminal failure |
| Validation | ReadyForValidation |
ValidationComplete |
None |
The problem being solved
Link services communicate through Kafka: for example, when patient data has been normalized, a
ResourcesNormalized message tells MeasureEval to evaluate the patient against a measure.
Processing such a message can fail for two very different reasons:
- Transient failures — a dependency is briefly down (MongoDB, Redis, blob storage, another service). The same message would succeed if simply tried again a bit later.
- Permanent failures ("poison" messages) — the message content itself is bad (malformed JSON, an unparseable FHIR resource, a validation error). Retrying can never succeed.
Without deliberate handling, the first kind loses real patient data over a 30-second network blip, and the second kind can jam a consumer forever, re-failing on the same message. The design below retries transient failures automatically on a fixed schedule and moves permanent failures aside — into an error topic — where they are preserved for investigation instead of blocking everything behind them.
What happens when a message fails
Each consumed topic has two companion topics, created alongside it:
<Topic>-Retry— a holding area for messages awaiting another attempt. A dedicated listener re-reads them, waits until the message's scheduled due-time, and feeds it back through the exact same processing logic.<Topic>-Error— the dead-letter topic. Messages that will never be retried again end up here, preserved with their original content and headers. Nothing consumes this topic; it is an archive for humans (and a future replay tool).
flowchart LR
Main["ResourcesNormalized<br>(main topic)"] --> Consumer["Consumer<br>(processing logic)"]
Consumer -->|success| Ack["done — offset committed"]
Consumer -->|failure| Decide{"classify the<br>failure"}
Decide -->|"transient +<br>attempts left"| Retry["ResourcesNormalized-Retry"]
Decide -->|"poison, or attempts<br>exhausted, or retries disabled"| Error["ResourcesNormalized-Error<br>(dead letter)"]
Retry -->|"after the backoff delay"| Consumer
classDef error fill:#f8d7da,stroke:#c00,color:#000;
class Error error;
Step by step, when processing throws an exception:
- The failure is classified (see the next section): can a retry ever help?
- Transient, attempts remaining → the message is republished to
-Retry, stamped with its attempt count and an absolute "not before" timestamp (now + the configured delay). The retry listener redelivers it after that delay and the full processing logic runs again. - Poison, attempts exhausted, or retries disabled → the message is republished to
-Error. This is terminal: it will never be redelivered automatically. - Only after the republish succeeds is the original message acknowledged. If the service crashes mid-flight — or the republish itself fails — the offset stays uncommitted and Kafka redelivers the original message. A message is never silently lost.
One special case: a message whose bytes can't even be deserialized fails before the
processing logic runs at all. Such messages skip the retry ladder (deserialization can never
succeed on retry) and go straight to -Error.
The diagram uses MeasureEval's ResourcesNormalized as the example; the flow is identical for
its EvaluationRequested topic and for Validation's ReadyForValidation
(ReadyForValidation-Retry / ReadyForValidation-Error).
One operational prerequisite: the -Retry and -Error topics are not auto-created — they
must exist in the broker (they are part of the standard topic provisioning). In a fresh
environment where they are missing, failed messages have nowhere to go.
Transient vs poison: how failures are classified
| Outcome | When | Where the message goes |
|---|---|---|
| Retry | Any exception not classified as poison, with attempts remaining | -Retry, redelivered after the scheduled delay |
| Dead-letter: poison | The exception (anywhere in its cause chain) is a content problem that can never succeed — in both services: FHIR parse errors, validation errors, message-format and deserialization errors | -Error, immediately, without burning retry attempts |
| Dead-letter: exhausted | All attempts used up | -Error |
| Dead-letter: disabled | Retries are turned off by configuration | -Error on the first failure |
The default is deliberately retry: anything not on the service's explicit poison list — including infrastructure outages like "cannot connect to Redis" — is treated as transient. The flip side of that rule for developers: processing code must let such exceptions escape unhandled. Catching and swallowing a Redis outage would not prevent the failure — it would convert it into a silently wrong result (e.g. evaluating an empty patient bundle and reporting a false "not reportable") with no retry and no dead letter.
The retry schedule
The delays come from configuration (next section). Both services deploy the same schedule:
| Attempt | Delay before it | Elapsed since first failure |
|---|---|---|
| 1 (original) | — | 0 |
| 2 | 20 s | ~20 s |
| 3 | 60 s | ~1 min 20 s |
| 4 (last) | 120 s | ~3 min 20 s |
| dead letter | — | ~3 min 20 s |
So a dependency outage shorter than ~3.5 minutes self-heals with no data loss and no manual
action: whichever attempt lands after the dependency returns completes the real work, produces
the normal downstream event, and the message never reaches -Error.
Configuration
Two properties control everything; both live under spring.kafka.retry (in
application.yml, overridable per environment via Azure App Configuration):
spring:
kafka:
retry:
consumer-retry-duration: [PT20S, PT60S, PT120S] # ISO-8601 durations
disable-retry-consumer: false
consumer-retry-durationis the single source of truth for both how many retries happen and when: N entries mean N retries (N+1 total attempts), and entry K is the delay before retry K. An empty list means the first failure dead-letters. This mirrors the .NET services'ConsumerSettings:ConsumerRetryDuration, so both stacks are tuned the same way.disable-retry-consumer: truesends every failure straight to-Errorand doesn't start the retry listener at all.- Changing depth, timing, or turning retries off is a configuration change and a restart — no code change.
Guarantees — and their consequences
- At-least-once delivery. A message is acknowledged only after it has been fully processed
or durably parked in
-Retry/-Error. The trade-off: a crash or restart at the wrong moment can deliver a message twice. This is by design (the alternative is losing it), and it means duplicate downstream events (e.g. twoMeasureReportGeneratedfor one input) are possible and expected under failure conditions — downstream consumers must tolerate them. - Dead letters are terminal today. Nothing consumes
-Error; there is no automatic replay. Recovering a dead-lettered message is currently a manual operation. (A cross-service replay capability is a known future work item.) - Cached data survives retries. Messages like
ResourcesNormalizeddon't carry the patient data themselves — it sits in a shared resource cache (Redis / blob storage) keyed by correlation id. The rule: a failed attempt must leave the cache untouched, because the retry needs it. The cache is released at exactly two moments — successful processing, or after the message is durably dead-lettered — never in between. Validation carries no per-message side resources, so it has no cleanup step at all. - Dependency clients fail fast on the consumer path. Blob-storage clients used during message processing are capped at 2 quick tries (~10 s) instead of the SDK default of ~3 minutes of hidden internal retrying. One system owns retries — the visible, configurable Kafka ladder — so outages surface in seconds and the schedule above is the real schedule.
For QA: what to expect and how to verify
Expected behavior by scenario
| Scenario | Expected behavior |
|---|---|
| Malformed / poison message | Goes to -Error immediately (no retry attempts); normal traffic continues unaffected |
| Dependency down briefly (< ~3.5 min) | Retry ladder runs (20 s / 60 s / 120 s); once the dependency is restored, the next attempt succeeds and produces the normal downstream event; -Error stays empty |
| Dependency down past the ladder (~3.5 min+) | Message lands in -Error after 4 attempts; downstream event is not produced |
| Service restarted mid-ladder | Message is redelivered and continues; a duplicate downstream event is possible and acceptable |
| Retries disabled via config | First failure goes straight to -Error |
Which dependency outages actually trigger the ladder differs per service — the ladder only engages when a failure makes processing throw, and each service has different hard dependencies:
- MeasureEval retries on: Redis down (its resource cache is a hard dependency), MongoDB
down, blob storage (Azurite locally) down. Data-only triggers, needing no infrastructure
change: an
EvaluationRequestedreferencing a nonexistent previous report, or an unknown measure id. - Validation retries on: blob storage down while the message carries a valid
payloadUri, or its SQL Server database down (Validation has no MongoDB dependency). - Validation trap: stopping Redis does not produce retries in Validation. Its Redis
layer is a read-through cache whose errors are deliberately swallowed and treated as cache
misses — processing continues against the backing store. Expect warning logs and a degraded
health check, not
-Retrytraffic. This is correct behavior, not a bug.
What to observe
Three signals tell the whole story per test message:
- Service logs — each retry hop and each dead-letter decision is logged. Against local
docker-compose (
link-measureeval/link-validation):docker logs link-measureeval -f | grep -E "Retry attempt|Max retry|Cache cleanup" docker logs link-validation -f | grep -E "Retry attempt|Max retry" - The
-Errortopic offset — did the dead-letter count go up? (SubstituteReadyForValidation-Errorwhen testing Validation.)docker exec kafka-broker /opt/bitnami/kafka/bin/kafka-get-offsets.sh \ --bootstrap-server kafka_b:9092 --command-config /tmp/client.properties \ --topic ResourcesNormalized-Error - The downstream event —
MeasureReportGenerated(MeasureEval) /ValidationComplete(Validation) must appear on success and only on success.
Inducing failures: stopping dependency containers
With the docker-compose stack running, use the levers from the list above:
docker stop mongo (MeasureEval), docker stop mssql (Validation), docker stop azurite
(MeasureEval's resource cache; Validation's payload reads), docker stop redis_cache
(MeasureEval only — see the Validation trap above). Restart the dependency mid-ladder to
verify recovery.
docker stop and docker pause are two different scenarios, and both are worth running:
stop gives connection-refused (fast, clean failures); pause gives a TCP black hole —
connections hang until a timeout fires — which exercises the configured timeout bounds (the
2 s Redis command timeout, Mongo's serverSelectionTimeoutMS) and is closer to a real
network outage.
Inducing failures: breaking configuration
Config-based faults complement container stops: they model real misconfiguration incidents, and they target one service without disturbing other consumers of the same dependency.
Mechanics (local docker-compose): add or override the property in the service's
environment: block in docker-compose.yml, in dotted form exactly like the existing
entries (e.g. spring.data.mongodb.uri: ...), then recreate the service with
docker compose up -d measureeval (or validation). In deployed environments the same
property paths go into Azure App Configuration plus a service restart.
The golden rule: a broken config only exercises the retry ladder if the service still boots and the failure happens at message-processing time. Some breakages fail at startup instead — a crashed or never-healthy container runs no consumers and produces zero retries, which reads like a failed test when it's really a failed setup. The traps are flagged below.
MeasureEval:
| Lever | Property to break | Behavior |
|---|---|---|
| Redis port | spring.data.redis.port: 6380 (valid but unused port) |
✅ Boots fine; every cache read throws at processing time → ladder. Use a valid-but-dead port — an out-of-range value like 66666 fails property binding and kills startup |
| Redis password | spring.data.redis.password: wrong |
✅ Runtime (the client connects lazily) — auth failure at first cache read → ladder |
| Mongo credentials | spring.data.mongodb.uri: mongodb://linkuser:wrongpass@mongo:27017/link-measureeval?authSource=admin&serverSelectionTimeoutMS=4000 |
✅ Runtime — every Mongo access fails with AuthenticationFailed (error 18) → ladder. Works even though the local compose Mongo has no auth enabled: the server still validates the SCRAM handshake against its user catalog, and the user doesn't exist |
| Mongo unreachable | dead port in the URI, e.g. mongodb://mongo:27018/...?serverSelectionTimeoutMS=4000 |
✅ Runtime; equivalent to stopping Mongo but scoped to this one service |
| Blob container | resource-cache.blob-storage.blob-container-name: no-such-container |
✅ Runtime — cache blob reads 404 (ContainerNotFound) → ladder; the fail-fast client (2 tries, 10 s cap) surfaces it in seconds. Breaking the account key inside the connection string gives the 403 variant |
Mongo URI rules that make or break the test: keep the URI well-formed (a malformed one —
e.g. a port inside an +srv URI — is a startup crash, not a runtime failure), and always
keep serverSelectionTimeoutMS=4000 — the driver's default is 30 s of server selection per
attempt, which stalls every rung and makes the ladder look hung.
Validation:
| Lever | Property to break | Behavior |
|---|---|---|
| SQL Server password | spring.datasource.password: wrong |
❌ Startup killer, not a lever — Hibernate opens a connection during boot (ddl-auto: update), so the container never becomes healthy and no consumer runs. For DB-driven retries, break the DB after startup: docker stop mssql or docker pause mssql |
| Blob container | internal-blob-storage.blob-container-name: no-such-container |
✅ Runtime — the report-payload read fails (for messages with a valid payloadUri) → ladder |
| Redis (any breakage) | port, password, or stopped | ❌ Deliberately no retries — see the Validation trap above (cache errors are swallowed as misses) |
Config faults are sticky — know which command heals them. A container keeps the
environment it was created with: reverting docker-compose.yml does not change the
running service, so the fault stays live until the container is recreated. The two commands
do opposite things:
| Command | Effect |
|---|---|
docker restart link-measureeval |
Keeps the injected fault (restart preserves the container's env) |
docker compose up -d measureeval |
Removes it (recreates the container from the current file) — this is also the "dependency restored mid-ladder" recovery step; the in-flight message is parked in -Retry and survives the recreation |
Two observation caveats while a config fault is live: the container may show unhealthy
purely from its health-check probes hitting the broken dependency — that is not ladder
traffic, and the resulting log noise (on HTTP threads) is not retry activity. Ladder behavior
only appears once a message is actually consumed, so always produce a test record and watch
for the Retry attempt log lines.
Producing test records
A poison message is any record whose value isn't valid JSON for the topic's type. To produce test records, use the Kafka console producer inside the broker container from Git Bash, not PowerShell (PowerShell pipes prepend a byte-order mark that corrupts the message key):
printf '%s\n' '{"facilityId":"MyFacility","patientId":"p1"}|{"queryType":"INITIAL",...}' |
MSYS_NO_PATHCONV=1 docker exec -i kafka-broker \
/opt/bitnami/kafka/bin/kafka-console-producer.sh \
--bootstrap-server kafka_b:9092 --producer.config /tmp/client.properties \
--topic ResourcesNormalized --property parse.key=true --property 'key.separator=|'
For developers: how it's built and how to adopt it
The mechanism lives in the shared Java module and is service-agnostic; each service plugs in its own topics and its own poison list.
Background on why it's hand-rolled: these services process messages on a background executor
(long-running CQL evaluation would otherwise exceed Kafka's poll interval and trigger consumer
rebalancing). Exceptions therefore surface on the executor thread, which Spring Kafka's native
retry machinery cannot see — so the services own the retry policy (classification, attempt
counting, schedule, destination) and reuse Spring only for the retry plumbing (the listener
that delays -Retry redeliveries).
| Shared component | Role |
|---|---|
AbstractAsyncConsumer |
Executor hand-off; ack on success, route-then-ack on failure |
RetryTopicRecoverer |
Makes the routing decision, stamps attempt/due-time headers, publishes |
RetryTopicRecovererFactory |
Assembles the above; takes the service's poison set and an optional terminal-cleanup hook as parameters |
KafkaRetryConfig |
Binds the two spring.kafka.retry properties; derives attempt count and per-attempt delay |
Relationships
flowchart LR nlogging_error_handling_6DA934E6["Design: Logging & Error Handling"]