RANGER-5752:Audit dispatcher consumer dies on TopicAuthorizationException - container stays healthy, audits stop flowing - #1231
Open
rameeshm wants to merge 1 commit into
Open
RANGER-5752:Audit dispatcher consumer dies on TopicAuthorizationException - container stays healthy, audits stop flowing#1231rameeshm wants to merge 1 commit into
rameeshm wants to merge 1 commit into
Conversation
…tion - container stays healthy, audits stop flowing
kumaab
approved these changes
Sep 11, 2026
| dispatcherWorkers.remove(workerId); | ||
| // Start a new worker with the same ID | ||
| startWorker(workerId); | ||
| LOG.info("Successfully restarted worker '{}'", workerId); |
Contributor
There was a problem hiding this comment.
It may be helpful to add the number of restarts happened so far, to help with diagnostics.
mneethiraj
reviewed
Sep 11, 2026
| public static final String PROP_DISPATCHER_MAX_POLL_INTERVAL_MS = "max.poll.interval.ms"; | ||
| public static final String PROP_DISPATCHER_HEARTBEAT_INTERVAL_MS = "heartbeat.interval.ms"; | ||
| public static final String PROP_DISPATCHER_PARTITION_ASSIGNMENT_STRATEGY = "partition.assignment.strategy"; | ||
| public static final String PROP_DISPATCHER_AUTH_RETRY_DELAY_MS = "auth.retry.delay.ms"; |
Contributor
There was a problem hiding this comment.
I suggest replacing "auth" with "authz" - to make it clear this configuration is about retry on authorization failure (instead of authentication failure).
PROP_DISPATCHER_AUTH_RETRY_DELAY_MS => PROP_DISPATCHER_AUTHZ_RETRY_DELAY_MS
auth.retry.delay.ms => authz.retry.delay.ms
| protected int dispatcherThreadCount = 1; | ||
| protected String offsetCommitStrategy = AuditServerConstants.DEFAULT_OFFSET_COMMIT_STRATEGY; | ||
| protected long offsetCommitInterval = AuditServerConstants.DEFAULT_OFFSET_COMMIT_INTERVAL_MS; | ||
| protected long authRetryDelayMs = AuditServerConstants.DEFAULT_DISPATCHER_AUTH_RETRY_DELAY_MS; |
Contributor
There was a problem hiding this comment.
Consider marking members authRetryDelayMs and pollErrorRetryDelayMs as final, and initializing them only in the constructor at lines 105 and 106 below.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
What changes were proposed in this pull request?
The patch addresses the issue where the AuditDispatcherBase.DispatcherWorker dies when Kafka is not ready yet.
Retry Logic:
It wraps the poll() and processing logic in a try-catch block. When an AuthorizationException or AuthenticationException is thrown, the worker now sleeps for a configurable delay (auth.retry.delay.ms) and continues the loop instead of exiting.
Worker Health Monitoring:
It adds monitorAndRestartWorkers() to the main thread's loop. If a worker thread terminates unexpectedly (e.g., due to an unhandled Error), the main thread will detect that its Future is done and restart the worker.
How was this patch tested?
Patch tested in docker with scenarios to repro the situation and see how resilient is AuditDispatcher is.