Introduce Micrometer API to add tracing.

Closes #1212.
This commit is contained in:
Greg L. Turnquist
2022-01-19 09:55:10 -06:00
parent 9686749178
commit c1e3dbf0da
13 changed files with 1121 additions and 2 deletions

View File

@@ -354,7 +354,7 @@
<id>external-cassandra</id>
<properties>
<build.cassandra.mode>external</build.cassandra.mode>
<build.cassandra.native_transport_port>9042 </build.cassandra.native_transport_port>
<build.cassandra.native_transport_port>9042</build.cassandra.native_transport_port>
<build.cassandra.rpc_port>9160</build.cassandra.rpc_port>
<build.cassandra.ssl_storage_port>7001</build.cassandra.ssl_storage_port>
<build.cassandra.storage_port>7000</build.cassandra.storage_port>
@@ -369,7 +369,7 @@
<properties>
<build.cassandra.mode>testcontainers</build.cassandra.mode>
<build.cassandra.native_transport_port>0 </build.cassandra.native_transport_port>
<build.cassandra.native_transport_port>0</build.cassandra.native_transport_port>
<build.cassandra.rpc_port>0</build.cassandra.rpc_port>
<build.cassandra.ssl_storage_port>0</build.cassandra.ssl_storage_port>
<build.cassandra.storage_port>0</build.cassandra.storage_port>

View File

@@ -71,6 +71,36 @@
<artifactId>java-driver-query-builder</artifactId>
</dependency>
<!-- Micrometer -->
<dependency>
<groupId>io.micrometer</groupId>
<artifactId>micrometer-observation</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>io.micrometer</groupId>
<artifactId>micrometer-tracing-api</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>io.micrometer</groupId>
<artifactId>micrometer-test</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>io.micrometer</groupId>
<artifactId>micrometer-tracing-test</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>io.micrometer</groupId>
<artifactId>micrometer-tracing-integration-test</artifactId>
<scope>test</scope>
</dependency>
<!-- Reactor -->
<dependency>
<groupId>io.projectreactor</groupId>

View File

@@ -0,0 +1,127 @@
/*
* Copyright 2013-2022 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.data.cassandra.observability;
import io.micrometer.common.docs.TagKey;
import io.micrometer.observation.docs.DocumentedObservation;
/**
* Cassandra-based implementation of {@link DocumentedObservation}.
*
* @author Mark Paluch
* @author Marcin Grzejszczak
* @author Greg Turnquist
* @since 4.0.0
*/
enum CassandraObservation implements DocumentedObservation {
/**
* Create an {@link io.micrometer.observation.Observation} for Cassandra-based queries.
*/
CASSANDRA_QUERY_OBSERVATION {
@Override
public String getName() {
return "spring.data.cassandra.query";
}
@Override
public String getContextualName() {
return "query";
}
@Override
public TagKey[] getLowCardinalityTagKeys() {
return LowCardinalityTags.values();
}
@Override
public TagKey[] getHighCardinalityTagKeys() {
return HighCardinalityTags.values();
}
@Override
public String getPrefix() {
return "spring.data.cassandra.";
}
};
enum LowCardinalityTags implements TagKey {
/**
* Name of the Cassandra keyspace.
*/
KEYSPACE_NAME {
@Override
public String getKey() {
return "spring.data.cassandra.keyspace";
}
},
/**
* Cassandra session
*/
SESSION_NAME {
@Override
public String getKey() {
return "spring.data.cassandra.sessionName";
}
},
/**
* The method name
*/
METHOD_NAME {
@Override
public String getKey() {
return "spring.data.cassandra.methodName";
}
},
/**
* Cassandra URL
*/
URL {
@Override
public String getKey() {
return "spring.data.cassandra.url";
}
},
/**
* A tag containing error that occurred for the given node.
*/
NODE_ERROR_TAG {
@Override
public String getKey() {
return "spring.data.cassandra.node[%s].error";
}
}
}
enum HighCardinalityTags implements TagKey {
/**
* A tag containing Cassandra CQL.
*/
CQL_TAG {
@Override
public String getKey() {
return "spring.data.cassandra.cql";
}
}
}
}

View File

