This is an automated email from the ASF dual-hosted git repository. tenthe pushed a commit to branch debug-adapter-monitoring-metrics in repository https://gitbox.apache.org/repos/asf/streampipes.git
commit 0284967cb3101893a1107b718652bd2ae895ad25 Author: Philipp Zehnder <[email protected]> AuthorDate: Mon Jun 29 17:08:39 2026 +0200 Add debug logging for adapter monitoring metrics --- .../api/monitoring/SpMonitoringManager.java | 27 +++++++++++++++ .../elements/SendToBrokerAdapterSink.java | 1 - .../monitoring/MonitoringManagement.java | 32 +++++++++++++++++- .../health/monitoring/AdapterHealthCheck.java | 39 +++++++++++++++++----- .../pipeline/ExtensionsLogProvider.java | 33 ++++++++++++++++++ .../pipeline/ExtensionsServiceLogExecutor.java | 38 ++++++++++++++++++++- .../rest/impl/AdapterMonitoringResource.java | 7 ++++ 7 files changed, 165 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..3699ddd59c 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,13 @@ public enum SpMonitoringManager { public SpEndpointMonitoringInfo getMonitoringInfo() { var logInfos = makeLogInfos(); + LOG.debug("Providing extension monitoring snapshot: resourceCount={}, totalOutputCounter={}, " + + "latestOutputTimestamp={}, thread={}", + metricsInfos.size(), + totalOutputCounter(), + latestOutputTimestamp(), + Thread.currentThread().getName()); + return new SpEndpointMonitoringInfo(logInfos, metricsInfos); } @@ -120,4 +132,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..49c7b499a7 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,34 @@ public class MonitoringManagement { public SpEndpointMonitoringInfo getMonitoringInfos() { try { - return monitoringManager.getMonitoringInfo(); + var monitoringInfo = monitoringManager.getMonitoringInfo(); + 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..b30007e5a3 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,27 @@ 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; + + for (AdapterDescription adapterDescription : runningAdapterDescriptions) { + var metricsEntry = updateTotalEventsPublished(adapterMetrics, + adapterDescription.getElementId(), + adapterDescription.getName()); + totalEventsPublished += metricsEntry.getMessagesOut().getCounter(); + latestEventTimestamp = Math.max(latestEventTimestamp, metricsEntry.getMessagesOut().getLastTimestamp()); + } + + 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 +132,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..689f310d09 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,28 @@ 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) { + 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()); + 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 +138,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..3024d1bae9 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,23 @@ public class ExtensionsServiceLogExecutor implements Runnable { public void triggerUpdate() { List<SpServiceRegistration> serviceEndpoints = getActiveExtensionsEndpoints(); + LOG.debug("Monitoring fetch triggered: serviceCount={}, thread={}", + serviceEndpoints.size(), + Thread.currentThread().getName()); serviceEndpoints.forEach(serviceEndpoint -> { try { + 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)); + 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 +93,16 @@ public class ExtensionsServiceLogExecutor implements Runnable { } SpEndpointMonitoringInfo monitoringInfo = parseLogResponse(response.responseBody()); + 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 +169,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))
