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
6 changes: 6 additions & 0 deletions pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -479,6 +479,12 @@ flexible messaging model and an intuitive client API.</description>
<artifactId>aspectjweaver</artifactId>
<version>${aspectj.version}</version>
</dependency>

<dependency>
<groupId>com.ea.agentloader</groupId>
<artifactId>ea-agent-loader</artifactId>
<version>1.0.2</version>
</dependency>

</dependencies>
</dependencyManagement>
Expand Down
70 changes: 70 additions & 0 deletions pulsar-broker/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -215,6 +215,22 @@
<artifactId>java-semver</artifactId>
</dependency>

<!-- aspectJ dependencies -->
<dependency>
<groupId>org.aspectj</groupId>
<artifactId>aspectjrt</artifactId>
</dependency>

<dependency>
<groupId>org.aspectj</groupId>
<artifactId>aspectjweaver</artifactId>
</dependency>

<dependency>
<groupId>com.ea.agentloader</groupId>
<artifactId>ea-agent-loader</artifactId>
</dependency>

<dependency>
<groupId>org.mockito</groupId>
<artifactId>mockito-core</artifactId>
Expand All @@ -236,6 +252,31 @@

<build>
<plugins>
<plugin>
<groupId>org.codehaus.mojo</groupId>
<artifactId>aspectj-maven-plugin</artifactId>
<version>1.10</version>
<configuration>
<complianceLevel>1.8</complianceLevel>
<source>1.8</source>
<target>1.8</target>
<showWeaveInfo>true</showWeaveInfo>
<weaveDependencies>
<weaveDependency>
<groupId>org.apache.zookeeper</groupId>
<artifactId>zookeeper</artifactId>
</weaveDependency>
</weaveDependencies>
</configuration>
<executions>
<execution>
<phase>process-sources</phase>
<goals>
<goal>compile</goal>
</goals>
</execution>
</executions>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-surefire-plugin</artifactId>
Expand Down Expand Up @@ -283,6 +324,35 @@
</executions>
</plugin>
</plugins>
<pluginManagement>
<plugins>
<!--This plugin's configuration is used to store Eclipse m2e settings only. It has no influence on the Maven build itself.-->
<plugin>
<groupId>org.eclipse.m2e</groupId>
<artifactId>lifecycle-mapping</artifactId>
<version>1.0.0</version>
<configuration>
<lifecycleMappingMetadata>
<pluginExecutions>
<pluginExecution>
<pluginExecutionFilter>
<groupId>org.codehaus.mojo</groupId>
<artifactId>aspectj-maven-plugin</artifactId>
<versionRange>[1.10,)</versionRange>
<goals>
<goal>compile</goal>
</goals>
</pluginExecutionFilter>
<action>
<ignore></ignore>
</action>
</pluginExecution>
</pluginExecutions>
</lifecycleMappingMetadata>
</configuration>
</plugin>
</plugins>
</pluginManagement>
<resources>
<resource>
<directory>src/main/resources</directory>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,10 +26,13 @@
import org.apache.pulsar.broker.PulsarServerException;
import org.apache.pulsar.broker.PulsarService;
import org.apache.pulsar.broker.ServiceConfiguration;
import org.aspectj.weaver.loadtime.Agent;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.bridge.SLF4JBridgeHandler;

import com.ea.agentloader.AgentLoader;

