# How to measure streaming time on each record?

**URL:** <https://forum.confluent.io/t/how-to-measure-streaming-time-on-each-record/3629>\
**Category:** Kafka Streams\
**Created:** [11 December 2021 15:05 UTC](https://forum.confluent.io/t/how-to-measure-streaming-time-on-each-record/3629 "2021-12-11T15:05:21Z")\
**Posts on this page:** 13\
**Page:** 1

<div class="post-metadata">

**Author:** ![programista4k](https://sea1.discourse-cdn.com/flex019/user_avatar/forum.confluent.io/programista4k/32/767_2.png) [@programista4k](https://forum.confluent.io/u/programista4k)\
**Post date:** [11 December 2021 15:05 UTC](https://forum.confluent.io/t/how-to-measure-streaming-time-on-each-record/3629/1 "2021-12-11T15:05:21Z")

</div>

Hello frens,  
I’ve add some metrics to my cryptocurrency streaming application (it source all cryptocurrency transactions from 60 markets and aggregates average prices on it to make a visualisation) and the results worries me.

It takes about 20 ms on my old computer from the moment when service1 sends a message to Kafka to the moment when service2 consumes it.

but simmilar measurement on Kafka Streams shows me about 150 ms latency!

Fragment of the code:

```auto
// service 1:
kStream.process() -> new Processor() {
...
@Override
public void process(K k, V v) {
    timerZero.record(now() - v.getEventTime());
}
...
}
...
kStream.to("myTopic", Produced.with(myKeySerde, myValSerde);
kStream.process(() -> new Processor() {
...
@Override
public void process(K k, V v){
    timerStart.record(now() - v.getEventTime());
}
...
}

// service 2:
streamsBuilder.stream("myTopic", Consumed.with(myKeySerde, myValSerce))
.process(() -> new Processor() {
...
@Override
public void process(K k, V v){
    timerStop.record(now() - v.getEventTime());
}
...
}

```

example results:

**timerZero: 20 ms. after event Time**  
**timerStart: 20 ms. after eventTime**  
**timerStop: 120 ms. after eventTime**!!!

Is it a good way to measure it?

---

<div class="post-metadata">

**Author:** ![an0r0c](https://sea1.discourse-cdn.com/flex019/user_avatar/forum.confluent.io/an0r0c/32/1296_2.png) [@an0r0c](https://forum.confluent.io/u/an0r0c)\
**Post date:** [11 December 2021 18:15 UTC](https://forum.confluent.io/t/how-to-measure-streaming-time-on-each-record/3629/2 "2021-12-11T18:15:30Z")

</div>

I would take a look at [GitHub - opentracing-contrib/java-kafka-client: OpenTracing Instrumentation for Apache Kafka Client](https://github.com/opentracing-contrib/java-kafka-client)  
This + i.e. Jaeger gives you a very good visualization where you “loose” time.

You can use Jaeger in a docker Container without external persistence for an easy Start.  
In my experience you can setup this in about 2 hours locally (incl Research of how exactly)

---

<div class="post-metadata">

**Author:** ![programista4k](https://sea1.discourse-cdn.com/flex019/user_avatar/forum.confluent.io/programista4k/32/767_2.png) [@programista4k](https://forum.confluent.io/u/programista4k)\
**Post date:** [11 December 2021 19:08 UTC](https://forum.confluent.io/t/how-to-measure-streaming-time-on-each-record/3629/3 "2021-12-11T19:08:51Z")

</div>

so you mean measure times in the interceptors/callbacks of consumer and producer

---

<div class="post-metadata">

**Author:** ![an0r0c](https://sea1.discourse-cdn.com/flex019/user_avatar/forum.confluent.io/an0r0c/32/1296_2.png) [@an0r0c](https://forum.confluent.io/u/an0r0c)\
**Post date:** [11 December 2021 19:33 UTC](https://forum.confluent.io/t/how-to-measure-streaming-time-on-each-record/3629/4 "2021-12-11T19:33:35Z")

</div>

Yes, this is imho an easy way to go. For sure it won’t give you details in the steps between one producer/consumer in your Pipeline but I assume you are not doing compute intensive stuff there

To my experience when it comes down to \<100 Milliseconds e2e latency in kstreams you need to tune a lot to achieve this.  
Each producer/consumer adds latency due Network, polling …

---

<div class="post-metadata">

**Author:** ![programista4k](https://sea1.discourse-cdn.com/flex019/user_avatar/forum.confluent.io/programista4k/32/767_2.png) [@programista4k](https://forum.confluent.io/u/programista4k)\
**Post date:** [11 December 2021 19:37 UTC](https://forum.confluent.io/t/how-to-measure-streaming-time-on-each-record/3629/5 "2021-12-11T19:37:52Z")

</div>

Yeah, but assuming I have measured times well then why using normal kafka write/read takes 20ms and in kafka streams 150 ms…

---

<div class="post-metadata">

**Author:** ![an0r0c](https://sea1.discourse-cdn.com/flex019/user_avatar/forum.confluent.io/an0r0c/32/1296_2.png) [@an0r0c](https://forum.confluent.io/u/an0r0c)\
**Post date:** [11 December 2021 19:48 UTC](https://forum.confluent.io/t/how-to-measure-streaming-time-on-each-record/3629/6 "2021-12-11T19:48:15Z")

</div>

Kstreams has some configs like commit.interval.ms which - depending on your topology- can slow things down in terms of latency to usually enable higher throughput.

I would have a look on linger.ms and commit.interval as a first step

---

<div class="post-metadata">

**Author:** ![programista4k](https://sea1.discourse-cdn.com/flex019/user_avatar/forum.confluent.io/programista4k/32/767_2.png) [@programista4k](https://forum.confluent.io/u/programista4k)\
**Post date:** [11 December 2021 22:12 UTC](https://forum.confluent.io/t/how-to-measure-streaming-time-on-each-record/3629/7 "2021-12-11T22:12:15Z")

</div>

> [@an0r0c](#):
>
> commit.interval.ms

this is the solution. After setting `commit.interval.ms=0` it takes 10 ms. now!

---

<div class="post-metadata">

**Author:** ![mjsax](https://sea1.discourse-cdn.com/flex019/user_avatar/forum.confluent.io/mjsax/32/3113_2.png) [@mjsax](https://forum.confluent.io/u/mjsax)\
**Post date:** [13 December 2021 23:16 UTC](https://forum.confluent.io/t/how-to-measure-streaming-time-on-each-record/3629/8 "2021-12-13T23:16:30Z")

</div>

Just be careful with this setting. If you have higher throughput, you could put quite some load on the broker if you commit all the time.

---

<div class="post-metadata">

**Author:** ![an0r0c](https://sea1.discourse-cdn.com/flex019/user_avatar/forum.confluent.io/an0r0c/32/1296_2.png) [@an0r0c](https://forum.confluent.io/u/an0r0c)\
**Post date:** [15 December 2021 17:34 UTC](https://forum.confluent.io/t/how-to-measure-streaming-time-on-each-record/3629/9 "2021-12-15T17:34:09Z")

</div>

@mjsax: I already had some challenges in finding the “best fit” in terms to optimize latency vs. throughput and load in some topologies. Do you have some best practice material about what parameters at all might be relevant and where to be careful (trade-offs) - I think that would be quite interesting. i.e. smaller commit.interval imho reduces throughput per partition and increases load on broker.  
I just did it based on my experience (incl. try-and-error ^^) and on some google research and would be interested on some more insights on that topic.

---

<div class="post-metadata">

**Author:** ![mjsax](https://sea1.discourse-cdn.com/flex019/user_avatar/forum.confluent.io/mjsax/32/3113_2.png) [@mjsax](https://forum.confluent.io/u/mjsax)\
**Post date:** [15 December 2021 21:09 UTC](https://forum.confluent.io/t/how-to-measure-streaming-time-on-each-record/3629/10 "2021-12-15T21:09:37Z")

</div>

I am not aware of any comprehensive single document from the top of my head.

---

<div class="post-metadata">

**Author:** ![programista4k](https://sea1.discourse-cdn.com/flex019/user_avatar/forum.confluent.io/programista4k/32/767_2.png) [@programista4k](https://forum.confluent.io/u/programista4k)\
**Post date:** [15 December 2021 23:06 UTC](https://forum.confluent.io/t/how-to-measure-streaming-time-on-each-record/3629/11 "2021-12-15T23:06:38Z")

</div>

hmm… following Sir @an0r0c 's advice I have created my own Producer with callback. The callback sends metrics. So after every message sent to Kafka the metrics to Prometheus are sent too.

```auto
class MyProducer implements Producer<String, TracedRecord> {
...
    @Override
    public Future<RecordMetadata> send(ProducerRecord<String, TracedRecord> producerRecord) {
        return send(producerRecord, ((recordMetadata, e) -> {
            timer.record(clock.millis() - producerRecord.value().getBirthTimestamp(), TimeUnit.MILLISECONDS);
        }));
    }
...

```

To make Kafka Streams use it I must create `KafkaClientSupplier` which will return this producer. But look at the types of the producer, it’s byte ! The KafkaClientSupplier is not generic!

```auto
public interface KafkaClientSupplier {
...
    Producer<byte[], byte[]> getProducer(Map<String, Object> var1);
...
}

```

And my Producer is `Producer<String, TracedRecord>`. It must be like this because I measure the diff between `clock.millis()` and `tradeRecord.getBirthTimestamp()`.  
Deserializing byte to TraceRecord in every callback (so for every message) will be overkill, frens.

 ![comment_1625042959lsDB470io1HQetBRi7A58E,w1200h627f](https://us1.discourse-cdn.com/flex019/uploads/confluentcommunity/original/2X/b/bb9a5e60f67e715402b723aece61529106d3c5e6.jpeg)

---

<div class="post-metadata">

**Author:** ![an0r0c](https://sea1.discourse-cdn.com/flex019/user_avatar/forum.confluent.io/an0r0c/32/1296_2.png) [@an0r0c](https://forum.confluent.io/u/an0r0c)\
**Post date:** [16 December 2021 05:30 UTC](https://forum.confluent.io/t/how-to-measure-streaming-time-on-each-record/3629/12 "2021-12-16T05:30:48Z")

</div>

Do you really need to deserialize it? It’s a will until I used that last time but I think opentracing also just adds headers to the record therefore deserialization shouldn’t be necessary.

---

<div class="post-metadata">

**Author:** ![system](https://us1.discourse-cdn.com/flex019/uploads/confluentcommunity/original/1X/c49438c90c9df282e9996fdf6971be890c71b65a.svg) [@system](https://forum.confluent.io/u/system)\
**Post date:** [23 December 2021 05:31 UTC](https://forum.confluent.io/t/how-to-measure-streaming-time-on-each-record/3629/13 "2021-12-23T05:31:11Z")

</div>

This topic was automatically closed 7 days after the last reply. New replies are no longer allowed.
