diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/aot/CassandraRuntimeHints.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/aot/CassandraRuntimeHints.java index 0ff6169d5..a1370a8eb 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/aot/CassandraRuntimeHints.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/aot/CassandraRuntimeHints.java @@ -31,7 +31,7 @@ import org.springframework.data.cassandra.core.mapping.event.ReactiveBeforeSaveC import org.springframework.data.cassandra.observability.CassandraObservationSupplier; import org.springframework.data.cassandra.repository.support.SimpleCassandraRepository; import org.springframework.data.cassandra.repository.support.SimpleReactiveCassandraRepository; -import org.springframework.data.repository.util.ReactiveWrappers; +import org.springframework.data.util.ReactiveWrappers; import org.springframework.lang.Nullable; import org.springframework.util.ClassUtils; @@ -91,6 +91,19 @@ class CassandraRuntimeHints implements RuntimeHintsRegistrar { } hints.proxies().registerJdkProxy(CqlSession.class, SpringProxy.class, Advised.class, DecoratingProxy.class); + Class observationDecorated; + try { + observationDecorated = Class.forName( + "org.springframework.data.cassandra.observability.CqlSessionObservationInterceptor.ObservationDecoratedProxy", + false, classLoader); + } catch (Exception e) { + observationDecorated = null; + } + + if (observationDecorated != null) { + hints.proxies().registerJdkProxy(CqlSession.class, SpringProxy.class, Advised.class, DecoratingProxy.class, + observationDecorated); + } } } } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultBridgedReactiveSession.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultBridgedReactiveSession.java index c595491cf..c9f6b593d 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultBridgedReactiveSession.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/cql/session/DefaultBridgedReactiveSession.java @@ -27,6 +27,8 @@ import java.util.concurrent.CompletionStage; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; +import org.springframework.aop.TargetSource; +import org.springframework.aop.framework.AopProxyUtils; import org.springframework.data.cassandra.ReactiveResultSet; import org.springframework.data.cassandra.ReactiveSession; import org.springframework.util.Assert; @@ -54,7 +56,7 @@ import com.datastax.oss.driver.api.core.metadata.Metadata; * Calls are deferred until a subscriber subscribes to the resulting {@link org.reactivestreams.Publisher}. The calls * are executed by subscribing to {@link CompletionStage} and returning the result as calls complete. *

- * Elements are emitted on netty EventLoop threads. {@link AsyncResultSet} allows {@link AsyncResultSet#fetchNextPage()} + * Elements are emitted on netty EventLoop threads. {@link AsyncResultSet} allows {@link AsyncResultSet#fetchNextPage() * asynchronous requesting} of subsequent pages. The next page is requested after emitting all elements of the previous * page. However, this is an intermediate solution until Datastax can provide a fully reactive driver. *

@@ -85,6 +87,18 @@ public class DefaultBridgedReactiveSession implements ReactiveSession { Assert.notNull(session, "Session must not be null"); + // potentially unwrap a ObservationDecoratedProxy as reactive observability + // requires its own approach to span creation. We do not want to participate in + // async API spans but rather drive our own spans. + if (session instanceof TargetSource) { + Class[] interfaces = session.getClass().getInterfaces(); + for (Class anInterface : interfaces) { + if (anInterface.getName().endsWith("ObservationDecoratedProxy")) { + session = (CqlSession) AopProxyUtils.getSingletonTarget(session); + } + } + } + this.session = session; } @@ -239,7 +253,6 @@ public class DefaultBridgedReactiveSession implements ReactiveSession { return Flux.fromIterable(resultSet.currentPage()); } - @Override public ColumnDefinitions getColumnDefinitions() { return this.resultSet.getColumnDefinitions(); diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/observability/CqlSessionObservationInterceptor.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/observability/CqlSessionObservationInterceptor.java index de03fd3bc..cc2d892a2 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/observability/CqlSessionObservationInterceptor.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/observability/CqlSessionObservationInterceptor.java @@ -25,6 +25,7 @@ import java.util.function.Function; import org.aopalliance.intercept.MethodInterceptor; import org.aopalliance.intercept.MethodInvocation; +import org.springframework.aop.TargetSource; import com.datastax.oss.driver.api.core.CqlIdentifier; import com.datastax.oss.driver.api.core.CqlSession; @@ -183,4 +184,11 @@ final class CqlSessionObservationInterceptor implements MethodInterceptor { return observation.start(); } + /** + * Marker interface for components that want to participate in observation but do not want to work with a + * {@code CqlSession} that is already decorated for observation. + */ + public interface ObservationDecoratedProxy extends TargetSource { + + } } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/observability/ObservableCqlSessionFactory.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/observability/ObservableCqlSessionFactory.java index 9aaea087d..021bc6053 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/observability/ObservableCqlSessionFactory.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/observability/ObservableCqlSessionFactory.java @@ -17,7 +17,10 @@ package org.springframework.data.cassandra.observability; import io.micrometer.observation.ObservationRegistry; +import org.springframework.aop.RawTargetAccess; +import org.springframework.aop.TargetSource; import org.springframework.aop.framework.ProxyFactory; +import org.springframework.data.cassandra.observability.CqlSessionObservationInterceptor.ObservationDecoratedProxy; import org.springframework.util.Assert; import com.datastax.oss.driver.api.core.CqlSession; @@ -65,7 +68,10 @@ public final class ObservableCqlSessionFactory { proxyFactory.setTarget(session); proxyFactory.addAdvice(new CqlSessionObservationInterceptor(session, remoteServiceName, observationRegistry)); proxyFactory.addInterface(CqlSession.class); + proxyFactory.addInterface(ObservationDecoratedProxy.class); return (CqlSession) proxyFactory.getProxy(); } + + } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/observability/ObservableReactiveSession.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/observability/ObservableReactiveSession.java index fb2ada906..0679dd15a 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/observability/ObservableReactiveSession.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/observability/ObservableReactiveSession.java @@ -180,4 +180,5 @@ public class ObservableReactiveSession implements ReactiveSession { private static Observation getParentObservation(ContextView contextView) { return contextView.getOrDefault(ObservationThreadLocalAccessor.KEY, null); } + } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/observability/ObservableReactiveSessionFactoryBean.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/observability/ObservableReactiveSessionFactoryBean.java index 16d19d8dc..04caf2f5d 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/observability/ObservableReactiveSessionFactoryBean.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/observability/ObservableReactiveSessionFactoryBean.java @@ -17,6 +17,8 @@ package org.springframework.data.cassandra.observability; import io.micrometer.observation.ObservationRegistry; +import org.springframework.aop.TargetSource; +import org.springframework.aop.framework.AopProxyUtils; import org.springframework.beans.factory.config.AbstractFactoryBean; import org.springframework.data.cassandra.ReactiveSession; import org.springframework.data.cassandra.core.cql.session.DefaultBridgedReactiveSession; @@ -25,6 +27,7 @@ import org.springframework.util.Assert; import org.springframework.util.ObjectUtils; import com.datastax.oss.driver.api.core.CqlSession; +import com.datastax.oss.driver.api.core.CqlSessionBuilder; /** * Factory bean to construct a {@link ReactiveSession} integrated with given {@link ObservationRegistry}. The required @@ -41,10 +44,30 @@ public class ObservableReactiveSessionFactoryBean extends AbstractFactoryBean