public class PulsarBrokerStarter {

private static ServiceConfiguration loadConfig(String configFile) throws Exception {
Expand All @@ -53,6 +56,9 @@ public static void main(String[] args) throws Exception {
String configFile = args[0];
ServiceConfiguration config = loadConfig(configFile);

// load aspectj-weaver agent for instrumentation
AgentLoader.loadAgentClass(Agent.class.getName(), null);

@SuppressWarnings("resource")
final PulsarService service = new PulsarService(config);
Runtime.getRuntime().addShutdownHook(service.getShutdownService());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,6 @@
import static org.apache.commons.lang3.StringUtils.isBlank;

import java.io.FileInputStream;
import java.net.URI;
import java.net.URL;

import org.apache.pulsar.broker.PulsarService;
Expand All @@ -33,11 +32,13 @@
import org.apache.pulsar.common.policies.data.ClusterData;
import org.apache.pulsar.common.policies.data.PropertyAdmin;
import org.apache.pulsar.zookeeper.LocalBookkeeperEnsemble;
import org.aspectj.weaver.loadtime.Agent;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import com.beust.jcommander.JCommander;
import com.beust.jcommander.Parameter;
import com.ea.agentloader.AgentLoader;
import com.google.common.collect.Lists;
import com.google.common.collect.Sets;

Expand Down Expand Up @@ -155,6 +156,9 @@ void start() throws Exception {
return;
}

// load aspectj-weaver agent for instrumentation
AgentLoader.loadAgentClass(Agent.class.getName(), null);

// Start Broker
broker = new PulsarService(config);
broker.start();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,18 +33,15 @@

import org.apache.bookkeeper.mledger.ManagedLedgerFactory;
import org.apache.bookkeeper.util.OrderedSafeExecutor;
import org.apache.pulsar.broker.PulsarServerException;
import org.apache.pulsar.broker.ServiceConfiguration;
import org.apache.pulsar.broker.ServiceConfigurationUtils;
import org.apache.pulsar.broker.admin.AdminResource;
import org.apache.pulsar.broker.cache.ConfigurationCacheService;
import org.apache.pulsar.broker.cache.LocalZooKeeperCacheService;
import org.apache.pulsar.broker.loadbalance.LeaderElectionService;
import org.apache.pulsar.broker.loadbalance.LeaderElectionService.LeaderListener;
import org.apache.pulsar.broker.loadbalance.LoadManager;
import org.apache.pulsar.broker.loadbalance.LoadReportUpdaterTask;
import org.apache.pulsar.broker.loadbalance.LoadResourceQuotaUpdaterTask;
import org.apache.pulsar.broker.loadbalance.LoadSheddingTask;
import org.apache.pulsar.broker.loadbalance.LeaderElectionService.LeaderListener;
import org.apache.pulsar.broker.loadbalance.impl.SimpleLoadManagerImpl;
import org.apache.pulsar.broker.namespace.NamespaceService;
import org.apache.pulsar.broker.service.BrokerService;
Expand All @@ -67,7 +64,7 @@
import org.apache.pulsar.zookeeper.LocalZooKeeperConnectionService;
import org.apache.pulsar.zookeeper.ZooKeeperCache;
import org.apache.pulsar.zookeeper.ZooKeeperClientFactory;
import org.apache.pulsar.zookeeper.ZookeeperClientFactoryImpl;
import org.apache.pulsar.zookeeper.ZookeeperBkClientFactoryImpl;
import org.apache.zookeeper.ZooKeeper;
import org.eclipse.jetty.servlet.ServletHolder;
import org.slf4j.Logger;
Expand Down Expand Up @@ -564,7 +561,7 @@ public LocalZooKeeperCacheService getLocalZkCacheService() {

public ZooKeeperClientFactory getZooKeeperClientFactory() {
if (zkClientFactory == null) {
zkClientFactory = new ZookeeperClientFactoryImpl();
zkClientFactory = new ZookeeperBkClientFactoryImpl();
}
// Return default factory
return zkClientFactory;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,7 @@
import org.apache.pulsar.broker.service.persistent.PersistentTopic;
import org.apache.pulsar.broker.stats.ClusterReplicationMetrics;
import org.apache.pulsar.broker.web.PulsarWebResource;
import org.apache.pulsar.broker.zookeeper.aspectj.ClientCnxnAspect;
import org.apache.pulsar.client.api.ClientConfiguration;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.PulsarClientException;
Expand Down Expand Up @@ -189,6 +190,10 @@ public BrokerService(PulsarService pulsar) throws Exception {

this.multiLayerTopicsMap = new ConcurrentOpenHashMap<>();
this.pulsarStats = new PulsarStats(pulsar);
// register listener to capture zk-latency
ClientCnxnAspect.addListener((eventType, latencyMs) -> {
this.pulsarStats.recordZkLatencyTimeValue(eventType, latencyMs);
});
this.offlineTopicStatCache = new ConcurrentOpenHashMap<>();

final DefaultThreadFactory acceptorThreadFactory = new DefaultThreadFactory("pulsar-acceptor");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
import org.apache.pulsar.broker.stats.BrokerOperabilityMetrics;
import org.apache.pulsar.broker.stats.ClusterReplicationMetrics;
import org.apache.pulsar.broker.stats.NamespaceStats;
import org.apache.pulsar.broker.zookeeper.aspectj.ClientCnxnAspect.EventType;
import org.apache.pulsar.common.naming.NamespaceBundle;
import org.apache.pulsar.common.stats.Metrics;
import org.apache.pulsar.common.util.collections.ConcurrentOpenHashMap;
Expand Down Expand Up @@ -190,4 +191,16 @@ public void recordTopicLoadTimeValue(String topic, long topicLoadLatencyMs) {
log.warn("Exception while recording topic load time for topic {}, {}", topic, ex.getMessage());
}
}

public void recordZkLatencyTimeValue(EventType eventType, long latencyMs) {
try {
if (EventType.write.equals(eventType)) {
brokerOperabilityMetrics.recordZkWriteLatencyTimeValue(latencyMs);
} else if (EventType.read.equals(eventType)) {
brokerOperabilityMetrics.recordZkReadLatencyTimeValue(latencyMs);
}
} catch (Exception ex) {
log.warn("Exception while recording zk-latency {}, {}", eventType, ex.getMessage());
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.TimeUnit;

import org.apache.pulsar.common.stats.Metrics;

Expand All @@ -31,13 +32,17 @@
public class BrokerOperabilityMetrics {
private final List<Metrics> metricsList;
private final String localCluster;
private final TopicLoadStats topicLoadStats;
private final DimensionStats topicLoadStats;
private final DimensionStats zkWriteLatencyStats;
private final DimensionStats zkReadLatencyStats;
private final String brokerName;

public BrokerOperabilityMetrics(String localCluster, String brokerName) {
this.metricsList = new ArrayList<>();
this.localCluster = localCluster;
this.topicLoadStats = new TopicLoadStats();
this.topicLoadStats = new DimensionStats("topic_load_times", 60);
this.zkWriteLatencyStats = new DimensionStats("zk_write_latency", 60);
this.zkReadLatencyStats = new DimensionStats("zk_read_latency", 60);
this.brokerName = brokerName;
}

Expand All @@ -48,25 +53,37 @@ public List<Metrics> getMetrics() {

private void generate() {
metricsList.add(getTopicLoadMetrics());
metricsList.add(getZkWriteLatencyMetrics());
metricsList.add(getZkReadLatencyMetrics());
}

Metrics getTopicLoadMetrics() {
return getDimensionMetrics("topic_load_times", "topic_load", topicLoadStats);
}

Metrics getZkWriteLatencyMetrics() {
return getDimensionMetrics("zk_write_latency", "zk_write_latency", zkWriteLatencyStats);
}

Metrics getZkReadLatencyMetrics() {
return getDimensionMetrics("zk_read_latency", "zk_read_latency", zkReadLatencyStats);
}

Metrics getDimensionMetrics(String metricsName, String dimensionName, DimensionStats stats) {
Map<String, String> dimensionMap = Maps.newHashMap();
dimensionMap.put("broker", brokerName);
dimensionMap.put("cluster", localCluster);
dimensionMap.put("metric", "topic_load_times");
dimensionMap.put("metric", metricsName);
Metrics dMetrics = Metrics.create(dimensionMap);

topicLoadStats.updateStats();

dMetrics.put("brk_topic_load_time_mean_ms", topicLoadStats.meanTopicLoadMs);
dMetrics.put("brk_topic_load_time_median_ms", topicLoadStats.medianTopicLoadMs);
dMetrics.put("brk_topic_load_time_95percentile_ms", topicLoadStats.topicLoad95Ms);
dMetrics.put("brk_topic_load_time_99_percentile_ms", topicLoadStats.topicLoad99Ms);
dMetrics.put("brk_topic_load_time_99_9_percentile_ms", topicLoadStats.topicLoad999Ms);
dMetrics.put("brk_topic_load_time_99_99_percentile_ms", topicLoadStats.topicsLoad9999Ms);
dMetrics.put("brk_topic_load_rate_s",
(1000 * topicLoadStats.topicLoadCounts) / topicLoadStats.elapsedIntervalMs);
dMetrics.put("brk_" + dimensionName + "_time_mean_ms", stats.getMeanDimension());
dMetrics.put("brk_" + dimensionName + "_time_median_ms", stats.getMedianDimension());
dMetrics.put("brk_" + dimensionName + "_time_75percentile_ms", stats.getDimension75());
dMetrics.put("brk_" + dimensionName + "_time_95percentile_ms", stats.getDimension95());
dMetrics.put("brk_" + dimensionName + "_time_99_percentile_ms", stats.getDimension99());
dMetrics.put("brk_" + dimensionName + "_time_99_9_percentile_ms", stats.getDimension999());
dMetrics.put("brk_" + dimensionName + "_time_99_99_percentile_ms", stats.getDimension9999());
dMetrics.put("brk_" + dimensionName + "_rate_s", stats.getDimensionCount());

return dMetrics;
}
Expand All @@ -76,6 +93,14 @@ public void reset() {
}

public void recordTopicLoadTimeValue(long topicLoadLatencyMs) {
topicLoadStats.recordTopicLoadTimeValue(topicLoadLatencyMs);
topicLoadStats.recordDimensionTimeValue(topicLoadLatencyMs, TimeUnit.MILLISECONDS);
}

public void recordZkWriteLatencyTimeValue(long topicLoadLatencyMs) {
zkWriteLatencyStats.recordDimensionTimeValue(topicLoadLatencyMs, TimeUnit.MILLISECONDS);
}

public void recordZkReadLatencyTimeValue(long topicLoadLatencyMs) {
zkReadLatencyStats.recordDimensionTimeValue(topicLoadLatencyMs, TimeUnit.MILLISECONDS);
}
}
Loading