Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -94,13 +94,16 @@ public MetadataFileContents(String version, String metricsSnapshot) {
public static Optional<Pair<DiagnosticsManager, MetricsSnapshotReporter>> buildDiagnosticsManager(String jobName,
String jobId, JobModel jobModel, String containerId, Optional<String> execEnvContainerId, Config config) {

JobConfig jobConfig = new JobConfig(config);
Optional<Pair<DiagnosticsManager, MetricsSnapshotReporter>> 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 = jobConfig.getThreadPoolSize();

// Diagnostic stream, producer, and reporter related parameters
String diagnosticsReporterName = MetricsConfig.METRICS_SNAPSHOT_REPORTER_NAME_FOR_DIAGNOSTICS;
Expand Down Expand Up @@ -129,7 +132,7 @@ public static Optional<Pair<DiagnosticsManager, MetricsSnapshotReporter>> 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()));

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -65,9 +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 int containerThreadPoolSize;
private final Map<String, ContainerModel> containerModels;
private boolean jobParamsEmitted = false;

Expand All @@ -79,29 +81,55 @@ 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<String, ContainerModel> containerModels,
Integer containerMemoryMb, Integer containerNumCores, Integer numStoresWithChangelog, 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,
public DiagnosticsManager(String jobName,
String jobId,
Map<String, ContainerModel> 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,
terminationDuration, Executors.newSingleThreadScheduledExecutor(
new ThreadFactoryBuilder().setNameFormat(PUBLISH_THREAD_NAME).setDaemon(true).build()));
}

@VisibleForTesting
DiagnosticsManager(String jobName, String jobId, Map<String, ContainerModel> containerModels,
int containerMemoryMb, int containerNumCores, int numStoresWithChangelog, String containerId,
String executionEnvContainerId, String taskClassVersion, String samzaVersion, String hostname,
SystemStream diagnosticSystemStream, SystemProducer systemProducer, Duration terminationDuration,
DiagnosticsManager(String jobName,
String jobId,
Map<String, ContainerModel> 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,
ScheduledExecutorService executorService) {
this.jobName = jobName;
this.jobId = jobId;
this.containerModels = containerModels;
this.containerMemoryMb = containerMemoryMb;
this.containerNumCores = containerNumCores;
this.numStoresWithChangelog = numStoresWithChangelog;
this.maxHeapSizeBytes = maxHeapSizeBytes;
this.containerThreadPoolSize = containerThreadPoolSize;
this.containerId = containerId;
this.executionEnvContainerId = executionEnvContainerId;
this.taskClassVersion = taskClassVersion;
Expand Down Expand Up @@ -185,6 +213,8 @@ public void run() {
diagnosticsStreamMessage.addContainerNumCores(containerNumCores);
diagnosticsStreamMessage.addNumStoresWithChangelog(numStoresWithChangelog);
diagnosticsStreamMessage.addContainerModels(containerModels);
diagnosticsStreamMessage.addMaxHeapSize(maxHeapSizeBytes);
diagnosticsStreamMessage.addContainerThreadPoolSize(containerThreadPoolSize);
}

// Add stop event list to the message
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,9 @@ public class DiagnosticsStreamMessage {
private static final String STOP_EVENT_LIST_METRIC_NAME = "stopEvents";
private static final String CONTAINER_MB_METRIC_NAME = "containerMemoryMb";
private static final String CONTAINER_NUM_CORES_METRIC_NAME = "containerNumCores";
public static final String CONTAINER_NUM_STORES_WITH_CHANGELOG_METRIC_NAME = "numStoresWithChangelog";
private static final String CONTAINER_NUM_STORES_WITH_CHANGELOG_METRIC_NAME = "numStoresWithChangelog";
private static final String CONTAINER_MAX_CONFIGURED_HEAP_METRIC_NAME = "maxHeap";
private static final String CONTAINER_THREAD_POOL_SIZE_METRIC_NAME = "containerThreadPoolSize";
private static final String CONTAINER_MODELS_METRIC_NAME = "containerModels";

private final MetricsHeader metricsHeader;
Expand Down Expand Up @@ -97,6 +99,22 @@ public void addNumStoresWithChangelog(Integer numStoresWithChangelog) {
numStoresWithChangelog);
}

/**
* Add the configured max heap size in bytes.
* @param maxHeapSize the parameter value.
*/
public void addMaxHeapSize(Long maxHeapSize) {
addToMetricsMessage(GROUP_NAME_FOR_DIAGNOSTICS_MANAGER, CONTAINER_MAX_CONFIGURED_HEAP_METRIC_NAME, maxHeapSize);
}

/**
* Add the configured container thread pool size.
* @param threadPoolSize the parameter value.
*/
public void addContainerThreadPoolSize(Integer threadPoolSize) {
addToMetricsMessage(GROUP_NAME_FOR_DIAGNOSTICS_MANAGER, CONTAINER_THREAD_POOL_SIZE_METRIC_NAME, threadPoolSize);
}

/**
* Add a map of container models (indexed by containerID) to the message.
* @param containerModelMap the container models map
Expand Down Expand Up @@ -185,6 +203,14 @@ public Integer getNumStoresWithChangelog() {
CONTAINER_NUM_STORES_WITH_CHANGELOG_METRIC_NAME);
}

public Long getMaxHeapSize() {
Comment thread
rmatharu-zz marked this conversation as resolved.
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<String, ContainerModel> getContainerModels() {
return deserializeContainerModelMap((String) getFromMetricsMessage(GROUP_NAME_FOR_DIAGNOSTICS_MANAGER, CONTAINER_MODELS_METRIC_NAME));
}
Expand All @@ -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<ProcessorStopEvent>) diagnosticsManagerGroupMap.get(STOP_EVENT_LIST_METRIC_NAME));
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<String, ContainerModel> containerModels = TestDiagnosticsStreamMessage.getSampleContainerModels();
Expand All @@ -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);

Expand Down Expand Up @@ -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());
Expand Down