Skip to content

Commit e598837

Browse files
authored
Standardized Kafka telemetry (#623)
* Improved and cleaned and standardized up KafkaPublisher & KafkaListener telemetry * Replaced with ONLY Meter.Provider * Added TimerMeter * Added CounterMeter and GaugeLongMeter and cleaned up full Kafka telemetry * Cleanup * Cleanup * Final Cleanup * Improve telemetry * Kafka Telemetry final cleanup * Log improved * Standardized telemetry for producer and consumer * Cleanup * Fix * Standardized telemetry * Cleanup
1 parent 3f1973a commit e598837

62 files changed

Lines changed: 2085 additions & 906 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

common/src/main/java/io/koraframework/common/util/TimeUtils.java

Lines changed: 10 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -11,15 +11,19 @@ public static long started() {
1111
return System.nanoTime();
1212
}
1313

14-
public static Duration took(long started) {
15-
return Duration.ofNanos(System.nanoTime() - started).truncatedTo(ChronoUnit.MILLIS);
14+
public static Duration took(long startedNanos) {
15+
return Duration.ofNanos(System.nanoTime() - startedNanos).truncatedTo(ChronoUnit.MILLIS);
1616
}
1717

18-
public static String tookForLogging(long started) {
19-
return durationForLogging(System.nanoTime() - started);
18+
public static String tookForLogging(long startedNanos) {
19+
return durationForLogging(System.nanoTime() - startedNanos);
2020
}
2121

22-
public static String durationForLogging(long duration) {
23-
return Duration.ofNanos(duration).truncatedTo(ChronoUnit.MILLIS).toString().substring(2).toLowerCase();
22+
public static String durationForLogging(long durationNanos) {
23+
return Duration.ofNanos(durationNanos).truncatedTo(ChronoUnit.MILLIS).toString().substring(2).toLowerCase();
24+
}
25+
26+
public static String durationForLogging(Duration duration) {
27+
return duration.truncatedTo(ChronoUnit.MILLIS).toString().substring(2).toLowerCase();
2428
}
2529
}

kafka/kafka-annotation-processor/src/main/java/io/koraframework/kafka/annotation/processor/consumer/KafkaConsumerContainerGenerator.java

Lines changed: 8 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33
import com.palantir.javapoet.MethodSpec;
44
import com.palantir.javapoet.ParameterSpec;
55
import com.palantir.javapoet.ParameterizedTypeName;
6+
import io.koraframework.annotation.processor.common.AnnotationUtils;
67
import org.jspecify.annotations.Nullable;
78
import io.koraframework.annotation.processor.common.CommonClassNames;
89
import io.koraframework.annotation.processor.common.TagUtils;
@@ -64,8 +65,10 @@ public MethodSpec generate(Elements elements, ExecutableElement executableElemen
6465
.addAnnotation(Nullable.class)
6566
.build());
6667

68+
var configPath = AnnotationUtils.parseAnnotationValueWithoutDefault(listenerAnnotation, "value");
6769
var consumerName = ((TypeElement) executableElement.getEnclosingElement()).getQualifiedName() + "#" + executableElement.getSimpleName();
68-
methodBuilder.addStatement("var telemetry = telemetryFactory.get($S, config.driverProperties(), config.telemetry())", consumerName);
70+
methodBuilder.addStatement("var telemetry = telemetryFactory.get($S, $S, config.driverProperties(), config.telemetry())",
71+
configPath, consumerName);
6972

7073
var consumerParameter = parameters.stream().filter(r -> r instanceof ConsumerParameter.Consumer).map(ConsumerParameter.Consumer.class::cast).findFirst();
7174
if (handlerTypeName.rawType().equals(recordHandler)) {
@@ -77,11 +80,11 @@ public MethodSpec generate(Elements elements, ExecutableElement executableElemen
7780
methodBuilder.beginControlFlow("if (config.topics() == null || config.topics().size() != 1)"); // todo allow list?
7881
methodBuilder.addStatement("throw new java.lang.IllegalArgumentException($S + config.topics())", "@KafkaListener require to specify 1 topic to subscribe when groupId is null, but received: ");
7982
methodBuilder.endControlFlow();
80-
methodBuilder.addCode("return new $T<>($S, config, config.topics().get(0), keyDeserializer, valueDeserializer, telemetry, wrappedHandler);",
81-
kafkaAssignConsumerContainer, consumerName);
83+
methodBuilder.addCode("return new $T<>($S, $S, config, config.topics().get(0), keyDeserializer, valueDeserializer, telemetry, wrappedHandler);",
84+
kafkaAssignConsumerContainer, configPath, consumerName);
8285
methodBuilder.addCode("$<\n} else {$>\n");
83-
methodBuilder.addCode("return new $T<>($S, config, keyDeserializer, valueDeserializer, wrappedHandler, telemetry, rebalanceListener);",
84-
kafkaSubscribeConsumerContainer, consumerName);
86+
methodBuilder.addCode("return new $T<>($S, $S, config, keyDeserializer, valueDeserializer, wrappedHandler, telemetry, rebalanceListener);",
87+
kafkaSubscribeConsumerContainer, configPath, consumerName);
8588
methodBuilder.addCode("$<\n}\n");
8689
return methodBuilder.build();
8790
}

kafka/kafka-annotation-processor/src/main/java/io/koraframework/kafka/annotation/processor/producer/KafkaPublisherGenerator.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -177,8 +177,8 @@ public void generatePublisherImplementation(TypeElement publisher, List<Executab
177177
.addParameter(publisherTelemetryConfig, "telemetryConfig")
178178
.addParameter(ClassName.get(Properties.class), "driverProperties")
179179
.addParameter(topicConfigTypeName, "topicConfig")
180-
.addStatement("var telemetry = telemetryFactory.get($S, telemetryConfig, driverProperties);", configPath)
181-
.addStatement("super(driverProperties, telemetryConfig, telemetry)")
180+
.addStatement("var telemetry = telemetryFactory.get($S, $S, telemetryConfig, driverProperties);", configPath, publisher.getQualifiedName().toString())
181+
.addStatement("super($S, $S, driverProperties, telemetryConfig, telemetry)", configPath, publisher.getQualifiedName().toString())
182182
.addStatement("this.topicConfig = topicConfig");
183183
record TypeWithTag(TypeName typeName, String tag) {}
184184
var parameters = new HashMap<TypeWithTag, String>();

kafka/kafka-annotation-processor/src/test/java/io/koraframework/kafka/annotation/processor/consumer/AbstractKafkaListenerAnnotationProcessorTest.java

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -13,8 +13,8 @@
1313
import io.koraframework.kafka.common.consumer.containers.handlers.KafkaRecordHandler;
1414
import io.koraframework.kafka.common.consumer.containers.handlers.KafkaRecordsHandler;
1515
import io.koraframework.kafka.common.consumer.telemetry.KafkaConsumerTelemetryFactory;
16-
import io.koraframework.kafka.common.consumer.telemetry.NoopKafkaConsumerPollObservation;
17-
import io.koraframework.kafka.common.consumer.telemetry.NoopKafkaConsumerRecordObservation;
16+
import io.koraframework.kafka.common.consumer.telemetry.impl.NoopKafkaConsumerPollObservation;
17+
import io.koraframework.kafka.common.consumer.telemetry.impl.NoopKafkaConsumerRecordObservation;
1818
import io.koraframework.kafka.common.exceptions.RecordKeyDeserializationException;
1919
import io.koraframework.kafka.common.exceptions.RecordValueDeserializationException;
2020
import org.apache.kafka.clients.consumer.Consumer;
@@ -337,11 +337,11 @@ protected RecordsHandlerAssertions(Class<?> controllerClass, Class<?> moduleClas
337337
}
338338

339339
public void handle(ConsumerRecords<K, V> record) {
340-
moduleHandler.handle(consumer, new NoopKafkaConsumerPollObservation(), record);
340+
moduleHandler.handle(consumer, NoopKafkaConsumerPollObservation.INSTANCE, record);
341341
}
342342

343343
public void handle(ConsumerRecord<K, V> record, ThrowingConsumer<InvocationAssertions<K, V>> verifier) {
344-
moduleHandler.handle(consumer, new NoopKafkaConsumerPollObservation(), new ConsumerRecords<>(Map.of(
344+
moduleHandler.handle(consumer, NoopKafkaConsumerPollObservation.INSTANCE, new ConsumerRecords<>(Map.of(
345345
new TopicPartition("test", 1),
346346
List.of(record)
347347
)));
@@ -350,7 +350,7 @@ public void handle(ConsumerRecord<K, V> record, ThrowingConsumer<InvocationAsser
350350

351351
public void handle(ConsumerRecord<K, V> record, Class<? extends Throwable> expectedError) {
352352
assertThatThrownBy(() -> {
353-
moduleHandler.handle(consumer, new NoopKafkaConsumerPollObservation(), new ConsumerRecords<>(Map.of(
353+
moduleHandler.handle(consumer, NoopKafkaConsumerPollObservation.INSTANCE, new ConsumerRecords<>(Map.of(
354354
new TopicPartition("test", 1),
355355
List.of(record)
356356
)));

kafka/kafka-symbol-processor/src/main/kotlin/io/koraframework/kafka/symbol/processor/consumer/KafkaContainerGenerator.kt

Lines changed: 8 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import io.koraframework.kafka.symbol.processor.KafkaClassNames
1010
import io.koraframework.kafka.symbol.processor.KafkaUtils.consumerTag
1111
import io.koraframework.kafka.symbol.processor.KafkaUtils.containerFunName
1212
import io.koraframework.kafka.symbol.processor.consumer.KafkaHandlerGenerator.HandlerFunction
13+
import io.koraframework.ksp.common.AnnotationUtils.findValueNoDefault
1314
import io.koraframework.ksp.common.CommonClassNames
1415
import io.koraframework.ksp.common.KotlinPoetUtils.controlFlow
1516
import io.koraframework.ksp.common.TagUtils.addTag
@@ -36,8 +37,10 @@ class KafkaContainerGenerator {
3637
.addTag(consumerTag)
3738
.returns(CommonClassNames.lifecycle)
3839

40+
val configPath = listenerAnnotation.findValueNoDefault<String>("value")!!
3941
val consumerName = functionDeclaration.parentDeclaration?.qualifiedName?.asString() + "." + functionDeclaration.simpleName.asString()
40-
funBuilder.addStatement("val telemetry = telemetryFactory.get(%S, config.driverProperties(), config.telemetry())", consumerName)
42+
funBuilder.addStatement("val telemetry = telemetryFactory.get(%S, %S, config.driverProperties(), config.telemetry())",
43+
configPath, consumerName)
4144
if (handlerType.rawType == KafkaClassNames.recordHandler) {
4245
funBuilder.addStatement("val wrappedHandler = %T.wrapHandlerRecord(%L, handler)", KafkaClassNames.handlerWrapper, consumerParameter == null)
4346
} else {
@@ -47,11 +50,11 @@ class KafkaContainerGenerator {
4750
addStatement("val topics = config.topics()")
4851
addStatement("require(topics != null)")
4952
addStatement("require(topics.size == 1)")
50-
addStatement("return %T(%S, config, topics[0], keyDeserializer, valueDeserializer, telemetry, wrappedHandler)",
51-
KafkaClassNames.kafkaAssignConsumerContainer, consumerName)
53+
addStatement("return %T(%S, %S, config, topics[0], keyDeserializer, valueDeserializer, telemetry, wrappedHandler)",
54+
KafkaClassNames.kafkaAssignConsumerContainer, configPath, consumerName)
5255
nextControlFlow("else")
53-
addStatement("return %T(%S, config, keyDeserializer, valueDeserializer, wrappedHandler, telemetry, rebalanceListener)",
54-
KafkaClassNames.kafkaSubscribeConsumerContainer, consumerName)
56+
addStatement("return %T(%S, %S, config, keyDeserializer, valueDeserializer, wrappedHandler, telemetry, rebalanceListener)",
57+
KafkaClassNames.kafkaSubscribeConsumerContainer, configPath, consumerName)
5558
}
5659
return funBuilder.build()
5760
}

kafka/kafka-symbol-processor/src/main/kotlin/io/koraframework/kafka/symbol/processor/producer/KafkaPublisherGenerator.kt

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -204,9 +204,11 @@ class KafkaPublisherGenerator(val env: SymbolProcessorEnvironment, val resolver:
204204
val b = classDeclaration.extendsKeepAop(implementationName, resolver)
205205
.generated(KafkaPublisherSymbolProcessor::class)
206206
.superclass(KafkaClassNames.abstractPublisher)
207+
.addSuperclassConstructorParameter("%S", configPath)
208+
.addSuperclassConstructorParameter("%S", classDeclaration.qualifiedName!!.asString())
207209
.addSuperclassConstructorParameter("driverProperties")
208210
.addSuperclassConstructorParameter("telemetryConfig")
209-
.addSuperclassConstructorParameter("telemetryFactory.get(%S, telemetryConfig, driverProperties)", configPath)
211+
.addSuperclassConstructorParameter("telemetryFactory.get(%S, %S, telemetryConfig, driverProperties)", configPath, classDeclaration.qualifiedName!!.asString())
210212
.apply { topicConfig?.let { addProperty(PropertySpec.builder("topicConfig", it, KModifier.PRIVATE, KModifier.FINAL).initializer("topicConfig").build()) } }
211213

212214
val constructorBuilder = FunSpec.constructorBuilder()

kafka/kafka/src/main/java/io/koraframework/kafka/common/KafkaDeserializersModule.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@
1414
* Default Kafka deserializes provided by module for base types
1515
*/
1616
public interface KafkaDeserializersModule {
17+
1718
@DefaultComponent
1819
default Deserializer<String> stringDeserializer() {
1920
return new StringDeserializer();
Lines changed: 18 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -1,35 +1,31 @@
11
package io.koraframework.kafka.common;
22

3+
import io.koraframework.common.DefaultComponent;
4+
import io.koraframework.kafka.common.consumer.telemetry.impl.DefaultKafkaConsumerLoggerFactory;
5+
import io.koraframework.kafka.common.consumer.telemetry.impl.DefaultKafkaConsumerMetricsFactory;
6+
import io.koraframework.kafka.common.consumer.telemetry.impl.DefaultKafkaConsumerTelemetryFactory;
7+
import io.koraframework.kafka.common.producer.telemetry.impl.DefaultKafkaPublisherLoggerFactory;
8+
import io.koraframework.kafka.common.producer.telemetry.impl.DefaultKafkaPublisherMetricsFactory;
9+
import io.koraframework.kafka.common.producer.telemetry.impl.DefaultKafkaPublisherTelemetryFactory;
310
import io.micrometer.core.instrument.MeterRegistry;
4-
import io.micrometer.core.instrument.composite.CompositeMeterRegistry;
511
import io.opentelemetry.api.trace.Tracer;
6-
import io.opentelemetry.api.trace.TracerProvider;
712
import org.jspecify.annotations.Nullable;
8-
import io.koraframework.common.DefaultComponent;
9-
import io.koraframework.kafka.common.consumer.telemetry.DefaultKafkaConsumerTelemetryFactory;
10-
import io.koraframework.kafka.common.producer.telemetry.DefaultKafkaPublisherTelemetryFactory;
1113

1214
public interface KafkaModule extends KafkaDeserializersModule, KafkaSerializersModule {
13-
@DefaultComponent
14-
default DefaultKafkaConsumerTelemetryFactory defaultKafkaConsumerTelemetryFactory(@Nullable Tracer tracer, @Nullable MeterRegistry meterRegistry) {
15-
if (tracer == null) {
16-
tracer = TracerProvider.noop().get("kafkaPublisherTelemetry");
17-
}
18-
if (meterRegistry == null) {
19-
meterRegistry = new CompositeMeterRegistry();
20-
}
2115

22-
return new DefaultKafkaConsumerTelemetryFactory(tracer, meterRegistry);
16+
@DefaultComponent
17+
default DefaultKafkaConsumerTelemetryFactory defaultKafkaConsumerTelemetryFactory(@Nullable Tracer tracer,
18+
@Nullable MeterRegistry meterRegistry,
19+
@Nullable DefaultKafkaConsumerLoggerFactory loggerFactory,
20+
@Nullable DefaultKafkaConsumerMetricsFactory metricsFactory) {
21+
return new DefaultKafkaConsumerTelemetryFactory(tracer, meterRegistry, loggerFactory, metricsFactory);
2322
}
2423

2524
@DefaultComponent
26-
default DefaultKafkaPublisherTelemetryFactory defaultKafkaProducerTelemetryFactory(@Nullable Tracer tracer, @Nullable MeterRegistry meterRegistry) {
27-
if (tracer == null) {
28-
tracer = TracerProvider.noop().get("kafkaPublisherTelemetry");
29-
}
30-
if (meterRegistry == null) {
31-
meterRegistry = new CompositeMeterRegistry();
32-
}
33-
return new DefaultKafkaPublisherTelemetryFactory(tracer, meterRegistry);
25+
default DefaultKafkaPublisherTelemetryFactory defaultKafkaPublisherTelemetryFactory(@Nullable Tracer tracer,
26+
@Nullable MeterRegistry meterRegistry,
27+
@Nullable DefaultKafkaPublisherLoggerFactory loggerFactory,
28+
@Nullable DefaultKafkaPublisherMetricsFactory metricsFactory) {
29+
return new DefaultKafkaPublisherTelemetryFactory(tracer, meterRegistry, loggerFactory, metricsFactory);
3430
}
3531
}

kafka/kafka/src/main/java/io/koraframework/kafka/common/KafkaSerializersModule.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@
1414
* Default Kafka serializes provided by module for base types
1515
*/
1616
public interface KafkaSerializersModule {
17+
1718
@DefaultComponent
1819
default Serializer<String> stringSerializer() {
1920
return new StringSerializer();

kafka/kafka/src/main/java/io/koraframework/kafka/common/KafkaUtils.java

Lines changed: 2 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -1,36 +1,15 @@
11
package io.koraframework.kafka.common;
22

3-
import org.apache.kafka.clients.CommonClientConfigs;
4-
import org.jspecify.annotations.NullMarked;
5-
import io.koraframework.kafka.common.consumer.KafkaListenerConfig;
6-
73
import java.util.concurrent.ThreadFactory;
84
import java.util.concurrent.atomic.AtomicInteger;
95

10-
@NullMarked
116
public final class KafkaUtils {
127

138
private KafkaUtils() {}
149

15-
public static String getConsumerPrefix(KafkaListenerConfig config) {
16-
final Object groupId = config.driverProperties().get(CommonClientConfigs.GROUP_ID_CONFIG);
17-
if (groupId != null) {
18-
return groupId.toString();
19-
}
20-
21-
if (config.topics() != null) {
22-
return String.join(";", config.topics());
23-
} else if (config.topicsPattern() != null) {
24-
return config.topicsPattern().toString();
25-
} else if (config.partitions() != null) {
26-
return String.join(";", config.partitions());
27-
} else {
28-
return "unknown";
29-
}
30-
}
31-
3210
public static class NamedThreadFactory implements ThreadFactory {
33-
private static final String CONSUMER_PREFIX = "kafka-consumer-";
11+
12+
private static final String CONSUMER_PREFIX = "kafka-listener-";
3413

3514
private final AtomicInteger threadNumber = new AtomicInteger(1);
3615
private final String namePrefix;

0 commit comments

Comments
 (0)