From 462c6de891b01c20bddaeac78bf2ecc4a21265a0 Mon Sep 17 00:00:00 2001 From: Ray Matharu Date: Mon, 5 Aug 2019 14:15:37 -0700 Subject: [PATCH 1/3] Adding two additional params to DiagnosticsStreamMessage --- .../apache/samza/util/DiagnosticsUtil.java | 4 ++- .../samza/diagnostics/DiagnosticsManager.java | 18 +++++++---- .../diagnostics/DiagnosticsStreamMessage.java | 30 ++++++++++++++++++- .../diagnostics/TestDiagnosticsManager.java | 6 +++- 4 files changed, 49 insertions(+), 9 deletions(-) diff --git a/samza-core/src/main/java/org/apache/samza/util/DiagnosticsUtil.java b/samza-core/src/main/java/org/apache/samza/util/DiagnosticsUtil.java index 4bc0f240b4..5a45d4c980 100644 --- a/samza-core/src/main/java/org/apache/samza/util/DiagnosticsUtil.java +++ b/samza-core/src/main/java/org/apache/samza/util/DiagnosticsUtil.java @@ -101,6 +101,8 @@ public static Optional> buildD ClusterManagerConfig clusterManagerConfig = new ClusterManagerConfig(config); int containerMemoryMb = clusterManagerConfig.getContainerMemoryMb(); int containerNumCores = clusterManagerConfig.getNumCores(); + long maxHeapSizeBytes = Runtime.getRuntime().maxMemory(); + int containerThreadPoolSize = new JobConfig(config).getThreadPoolSize(); // Diagnostic stream, producer, and reporter related parameters String diagnosticsReporterName = MetricsConfig.METRICS_SNAPSHOT_REPORTER_NAME_FOR_DIAGNOSTICS; @@ -129,7 +131,7 @@ public static Optional> buildD systemFactory.getProducer(diagnosticsSystemStream.getSystem(), config, new MetricsRegistryMap()); DiagnosticsManager diagnosticsManager = new DiagnosticsManager(jobName, jobId, jobModel.getContainers(), containerMemoryMb, containerNumCores, - new StorageConfig(config).getNumStoresWithChangelog(), containerId, execEnvContainerId.orElse(""), + new StorageConfig(config).getNumStoresWithChangelog(), maxHeapSizeBytes, containerThreadPoolSize, containerId, execEnvContainerId.orElse(""), taskClassVersion, samzaVersion, hostName, diagnosticsSystemStream, systemProducer, Duration.ofMillis(new TaskConfig(config).getShutdownMs())); diff --git a/samza-core/src/main/scala/org/apache/samza/diagnostics/DiagnosticsManager.java b/samza-core/src/main/scala/org/apache/samza/diagnostics/DiagnosticsManager.java index aa41940ae3..9d8c15808f 100644 --- a/samza-core/src/main/scala/org/apache/samza/diagnostics/DiagnosticsManager.java +++ b/samza-core/src/main/scala/org/apache/samza/diagnostics/DiagnosticsManager.java @@ -68,6 +68,8 @@ public class DiagnosticsManager { private final Integer containerMemoryMb; private final Integer containerNumCores; private final Integer numStoresWithChangelog; + private final long maxHeapSizeBytes; + private final Integer containerThreadPoolSize; private final Map containerModels; private boolean jobParamsEmitted = false; @@ -80,20 +82,20 @@ public class DiagnosticsManager { private final SystemStream diagnosticSystemStream; public DiagnosticsManager(String jobName, String jobId, Map containerModels, - Integer containerMemoryMb, Integer containerNumCores, Integer numStoresWithChangelog, String containerId, - String executionEnvContainerId, String taskClassVersion, String samzaVersion, String hostname, + Integer containerMemoryMb, Integer containerNumCores, Integer numStoresWithChangelog, Long maxHeapSizeBytes, Integer containerThreadPoolSize, + String containerId, String executionEnvContainerId, String taskClassVersion, String samzaVersion, String hostname, SystemStream diagnosticSystemStream, SystemProducer systemProducer, Duration terminationDuration) { - this(jobName, jobId, containerModels, containerMemoryMb, containerNumCores, numStoresWithChangelog, containerId, - executionEnvContainerId, taskClassVersion, samzaVersion, hostname, diagnosticSystemStream, systemProducer, + this(jobName, jobId, containerModels, containerMemoryMb, containerNumCores, numStoresWithChangelog, maxHeapSizeBytes, containerThreadPoolSize, + containerId, executionEnvContainerId, taskClassVersion, samzaVersion, hostname, diagnosticSystemStream, systemProducer, terminationDuration, Executors.newSingleThreadScheduledExecutor( new ThreadFactoryBuilder().setNameFormat(PUBLISH_THREAD_NAME).setDaemon(true).build())); } @VisibleForTesting DiagnosticsManager(String jobName, String jobId, Map containerModels, - int containerMemoryMb, int containerNumCores, int numStoresWithChangelog, String containerId, - String executionEnvContainerId, String taskClassVersion, String samzaVersion, String hostname, + int containerMemoryMb, int containerNumCores, int numStoresWithChangelog, Long maxHeapSizeBytes, Integer containerThreadPoolSize, + String containerId, String executionEnvContainerId, String taskClassVersion, String samzaVersion, String hostname, SystemStream diagnosticSystemStream, SystemProducer systemProducer, Duration terminationDuration, ScheduledExecutorService executorService) { this.jobName = jobName; @@ -102,6 +104,8 @@ public DiagnosticsManager(String jobName, String jobId, Map getContainerModels() { return deserializeContainerModelMap((String) getFromMetricsMessage(GROUP_NAME_FOR_DIAGNOSTICS_MANAGER, CONTAINER_MODELS_METRIC_NAME)); } @@ -210,6 +236,8 @@ public static DiagnosticsStreamMessage convertToDiagnosticsStreamMessage(Metrics diagnosticsStreamMessage.addContainerMb((Integer) diagnosticsManagerGroupMap.get(CONTAINER_MB_METRIC_NAME)); diagnosticsStreamMessage.addNumStoresWithChangelog((Integer) diagnosticsManagerGroupMap.get(CONTAINER_NUM_STORES_WITH_CHANGELOG_METRIC_NAME)); diagnosticsStreamMessage.addContainerModels(deserializeContainerModelMap((String) diagnosticsManagerGroupMap.get(CONTAINER_MODELS_METRIC_NAME))); + diagnosticsStreamMessage.addMaxHeapSize((Long) diagnosticsManagerGroupMap.get(CONTAINER_MAX_CONFIGURED_HEAP_METRIC_NAME)); + diagnosticsStreamMessage.addContainerThreadPoolSize((Integer) diagnosticsManagerGroupMap.get(CONTAINER_THREAD_POOL_SIZE_METRIC_NAME)); diagnosticsStreamMessage.addProcessorStopEvents((List) diagnosticsManagerGroupMap.get(STOP_EVENT_LIST_METRIC_NAME)); } diff --git a/samza-core/src/test/java/org/apache/samza/diagnostics/TestDiagnosticsManager.java b/samza-core/src/test/java/org/apache/samza/diagnostics/TestDiagnosticsManager.java index c69b278d45..33d16e3300 100644 --- a/samza-core/src/test/java/org/apache/samza/diagnostics/TestDiagnosticsManager.java +++ b/samza-core/src/test/java/org/apache/samza/diagnostics/TestDiagnosticsManager.java @@ -53,6 +53,8 @@ public class TestDiagnosticsManager { private String samzaVersion = "1.3.0"; private String hostname = "sample host name"; private int containerMb = 1024; + private int containerThreadPoolSize = 2; + private long maxHeapSize = 900; private int numStoresWithChangelog = 2; private int containerNumCores = 2; private Map containerModels = TestDiagnosticsStreamMessage.getSampleContainerModels(); @@ -73,7 +75,7 @@ public void setup() { }); this.diagnosticsManager = - new DiagnosticsManager(jobName, jobId, containerModels, containerMb, containerNumCores, numStoresWithChangelog, + new DiagnosticsManager(jobName, jobId, containerModels, containerMb, containerNumCores, numStoresWithChangelog, maxHeapSize, containerThreadPoolSize, "0", executionEnvContainerId, taskClassVersion, samzaVersion, hostname, diagnosticsSystemStream, mockSystemProducer, Duration.ofSeconds(1), mockExecutorService); @@ -202,6 +204,8 @@ private void validateOutgoingMessageEnvelope(OutgoingMessageEnvelope outgoingMes DiagnosticsStreamMessage.convertToDiagnosticsStreamMessage(metricsSnapshot); Assert.assertEquals(containerMb, diagnosticsStreamMessage.getContainerMb().intValue()); + Assert.assertEquals(maxHeapSize, diagnosticsStreamMessage.getMaxHeapSize().longValue()); + Assert.assertEquals(containerThreadPoolSize, diagnosticsStreamMessage.getContainerThreadPoolSize().intValue()); Assert.assertEquals(exceptionEventList, diagnosticsStreamMessage.getExceptionEvents()); Assert.assertEquals(diagnosticsStreamMessage.getProcessorStopEvents(), Arrays.asList(new ProcessorStopEvent("0", executionEnvContainerId, hostname, 101))); Assert.assertEquals(containerModels, diagnosticsStreamMessage.getContainerModels()); From 1d6dceceab7fa05374d3e74206e103ea06ab22a9 Mon Sep 17 00:00:00 2001 From: Ray Matharu Date: Tue, 6 Aug 2019 13:56:21 -0700 Subject: [PATCH 2/3] Addressing review comments --- .../apache/samza/util/DiagnosticsUtil.java | 5 +- .../samza/diagnostics/DiagnosticsManager.java | 48 ++++++++++++++----- .../diagnostics/DiagnosticsStreamMessage.java | 8 ++-- 3 files changed, 43 insertions(+), 18 deletions(-) diff --git a/samza-core/src/main/java/org/apache/samza/util/DiagnosticsUtil.java b/samza-core/src/main/java/org/apache/samza/util/DiagnosticsUtil.java index 5a45d4c980..a3245a17bf 100644 --- a/samza-core/src/main/java/org/apache/samza/util/DiagnosticsUtil.java +++ b/samza-core/src/main/java/org/apache/samza/util/DiagnosticsUtil.java @@ -94,15 +94,16 @@ public MetadataFileContents(String version, String metricsSnapshot) { public static Optional> buildDiagnosticsManager(String jobName, String jobId, JobModel jobModel, String containerId, Optional execEnvContainerId, Config config) { + JobConfig jobConfig = new JobConfig(config); Optional> diagnosticsManagerReporterPair = Optional.empty(); - if (new JobConfig(config).getDiagnosticsEnabled()) { + if (jobConfig.getDiagnosticsEnabled()) { ClusterManagerConfig clusterManagerConfig = new ClusterManagerConfig(config); int containerMemoryMb = clusterManagerConfig.getContainerMemoryMb(); int containerNumCores = clusterManagerConfig.getNumCores(); long maxHeapSizeBytes = Runtime.getRuntime().maxMemory(); - int containerThreadPoolSize = new JobConfig(config).getThreadPoolSize(); + int containerThreadPoolSize = jobConfig.getThreadPoolSize(); // Diagnostic stream, producer, and reporter related parameters String diagnosticsReporterName = MetricsConfig.METRICS_SNAPSHOT_REPORTER_NAME_FOR_DIAGNOSTICS; diff --git a/samza-core/src/main/scala/org/apache/samza/diagnostics/DiagnosticsManager.java b/samza-core/src/main/scala/org/apache/samza/diagnostics/DiagnosticsManager.java index 9d8c15808f..624683d1c0 100644 --- a/samza-core/src/main/scala/org/apache/samza/diagnostics/DiagnosticsManager.java +++ b/samza-core/src/main/scala/org/apache/samza/diagnostics/DiagnosticsManager.java @@ -65,11 +65,11 @@ public class DiagnosticsManager { private final Instant resetTime; // Job-related params - private final Integer containerMemoryMb; - private final Integer containerNumCores; - private final Integer numStoresWithChangelog; + private final int containerMemoryMb; + private final int containerNumCores; + private final int numStoresWithChangelog; private final long maxHeapSizeBytes; - private final Integer containerThreadPoolSize; + private final int containerThreadPoolSize; private final Map containerModels; private boolean jobParamsEmitted = false; @@ -81,10 +81,22 @@ public class DiagnosticsManager { private final Duration terminationDuration; // duration to wait when terminating the scheduler private final SystemStream diagnosticSystemStream; - public DiagnosticsManager(String jobName, String jobId, Map containerModels, - Integer containerMemoryMb, Integer containerNumCores, Integer numStoresWithChangelog, Long maxHeapSizeBytes, Integer containerThreadPoolSize, - String containerId, String executionEnvContainerId, String taskClassVersion, String samzaVersion, String hostname, - SystemStream diagnosticSystemStream, SystemProducer systemProducer, Duration terminationDuration) { + public DiagnosticsManager(String jobName, + String jobId, + Map containerModels, + int containerMemoryMb, + int containerNumCores, + int numStoresWithChangelog, + long maxHeapSizeBytes, + int containerThreadPoolSize, + String containerId, + String executionEnvContainerId, + String taskClassVersion, + String samzaVersion, + String hostname, + SystemStream diagnosticSystemStream, + SystemProducer systemProducer, + Duration terminationDuration) { this(jobName, jobId, containerModels, containerMemoryMb, containerNumCores, numStoresWithChangelog, maxHeapSizeBytes, containerThreadPoolSize, containerId, executionEnvContainerId, taskClassVersion, samzaVersion, hostname, diagnosticSystemStream, systemProducer, @@ -93,10 +105,22 @@ public DiagnosticsManager(String jobName, String jobId, Map containerModels, - int containerMemoryMb, int containerNumCores, int numStoresWithChangelog, Long maxHeapSizeBytes, Integer containerThreadPoolSize, - String containerId, String executionEnvContainerId, String taskClassVersion, String samzaVersion, String hostname, - SystemStream diagnosticSystemStream, SystemProducer systemProducer, Duration terminationDuration, + DiagnosticsManager(String jobName, + String jobId, + Map containerModels, + int containerMemoryMb, + int containerNumCores, + int numStoresWithChangelog, + Long maxHeapSizeBytes, + Integer containerThreadPoolSize, + String containerId, + String executionEnvContainerId, + String taskClassVersion, + String samzaVersion, + String hostname, + SystemStream diagnosticSystemStream, + SystemProducer systemProducer, + Duration terminationDuration, ScheduledExecutorService executorService) { this.jobName = jobName; this.jobId = jobId; diff --git a/samza-core/src/main/scala/org/apache/samza/diagnostics/DiagnosticsStreamMessage.java b/samza-core/src/main/scala/org/apache/samza/diagnostics/DiagnosticsStreamMessage.java index bb41088778..81642d5545 100644 --- a/samza-core/src/main/scala/org/apache/samza/diagnostics/DiagnosticsStreamMessage.java +++ b/samza-core/src/main/scala/org/apache/samza/diagnostics/DiagnosticsStreamMessage.java @@ -203,14 +203,14 @@ public Integer getNumStoresWithChangelog() { CONTAINER_NUM_STORES_WITH_CHANGELOG_METRIC_NAME); } - public Integer getContainerThreadPoolSize() { - return (Integer) getFromMetricsMessage(GROUP_NAME_FOR_DIAGNOSTICS_MANAGER, CONTAINER_THREAD_POOL_SIZE_METRIC_NAME); - } - public Long getMaxHeapSize() { return (Long) getFromMetricsMessage(GROUP_NAME_FOR_DIAGNOSTICS_MANAGER, CONTAINER_MAX_CONFIGURED_HEAP_METRIC_NAME); } + public Integer getContainerThreadPoolSize() { + return (Integer) getFromMetricsMessage(GROUP_NAME_FOR_DIAGNOSTICS_MANAGER, CONTAINER_THREAD_POOL_SIZE_METRIC_NAME); + } + public Map getContainerModels() { return deserializeContainerModelMap((String) getFromMetricsMessage(GROUP_NAME_FOR_DIAGNOSTICS_MANAGER, CONTAINER_MODELS_METRIC_NAME)); } From 2f803b1decf7ab1761f4786925f7765b17f64add Mon Sep 17 00:00:00 2001 From: Ray Matharu Date: Tue, 6 Aug 2019 14:32:34 -0700 Subject: [PATCH 3/3] minor --- .../org/apache/samza/diagnostics/DiagnosticsManager.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/samza-core/src/main/scala/org/apache/samza/diagnostics/DiagnosticsManager.java b/samza-core/src/main/scala/org/apache/samza/diagnostics/DiagnosticsManager.java index 624683d1c0..ed5179e5a7 100644 --- a/samza-core/src/main/scala/org/apache/samza/diagnostics/DiagnosticsManager.java +++ b/samza-core/src/main/scala/org/apache/samza/diagnostics/DiagnosticsManager.java @@ -111,8 +111,8 @@ public DiagnosticsManager(String jobName, int containerMemoryMb, int containerNumCores, int numStoresWithChangelog, - Long maxHeapSizeBytes, - Integer containerThreadPoolSize, + long maxHeapSizeBytes, + int containerThreadPoolSize, String containerId, String executionEnvContainerId, String taskClassVersion,