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 ed67499d14 Add debug logging for adapter monitoring metrics (#4654)
ed67499d14 is described below
commit ed67499d14b99fea64c03ea6ed27305773cfaa32
Author: Philipp Zehnder <[email protected]>
AuthorDate: Tue Jun 30 14:56:53 2026 +0200
Add debug logging for adapter monitoring metrics (#4654)
---
.../api/monitoring/SpMonitoringManager.java | 29 ++++++++++++++
.../elements/SendToBrokerAdapterSink.java | 1 -
.../monitoring/MonitoringManagement.java | 34 +++++++++++++++-
.../health/monitoring/AdapterHealthCheck.java | 44 ++++++++++++++++-----
.../pipeline/ExtensionsLogProvider.java | 37 +++++++++++++++++
.../pipeline/ExtensionsServiceLogExecutor.java | 46 +++++++++++++++++++++-
.../rest/impl/AdapterMonitoringResource.java | 7 ++++
7 files changed, 186 insertions(+), 12 deletions(-)
diff --git
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/monitoring/SpMonitoringManager.java
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/monitoring/SpMonitoringManager.java
index 7dd1c92228..e3fb7d79ba 100644
---
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/monitoring/SpMonitoringManager.java
+++
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/monitoring/SpMonitoringManager.java
@@ -22,6 +22,9 @@ import
org.apache.streampipes.model.monitoring.SpEndpointMonitoringInfo;
import org.apache.streampipes.model.monitoring.SpLogEntry;
import org.apache.streampipes.model.monitoring.SpMetricsEntry;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@@ -30,6 +33,8 @@ public enum SpMonitoringManager {
INSTANCE;
+ private static final Logger LOG =
LoggerFactory.getLogger(SpMonitoringManager.class);
+
private final Map<String, FixedSizeList<SpLogEntry>> logInfos;
private final Map<String, SpMetricsEntry> metricsInfos;
@@ -90,6 +95,15 @@ public enum SpMonitoringManager {
public SpEndpointMonitoringInfo getMonitoringInfo() {
var logInfos = makeLogInfos();
+ if (LOG.isDebugEnabled()) {
+ LOG.debug("Providing extension monitoring snapshot: resourceCount={},
totalOutputCounter={}, "
+ + "latestOutputTimestamp={}, thread={}",
+ metricsInfos.size(),
+ totalOutputCounter(),
+ latestOutputTimestamp(),
+ Thread.currentThread().getName());
+ }
+
return new SpEndpointMonitoringInfo(logInfos, metricsInfos);
}
@@ -120,4 +134,19 @@ public enum SpMonitoringManager {
this.metricsInfos.put(resourceId, new SpMetricsEntry());
}
+ private long totalOutputCounter() {
+ return metricsInfos.values()
+ .stream()
+ .mapToLong(metricsEntry -> metricsEntry.getMessagesOut().getCounter())
+ .sum();
+ }
+
+ private long latestOutputTimestamp() {
+ return metricsInfos.values()
+ .stream()
+ .mapToLong(metricsEntry ->
metricsEntry.getMessagesOut().getLastTimestamp())
+ .max()
+ .orElse(0);
+ }
+
}
diff --git
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java
index 859035593f..04910de464 100644
---
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java
+++
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java
@@ -108,4 +108,3 @@ public class SendToBrokerAdapterSink implements
IAdapterPipelineElement {
}
}
-
diff --git
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/monitoring/MonitoringManagement.java
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/monitoring/MonitoringManagement.java
index 91c1fa3c0f..04f94913c4 100644
---
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/monitoring/MonitoringManagement.java
+++
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/monitoring/MonitoringManagement.java
@@ -21,8 +21,13 @@ package
org.apache.streampipes.extensions.management.monitoring;
import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
import org.apache.streampipes.model.monitoring.SpEndpointMonitoringInfo;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
public class MonitoringManagement {
+ private static final Logger LOG =
LoggerFactory.getLogger(MonitoringManagement.class);
+
private final SpMonitoringManager monitoringManager;
public MonitoringManagement() {
@@ -35,9 +40,36 @@ public class MonitoringManagement {
public SpEndpointMonitoringInfo getMonitoringInfos() {
try {
- return monitoringManager.getMonitoringInfo();
+ var monitoringInfo = monitoringManager.getMonitoringInfo();
+ if (LOG.isDebugEnabled()) {
+ LOG.debug("Returning extension monitoring response: resourceCount={},
totalOutputCounter={}, "
+ + "latestOutputTimestamp={}, thread={}",
+ monitoringInfo.getMetricsInfos().size(),
+ totalOutputCounter(monitoringInfo),
+ latestOutputTimestamp(monitoringInfo),
+ Thread.currentThread().getName());
+ }
+
+ return monitoringInfo;
} finally {
monitoringManager.clearAllLogs();
}
}
+
+ private long totalOutputCounter(SpEndpointMonitoringInfo monitoringInfo) {
+ return monitoringInfo.getMetricsInfos()
+ .values()
+ .stream()
+ .mapToLong(metricsEntry -> metricsEntry.getMessagesOut().getCounter())
+ .sum();
+ }
+
+ private long latestOutputTimestamp(SpEndpointMonitoringInfo monitoringInfo) {
+ return monitoringInfo.getMetricsInfos()
+ .values()
+ .stream()
+ .mapToLong(metricsEntry ->
metricsEntry.getMessagesOut().getLastTimestamp())
+ .max()
+ .orElse(0);
+ }
}
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 aaac997e51..9c5ea287fb 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
@@ -24,6 +24,7 @@ import
org.apache.streampipes.health.monitoring.model.HealthCheckData;
import org.apache.streampipes.loadbalance.pipeline.ExtensionsLogProvider;
import org.apache.streampipes.model.connect.adapter.AdapterDescription;
import org.apache.streampipes.model.health.AdapterInstanceState;
+import org.apache.streampipes.model.monitoring.SpMetricsEntry;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -69,6 +70,11 @@ public class AdapterHealthCheck {
.noneMatch(r -> r.getElementId().equals(entry.getElementId()))
)
.toList();
+ LOG.debug("Adapter monitoring candidates: runningAdapters={},
adaptersToRecover={}, "
+ + "adaptersToMonitor={}",
+ healthCheckData.activeResources().runningAdapters().size(),
+ allAdaptersToRecover.size(),
+ adaptersToMonitor.size());
if (!adaptersToMonitor.isEmpty()) {
updateMonitoringMetrics(adaptersToMonitor);
@@ -96,15 +102,32 @@ public class AdapterHealthCheck {
protected void updateMonitoringMetrics(List<AdapterDescription>
runningAdapterDescriptions) {
var adapterMetrics =
AdapterMetricsManager.getInstance().getAdapterMetrics();
- runningAdapterDescriptions
- .forEach(adapterDescription ->
updateTotalEventsPublished(adapterMetrics,
-
adapterDescription.getElementId(),
-
adapterDescription.getName()));
- LOG.debug("Monitoring {} adapter instances", adapterMetrics.size());
+ var totalEventsPublished = 0L;
+ var latestEventTimestamp = 0L;
+ var debugEnabled = LOG.isDebugEnabled();
+
+ for (AdapterDescription adapterDescription : runningAdapterDescriptions) {
+ var metricsEntry = updateTotalEventsPublished(adapterMetrics,
+
adapterDescription.getElementId(),
+
adapterDescription.getName());
+ if (debugEnabled) {
+ totalEventsPublished += metricsEntry.getMessagesOut().getCounter();
+ latestEventTimestamp = Math.max(latestEventTimestamp,
metricsEntry.getMessagesOut().getLastTimestamp());
+ }
+ }
+
+ if (debugEnabled) {
+ LOG.debug("Monitoring {} adapter instances, totalEventsPublished={},
latestEventTimestamp={}",
+ adapterMetrics.size(),
+ totalEventsPublished,
+ latestEventTimestamp);
+ }
}
- private void updateTotalEventsPublished(AdapterMetrics adapterMetrics,
String adapterId,
- String adapterName) {
+ private SpMetricsEntry updateTotalEventsPublished(
+ AdapterMetrics adapterMetrics,
+ String adapterId,
+ String adapterName) {
// Check if the adapter is already registered; if not, register it first.
// This step is crucial, especially when the StreamPipes Core service is
restarted,
@@ -114,8 +137,11 @@ public class AdapterHealthCheck {
adapterMetrics.register(adapterId, adapterName);
}
- adapterMetrics.updateTotalEventsPublished(adapterId, adapterName,
ExtensionsLogProvider.INSTANCE
- .getMetricInfosForResource(adapterId).getMessagesOut().getCounter());
+ var metricsEntry =
ExtensionsLogProvider.INSTANCE.getMetricInfosForResource(adapterId);
+ var counter = metricsEntry.getMessagesOut().getCounter();
+
+ adapterMetrics.updateTotalEventsPublished(adapterId, adapterName, counter);
+ return metricsEntry;
}
diff --git
a/streampipes-load-balancer/src/main/java/org/apache/streampipes/loadbalance/pipeline/ExtensionsLogProvider.java
b/streampipes-load-balancer/src/main/java/org/apache/streampipes/loadbalance/pipeline/ExtensionsLogProvider.java
index 7193a7535e..54a40025f5 100644
---
a/streampipes-load-balancer/src/main/java/org/apache/streampipes/loadbalance/pipeline/ExtensionsLogProvider.java
+++
b/streampipes-load-balancer/src/main/java/org/apache/streampipes/loadbalance/pipeline/ExtensionsLogProvider.java
@@ -24,6 +24,9 @@ import org.apache.streampipes.model.monitoring.SpMetricsEntry;
import org.apache.streampipes.model.pipeline.Pipeline;
import org.apache.streampipes.storage.api.pipeline.IPipelineStorage;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
@@ -36,13 +39,32 @@ public enum ExtensionsLogProvider {
INSTANCE;
+ private static final Logger LOG =
LoggerFactory.getLogger(ExtensionsLogProvider.class);
private static final int MAX_ITEMS = 50;
private final Map<String, List<SpLogEntry>> allLogInfos = new HashMap<>();
private final Map<String, SpMetricsEntry> allMetricsInfos = new HashMap<>();
public void addMonitoringInfos(SpEndpointMonitoringInfo monitoringInfo) {
+ if (LOG.isDebugEnabled()) {
+ LOG.debug("Updating core monitoring cache: incomingResourceCount={},
cachedResourceCountBefore={}, "
+ + "incomingTotalOutputCounter={},
incomingLatestOutputTimestamp={}, thread={}",
+ monitoringInfo.getMetricsInfos().size(),
+ allMetricsInfos.size(),
+ totalOutputCounter(monitoringInfo.getMetricsInfos()),
+ latestOutputTimestamp(monitoringInfo.getMetricsInfos()),
+ Thread.currentThread().getName());
+ }
+
allMetricsInfos.putAll(monitoringInfo.getMetricsInfos());
+ if (LOG.isDebugEnabled()) {
+ LOG.debug("Updated core monitoring cache: cachedResourceCountAfter={},
cachedTotalOutputCounter={}, "
+ + "cachedLatestOutputTimestamp={}",
+ allMetricsInfos.size(),
+ totalOutputCounter(allMetricsInfos),
+ latestOutputTimestamp(allMetricsInfos));
+ }
+
monitoringInfo.getLogInfos().forEach((key, value) -> {
if (!allLogInfos.containsKey(key)) {
allLogInfos.put(key, new ArrayList<>());
@@ -120,6 +142,21 @@ public enum ExtensionsLogProvider {
public Map<String, SpMetricsEntry> getAllMetricsInfos() {
return this.allMetricsInfos;
}
+
+ private long totalOutputCounter(Map<String, SpMetricsEntry> metricsInfos) {
+ return metricsInfos.values()
+ .stream()
+ .mapToLong(metricsEntry -> metricsEntry.getMessagesOut().getCounter())
+ .sum();
+ }
+
+ private long latestOutputTimestamp(Map<String, SpMetricsEntry> metricsInfos)
{
+ return metricsInfos.values()
+ .stream()
+ .mapToLong(metricsEntry ->
metricsEntry.getMessagesOut().getLastTimestamp())
+ .max()
+ .orElse(0);
+ }
private List<String> collectPipelineElementIds(Pipeline pipeline) {
if (pipeline != null){
diff --git
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/ExtensionsServiceLogExecutor.java
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/ExtensionsServiceLogExecutor.java
index 579e964b27..5bac74bd1f 100644
---
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/ExtensionsServiceLogExecutor.java
+++
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/ExtensionsServiceLogExecutor.java
@@ -68,11 +68,29 @@ public class ExtensionsServiceLogExecutor implements
Runnable {
public void triggerUpdate() {
List<SpServiceRegistration> serviceEndpoints =
getActiveExtensionsEndpoints();
+ if (LOG.isDebugEnabled()) {
+ LOG.debug("Monitoring fetch triggered: serviceCount={}, thread={}",
+ serviceEndpoints.size(),
+ Thread.currentThread().getName());
+ }
serviceEndpoints.forEach(serviceEndpoint -> {
try {
+ if (LOG.isDebugEnabled()) {
+ LOG.debug("Fetching monitoring info from extension service:
serviceId={}, serviceUrl={}, thread={}",
+ serviceEndpoint.getSvcId(),
+ serviceEndpoint.getServiceUrl(),
+ Thread.currentThread().getName());
+ }
+
var target =
ExtensionServiceRequestTargets.serviceHealth(serviceEndpoint, LOG_PATH);
var response =
extensionRequestManager.request(ExtensionServiceRequests.serviceHealth(target,
resourceManager));
+ if (LOG.isDebugEnabled()) {
+ LOG.debug("Monitoring fetch response from extension service:
serviceId={}, status={}, success={}",
+ serviceEndpoint.getSvcId(),
+ response.statusCode(),
+ response.isSuccess());
+ }
if (!response.isSuccess()) {
LOG.info("Could not fetch log info from endpoint {} (status {})",
@@ -81,9 +99,18 @@ public class ExtensionsServiceLogExecutor implements
Runnable {
}
SpEndpointMonitoringInfo monitoringInfo =
parseLogResponse(response.responseBody());
+ if (LOG.isDebugEnabled()) {
+ LOG.debug("Fetched monitoring info from extension service:
serviceId={}, resourceCount={}, "
+ + "totalOutputCounter={}, latestOutputTimestamp={}",
+ serviceEndpoint.getSvcId(),
+ monitoringInfo.getMetricsInfos().size(),
+ totalOutputCounter(monitoringInfo),
+ latestOutputTimestamp(monitoringInfo));
+ }
+
ExtensionsLogProvider.INSTANCE.addMonitoringInfos(monitoringInfo);
} catch (IOException e) {
- LOG.info("Could not fetch log info from endpoint {}", serviceEndpoint);
+ LOG.info("Could not fetch log info from endpoint {}", serviceEndpoint,
e);
}
});
@@ -150,4 +177,21 @@ public class ExtensionsServiceLogExecutor implements
Runnable {
throws JsonProcessingException {
return JacksonSerializer.getObjectMapper().readValue(response,
SpEndpointMonitoringInfo.class);
}
+
+ private long totalOutputCounter(SpEndpointMonitoringInfo monitoringInfo) {
+ return monitoringInfo.getMetricsInfos()
+ .values()
+ .stream()
+ .mapToLong(metricsEntry -> metricsEntry.getMessagesOut().getCounter())
+ .sum();
+ }
+
+ private long latestOutputTimestamp(SpEndpointMonitoringInfo monitoringInfo) {
+ return monitoringInfo.getMetricsInfos()
+ .values()
+ .stream()
+ .mapToLong(metricsEntry ->
metricsEntry.getMessagesOut().getLastTimestamp())
+ .max()
+ .orElse(0);
+ }
}
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/AdapterMonitoringResource.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/AdapterMonitoringResource.java
index f3c9630b1a..de66e0b567 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/AdapterMonitoringResource.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/AdapterMonitoringResource.java
@@ -29,6 +29,8 @@ import org.apache.streampipes.model.monitoring.SpMetricsEntry;
import org.apache.streampipes.resource.management.SpResourceManager;
import
org.apache.streampipes.resource.management.permission.SpPermissionEvaluator;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity;
import org.springframework.security.access.prepost.PreAuthorize;
@@ -46,6 +48,8 @@ import java.util.Map;
@RequestMapping("/api/v2/adapter-monitoring")
public class AdapterMonitoringResource extends AbstractMonitoringResource {
+ private static final Logger LOG =
LoggerFactory.getLogger(AdapterMonitoringResource.class);
+
private final ExtensionServiceRequestManager extensionServiceRequestManager;
public AdapterMonitoringResource(ExtensionServiceRequestManager
extensionServiceRequestManager,
@@ -94,6 +98,9 @@ public class AdapterMonitoringResource extends
AbstractMonitoringResource {
public ResponseEntity<Map<String, SpMetricsEntry>> getMetricsInfos(
@RequestParam(value = "filter") List<String> elementIds
) {
+ LOG.debug("Manual adapter monitoring refresh requested from REST endpoint:
filters={}, thread={}",
+ elementIds,
+ Thread.currentThread().getName());
new ExtensionsServiceLogExecutor(extensionServiceRequestManager,
resourceManager).triggerUpdate();
var filteredElementIds = elementIds.stream()
.map(a -> resourceManager.manageAdapters().getDb().getElementById(a))