This is an automated email from the ASF dual-hosted git repository.
tenthe pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git
The following commit(s) were added to refs/heads/dev by this push:
new a85a7fee0d Optimize adapter start and stop lifecycle (#4656)
a85a7fee0d is described below
commit a85a7fee0dc4ac8b0d5dcca370d9b46deb7532a8
Author: Philipp Zehnder <[email protected]>
AuthorDate: Wed Jul 1 08:08:18 2026 +0200
Optimize adapter start and stop lifecycle (#4656)
---
.../management/AdapterMasterManagement.java | 29 +++++++------
.../management/AdapterMigrationManager.java | 2 +-
.../management/AdapterUpdateManagement.java | 4 +-
.../management/management/WorkerRestClient.java | 47 +++++++++++-----------
.../export/resolver/AdapterResolver.java | 2 +-
.../health/monitoring/AdapterHealthCheck.java | 2 +-
.../rest/impl/connect/AdapterResource.java | 4 +-
.../rest/impl/connect/CompactAdapterResource.java | 2 +-
8 files changed, 49 insertions(+), 43 deletions(-)
diff --git
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMasterManagement.java
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMasterManagement.java
index 53b7c66fad..8e87a97401 100644
---
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMasterManagement.java
+++
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMasterManagement.java
@@ -120,7 +120,7 @@ public class AdapterMasterManagement {
try {
// Stop stream adapter
try {
- stopStreamAdapter(elementId, true);
+ stopAdapter(elementId, true);
} catch (AdapterException e) {
LOG.info("Could not stop adapter: " + elementId, e);
}
@@ -144,9 +144,12 @@ public class AdapterMasterManagement {
return adapterResourceManager.getDb().findAll();
}
- public void stopStreamAdapter(String elementId, boolean forceStop) throws
AdapterException {
+ public void stopAdapter(String elementId, boolean forceStop) throws
AdapterException {
+ stopAdapter(getAdapter(elementId), forceStop);
+ }
+
+ public void stopAdapter(AdapterDescription ad, boolean forceStop) throws
AdapterException {
LoadManager.tryLockForAdapter();
- AdapterDescription ad =
adapterResourceManager.getDb().getElementById(elementId);
try {
try {
var service = extensionsServiceStorage.findAll().stream()
@@ -154,11 +157,11 @@ public class AdapterMasterManagement {
if (ad.getSelectedServiceId() != null) {
return svc.getSvcId().equals(ad.getSelectedServiceId());
} else {
- return svc.getServiceUrl().equals(ad.getSelectedEndpointUrl());
+ return svc.getServiceUrl().equals(ad.getSelectedEndpointUrl());
}
})
.findFirst().orElseThrow(AdapterException::new);
- workerRestClient.stopStreamAdapter(service, ad);
+ workerRestClient.stopAdapter(service, ad);
} catch (AdapterException e) {
if (!forceStop) {
throw new AdapterException("Could not stop adapter", e);
@@ -168,7 +171,7 @@ public class AdapterMasterManagement {
adapterResourceManager.getDb().updateElement(ad);
}
}
- ExtensionsLogProvider.INSTANCE.reset(elementId);
+ ExtensionsLogProvider.INSTANCE.reset(ad.getElementId());
// remove the adapter from the metrics manager so that
// no metrics for this adapter are exposed anymore
@@ -182,11 +185,13 @@ public class AdapterMasterManagement {
}
}
- public void startStreamAdapter(String elementId) throws AdapterException {
+ public void startAdapter(String elementId) throws AdapterException {
+ startAdapter(getAdapter(elementId));
+ }
+
+ public void startAdapter(AdapterDescription ad) throws AdapterException {
LoadManager.tryLockForAdapter();
try {
- var ad = adapterResourceManager.getDb().getElementById(elementId);
-
try {
// Find endpoint to start adapter on
var service = new ExtensionsServiceEndpointGenerator()
@@ -196,16 +201,16 @@ public class AdapterMasterManagement {
// Update selected endpoint URL of adapter
ad.setSelectedEndpointUrl(service.getServiceUrl());
ad.setSelectedServiceId(service.getSvcId());
- adapterResourceManager.getDb().updateElement(ad);
+ ad = adapterResourceManager.getDb().updateElement(ad);
// Invoke adapter instance
- workerRestClient.invokeStreamAdapter(service, elementId);
+ workerRestClient.invokeStreamAdapter(service, ad);
// register the adapter at the metrics manager so that the
AdapterHealthCheck
// can send metrics
adapterMetrics.register(ad.getElementId(), ad.getName());
- LOG.info("Started adapter " + elementId + " on: " +
ad.getSelectedServiceId());
+ LOG.info("Started adapter " + ad.getElementId() + " on: " +
ad.getSelectedServiceId());
} catch (NoServiceEndpointsAvailableException e) {
throw new AdapterException("Could not start adapter due to unavailable
service endpoint",
e);
diff --git
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMigrationManager.java
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMigrationManager.java
index 69f4ece820..0bfcb89a83 100644
---
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMigrationManager.java
+++
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMigrationManager.java
@@ -101,7 +101,7 @@ public class AdapterMigrationManager extends
AbstractMigrationManager implements
migrationResult.element().getElementId()
);
try {
- workerRestClient.stopStreamAdapter(service, adapterDescription);
+ workerRestClient.stopAdapter(service, adapterDescription);
} catch (AdapterException e) {
LOG.error("Stopping adapter failed: {}",
StringUtils.join(e.getStackTrace(), "\n"));
}
diff --git
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterUpdateManagement.java
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterUpdateManagement.java
index 01811ced18..2c3d9f2504 100644
---
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterUpdateManagement.java
+++
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterUpdateManagement.java
@@ -52,7 +52,7 @@ public class AdapterUpdateManagement {
boolean shouldRestart = ad.isRunning();
if (ad.isRunning()) {
- this.adapterMasterManagement.stopStreamAdapter(ad.getElementId(), true);
+ this.adapterMasterManagement.stopAdapter(ad.getElementId(), true);
}
// update data source in database
@@ -61,7 +61,7 @@ public class AdapterUpdateManagement {
pipelineUpdateCoordinator.updatePipelines(ad);
if (shouldRestart) {
- this.adapterMasterManagement.startStreamAdapter(ad.getElementId());
+ this.adapterMasterManagement.startAdapter(ad.getElementId());
}
}
diff --git
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/WorkerRestClient.java
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/WorkerRestClient.java
index 27df7db7a1..1416ebf564 100644
---
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/WorkerRestClient.java
+++
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/WorkerRestClient.java
@@ -42,7 +42,6 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
-import java.util.List;
/**
* This client can be used to interact with the adapter workers executing the
adapter instances
@@ -61,19 +60,23 @@ public class WorkerRestClient {
public void invokeStreamAdapter(SpServiceRegistration service,
String elementId) throws AdapterException {
- var adapterStreamDescription = getAndDecryptAdapter(elementId);
+ invokeStreamAdapter(service,
getAdapterStorage().getElementById(elementId));
+ }
+
+ public void invokeStreamAdapter(SpServiceRegistration service,
+ AdapterDescription adapterStreamDescription)
throws AdapterException {
+ var decryptedAdapterStreamDescription =
cloneAndDecryptAdapter(adapterStreamDescription);
var requestTarget = ExtensionServiceRequestTargets.adapterStart(service);
- startAdapter(requestTarget, adapterStreamDescription);
- updateStreamAdapterStatus(adapterStreamDescription.getElementId(), true);
+ startAdapter(requestTarget, decryptedAdapterStreamDescription);
+ updateStreamAdapterStatus(decryptedAdapterStreamDescription, true);
}
- public void stopStreamAdapter(SpServiceRegistration service,
- AdapterDescription adapterStreamDescription)
throws AdapterException {
+ public void stopAdapter(SpServiceRegistration service,
+ AdapterDescription adapterStreamDescription) throws
AdapterException {
+ var decryptedAdapterStreamDescription =
cloneAndDecryptAdapter(adapterStreamDescription);
var requestTarget = ExtensionServiceRequestTargets.adapterStop(service);
- var ad =
getAdapterDescriptionById(resourceManager.manageAdapters().getDb(),
adapterStreamDescription.getElementId());
-
- stopAdapter(requestTarget, ad);
- updateStreamAdapterStatus(adapterStreamDescription.getElementId(), false);
+ stopAdapter(requestTarget, decryptedAdapterStreamDescription);
+ updateStreamAdapterStatus(decryptedAdapterStreamDescription, false);
}
private void startAdapter(ExtensionServiceRequestTarget requestTarget,
@@ -200,22 +203,14 @@ public class WorkerRestClient {
}
- private AdapterDescription getAdapterDescriptionById(IAdapterStorage
adapterStorage,
- String id) {
- AdapterDescription adapterDescription = null;
- List<AdapterDescription> allAdapters = adapterStorage.findAll();
- for (AdapterDescription a : allAdapters) {
- if (a.getElementId().endsWith(id)) {
- adapterDescription = a;
- }
- }
-
- return adapterDescription;
- }
-
private void updateStreamAdapterStatus(String adapterId,
boolean running) {
var adapter = getAndDecryptAdapter(adapterId);
+ updateStreamAdapterStatus(adapter, running);
+ }
+
+ private void updateStreamAdapterStatus(AdapterDescription adapter,
+ boolean running) {
adapter.setRunning(running);
encryptAndUpdateAdapter(adapter);
}
@@ -232,6 +227,12 @@ public class WorkerRestClient {
return adapter;
}
+ private AdapterDescription cloneAndDecryptAdapter(AdapterDescription
adapter) {
+ AdapterDescription decryptedDescription = new
Cloner().adapterDescription(adapter);
+ SecretProvider.getDecryptionService().apply(decryptedDescription);
+ return decryptedDescription;
+ }
+
private IAdapterStorage getAdapterStorage() {
return resourceManager.manageAdapters().getDb();
}
diff --git
a/streampipes-data-export/src/main/java/org/apache/streampipes/export/resolver/AdapterResolver.java
b/streampipes-data-export/src/main/java/org/apache/streampipes/export/resolver/AdapterResolver.java
index f2ba42bc6e..1de104d39e 100644
---
a/streampipes-data-export/src/main/java/org/apache/streampipes/export/resolver/AdapterResolver.java
+++
b/streampipes-data-export/src/main/java/org/apache/streampipes/export/resolver/AdapterResolver.java
@@ -104,7 +104,7 @@ public class AdapterResolver extends
AbstractResolver<AdapterDescription> {
new WorkerRestClient(extensionServiceRequestManager,
resourceManager),
getNoSqlStore().getExtensionsServiceStorage(),
extensionServiceRequestManager
- ).stopStreamAdapter(resourceId, true);
+ ).stopAdapter(resourceId, true);
} catch (AdapterException e) {
LOG.warn("Error when stopping adapter with id {} and name {}",
resourceId, existingAdapter.getName());
}
diff --git
a/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/AdapterHealthCheck.java
b/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/AdapterHealthCheck.java
index 4fb8cb5bb2..aedb4b34be 100644
---
a/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/AdapterHealthCheck.java
+++
b/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/AdapterHealthCheck.java
@@ -188,7 +188,7 @@ public class AdapterHealthCheck implements HealthCheck {
try {
if (adapterDescription.isRunning()) {
LOG.debug("Start recovering adapter {} ",
adapterDescription.getElementId());
-
this.healthCheckData.resourceProvider().adapterMasterManagement().startStreamAdapter(adapterDescription.getElementId());
+
this.healthCheckData.resourceProvider().adapterMasterManagement().startAdapter(adapterDescription.getElementId());
LOG.info("Adapter {} is recovered",
adapterDescription.getElementId());
}
} catch (AdapterException e) {
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java
index bc2103895a..84d86461d0 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java
@@ -230,7 +230,7 @@ public class AdapterResource extends
AbstractAdapterResource<AdapterMasterManage
try {
var adapter = getAdapterDescription(elementId);
if (checkAdapterPermission(adapter, "WRITE")) {
- managementService.stopStreamAdapter(elementId, forceStop);
+ managementService.stopAdapter(adapter, forceStop);
return ok(Notifications.success("Adapter stopped"));
} else {
return unauthorized();
@@ -247,7 +247,7 @@ public class AdapterResource extends
AbstractAdapterResource<AdapterMasterManage
try {
var adapterDescription = getAdapterDescription(elementId);
if (checkAdapterPermission(adapterDescription, "WRITE")) {
- managementService.startStreamAdapter(elementId);
+ managementService.startAdapter(adapterDescription);
return ok(Notifications.success("Adapter started"));
} else {
return unauthorized();
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/CompactAdapterResource.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/CompactAdapterResource.java
index 9d2cc76625..e8f88ae5e9 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/CompactAdapterResource.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/CompactAdapterResource.java
@@ -144,7 +144,7 @@ public class CompactAdapterResource extends
AbstractAdapterResource<AdapterMaste
}
if (compactAdapter.createOptions()
.start()) {
- managementService.startStreamAdapter(adapterId);
+ managementService.startAdapter(adapterId);
}
}
return ok(Notifications.success(adapterId));