@@ -0,0 +1,57 @@
/*
* Copyright 2013-2022 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.data.cassandra.observability;
import io.micrometer.observation.Observation;
import org.springframework.lang.Nullable;
import com.datastax.oss.driver.api.core.CqlSession;
import com.datastax.oss.driver.api.core.cql.Statement;
/**
* A {@link Observation.Context} for {@link CqlSession}.
*
* @author Greg Turnquist
* @since 4.0.0
*/
public class CqlSessionContext extends Observation.Context {
private final @Nullable Statement<?> statement;
private final String methodName;
private final @Nullable CqlSession delegateSession;
public CqlSessionContext(@Nullable Statement<?> statement, String methodName, @Nullable CqlSession delegateSession) {
this.statement = statement;
this.methodName = methodName;
this.delegateSession = delegateSession;
}
@Nullable
public Statement<?> getStatement() {
return statement;
}
public String getMethodName() {
return methodName;
}
@Nullable
public CqlSession getDelegateSession() {
return delegateSession;
}
}

View File

@@ -0,0 +1,32 @@
/*
* Copyright 2013-2022 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.data.cassandra.observability;
import io.micrometer.observation.Observation;
/**
* {@link Observation.TagsProvider} for Cassandra.
*
* @author Greg Turnquist
* @since 4.0.0
*/
public interface CqlSessionTagsProvider extends Observation.TagsProvider<CqlSessionContext> {
@Override
default boolean supportsContext(Observation.Context context) {
return context instanceof CqlSessionContext;
}
}

View File

@@ -0,0 +1,55 @@
/*
* Copyright 2013-2022 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.data.cassandra.observability;
import io.micrometer.observation.ObservationRegistry;
import org.springframework.beans.BeansException;
import org.springframework.beans.factory.config.BeanPostProcessor;
import com.datastax.oss.driver.api.core.CqlSession;
/**
* {@link BeanPostProcessor} to automatically wrap all {@link CqlSession}s with a {@link CqlSessionTracingInterceptor}.
*
* @author Marcin Grzejszczak
* @author Mark Paluch
* @author Greg Turnquist
* @since 4.0.0
*/
public class CqlSessionTracingBeanPostProcessor implements BeanPostProcessor {
private final ObservationRegistry observationRegistry;
private final CqlSessionTagsProvider tagsProvider;
public CqlSessionTracingBeanPostProcessor(ObservationRegistry observationRegistry,
CqlSessionTagsProvider tagsProvider) {
this.observationRegistry = observationRegistry;
this.tagsProvider = tagsProvider;
}
@Override
public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException {
if (bean instanceof CqlSession) {
return CqlSessionTracingFactory.wrap((CqlSession) bean, this.observationRegistry, this.tagsProvider);
}
return bean;
}
}

View File

@@ -0,0 +1,56 @@
/*
* Copyright 2013-2022 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.data.cassandra.observability;
import io.micrometer.observation.ObservationRegistry;
import org.springframework.aop.framework.ProxyFactory;
import com.datastax.oss.driver.api.core.CqlSession;
/**
* Factory to wrap a {@link CqlSession} with a {@link CqlSessionTracingInterceptor}.
*
* @author Mark Paluch
* @author Greg Turnquist
* @since 4.0.0
*/
public final class CqlSessionTracingFactory {
private CqlSessionTracingFactory() {
throw new IllegalStateException("Can't instantiate a utility class.");
}
/**
* Wrap the {@link CqlSession} with a {@link CqlSessionTracingInterceptor}.
*
* @param session
* @param observationRegistry
* @param tagsProvider
* @return
*/
public static CqlSession wrap(CqlSession session, ObservationRegistry observationRegistry,
CqlSessionTagsProvider tagsProvider) {
ProxyFactory proxyFactory = new ProxyFactory();
proxyFactory.setTarget(session);
proxyFactory.addAdvice(new CqlSessionTracingInterceptor(session, observationRegistry, tagsProvider));
proxyFactory.addInterface(CqlSession.class);
return (CqlSession) proxyFactory.getProxy();
}
}

View File

@@ -0,0 +1,165 @@
/*
* Copyright 2013-2022 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.data.cassandra.observability;
import io.micrometer.observation.Observation;
import io.micrometer.observation.ObservationRegistry;
import java.lang.reflect.Method;
import java.util.Arrays;
import java.util.Map;
import java.util.function.Function;
import org.aopalliance.intercept.MethodInterceptor;
import org.aopalliance.intercept.MethodInvocation;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import com.datastax.oss.driver.api.core.CqlSession;
import com.datastax.oss.driver.api.core.cql.ResultSet;
import com.datastax.oss.driver.api.core.cql.SimpleStatement;
import com.datastax.oss.driver.api.core.cql.Statement;
/**
* A {@link MethodInterceptor} that wraps calls around {@link CqlSession} in a trace representation. This interceptor
* wraps statements for {@code execute} and {@code prepare} (including their asynchronous variants) only. Graph and
* reactive {@link CqlSession} method are called as-is.
*
* @author Mark Paluch
* @author Marcin Grzejszczak
* @author Greg Turnquist
* @since 4.0.0
*/
final class CqlSessionTracingInterceptor
implements MethodInterceptor, Observation.TagsProviderAware<CqlSessionTagsProvider> {
private static final Log log = LogFactory.getLog(CqlSessionTracingInterceptor.class);
private final CqlSession delegateSession;
private final ObservationRegistry observationRegistry;
private CqlSessionTagsProvider tagsProvider;
CqlSessionTracingInterceptor(CqlSession delegateSession, ObservationRegistry observationRegistry,
CqlSessionTagsProvider tagsProvider) {
this.delegateSession = delegateSession;
this.observationRegistry = observationRegistry;
this.tagsProvider = tagsProvider;
}
@Nullable
@Override
public Object invoke(@NotNull MethodInvocation invocation) throws Throwable {
Method method = invocation.getMethod();
Object[] args = invocation.getArguments();
if (method.getName().equals("execute") && args.length > 0) {
return tracedCall(createStatement(args), method.getName(), this.delegateSession::execute);
}
if (method.getName().equals("executeAsync") && args.length > 0) {
return tracedCall(createStatement(args), method.getName(), this.delegateSession::executeAsync);
}
if (method.getName().equals("prepare") && args.length > 0) {
return tracedCall(createStatement(args), method.getName(),
statement -> this.delegateSession.prepare((SimpleStatement) statement));
}
if (method.getName().equals("prepareAsync") && args.length > 0) {
return tracedCall(createStatement(args), method.getName(),
statement -> this.delegateSession.prepareAsync((SimpleStatement) statement));
}
return invocation.proceed();
}
/**
* Apply tracing to a {@link Statement}.
*
* @param statement original CQL {@link Statement}
* @param statementExecutor function that transforms a {@link Statement} into a resulting {@link Object}
* @return {@link ResultSet}
*/
private Object tracedCall(Statement<?> statement, String methodName,
Function<Statement<?>, Object> statementExecutor) {
if (this.observationRegistry.getCurrentObservation() == null) {
return null;
}
Observation observation = childObservation(statement, methodName, this.delegateSession);
if (log.isDebugEnabled()) {
log.debug("Created a new child observation before query [" + observation + "]");
}
try (Observation.Scope scope = observation.openScope()) {
return statementExecutor.apply(statement);
} catch (Exception e) {
observation.error(e);
throw e;
} finally {
observation.stop();
}
}
/**
* Convert list of arguments into a {@link Statement}.
*
* @param args
* @return CQL statement
*/
private static Statement<?> createStatement(Object[] args) {
if (args[0] instanceof Statement) {
return (Statement<?>) args[0];
}
if (args[0] instanceof String & args.length == 1) {
return SimpleStatement.newInstance((String) args[0]);
}
if (args[0]instanceof String query && args.length == 2) {
return args[1] instanceof Map //
? SimpleStatement.newInstance(query, (Map) args[1]) //
: SimpleStatement.newInstance(query, (Object[]) args[1]);
}
throw new IllegalArgumentException(String.format("Unsupported arguments %s", Arrays.toString(args)));
}
private Observation childObservation(Statement<?> statement, String methodName, CqlSession delegateSession) {
CqlSessionContext observationContext = new CqlSessionContext(statement, methodName, delegateSession);
return CassandraObservation.CASSANDRA_QUERY_OBSERVATION //
.observation(this.observationRegistry, observationContext) //
.contextualName(CassandraObservation.CASSANDRA_QUERY_OBSERVATION.getContextualName()) //
.tagsProvider(this.tagsProvider) //
.start();
}
@Override
public void setTagsProvider(CqlSessionTagsProvider tagsProvider) {
this.tagsProvider = tagsProvider;
}
}

View File

@@ -0,0 +1,102 @@
/*
* Copyright 2013-2022 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.data.cassandra.observability;
import io.micrometer.common.Tag;
import io.micrometer.observation.Observation;
import io.micrometer.tracing.Span;
import io.micrometer.tracing.Tracer;
import io.micrometer.tracing.handler.TracingObservationHandler;
import java.net.URI;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
/**
* A {@link TracingObservationHandler} for {@link CqlSessionContext}.
*
* @author Greg Turnquist
* @since 4.0.0
*/
public class CqlSessionTracingObservationHandler implements TracingObservationHandler<CqlSessionContext> {
private static final Log log = LogFactory.getLog(CqlSessionTracingObservationHandler.class);
private final Tracer tracer;
public CqlSessionTracingObservationHandler(Tracer tracer) {
this.tracer = tracer;
}
@Override
public void onStart(CqlSessionContext context) {
Span.Builder builder = this.tracer.spanBuilder() //
.name(context.getContextualName()) //
.kind(Span.Kind.CLIENT);
getTracingContext(context).setSpan(builder.start());
}
@Override
public void onStop(CqlSessionContext context) {
Span span = getRequiredSpan(context);
tagSpan(context, span);
String sessionName = null;
String url = null;
for (Tag tag : context.getLowCardinalityTags()) {
if (tag.getKey().equals(CassandraObservation.LowCardinalityTags.SESSION_NAME.getKey())) {
sessionName = tag.getValue();
}
if (tag.getKey().equals(CassandraObservation.LowCardinalityTags.URL.getKey())) {
url = tag.getValue();
}
}
if (sessionName != null) {
span.remoteServiceName("cassandra-" + sessionName);
}
if (url != null) {
try {
URI uri = URI.create(url);
span.remoteIpAndPort(uri.getHost(), uri.getPort());
} catch (Exception e) {
if (log.isDebugEnabled()) {
log.debug("Failed to parse the url [" + url + "]. Won't set this value on the span");
}
}
}
span.end();
}
@Override
public boolean supportsContext(Observation.Context context) {
return context instanceof CqlSessionContext;
}
@Override
public Tracer getTracer() {
return this.tracer;
}
}

View File

@@ -0,0 +1,107 @@
/*
* Copyright 2013-2022 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.data.cassandra.observability;
import io.micrometer.common.Tags;
import java.util.Optional;
import java.util.StringJoiner;
import com.datastax.oss.driver.api.core.CqlIdentifier;
import com.datastax.oss.driver.api.core.cql.BatchStatement;
import com.datastax.oss.driver.api.core.cql.BatchableStatement;
import com.datastax.oss.driver.api.core.cql.BoundStatement;
import com.datastax.oss.driver.api.core.cql.SimpleStatement;
import com.datastax.oss.driver.api.core.cql.Statement;
/**
* Default {@link CqlSessionTagsProvider} implementation.
*
* @author Greg Turnquist
* @since 4.0.0
*/
public class DefaultCassandraTagsProvider implements CqlSessionTagsProvider {
@Override
public Tags getLowCardinalityTags(CqlSessionContext context) {
Tags tags = Tags.of( //
CassandraObservation.LowCardinalityTags.SESSION_NAME
.of(Optional.ofNullable(context.getDelegateSession().getName()).orElse("unknown")),
CassandraObservation.LowCardinalityTags.KEYSPACE_NAME.of(
Optional.ofNullable(context.getStatement().getKeyspace()).map(CqlIdentifier::asInternal).orElse("unknown")),
CassandraObservation.LowCardinalityTags.METHOD_NAME.of(context.getMethodName()));
if (context.getStatement().getNode() != null) {
tags = tags.and(CassandraObservation.LowCardinalityTags.URL
.of(context.getStatement().getNode().getEndPoint().resolve().toString()));
}
return tags;
}
@Override
public Tags getHighCardinalityTags(CqlSessionContext context) {
return Tags.of(CassandraObservation.HighCardinalityTags.CQL_TAG.of(getCql(context.getStatement())));
}
/**
* Extract the CQL query from the delegate {@link Statement}.
*
* @return string-based CQL of the delegate
* @param statement
*/
private static String getCql(Statement<?> statement) {
String query = "";
if (statement instanceof SimpleStatement) {
query = getQuery(statement);
}
if (statement instanceof BatchStatement) {
StringJoiner joiner = new StringJoiner(";");
for (BatchableStatement<?> bs : (BatchStatement) statement) {
joiner.add(getQuery(bs));
}
query = joiner.toString();
}
return query;
}
/**
* Extract the query from a {@link Statement}.
*
* @param statement
* @return query
*/
private static String getQuery(Statement<?> statement) {
if (statement instanceof SimpleStatement) {
return ((SimpleStatement) statement).getQuery();
}
if (statement instanceof BoundStatement) {
return ((BoundStatement) statement).getPreparedStatement().getQuery();
}
return "";
}
}

View File

@@ -0,0 +1,87 @@
/*
* Copyright 2013-2022 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.data.cassandra.observability;
import static org.assertj.core.api.Assertions.*;
import static org.mockito.Mockito.*;
import io.micrometer.observation.ObservationRegistry;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.springframework.aop.Advisor;
import org.springframework.aop.framework.Advised;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringExtension;
import com.datastax.oss.driver.api.core.CqlSession;
/**
* Verify that {@link CqlSessionTracingBeanPostProcessor} properly wraps {@link CqlSession} beans registered in the app
* context.
*
* @author Greg Turnquist
* @since 4.0.0
*/
@ExtendWith(SpringExtension.class)
@ContextConfiguration
public class CqlSessionTracingBeanPostProcessorTests {
@Autowired CqlSession session;
@Test
void injectedCqlSessionShouldBeWrapped() throws Exception {
assertThat(session).isInstanceOf(CqlSession.class);
assertThat(session).isInstanceOf(Advised.class);
Advised advised = (Advised) session;
assertThat(advised.getAdvisors()).extracting(Advisor::getAdvice)
.satisfies(advice -> advice.getClass().equals(CqlSessionTracingInterceptor.class));
assertThat(advised.getTargetSource().getTarget()).isEqualTo(TestConfig.originalSession);
}
@Configuration
static class TestConfig {
static CqlSession originalSession = mock(CqlSession.class);
@Bean
CqlSession originalSession() {
return originalSession;
}
@Bean
ObservationRegistry meterRegistry() {
return ObservationRegistry.create();
}
@Bean
CqlSessionTagsProvider tagsProvider() {
return new DefaultCassandraTagsProvider();
}
@Bean
CqlSessionTracingBeanPostProcessor traceCqlSessionBeanPostProcessor(ObservationRegistry observationRegistry,
CqlSessionTagsProvider tagsProvider) {
return new CqlSessionTracingBeanPostProcessor(observationRegistry, tagsProvider);
}
}
}

View File

@@ -0,0 +1,182 @@
/*
* Copyright 2013-2022 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.data.cassandra.observability;
import static org.assertj.core.api.Assertions.*;
import static org.springframework.data.cassandra.observability.CassandraObservation.*;
import io.micrometer.common.Tag;
import io.micrometer.common.Tags;
import io.micrometer.core.instrument.MeterRegistry;
import io.micrometer.core.instrument.observation.TimerObservationHandler;
import io.micrometer.core.instrument.simple.SimpleMeterRegistry;
import io.micrometer.core.tck.MeterRegistryAssert;
import io.micrometer.observation.Observation;
import io.micrometer.observation.ObservationRegistry;
import io.micrometer.tracing.Span;
import io.micrometer.tracing.test.simple.SimpleSpan;
import io.micrometer.tracing.test.simple.SimpleTracer;
import io.micrometer.tracing.test.simple.SpanAssert;
import java.util.Deque;
import java.util.Map;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.annotation.AnnotationConfigApplicationContext;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.cassandra.config.AbstractSessionConfiguration;
import org.springframework.data.cassandra.test.util.CassandraExtension;
import org.springframework.data.cassandra.test.util.IntegrationTestsSupport;
import com.datastax.oss.driver.api.core.CqlSession;
import com.datastax.oss.driver.api.core.cql.Statement;
/**
* Verify that {@link CqlSessionTracingInterceptor} properly wraps {@link Statement} object with tracing.
*
* @author Greg Turnquist
* @since 4.0.0
*/
public class CqlSessionTracingTests extends IntegrationTestsSupport {
private static final String CREATE_KEYSPACE = "CREATE KEYSPACE ConfigTest " + "WITH "
+ "REPLICATION = { 'class' : 'SimpleStrategy', 'replication_factor' : 1 };";
private ConfigurableApplicationContext context;
private CqlSession session;
@BeforeEach
void setUp() {
this.context = new AnnotationConfigApplicationContext(Config.class);
this.session = context.getBean(CqlSession.class);
}
@AfterEach
void tearDown() {
this.context.close();
}
@Test
void shouldNotCreateAnyMetricsWhenThereIsNoObservation() {
MeterRegistry meterRegistry = new SimpleMeterRegistry();
ObservationRegistry observationRegistry = ObservationRegistry.create();
observationRegistry.observationConfig().observationHandler(new TimerObservationHandler(meterRegistry));
MeterRegistryAssert.then(meterRegistry).hasNoMetrics();
}
@Test
void tracingNoStatementsShouldProduceNoMetrics() {
MeterRegistry meterRegistry = new SimpleMeterRegistry();
ObservationRegistry observationRegistry = ObservationRegistry.create();
observationRegistry.observationConfig().observationHandler(new TimerObservationHandler(meterRegistry));
CqlSessionTagsProvider tagsProvider = new DefaultCassandraTagsProvider();
SimpleTracer tracer = new SimpleTracer();
observationRegistry.observationConfig().observationHandler(new CqlSessionTracingObservationHandler(tracer));
CqlSessionTracingFactory.wrap(session, observationRegistry, tagsProvider);
MeterRegistryAssert.then(meterRegistry).hasNoMetrics();
}
@Test
void shouldCreateObservationForCqlSessionOperations() {
MeterRegistry meterRegistry = new SimpleMeterRegistry();
ObservationRegistry observationRegistry = ObservationRegistry.create();
observationRegistry.observationConfig().observationHandler(new TimerObservationHandler(meterRegistry));
CqlSessionTagsProvider tagsProvider = new DefaultCassandraTagsProvider();
SimpleTracer tracer = new SimpleTracer();
observationRegistry.observationConfig().observationHandler(new CqlSessionTracingObservationHandler(tracer));
Observation.start("test", observationRegistry).scoped(() -> {
CqlSession traceSession = CqlSessionTracingFactory.wrap(session, observationRegistry, tagsProvider);
traceSession.execute(CREATE_KEYSPACE);
traceSession.executeAsync(CREATE_KEYSPACE);
traceSession.prepare(CREATE_KEYSPACE);
traceSession.prepareAsync(CREATE_KEYSPACE);
});
MeterRegistryAssert.then(meterRegistry).hasTimerWithNameAndTags(CASSANDRA_QUERY_OBSERVATION.getName(), Tags.of( //
LowCardinalityTags.SESSION_NAME.of("s5"), //
LowCardinalityTags.KEYSPACE_NAME.of("unknown"), //
Tag.of("error", "none") //
));
assertThat(tracer.getSpans()).hasSize(4);
assertThat(findSpan(tracer.getSpans(), LowCardinalityTags.METHOD_NAME.getKey(), "execute")).isNotNull();
assertThat(findSpan(tracer.getSpans(), LowCardinalityTags.METHOD_NAME.getKey(), "executeAsync")).isNotNull();
assertThat(findSpan(tracer.getSpans(), LowCardinalityTags.METHOD_NAME.getKey(), "prepare")).isNotNull();
assertThat(findSpan(tracer.getSpans(), LowCardinalityTags.METHOD_NAME.getKey(), "prepareAsync")).isNotNull();
tracer.getSpans().forEach(simpleSpan -> SpanAssert.then(simpleSpan) //
.hasRemoteServiceNameEqualTo("cassandra-s5") //
.hasNameEqualTo(CASSANDRA_QUERY_OBSERVATION.getContextualName()) //
.hasTag(LowCardinalityTags.SESSION_NAME.getKey(), "s5") //
.hasTag(LowCardinalityTags.KEYSPACE_NAME.getKey(), "unknown") //
.hasTag(HighCardinalityTags.CQL_TAG.getKey(), CREATE_KEYSPACE) //
.hasIpThatIsBlank() //
.hasPortEqualTo(0) //
.hasKindEqualTo(Span.Kind.CLIENT));
}
/**
* Find a {@link Span} with a specific key and value.
*
* @param spans
* @param key
* @param value
* @return
*/
private SimpleSpan findSpan(Deque<SimpleSpan> spans, String key, String value) {
return spans.stream() //
.filter(simpleSpan -> {
Map<String, String> tags = simpleSpan.getTags();
return tags.containsKey(key) && tags.get(key).equals(value);
}) //
.findAny() //
.orElse(null);
}
@Configuration
static class Config extends AbstractSessionConfiguration {
@Override
protected String getKeyspaceName() {
return "system";
}
@Override
protected int getPort() {
return CassandraExtension.getResources().getPort();
}
}
}

View File

@@ -0,0 +1,119 @@
/*
* Copyright 2013-2022 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.data.cassandra.observability;
import io.micrometer.core.instrument.MeterRegistry;
import io.micrometer.core.instrument.observation.TimerObservationHandler;
import io.micrometer.core.instrument.simple.SimpleMeterRegistry;
import io.micrometer.observation.ObservationHandler;
import io.micrometer.observation.ObservationRegistry;
import io.micrometer.tracing.test.SampleTestRunner;
import io.micrometer.tracing.test.reporter.BuildingBlocks;
import java.util.Deque;
import java.util.function.BiConsumer;
import org.junit.jupiter.api.Disabled;
import org.junit.jupiter.api.extension.ExtendWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.data.cassandra.support.AbstractTestJavaConfig;
import org.springframework.data.cassandra.test.util.CassandraExtension;
import org.springframework.data.cassandra.test.util.TestKeyspace;
import org.springframework.data.cassandra.test.util.TestKeyspaceName;
import org.springframework.test.context.junit.jupiter.SpringExtension;
import com.datastax.oss.driver.api.core.CqlSession;
/**
* Collection of tests that log metrics and tracing with an external tracing tool. Since this external tool must be up
* and running after the test is completed, this test is ONLY run manually.
*
* @author Greg Turnquist
* @since 4.0.0
*/
@Disabled("Run this manually to visually test spans in Zipkin")
@ExtendWith({ SpringExtension.class, CassandraExtension.class })
@TestKeyspaceName
public class ZipkinIntegrationTests extends SampleTestRunner {
private static final MeterRegistry METER_REGISTRY = new SimpleMeterRegistry();
private static final ObservationRegistry OBSERVATION_REGISTRY = ObservationRegistry.create();
static {
OBSERVATION_REGISTRY.observationConfig().observationHandler(new TimerObservationHandler(METER_REGISTRY));
}
@Autowired CqlSession session;
ZipkinIntegrationTests() {
super(SampleRunnerConfig.builder().build(), OBSERVATION_REGISTRY, METER_REGISTRY);
}
@Override
public BiConsumer<BuildingBlocks, Deque<ObservationHandler>> customizeObservationHandlers() {
return (buildingBlocks, observationHandlers) -> {
observationHandlers.addLast(new CqlSessionTracingObservationHandler(buildingBlocks.getTracer()));
};
}
@Override
public TracingSetup[] getTracingSetup() {
return new TracingSetup[] { TracingSetup.ZIPKIN_BRAVE };
}
@Override
public SampleTestRunnerConsumer yourCode() {
return (tracer, meterRegistry) -> {
session.execute("DROP KEYSPACE IF EXISTS ConfigTest");
session.execute("CREATE KEYSPACE ConfigTest " + "WITH "
+ "REPLICATION = { 'class' : 'SimpleStrategy', 'replication_factor' : 1 };");
session.execute("USE ConfigTest");
session.execute("CREATE TABLE IF NOT EXISTS person (id int, firstName text, lastName text, PRIMARY KEY(id));");
System.out.println(((SimpleMeterRegistry) meterRegistry).getMetersAsString());
};
}
static class TestConfig extends AbstractTestJavaConfig {
@TestKeyspace CqlSession session;
@Override
protected String getKeyspaceName() {
return "system";
}
@Bean
ObservationRegistry registry() {
return OBSERVATION_REGISTRY;
}
@Bean
CqlSessionTagsProvider tagsProvider() {
return new DefaultCassandraTagsProvider();
}
@Bean
CqlSessionTracingBeanPostProcessor traceCqlSessionBeanPostProcessor(ObservationRegistry observationRegistry,
CqlSessionTagsProvider tagsProvider) {
return new CqlSessionTracingBeanPostProcessor(observationRegistry, tagsProvider);
}
}
}