Class AxonServerQueryBus

java.lang.Object
org.axonframework.axonserver.connector.query.AxonServerQueryBus
All Implemented Interfaces:
Lifecycle, Distributed<QueryBus>, MessageDispatchInterceptorSupport<QueryMessage<?,?>>, MessageHandlerInterceptorSupport<QueryMessage<?,?>>, QueryBus

public class AxonServerQueryBus extends Object implements QueryBus, Distributed<QueryBus>, Lifecycle
Axon QueryBus implementation that connects to Axon Server to submit and receive queries and query responses. Delegates incoming queries to the provided localSegment.
Since:
4.0
Author:
Marc Gathier
  • Constructor Details

  • Method Details

    • streamingQuery

      public <Q, R> org.reactivestreams.Publisher<QueryResponseMessage<R>> streamingQuery(StreamingQueryMessage<Q,R> query)
      Description copied from interface: QueryBus
      Builds a Publisher of responses to the given query. The actual query is not dispatched until there is a subscription to the result. The query is dispatched to a single query handler. Implementations may opt for invoking several query handlers and then choosing a response from single one for performance or resilience reasons.

      When no handlers are available that can answer the given query, the return Publisher will be completed with a NoHandlerForQueryException.

      Specified by:
      streamingQuery in interface QueryBus
      Type Parameters:
      Q - the payload type of the streaming query
      R - the response type of the streaming query
      Parameters:
      query - the streaming query message
      Returns:
      a Publisher of responses
    • builder

      public static AxonServerQueryBus.Builder builder()
      Instantiate a Builder to be able to create an AxonServerQueryBus.

      The QueryPriorityCalculator is defaulted to QueryPriorityCalculator.defaultQueryPriorityCalculator(), the TargetContextResolver defaults to a lambda returning the AxonServerConfiguration.getContext() as the context, the ExecutorServiceBuilder defaults to ExecutorServiceBuilder.defaultQueryExecutorServiceBuilder(). The AxonServerConnectionManager and the QueryBusSpanFactory defaults to a DefaultQueryBusSpanFactory backed by a NoOpSpanFactory. The AxonServerConfiguration, the local QueryBus, the QueryUpdateEmitter, and the message and generic Serializers are hard requirements and as such should be provided.

      Returns:
      a Builder to be able to create a AxonServerQueryBus
    • registerLifecycleHandlers

      public void registerLifecycleHandlers(@Nonnull Lifecycle.LifecycleRegistry lifecycle)
      Description copied from interface: Lifecycle
      Registers the activities to be executed in the various phases of an application's lifecycle. This could either be at startup, shutdown, or both.
      Specified by:
      registerLifecycleHandlers in interface Lifecycle
      Parameters:
      lifecycle - the lifecycle instance to register the handlers with
      See Also:
    • start

      public void start()
      Start the Axon Server QueryBus implementation.
    • subscribe

      public <R> Registration subscribe(@Nonnull String queryName, @Nonnull Type responseType, @Nonnull MessageHandler<? super QueryMessage<?,R>> handler)
      Description copied from interface: QueryBus
      Subscribe the given handler to queries with the given queryName and responseType. Multiple handlers may subscribe to the same combination of queryName/responseType.
      Specified by:
      subscribe in interface QueryBus
      Parameters:
      queryName - the name of the query request to subscribe
      responseType - the type of response the subscribed component answers with
      handler - a handler that implements the query
      Returns:
      a handle to un-subscribe the query handler
    • query

      public <Q, R> CompletableFuture<QueryResponseMessage<R>> query(@Nonnull QueryMessage<Q,R> queryMessage)
      Description copied from interface: QueryBus
      Dispatch the given query to a single QueryHandler subscribed to the given query's queryName and responseType. This method returns all values returned by the Query Handler as a Collection. This may or may not be the exact collection as defined in the Query Handler.

      If the Query Handler defines a single return object (i.e. not a collection or array), that object is returned as the sole entry in a singleton collection.

      When no handlers are available that can answer the given query, the returned CompletableFuture will be completed with a NoHandlerForQueryException.

      Specified by:
      query in interface QueryBus
      Type Parameters:
      Q - the payload type of the query
      R - the response type of the query
      Parameters:
      queryMessage - the query
      Returns:
      a CompletableFuture that resolves when the response is available
    • scatterGather

      public <Q, R> Stream<QueryResponseMessage<R>> scatterGather(@Nonnull QueryMessage<Q,R> queryMessage, long timeout, @Nonnull TimeUnit timeUnit)
      Description copied from interface: QueryBus
      Dispatch the given query to all QueryHandlers subscribed to the given query's queryName/responseType. Returns a stream of results which blocks until all handlers have processed the request or when the timeout occurs.

      If no handlers are available to provide a result, or when all available handlers throw an exception while attempting to do so, the returned Stream is empty.

      Note that any terminal operation (such as Stream.forEach(Consumer)) on the Stream may cause it to block until the timeout has expired, awaiting additional data to include in the stream.

      Specified by:
      scatterGather in interface QueryBus
      Type Parameters:
      Q - the payload type of the query
      R - the response type of the query
      Parameters:
      queryMessage - the query
      timeout - time to wait for results
      timeUnit - unit for the timeout
      Returns:
      stream of query results
    • subscriptionQuery

      @Deprecated public <Q, I, U> SubscriptionQueryResult<QueryResponseMessage<I>,SubscriptionQueryUpdateMessage<U>> subscriptionQuery(@Nonnull SubscriptionQueryMessage<Q,I,U> query, SubscriptionQueryBackpressure backPressure, int updateBufferSize)
      Deprecated.
      Dispatch the given query to a single QueryHandler subscribed to the given query's queryName/initialResponseType/updateResponseType. The result is lazily created and there will be no execution of the query handler before there is a subscription to the initial result. In order not to miss updates, the query bus will queue all updates which happen after the subscription query is done and once the subscription to the flux is made, these updates will be emitted.

      If there is an error during retrieving or consuming initial result, stream for incremental updates is NOT interrupted.

      If there is an error during emitting an update, subscription is cancelled causing further emits not reaching the destination.

      Provided backpressure mechanism will be used to deal with fast emitters.

      Specified by:
      subscriptionQuery in interface QueryBus
      Type Parameters:
      Q - the payload type of the query
      I - the response type of the query
      U - the incremental response types of the query
      Parameters:
      query - the query
      backPressure - the backpressure mechanism to be used for emitting updates
      updateBufferSize - the size of buffer which accumulates updates before subscription to the flux is made
      Returns:
      query result containing initial result and incremental updates
    • subscriptionQuery

      public <Q, I, U> SubscriptionQueryResult<QueryResponseMessage<I>,SubscriptionQueryUpdateMessage<U>> subscriptionQuery(@Nonnull SubscriptionQueryMessage<Q,I,U> query, int updateBufferSize)
      Description copied from interface: QueryBus
      Dispatch the given query to a single QueryHandler subscribed to the given query's queryName/initialResponseType/updateResponseType. The result is lazily created and there will be no execution of the query handler before there is a subscription to the initial result. In order not to miss updates, the query bus will queue all updates which happen after the subscription query is done and once the subscription to the flux is made, these updates will be emitted.

      If there is an error during retrieving or consuming initial result, stream for incremental updates is NOT interrupted.

      If there is an error during emitting an update, subscription is cancelled causing further emits not reaching the destination.

      Specified by:
      subscriptionQuery in interface QueryBus
      Type Parameters:
      Q - the payload type of the query
      I - the response type of the query
      U - the incremental response types of the query
      Parameters:
      query - the query
      updateBufferSize - the size of buffer which accumulates updates before subscription to the flux is made
      Returns:
      query result containing initial result and incremental updates
    • queryUpdateEmitter

      public QueryUpdateEmitter queryUpdateEmitter()
      Description copied from interface: QueryBus
      Gets the QueryUpdateEmitter associated with this QueryBus.
      Specified by:
      queryUpdateEmitter in interface QueryBus
      Returns:
      the associated QueryUpdateEmitter
    • localSegment

      public QueryBus localSegment()
      Description copied from interface: Distributed
      Return the message bus of type MessageBus which is regarded as the local segment for this implementation. Would return the message bus used to dispatch and handle messages in a local environment to bridge the gap in a distributed set up.
      Specified by:
      localSegment in interface Distributed<QueryBus>
      Returns:
      a MessageBus which is the local segment for this distributed message bus implementation
    • registerHandlerInterceptor

      public Registration registerHandlerInterceptor(@Nonnull MessageHandlerInterceptor<? super QueryMessage<?,?>> interceptor)
      Description copied from interface: MessageHandlerInterceptorSupport
      Register the given handlerInterceptor. After registration, the interceptor will be invoked for each handled Message on the messaging component that it was registered to, prior to invoking the message's handler.
      Specified by:
      registerHandlerInterceptor in interface MessageHandlerInterceptorSupport<QueryMessage<?,?>>
      Parameters:
      interceptor - The interceptor to register
      Returns:
      A Registration, which may be used to deregister the interceptor.
    • registerDispatchInterceptor

      @Nonnull public Registration registerDispatchInterceptor(@Nonnull MessageDispatchInterceptor<? super QueryMessage<?,?>> dispatchInterceptor)
      Description copied from interface: MessageDispatchInterceptorSupport
      Register the given DispatchInterceptor. After registration, the interceptor will be invoked for each Message dispatched on the messaging component that it was registered to.
      Specified by:
      registerDispatchInterceptor in interface MessageDispatchInterceptorSupport<QueryMessage<?,?>>
      Parameters:
      dispatchInterceptor - The interceptor to register
      Returns:
      A Registration, which may be used to deregister the interceptor.
    • disconnect

      public void disconnect()
      Disconnect the query bus from Axon Server, by unsubscribing all known query handlers and aborting all queries in progress.
    • shutdownDispatching

      public CompletableFuture<Void> shutdownDispatching()
      Shutdown the query bus asynchronously for dispatching queries to Axon Server. This process will wait for dispatched queries which have not received a response yet and will close off running subscription queries. This shutdown operation is performed in the Phase.OUTBOUND_QUERY_CONNECTORS phase.
      Returns:
      a completable future which is resolved once all query dispatching activities are completed