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 @@ -111,49 +111,7 @@ private void init(ProducerConfigurationData conf) {
}

try {
long now = System.nanoTime();
double elapsed = (now - oldTime) / 1e9;
oldTime = now;

long currentNumMsgsSent = numMsgsSent.sumThenReset();
long currentNumBytesSent = numBytesSent.sumThenReset();
long currentNumSendFailedMsgs = numSendFailed.sumThenReset();
long currentNumAcksReceived = numAcksReceived.sumThenReset();

totalMsgsSent.add(currentNumMsgsSent);
totalBytesSent.add(currentNumBytesSent);
totalSendFailed.add(currentNumSendFailedMsgs);
totalAcksReceived.add(currentNumAcksReceived);

synchronized (ds) {
latencyPctValues = ds.getQuantiles(PERCENTILES);
ds.reset();
}

sendMsgsRate = currentNumMsgsSent / elapsed;
sendBytesRate = currentNumBytesSent / elapsed;

if ((currentNumMsgsSent | currentNumSendFailedMsgs | currentNumAcksReceived
| currentNumMsgsSent) != 0) {

for (int i = 0; i < latencyPctValues.length; i++) {
if (Double.isNaN(latencyPctValues[i])) {
latencyPctValues[i] = 0;
}
}

log.info("[{}] [{}] Pending messages: {} --- Publish throughput: {} msg/s --- {} Mbit/s --- "
+ "Latency: med: {} ms - 95pct: {} ms - 99pct: {} ms - 99.9pct: {} ms - max: {} ms --- "
+ "Ack received rate: {} ack/s --- Failed messages: {}", producer.getTopic(),
producer.getProducerName(), producer.getPendingQueueSize(),
THROUGHPUT_FORMAT.format(sendMsgsRate),
THROUGHPUT_FORMAT.format(sendBytesRate / 1024 / 1024 * 8),
DEC.format(latencyPctValues[0]), DEC.format(latencyPctValues[2]),
DEC.format(latencyPctValues[3]), DEC.format(latencyPctValues[4]),
DEC.format(latencyPctValues[5]),
THROUGHPUT_FORMAT.format(currentNumAcksReceived / elapsed), currentNumSendFailedMsgs);
}

updateStats();
} catch (Exception e) {
log.error("[{}] [{}]: {}", producer.getTopic(), producer.getProducerName(), e.getMessage());
} finally {
Expand All @@ -171,6 +129,51 @@ Timeout getStatTimeout() {
return statTimeout;
}

protected void updateStats() {
long now = System.nanoTime();
double elapsed = (now - oldTime) / 1e9;
oldTime = now;

long currentNumMsgsSent = numMsgsSent.sumThenReset();
long currentNumBytesSent = numBytesSent.sumThenReset();
long currentNumSendFailedMsgs = numSendFailed.sumThenReset();
long currentNumAcksReceived = numAcksReceived.sumThenReset();

totalMsgsSent.add(currentNumMsgsSent);
totalBytesSent.add(currentNumBytesSent);
totalSendFailed.add(currentNumSendFailedMsgs);
totalAcksReceived.add(currentNumAcksReceived);

synchronized (ds) {
latencyPctValues = ds.getQuantiles(PERCENTILES);
ds.reset();
}

sendMsgsRate = currentNumMsgsSent / elapsed;
sendBytesRate = currentNumBytesSent / elapsed;

if ((currentNumMsgsSent | currentNumSendFailedMsgs | currentNumAcksReceived
| currentNumMsgsSent) != 0) {

for (int i = 0; i < latencyPctValues.length; i++) {
if (Double.isNaN(latencyPctValues[i])) {
latencyPctValues[i] = 0;
}
}

log.info("[{}] [{}] Pending messages: {} --- Publish throughput: {} msg/s --- {} Mbit/s --- "
+ "Latency: med: {} ms - 95pct: {} ms - 99pct: {} ms - 99.9pct: {} ms - max: {} ms --- "
+ "Ack received rate: {} ack/s --- Failed messages: {}", producer.getTopic(),
producer.getProducerName(), producer.getPendingQueueSize(),
THROUGHPUT_FORMAT.format(sendMsgsRate),
THROUGHPUT_FORMAT.format(sendBytesRate / 1024 / 1024 * 8),
DEC.format(latencyPctValues[0]), DEC.format(latencyPctValues[2]),
DEC.format(latencyPctValues[3]), DEC.format(latencyPctValues[4]),
DEC.format(latencyPctValues[5]),
THROUGHPUT_FORMAT.format(currentNumAcksReceived / elapsed), currentNumSendFailedMsgs);
}
}

@Override
public void updateNumMsgsSent(long numMsgs, long totalMsgsSize) {
numMsgsSent.add(numMsgs);
Expand Down Expand Up @@ -297,6 +300,7 @@ public double getSendLatencyMillisMax() {
}

public void cancelStatsTimeout() {
this.updateStats();
if (statTimeout != null) {
statTimeout.cancel();
statTimeout = null;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -202,4 +202,37 @@ public void testGetStats() throws Exception {
impl.getStats();
}

@Test
public void testGetStatsWithoutArriveUpdateInterval() throws Exception {
String topicName = "test-stats-without-arrive-interval";
ClientConfigurationData conf = new ClientConfigurationData();
conf.setServiceUrl("pulsar://localhost:6650");
conf.setStatsIntervalSeconds(100);

ThreadFactory threadFactory =
new DefaultThreadFactory("client-test-stats", Thread.currentThread().isDaemon());
EventLoopGroup eventLoopGroup = EventLoopUtil
.newEventLoopGroup(conf.getNumIoThreads(), false, threadFactory);

PulsarClientImpl clientImpl = new PulsarClientImpl(conf, eventLoopGroup);

ProducerConfigurationData producerConfData = new ProducerConfigurationData();
producerConfData.setMessageRoutingMode(MessageRoutingMode.CustomPartition);
producerConfData.setCustomMessageRouter(new CustomMessageRouter());

assertEquals(Long.parseLong("100"), clientImpl.getConfiguration().getStatsIntervalSeconds());

PartitionedProducerImpl<byte[]> impl = new PartitionedProducerImpl<>(
clientImpl, topicName, producerConfData,
1, null, null, null);

impl.getProducers().get(0).getStats().incrementSendFailed();
ProducerStatsRecorderImpl stats = impl.getStats();
assertEquals(stats.getTotalSendFailed(), 0);
// When close producer, the ProducerStatsRecorder will update stats immediately
impl.close();
stats = impl.getStats();
assertEquals(stats.getTotalSendFailed(), 1);
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -54,4 +54,24 @@ public void testIncrementNumAcksReceived() throws Exception {
Thread.sleep(1200);
assertEquals(1000.0, recorder.getSendLatencyMillisMax(), 0.5);
}

@Test
public void testGetStatsAndCancelStatsTimeoutWithoutArriveUpdateInterval() {
ClientConfigurationData conf = new ClientConfigurationData();
conf.setStatsIntervalSeconds(60);
PulsarClientImpl client = mock(PulsarClientImpl.class);
when(client.getConfiguration()).thenReturn(conf);
Timer timer = new HashedWheelTimer();
when(client.timer()).thenReturn(timer);
ProducerImpl<?> producer = mock(ProducerImpl.class);
when(producer.getTopic()).thenReturn("topic-test");
when(producer.getProducerName()).thenReturn("producer-test");
when(producer.getPendingQueueSize()).thenReturn(1);
ProducerConfigurationData producerConfigurationData = new ProducerConfigurationData();
ProducerStatsRecorderImpl recorder = new ProducerStatsRecorderImpl(client, producerConfigurationData, producer);
long latencyNs = TimeUnit.SECONDS.toNanos(1);
recorder.incrementNumAcksReceived(latencyNs);
recorder.cancelStatsTimeout();
assertEquals(1000.0, recorder.getSendLatencyMillisMax(), 0.5);
}
}