Skip to content

Commit d4deac7

Browse files
committed
feat: schedule polling and retries on managed executors instead of Timers
PollingEventSource and TimerEventSource each created a java.util.Timer, and PerResourcePollingEventSource a ScheduledThreadPoolExecutor of its own, so every polling event source and every controller cost a thread that the operator neither sized nor shut down. They now schedule on executors managed by ExecutorServiceManager, which builds them from ConfigurationService: - getScheduledExecutorService() backs the polling event sources - getRetryAndRescheduleExecutorService() backs the retry and reschedule timer of every controller, kept separate so that a slow poll cannot delay a retry Both are sized by concurrentScheduledTaskThreads() (4 by default) and concurrentRetryAndRescheduleThreads() (2 by default), overridable programmatically or through the josdk.scheduled-tasks.concurrent-threads and josdk.retry-and-reschedule.concurrent-threads keys. Both default to daemon threads, discard the tasks scheduled for later on shutdown so they don't hold up the termination of the operator, and drop cancelled tasks eagerly. The previous Executors.newScheduledThreadPool(0) was effectively single threaded and created non daemon threads. An event source resolves its executor on every start rather than at creation time, since the manager replaces its pools when the operator is restarted, and only shuts one down if it created it itself. That also fixes PerResourcePollingEventSource shutting down a user supplied executor and not registering its tasks again after a restart. PollingEventSource gains a constructor taking an EventSourceContext, the one without it is deprecated, and both polling configurations accept an executor to run a single event source on a pool of its own.
1 parent 2339b7d commit d4deac7

22 files changed

Lines changed: 804 additions & 75 deletions

File tree

‎docs/content/en/docs/documentation/eventing.md‎

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -222,6 +222,28 @@ is similar to `PerResourcePollingEventSource` except that, contrary to that even
222222
doesn't poll a specific API separately per resource, but periodically and independently of
223223
actually observed primary resources.
224224

225+
#### Threading of the polling event sources
226+
227+
Both polling event sources schedule their polls on an executor the operator shares between all of
228+
its polling event sources. A poll therefore only starts once one of that executor's threads is
229+
free, so if your operator registers many polling event sources, or if fetching your external
230+
resources is slow, size the pool accordingly with
231+
`ConfigurationServiceOverrider.withConcurrentScheduledTaskThreads` (4 threads by default):
232+
233+
```java
234+
Operator operator = new Operator(overrider -> overrider.withConcurrentScheduledTaskThreads(20));
235+
```
236+
237+
Retried and rescheduled reconciliations are triggered on a separate executor, so a slow poll can
238+
never delay them. It is sized with `withConcurrentRetryAndRescheduleThreads` (2 threads by
239+
default); few threads are needed there since triggering a reconciliation only enqueues an event
240+
for one of the reconciliation threads to pick up.
241+
242+
Use `ConfigurationServiceOverrider.withScheduledExecutorService` to replace the polling executor
243+
altogether, or the `withExecutorService` method of the event source's own configuration builder to
244+
poll a single event source on an executor of its own. An executor provided that way is not managed
245+
by the operator: it is your responsibility to shut it down.
246+
225247
#### Inbound event sources
226248

227249
[SimpleInboundEventSource](https://github.com/operator-framework/java-operator-sdk/blob/main/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/inbound/SimpleInboundEventSource.java)

‎docs/content/en/docs/documentation/operations/configuration.md‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -281,6 +281,13 @@ All operator-level keys are prefixed with `josdk.`.
281281
|---|---|---|
282282
| `josdk.workflow.executor-threads` | `Integer` | Thread pool size for workflow execution |
283283

284+
#### Scheduled Tasks
285+
286+
| Key | Type | Description |
287+
|---|---|---|
288+
| `josdk.scheduled-tasks.concurrent-threads` | `Integer` | Thread pool size for the operator's scheduled tasks, i.e. the polling event sources |
289+
| `josdk.retry-and-reschedule.concurrent-threads` | `Integer` | Thread pool size for triggering retried and rescheduled reconciliations |
290+
284291
#### Informer
285292

286293
| Key | Type | Description |

‎operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationService.java‎

Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121
import java.util.Set;
2222
import java.util.concurrent.ExecutorService;
2323
import java.util.concurrent.Executors;
24+
import java.util.concurrent.ScheduledExecutorService;
2425
import java.util.function.Consumer;
2526

2627
import org.slf4j.Logger;
@@ -64,6 +65,20 @@ public interface ConfigurationService {
6465
/** The default number of threads used to process dependent workflows */
6566
int DEFAULT_WORKFLOW_EXECUTOR_THREAD_NUMBER = DEFAULT_RECONCILIATION_THREADS_NUMBER;
6667

68+
/**
69+
* The default number of threads used to run the operator's scheduled tasks, i.e. the periodic
70+
* polls of {@link io.javaoperatorsdk.operator.processing.event.source.polling.PollingEventSource}
71+
* and {@link
72+
* io.javaoperatorsdk.operator.processing.event.source.polling.PerResourcePollingEventSource}
73+
*/
74+
int DEFAULT_SCHEDULED_TASK_THREADS_NUMBER = 4;
75+
76+
/**
77+
* The default number of threads used to trigger the operator's retried and rescheduled
78+
* reconciliations
79+
*/
80+
int DEFAULT_RETRY_AND_RESCHEDULE_THREADS_NUMBER = 2;
81+
6782
/**
6883
* Creates a new {@link ConfigurationService} instance used to configure an {@link
6984
* io.javaoperatorsdk.operator.Operator} instance, starting from the specified base configuration
@@ -219,6 +234,33 @@ default int concurrentWorkflowExecutorThreads() {
219234
return DEFAULT_WORKFLOW_EXECUTOR_THREAD_NUMBER;
220235
}
221236

237+
/**
238+
* Number of threads the operator can spin out to run its scheduled (i.e. periodic or delayed)
239+
* tasks with the default executor. These threads are shared by all the polling event sources of
240+
* the operator, so this number should be raised when many, or slow, polling event sources are
241+
* registered: a poll only starts once a thread is available, and a slow poll therefore delays the
242+
* polls of the other event sources.
243+
*
244+
* @return the maximum number of concurrent scheduled task threads
245+
* @since 5.6.0
246+
*/
247+
default int concurrentScheduledTaskThreads() {
248+
return DEFAULT_SCHEDULED_TASK_THREADS_NUMBER;
249+
}
250+
251+
/**
252+
* Number of threads the operator can spin out to trigger its retried and rescheduled
253+
* reconciliations with the default executor. These threads are shared by all the controllers of
254+
* the operator, but the tasks they run only enqueue an event for the reconciliation to happen on
255+
* a reconciliation thread, so few of them are needed.
256+
*
257+
* @return the maximum number of concurrent retry and reschedule threads
258+
* @since 5.6.0
259+
*/
260+
default int concurrentRetryAndRescheduleThreads() {
261+
return DEFAULT_RETRY_AND_RESCHEDULE_THREADS_NUMBER;
262+
}
263+
222264
/**
223265
* Override to provide a custom {@link Metrics} implementation
224266
*
@@ -249,6 +291,44 @@ default ExecutorService getWorkflowExecutorService() {
249291
return Executors.newFixedThreadPool(concurrentWorkflowExecutorThreads());
250292
}
251293

294+
/**
295+
* Override to provide a custom {@link ScheduledExecutorService} implementation to change how the
296+
* operator's scheduled (i.e. periodic or delayed) tasks are run. This executor is shared by all
297+
* the polling event sources of the operator. The retried and rescheduled reconciliations run on
298+
* an executor of their own, see {@link #getRetryAndRescheduleExecutorService()}, so that a slow
299+
* poll can't delay them.
300+
*
301+
* <p>Note that the default implementation lets the executor discard the tasks that were scheduled
302+
* for later when it is shut down, so that they don't delay the termination of the operator, and
303+
* that it creates daemon threads so that a never stopped operator doesn't keep the JVM alive.
304+
* Custom implementations are advised to do the same.
305+
*
306+
* @return the {@link ScheduledExecutorService} implementation to use to run scheduled tasks
307+
* @since 5.6.0
308+
*/
309+
default ScheduledExecutorService getScheduledExecutorService() {
310+
return Utils.daemonScheduledThreadPool(
311+
concurrentScheduledTaskThreads(), "josdk-scheduled-task");
312+
}
313+
314+
/**
315+
* Override to provide a custom {@link ScheduledExecutorService} implementation to change how the
316+
* operator's retried and rescheduled reconciliations are triggered. This executor is kept
317+
* separate from the one the polling event sources use, see {@link
318+
* #getScheduledExecutorService()}, so that a slow poll can't delay a retry.
319+
*
320+
* <p>The same notes as for {@link #getScheduledExecutorService()} apply to custom
321+
* implementations.
322+
*
323+
* @return the {@link ScheduledExecutorService} implementation to use to trigger retried and
324+
* rescheduled reconciliations
325+
* @since 5.6.0
326+
*/
327+
default ScheduledExecutorService getRetryAndRescheduleExecutorService() {
328+
return Utils.daemonScheduledThreadPool(
329+
concurrentRetryAndRescheduleThreads(), "josdk-retry-reschedule");
330+
}
331+
252332
/**
253333
* Determines whether the associated Kubernetes client should be closed when the associated {@link
254334
* io.javaoperatorsdk.operator.Operator} is stopped.

‎operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationServiceOverrider.java‎

Lines changed: 104 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121
import java.util.Optional;
2222
import java.util.Set;
2323
import java.util.concurrent.ExecutorService;
24+
import java.util.concurrent.ScheduledExecutorService;
2425
import java.util.function.Function;
2526

2627
import org.slf4j.Logger;
@@ -45,11 +46,15 @@ public class ConfigurationServiceOverrider {
4546
private Boolean checkCR;
4647
private Integer concurrentReconciliationThreads;
4748
private Integer concurrentWorkflowExecutorThreads;
49+
private Integer concurrentScheduledTaskThreads;
50+
private Integer concurrentRetryAndRescheduleThreads;
4851
private Cloner cloner;
4952
private Boolean closeClientOnStop;
5053
private KubernetesClient client;
5154
private ExecutorService executorService;
5255
private ExecutorService workflowExecutorService;
56+
private ScheduledExecutorService scheduledExecutorService;
57+
private ScheduledExecutorService retryAndRescheduleExecutorService;
5358
private LeaderElectionConfiguration leaderElectionConfiguration;
5459
private String clusterScopedEventNamespace;
5560
private EventRecorder eventRecorder;
@@ -86,6 +91,34 @@ public ConfigurationServiceOverrider withConcurrentWorkflowExecutorThreads(int t
8691
return this;
8792
}
8893

94+
/**
95+
* Sets the number of threads used to run the operator's scheduled (i.e. periodic or delayed)
96+
* tasks, which are shared by all its polling event sources.
97+
*
98+
* @param threadNumber the maximum number of concurrent scheduled task threads
99+
* @return this {@link ConfigurationServiceOverrider} for chained customization
100+
* @see ConfigurationService#concurrentScheduledTaskThreads()
101+
* @since 5.6.0
102+
*/
103+
public ConfigurationServiceOverrider withConcurrentScheduledTaskThreads(int threadNumber) {
104+
this.concurrentScheduledTaskThreads = threadNumber;
105+
return this;
106+
}
107+
108+
/**
109+
* Sets the number of threads used to trigger the operator's retried and rescheduled
110+
* reconciliations, which are shared by all its controllers.
111+
*
112+
* @param threadNumber the maximum number of concurrent retry and reschedule threads
113+
* @return this {@link ConfigurationServiceOverrider} for chained customization
114+
* @see ConfigurationService#concurrentRetryAndRescheduleThreads()
115+
* @since 5.6.0
116+
*/
117+
public ConfigurationServiceOverrider withConcurrentRetryAndRescheduleThreads(int threadNumber) {
118+
this.concurrentRetryAndRescheduleThreads = threadNumber;
119+
return this;
120+
}
121+
89122
@SuppressWarnings("rawtypes")
90123
public ConfigurationServiceOverrider withDependentResourceFactory(
91124
DependentResourceFactory dependentResourceFactory) {
@@ -119,6 +152,37 @@ public ConfigurationServiceOverrider withWorkflowExecutorService(
119152
return this;
120153
}
121154

155+
/**
156+
* Replaces the executor used to run the operator's scheduled (i.e. periodic or delayed) tasks,
157+
* which are shared by all its polling event sources.
158+
*
159+
* @param scheduledExecutorService the executor to run scheduled tasks on
160+
* @return this {@link ConfigurationServiceOverrider} for chained customization
161+
* @see ConfigurationService#getScheduledExecutorService()
162+
* @since 5.6.0
163+
*/
164+
public ConfigurationServiceOverrider withScheduledExecutorService(
165+
ScheduledExecutorService scheduledExecutorService) {
166+
this.scheduledExecutorService = scheduledExecutorService;
167+
return this;
168+
}
169+
170+
/**
171+
* Replaces the executor used to trigger the operator's retried and rescheduled reconciliations,
172+
* which is shared by all its controllers.
173+
*
174+
* @param retryAndRescheduleExecutorService the executor to trigger retried and rescheduled
175+
* reconciliations on
176+
* @return this {@link ConfigurationServiceOverrider} for chained customization
177+
* @see ConfigurationService#getRetryAndRescheduleExecutorService()
178+
* @since 5.6.0
179+
*/
180+
public ConfigurationServiceOverrider withRetryAndRescheduleExecutorService(
181+
ScheduledExecutorService retryAndRescheduleExecutorService) {
182+
this.retryAndRescheduleExecutorService = retryAndRescheduleExecutorService;
183+
return this;
184+
}
185+
122186
/**
123187
* Replaces the default {@link KubernetesClient} instance by the specified one. This is the
124188
* preferred mechanism to configure which client will be used to access the cluster.
@@ -312,6 +376,28 @@ public int concurrentWorkflowExecutorThreads() {
312376
original.concurrentWorkflowExecutorThreads());
313377
}
314378

379+
@Override
380+
public int concurrentScheduledTaskThreads() {
381+
return Utils.ensureValid(
382+
overriddenValueOrDefault(
383+
concurrentScheduledTaskThreads,
384+
ConfigurationService::concurrentScheduledTaskThreads),
385+
"maximum scheduled task threads",
386+
1,
387+
original.concurrentScheduledTaskThreads());
388+
}
389+
390+
@Override
391+
public int concurrentRetryAndRescheduleThreads() {
392+
return Utils.ensureValid(
393+
overriddenValueOrDefault(
394+
concurrentRetryAndRescheduleThreads,
395+
ConfigurationService::concurrentRetryAndRescheduleThreads),
396+
"maximum retry and reschedule threads",
397+
1,
398+
original.concurrentRetryAndRescheduleThreads());
399+
}
400+
315401
@Override
316402
public Metrics getMetrics() {
317403
return overriddenValueOrDefault(metrics, ConfigurationService::getMetrics);
@@ -340,6 +426,24 @@ public ExecutorService getWorkflowExecutorService() {
340426
}
341427
}
342428

429+
@Override
430+
public ScheduledExecutorService getScheduledExecutorService() {
431+
if (scheduledExecutorService != null) {
432+
return scheduledExecutorService;
433+
} else {
434+
return super.getScheduledExecutorService();
435+
}
436+
}
437+
438+
@Override
439+
public ScheduledExecutorService getRetryAndRescheduleExecutorService() {
440+
if (retryAndRescheduleExecutorService != null) {
441+
return retryAndRescheduleExecutorService;
442+
} else {
443+
return super.getRetryAndRescheduleExecutorService();
444+
}
445+
}
446+
343447
@Override
344448
public Optional<LeaderElectionConfiguration> getLeaderElectionConfiguration() {
345449
return leaderElectionConfiguration != null

‎operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ExecutorServiceManager.java‎

Lines changed: 33 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@ public class ExecutorServiceManager {
4242
private ExecutorService workflowExecutor;
4343
private ExecutorService cachingExecutorService;
4444
private ScheduledExecutorService scheduledExecutorService;
45+
private ScheduledExecutorService retryAndRescheduleExecutorService;
4546
private boolean started;
4647
private ConfigurationService configurationService;
4748

@@ -128,30 +129,54 @@ public ExecutorService cachingExecutorService() {
128129
return cachingExecutorService;
129130
}
130131

132+
/**
133+
* The executor the operator runs its scheduled (i.e. periodic or delayed) tasks on, shared by its
134+
* polling event sources. Note that it is only valid while the manager is started: it is shut down
135+
* by {@link #stop(Duration)} and replaced by a fresh one on the next {@link
136+
* #start(ConfigurationService)}, so callers should retrieve it when they start rather than hold
137+
* on to it.
138+
*
139+
* @return the executor to run scheduled tasks on
140+
*/
131141
public ScheduledExecutorService scheduledExecutorService() {
132142
return scheduledExecutorService;
133143
}
134144

145+
/**
146+
* The executor the operator triggers its retried and rescheduled reconciliations on, kept
147+
* separate from {@link #scheduledExecutorService()} so that a slow poll can't delay a retry. The
148+
* same lifecycle caveat as for {@link #scheduledExecutorService()} applies.
149+
*
150+
* @return the executor to trigger retried and rescheduled reconciliations on
151+
*/
152+
public ScheduledExecutorService retryAndRescheduleExecutorService() {
153+
return retryAndRescheduleExecutorService;
154+
}
155+
135156
public synchronized void start(ConfigurationService configurationService) {
136157
if (!started) {
137158
this.configurationService = configurationService; // used to lazy init workflow executor
138159
this.cachingExecutorService = Executors.newCachedThreadPool();
139-
this.scheduledExecutorService = Executors.newScheduledThreadPool(0);
160+
this.scheduledExecutorService = configurationService.getScheduledExecutorService();
161+
this.retryAndRescheduleExecutorService =
162+
configurationService.getRetryAndRescheduleExecutorService();
140163
this.executor = new InstrumentedExecutorService(configurationService.getExecutorService());
141164
started = true;
142165
}
143166
}
144167

145168
public synchronized void stop(Duration gracefulShutdownTimeout) {
146-
var parallelExec = Executors.newFixedThreadPool(4);
169+
var shutdowns =
170+
List.of(
171+
shutdown(executor, gracefulShutdownTimeout),
172+
shutdown(workflowExecutor, gracefulShutdownTimeout),
173+
shutdown(cachingExecutorService, gracefulShutdownTimeout),
174+
shutdown(scheduledExecutorService, gracefulShutdownTimeout),
175+
shutdown(retryAndRescheduleExecutorService, gracefulShutdownTimeout));
176+
var parallelExec = Executors.newFixedThreadPool(shutdowns.size());
147177
try {
148178
log.debug("Closing executor");
149-
parallelExec.invokeAll(
150-
List.of(
151-
shutdown(executor, gracefulShutdownTimeout),
152-
shutdown(workflowExecutor, gracefulShutdownTimeout),
153-
shutdown(cachingExecutorService, gracefulShutdownTimeout),
154-
shutdown(scheduledExecutorService, gracefulShutdownTimeout)));
179+
parallelExec.invokeAll(shutdowns);
155180
} catch (InterruptedException e) {
156181
log.debug("Exception closing executor: {}", e.getLocalizedMessage());
157182
Thread.currentThread().interrupt();

0 commit comments

Comments
 (0)