Savonitar commented on code in PR #28639:
URL: https://github.com/apache/flink/pull/28639#discussion_r3536780852
##########
flink-runtime/src/main/java/org/apache/flink/runtime/resourcemanager/ResourceManager.java:
##########
@@ -427,6 +430,14 @@ public CompletableFuture<RegistrationResponse>
registerJobMaster(
jobMasterIdFuture,
(JobMasterGateway jobMasterGateway, JobMasterId
leadingJobMasterId) -> {
if (Objects.equals(leadingJobMasterId,
jobMasterId)) {
+ // Register with the delegation token
manager first; a
+ // provider failure rejects this
registration so the job
+ // never starts without the tokens it
requires.
+ try {
+
delegationTokenManager.registerJob(jobId, jobConfiguration);
+ } catch (Exception e) {
+ return new
RegistrationResponse.Failure(e);
Review Comment:
> oncall guys
I'm something of an on-call guy myself...
Yes, it **was** logged in
https://github.com/apache/flink/pull/28639/changes#diff-d32f89982ba30269ab9438d993d9deb3b29d046792b6cc5ea35a277eaacfb287R629
(flink-runtime/src/main/java/org/apache/flink/runtime/security/token/DefaultDelegationTokenManager.java
registerJob()
```
LOG.error("Failed to register job {}", jobId, e);
```
And moreover, we already had the failing **provider** in stack trace. The
rollback path has its own `ERROR log too ("Failed to roll back registration of
job {}")`.
Update to: ^ Since then the failure surfaces improved further: the logs now
name the failing provider directly (`Failed to register job {} for provider
{}`,
https://github.com/apache/flink/commit/3c1dbd76bce6fbd345735c233c40b30d28824d6c),
and the JobMaster gets the failure back wrapped as a FlinkException naming the
job and the delegation token manager in the RegistrationResponse
(https://github.com/apache/flink/commit/485c10a47f6c884b15ec34ecfd2999f29cead82a).
##########
flink-runtime/src/main/java/org/apache/flink/runtime/resourcemanager/ResourceManager.java:
##########
@@ -427,6 +430,14 @@ public CompletableFuture<RegistrationResponse>
registerJobMaster(
jobMasterIdFuture,
(JobMasterGateway jobMasterGateway, JobMasterId
leadingJobMasterId) -> {
if (Objects.equals(leadingJobMasterId,
jobMasterId)) {
+ // Register with the delegation token
manager first; a
+ // provider failure rejects this
registration so the job
+ // never starts without the tokens it
requires.
+ try {
+
delegationTokenManager.registerJob(jobId, jobConfiguration);
+ } catch (Exception e) {
+ return new
RegistrationResponse.Failure(e);
Review Comment:
One addition to my note above: that holds for an `Exception`, but an `Error`
slips past both catches. The realistic trigger is a `NoClassDefFoundError` (or
other `LinkageError` from provider plugin code), the same failure class
`loadProviders` already has:
```
} catch (Exception | NoClassDefFoundError e) {
// The intentional general rule is that if a
provider's init method throws
// exception
// then stop the workload
LOG.error(
"Failed to initialize delegation token
provider {}",
provider.serviceName(),
e);
throw new FlinkRuntimeException(e);
}
```
^ Covered in
https://github.com/apache/flink/commit/485c10a47f6c884b15ec34ecfd2999f29cead82a.
Both job hooks now treat Exception | LinkageError the same.
##########
flink-runtime/src/main/java/org/apache/flink/runtime/security/token/DefaultDelegationTokenManager.java:
##########
@@ -310,64 +361,127 @@ public void start(Listener listener) throws Exception {
this.listener = checkNotNull(listener, "Listener must not be null");
synchronized (tokensUpdateFutureLock) {
checkState(tokensUpdateFuture == null, "Manager is already
started");
+ stopped = false;
Review Comment:
Agreed on the rename and the early returns, implemented in 9c1c49d. However,
I used the name `running` because it is used more often in the codebase.
And another part I did differently: I set `running = true` before the inline
first cycle in `start()`, exactly where `stopped = false` sat, not at the end
of the method (if I understand your suggestion correctly). The first obtain
runs inline and checks the flag twice, at the top of startTokensUpdate() and
again in maybeScheduleRenewal() when it schedules the periodic renewal. With
the flag flipped only at the end of start(), the inline cycle sees not-running,
skips the whole first obtain, and never schedules the renewal.
startShouldBeIdempotent and startTokensUpdateShouldScheduleRenewal lock that
ordering in, both fail if the flag moves after the cycle.
Same reason inverted on the stop() side: `running = false` flips before the
per-job unregistration below it, so a re-obtain racing shutdown cannot schedule
a cycle for a session that is shutting down. That cleanup also runs
unconditionally (no early return) because ResourceManager calls stop() on the
failed-startup path even when start() never ran.
##########
flink-runtime/src/main/java/org/apache/flink/runtime/security/token/DefaultDelegationTokenManager.java:
##########
@@ -416,13 +531,114 @@ long calculateRenewalDelay(Clock clock, long
nextRenewal) {
return renewalDelay;
}
+ @VisibleForTesting
+ void setClock(Clock clock) {
+ this.clock = clock;
+ }
+
/** Stops re-occurring token obtain task. */
@Override
public void stop() {
LOG.info("Stopping credential renewal");
- stopTokensUpdate();
+ synchronized (tokensUpdateFutureLock) {
+ // Mark stopped, cancel the pending cycle, and reset on-demand
re-obtain bookkeeping
+ // atomically, so a concurrent reobtainDelegationTokens() cannot
leave a live future
+ // orphaned after stop and a later start() does not inherit stale
state.
+ stopped = true;
+ stopTokensUpdate();
+ reobtainScheduled = false;
+ lastReobtainAtMillis = NO_PREVIOUS_REOBTAIN;
+ }
+
+ for (DelegationTokenProvider provider :
delegationTokenProviders.values()) {
Review Comment:
> Start called from another thread and start some/all of them
Providers have no start/restart hooks: init() runs once in the manager
constructor, and start() only drives the obtain logic, so there is nothing
start() could call to re-start a provider mid stop. An obtain cycle overlapping
stop() is possible, but that overlap is already covered by the stop() javadoc.
But another issue is possible. The manager instance is reused across
leadership sessions, so it is possible to have stop() then start() on the same
manager with the same provider instances and no re-init. A provider that closes
its resources in stop() (as the javadoc tells it to) is broken in the next
leadership term. The SPI javadoc says stop() is "called once during manager
shutdown" and the wiring does not honor that.
I see 2 options (taking into account the existing architecture/design):
1.Fan out provider.stop() only at process shutdown and let manager stop()
just stop scheduling. Providers keep a single init-to-stop lifecycle, matching
their single init.
Or
2.Keep per-session stop() but redefine the contract: stop() may be followed
by another start(), so providers must release only per-run resources and
re-acquire lazy. Smaller diff, but it changes the agreed "called once" approach
and pushes restart-stuff into every provider implementation.
I implemented option 1 in
[35495a6](https://github.com/apache/flink/pull/28639/changes/35495a68533d5da7f78ea83e2dec2ba077dd8c66).
`DelegationTokenManager.close()` is the terminal teardown: it ends the session
via stop(), then stops the providers exactly once, and a closed manager rejects
a later start(). It is called by the component that created the manager:
ClusterEntrypoint, MiniCluster, and YarnClusterDescriptor. The per-session
stop() keeps the providers usable for the next leadership session. It is a
single commit, so easy to rework if you prefer option 2.
WDYT?
##########
flink-runtime/src/main/java/org/apache/flink/runtime/security/token/DefaultDelegationTokenManager.java:
##########
@@ -416,13 +531,114 @@ long calculateRenewalDelay(Clock clock, long
nextRenewal) {
return renewalDelay;
}
+ @VisibleForTesting
+ void setClock(Clock clock) {
+ this.clock = clock;
+ }
+
/** Stops re-occurring token obtain task. */
@Override
public void stop() {
LOG.info("Stopping credential renewal");
- stopTokensUpdate();
+ synchronized (tokensUpdateFutureLock) {
+ // Mark stopped, cancel the pending cycle, and reset on-demand
re-obtain bookkeeping
+ // atomically, so a concurrent reobtainDelegationTokens() cannot
leave a live future
+ // orphaned after stop and a later start() does not inherit stale
state.
+ stopped = true;
+ stopTokensUpdate();
+ reobtainScheduled = false;
+ lastReobtainAtMillis = NO_PREVIOUS_REOBTAIN;
+ }
+
+ for (DelegationTokenProvider provider :
delegationTokenProviders.values()) {
Review Comment:
1. Renamed stop() to close() in [[FLINK-40019][core][runtime] Rename
DelegationTokenProvider.stop() to
close()](https://github.com/apache/flink/pull/28639/changes/4df526e2daaa3badbf2222fbd22343e02e4a151b)
. Agree it meets naming API better.
2.
> correct me if I'm wrong
> stop (together with start) can happen when HA kicks in during failover and
it makes sure that new leader re-obtained tokens.
> stop has nothing to do with freeing resources allocated by the provider
and this function called regularly.
yes, your understanding is correct 👍
##########
flink-runtime/src/main/java/org/apache/flink/runtime/security/token/DefaultDelegationTokenManager.java:
##########
@@ -310,64 +361,127 @@ public void start(Listener listener) throws Exception {
this.listener = checkNotNull(listener, "Listener must not be null");
synchronized (tokensUpdateFutureLock) {
checkState(tokensUpdateFuture == null, "Manager is already
started");
+ stopped = false;
}
startTokensUpdate();
}
@VisibleForTesting
void startTokensUpdate() {
- try {
- LOG.info("Starting tokens update task");
- DelegationTokenContainer container = new
DelegationTokenContainer();
- Optional<Long> nextRenewal =
obtainDelegationTokensAndGetNextRenewal(container);
-
- if (container.hasTokens()) {
-
delegationTokenReceiverRepository.onNewTokensObtained(container);
-
- LOG.info("Notifying listener about new tokens");
- checkNotNull(listener, "Listener must not be null");
-
listener.onNewTokensObtained(InstantiationUtil.serializeObject(container));
- LOG.info("Listener notified successfully");
- } else {
- LOG.warn("No tokens obtained so skipping notifications");
+ synchronized (tokensUpdateFutureLock) {
+ // The obtain cycle is starting: clear the dedupe flag so later
on-demand requests can
+ // schedule a fresh cycle.
+ reobtainScheduled = false;
+ // If stop() ran before this cycle (already handed to the IO
executor) began, skip the
+ // obtain/broadcast: the providers may already be stopped. Safe
via this lock's
+ // happens-before with stop(). The dedupe flag is cleared above,
so it is never stuck.
+ if (stopped) {
+ return;
}
+ }
+ // Serialize the obtain-and-broadcast so a re-obtain racing the
periodic renewal cannot run
+ // two cycles concurrently on the (multi-threaded) IO executor and
broadcast out of order.
+ synchronized (obtainLock) {
+ try {
+ LOG.info("Starting tokens update task");
+ DelegationTokenContainer container = new
DelegationTokenContainer();
+ Optional<Long> nextRenewal =
obtainDelegationTokensAndGetNextRenewal(container);
+
+ if (container.hasTokens()) {
+
delegationTokenReceiverRepository.onNewTokensObtained(container);
+
+ LOG.info("Notifying listener about new tokens");
+ checkNotNull(listener, "Listener must not be null");
+
listener.onNewTokensObtained(InstantiationUtil.serializeObject(container));
+ LOG.info("Listener notified successfully");
+ } else {
+ LOG.warn("No tokens obtained so skipping notifications");
+ }
- if (nextRenewal.isPresent()) {
- lastKnownNextRenewal = nextRenewal.get();
- currentRetryBackoff = renewalRetryInitialBackoff;
- long renewalDelay =
- calculateRenewalDelay(Clock.systemDefaultZone(),
nextRenewal.get());
- synchronized (tokensUpdateFutureLock) {
- tokensUpdateFuture =
- scheduledExecutor.schedule(
- () ->
ioExecutor.execute(this::startTokensUpdate),
- renewalDelay,
- TimeUnit.MILLISECONDS);
+ if (nextRenewal.isPresent()) {
+ lastKnownNextRenewal = nextRenewal.get();
+ currentRetryBackoff = renewalRetryInitialBackoff;
+ long renewalDelay = calculateRenewalDelay(clock,
nextRenewal.get());
+ maybeScheduleRenewal(renewalDelay);
+ LOG.info(
+ "Tokens update task started with {} delay",
+
TimeUtils.formatWithHighestUnit(Duration.ofMillis(renewalDelay)));
+ } else {
+ LOG.warn(
+ "Tokens update task not started because either no
tokens obtained or none of the tokens specified its renewal date");
}
- LOG.info(
- "Tokens update task started with {} delay",
-
TimeUtils.formatWithHighestUnit(Duration.ofMillis(renewalDelay)));
- } else {
+ } catch (InterruptedException e) {
+ // Ignore, may happen if shutting down.
+ LOG.debug("Interrupted", e);
+ } catch (Exception e) {
+ long delay = calculateRetryDelay(clock);
+ maybeScheduleRenewal(delay);
LOG.warn(
- "Tokens update task not started because either no
tokens obtained or none of the tokens specified its renewal date");
+ "Failed to update tokens, will try again in {}",
+
TimeUtils.formatWithHighestUnit(Duration.ofMillis(delay)),
+ e);
Review Comment:
You are right, the log could claim a delay that was never scheduled. Fixed
in 9c1c49d: maybeScheduleRenewal now returns the delay that actually took
effect and the log prints that.
##########
flink-runtime/src/main/java/org/apache/flink/runtime/security/token/DefaultDelegationTokenManager.java:
##########
@@ -416,13 +531,114 @@ long calculateRenewalDelay(Clock clock, long
nextRenewal) {
return renewalDelay;
}
+ @VisibleForTesting
+ void setClock(Clock clock) {
+ this.clock = clock;
+ }
+
/** Stops re-occurring token obtain task. */
@Override
public void stop() {
LOG.info("Stopping credential renewal");
- stopTokensUpdate();
+ synchronized (tokensUpdateFutureLock) {
+ // Mark stopped, cancel the pending cycle, and reset on-demand
re-obtain bookkeeping
+ // atomically, so a concurrent reobtainDelegationTokens() cannot
leave a live future
+ // orphaned after stop and a later start() does not inherit stale
state.
+ stopped = true;
+ stopTokensUpdate();
+ reobtainScheduled = false;
+ lastReobtainAtMillis = NO_PREVIOUS_REOBTAIN;
+ }
+
+ for (DelegationTokenProvider provider :
delegationTokenProviders.values()) {
+ try {
+ provider.stop();
+ } catch (Throwable t) {
+ LOG.error("Failed to stop delegation token provider {}",
provider.serviceName(), t);
+ }
+ }
LOG.info("Stopped credential renewal");
}
+
+ @Override
+ public void reobtainDelegationTokens() {
+ synchronized (tokensUpdateFutureLock) {
+ if (scheduledExecutor == null || ioExecutor == null) {
+ LOG.debug(
+ "A re-obtain of delegation tokens was requested but
the manager was "
+ + "constructed without executors (one-shot
obtain path); the "
+ + "request is ignored.");
+ return;
+ }
+ if (stopped) {
+ LOG.debug(
+ "A re-obtain of delegation tokens was requested after
the manager was "
+ + "stopped; the request is ignored.");
+ return;
+ }
+ // Dedupe: if an on-demand re-obtain is already scheduled and has
not started yet, the
+ // newly registered job(s) will be covered by it, so coalesce this
request into it.
+ if (reobtainScheduled) {
+ LOG.debug("A re-obtain of delegation tokens is already
scheduled; coalescing.");
+ return;
+ }
+ // Cooldown: bound how often on-demand re-obtains can run by
deferring this cycle until
+ // at least reobtainCooldownMillis have passed since the previous
on-demand re-obtain.
+ long now = clock.millis();
+ long delayMillis =
+ lastReobtainAtMillis == NO_PREVIOUS_REOBTAIN
+ ? 0L
+ : Math.max(0L, lastReobtainAtMillis +
reobtainCooldownMillis - now);
+ // Only bring the next cycle forward: if a cycle (e.g. the
periodic renewal) is still
+ // pending and scheduled to fire sooner than the cooldown-deferred
time, fire at that
+ // earlier time instead of pushing it later — otherwise a
short-lived token could expire
+ // before it is renewed. The nextScheduledAtMillis > now guard
skips an already-fired
+ // future that has not yet been re-armed, so this never bypasses
the cooldown.
+ if (tokensUpdateFuture != null
+ && nextScheduledAtMillis > now
+ && nextScheduledAtMillis - now < delayMillis) {
+ delayMillis = nextScheduledAtMillis - now;
+ }
+ lastReobtainAtMillis = now;
+ reobtainScheduled = true;
+ LOG.debug(
+ "Re-obtain of delegation tokens requested; scheduling an
obtain cycle in {}",
+
TimeUtils.formatWithHighestUnit(Duration.ofMillis(delayMillis)));
+ scheduleRenewalLocked(delayMillis);
+ }
+ }
+
+ @Override
+ public void registerJob(JobID jobId, Configuration jobConfiguration)
throws Exception {
+ try {
+ for (DelegationTokenProvider provider :
delegationTokenProviders.values()) {
+ provider.registerJob(jobId, jobConfiguration);
+ }
+ } catch (Exception e) {
+ // If any of the providers fail to register, then unregister the
job from them all.
+ // unregisterJob is idempotent, so it is safe to call it for
providers that were never
+ // (or only partially) registered for this job before the failure.
The rollback must
+ // never mask the original failure, so swallow any rollback
exception.
+ try {
+ unregisterJob(jobId);
+ } catch (Exception rollbackException) {
+ LOG.error("Failed to roll back registration of job {}", jobId,
rollbackException);
+ }
+ LOG.error("Failed to register job {}", jobId, e);
+ throw e;
+ }
+ }
+
+ @Override
+ public void unregisterJob(JobID jobId) throws Exception {
+ for (DelegationTokenProvider provider :
delegationTokenProviders.values()) {
+ try {
+ provider.unregisterJob(jobId);
+ } catch (Exception e) {
+ LOG.error("Failed to unregister job for provider {}",
provider.serviceName(), e);
Review Comment:
Good catch, added in
https://github.com/apache/flink/pull/28639/changes/3c1dbd76bce6fbd345735c233c40b30d28824d6c
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]