Skip to content

RANGER-5752:Audit dispatcher consumer dies on TopicAuthorizationException - container stays healthy, audits stop flowing - #1231

Open
rameeshm wants to merge 1 commit into
masterfrom
RANGER-5752-patch
Open

RANGER-5752:Audit dispatcher consumer dies on TopicAuthorizationException - container stays healthy, audits stop flowing#1231
rameeshm wants to merge 1 commit into
masterfrom
RANGER-5752-patch

Conversation

@rameeshm

Copy link
Copy Markdown
Contributor

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.

…tion - container stays healthy, audits stop flowing
dispatcherWorkers.remove(workerId);
// Start a new worker with the same ID
startWorker(workerId);
LOG.info("Successfully restarted worker '{}'", workerId);

@kumaab kumaab Sep 11, 2026

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It may be helpful to add the number of restarts happened so far, to help with diagnostics.

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";

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Consider marking members authRetryDelayMs and pollErrorRetryDelayMs as final, and initializing them only in the constructor at lines 105 and 106 below.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants