Skip to content

Commit 606f10b

Browse files
authored
Allowing stateful service loaders (#1616)
Service loaders might be stateful, so caching them might be an isssue. Since internally serviceloader class already caches the classpath search, the right solution is to store a reference to the ServiceLoader itself and call stream for every invocation. That way, we have the best of both world, fresh new instance of the service loader class for every invocation (supporting stateful) and avoid the classpath search for every invocation (the original performance issue to be fixed) Signed-off-by: Francisco Javier Tirado Sarti <ftirados@ibm.com>
1 parent abffe00 commit 606f10b

2 files changed

Lines changed: 21 additions & 22 deletions

File tree

impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowApplication.java

Lines changed: 12 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -114,7 +114,7 @@ public class WorkflowApplication implements AutoCloseable {
114114
private final WorkflowLifeCycleCloudEventFactory lifeCycleCloudEventFactory;
115115
private final ScheduledExecutorService schedulerExecutorService;
116116
private final Set<String> allowedCommands;
117-
private final Map<Class<?>, List<?>> serviceLoadedClasses = new ConcurrentHashMap<>();
117+
private final Map<Class<?>, ServiceLoader<?>> servicesLoaded = new ConcurrentHashMap<>();
118118

119119
private WorkflowApplication(Builder builder) {
120120
this.taskFactory = builder.taskFactory;
@@ -712,21 +712,19 @@ public Set<String> allowedCommands() {
712712

713713
@SuppressWarnings("unchecked")
714714
public <T extends Comparable<?>> List<T> serviceLoadedClasses(Class<T> clazz) {
715-
return (List<T>)
716-
serviceLoadedClasses.computeIfAbsent(
717-
clazz,
718-
c ->
719-
ServiceLoader.load(clazz).stream()
720-
.map(ServiceLoader.Provider::get)
721-
.sorted()
722-
.toList());
715+
ServiceLoader<?> serviceLoader = servicesLoaded.computeIfAbsent(clazz, ServiceLoader::load);
716+
return (List<T>) serviceLoader.stream().map(ServiceLoader.Provider::get).sorted().toList();
723717
}
724718

725719
public <T extends Comparable<?>> T serviceLoadedClass(Class<T> serviceClass) {
726-
List<T> list = serviceLoadedClasses(serviceClass);
727-
if (list.isEmpty()) {
728-
throw new IllegalStateException("No " + serviceClass + " implementation found");
729-
}
730-
return list.get(0);
720+
ServiceLoader<?> serviceLoader =
721+
servicesLoaded.computeIfAbsent(serviceClass, ServiceLoader::load);
722+
return (T)
723+
serviceLoader.stream()
724+
.map(ServiceLoader.Provider::get)
725+
.sorted()
726+
.findFirst()
727+
.orElseThrow(
728+
() -> new IllegalStateException("No " + serviceClass + " implementation found"));
731729
}
732730
}

impl/core/src/main/java/io/serverlessworkflow/impl/executors/DefaultTaskExecutorFactory.java

Lines changed: 9 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
import io.serverlessworkflow.api.types.CallTask;
1919
import io.serverlessworkflow.api.types.Task;
2020
import io.serverlessworkflow.api.types.TaskBase;
21+
import io.serverlessworkflow.impl.WorkflowApplication;
2122
import io.serverlessworkflow.impl.WorkflowDefinition;
2223
import io.serverlessworkflow.impl.WorkflowMutablePosition;
2324
import io.serverlessworkflow.impl.executors.CallTaskExecutor.CallTaskExecutorBuilder;
@@ -32,9 +33,7 @@
3233
import io.serverlessworkflow.impl.executors.SwitchExecutor.SwitchExecutorBuilder;
3334
import io.serverlessworkflow.impl.executors.TryExecutor.TryExecutorBuilder;
3435
import io.serverlessworkflow.impl.executors.WaitExecutor.WaitExecutorBuilder;
35-
import java.util.Collection;
36-
import java.util.ServiceLoader;
37-
import java.util.ServiceLoader.Provider;
36+
import java.util.List;
3837

3938
public class DefaultTaskExecutorFactory implements TaskExecutorFactory {
4039

@@ -46,9 +45,6 @@ public static TaskExecutorFactory get() {
4645

4746
protected DefaultTaskExecutorFactory() {}
4847

49-
private Collection<CallableTaskBuilder> callTasks =
50-
ServiceLoader.load(CallableTaskBuilder.class).stream().map(Provider::get).sorted().toList();
51-
5248
@Override
5349
public TaskExecutorBuilder<? extends TaskBase> getTaskExecutor(
5450
WorkflowMutablePosition position, Task task, WorkflowDefinition definition) {
@@ -57,7 +53,10 @@ public TaskExecutorBuilder<? extends TaskBase> getTaskExecutor(
5753
TaskBase taskBase = (TaskBase) callTask.get();
5854
if (taskBase != null) {
5955
return new CallTaskExecutorBuilder(
60-
position, taskBase, definition, findCallTask(taskBase.getClass()));
56+
position,
57+
taskBase,
58+
definition,
59+
findCallTask(taskBase.getClass(), definition.application()));
6160
}
6261
} else if (task.getSwitchTask() != null) {
6362
return new SwitchExecutorBuilder(position, task.getSwitchTask(), definition);
@@ -86,7 +85,9 @@ public TaskExecutorBuilder<? extends TaskBase> getTaskExecutor(
8685
}
8786

8887
@SuppressWarnings("unchecked")
89-
private <T extends TaskBase> CallableTaskBuilder<T> findCallTask(Class<T> clazz) {
88+
private <T extends TaskBase> CallableTaskBuilder<T> findCallTask(
89+
Class<T> clazz, WorkflowApplication app) {
90+
List<CallableTaskBuilder> callTasks = app.serviceLoadedClasses(CallableTaskBuilder.class);
9091
return (CallableTaskBuilder<T>)
9192
callTasks.stream()
9293
.filter(s -> s.accept(clazz))

0 commit comments

Comments
 (0)