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))

Reply via email to