Infrastructure
When it comes to dispatching queries, as explained in the Dispatching queries section, there are several implementations for actually sending query messages. The next sections provide an overview of the available implementations, as well as how to set up query dispatching infrastructure with Axon.
All query operations are async-native, returning CompletableFuture or Publisher types. This enables efficient asynchronous processing without blocking threads.
Query gateway
The query gateway is a convenient interface towards the query dispatching mechanism. While you are not required to use a gateway to dispatch queries, it is generally the easiest option to do so.
Axon provides a QueryGateway interface and the DefaultQueryGateway implementation.
The query gateway provides methods for:
-
Point-to-point queries:
query()for single results,queryMany()for multiple results (both returnCompletableFuture) -
Streaming queries:
streamingQuery()returns aPublisherfor streaming large result sets -
Subscription queries:
subscriptionQuery()returns aPublishercombining initial result and updates
All methods support an optional ProcessingContext parameter for correlation data propagation when dispatching from within a message handler.
The query gateway automatically handles message construction, dispatch interceptors, and result conversion. It’s configured with access to the query bus and dispatch interceptors.
Query bus
The query bus is the mechanism that dispatches queries to query handlers.
Queries are registered using the query name (based on MessageType).
Each query name can have only one registered handler.
DistributedQueryBus
The DistributedQueryBus, together with the AxonServerQueryBusConnector, is the default query bus when Axon Server is available.
It connects to Axon Server to send and receive queries in a distributed way, allowing queries to be handled by any node in the cluster.
The DistributedQueryBus combined with the AxonServerQueryBusConnector provides:
-
Distributed query handling: Queries are routed to the appropriate handler across connected applications
-
Load balancing: When the same application is deployed across multiple instances, Axon Server can distribute queries among those instances
-
Subscription query support: Full support for subscription queries with update distribution across nodes
-
Query prioritization: Support for query priority to ensure critical queries are processed first
-
Local handler shortcut: Direct (point-to-point) queries for which a local handler is registered can be executed locally without going through the connector, avoiding network overhead. A
LocalQueryDispatchPredicatedecides, per query, whether to take this shortcut. Subscription queries are never shortcut
Configuring the distributed query bus
The DistributedQueryBus can be configured using the DistributedQueryBusConfiguration class, which provides a fluent API for customizing query processing behavior.
Available configuration options:
-
Query threads: Number of threads used for query processing (default: 10)
-
Query queue capacity: Capacity of the priority queue for query processing tasks (default: 1000)
-
Custom executor service: Provide a custom
ExecutorServicefor query processing
Local handler shortcut
The local handler shortcut applies to direct (point-to-point) queries only. A LocalQueryDispatchPredicate lets a node prefer its own query handler for selected direct queries instead of routing them through the connector. The predicate is consulted for every outgoing direct query and receives the QueryMessage and the (nullable) ProcessingContext, so the decision can be based on the query’s payload, type, metadata, or a flag placed on the context by the dispatcher.
The shortcut only takes effect when the local segment actually subscribed a handler for the query. When it did not, the query is routed through the connector as usual, regardless of the predicate. This ensures a locally preferred query that this node cannot handle is still routed to a segment that can.
Short-cut queries are handled on the same bounded, priority-ordered worker pool as queries arriving from remote segments, because the shortcut reuses the same local-handling path. A burst of locally preferred queries is therefore subject to the same back-pressure and cannot flood the local segment beyond what it would already accept from remote traffic.
Subscription queries are never shortcut. They are always routed through the connector, because their update registrations must be coordinated across all nodes. The LocalQueryDispatchPredicate is not consulted for subscription queries.
To enable the shortcut, register a LocalQueryDispatchPredicate component:
import io.axoniq.framework.messaging.queryhandling.distributed.LocalQueryDispatchPredicate;
import org.axonframework.messaging.core.configuration.MessagingConfigurer;
public class AxonConfig {
public void configureLocalQueryShortcut(MessagingConfigurer configurer) {
configurer.componentRegistry(registry -> registry.registerComponent(
LocalQueryDispatchPredicate.class,
config -> (query, context) -> true // always prefer the local handler when subscribed
));
}
}
|
The
preferLocalQueryHandler setting
|
Configuration examples
-
Configuration API
-
Spring Boot
Customize the distributed query bus configuration:
import io.axoniq.framework.messaging.queryhandling.distributed.DistributedQueryBusConfiguration;
import io.axoniq.framework.messaging.queryhandling.distributed.LocalQueryDispatchPredicate;
import org.axonframework.messaging.core.configuration.MessagingConfigurer;
public class AxonConfig {
public void configureQueryBus(MessagingConfigurer configurer) {
// Customize the configuration
DistributedQueryBusConfiguration config = DistributedQueryBusConfiguration.DEFAULT
.queryThreads(20) // Set number of query processing threads
.queryQueueCapacity(2000); // Set queue capacity
// Register the custom configuration
configurer.componentRegistry(
cr -> cr.registerComponent(DistributedQueryBusConfiguration.class, c -> config)
.registerComponent(
LocalQueryDispatchPredicate.class,
c -> (query, context) -> true
)
);
}
}
To control the local handler shortcut, register a LocalQueryDispatchPredicate as shown above rather than using the deprecated preferLocalQueryHandler setting.
When using Spring Boot with Axon Server, configure the query bus through application properties:
# Configure query processing threads
axon.axonserver.query-threads=20
# Or configure programmatically
import io.axoniq.framework.messaging.queryhandling.distributed.DistributedQueryBusConfiguration;
import io.axoniq.framework.messaging.queryhandling.distributed.LocalQueryDispatchPredicate;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class AxonConfig {
@Bean
public DistributedQueryBusConfiguration queryBusConfiguration() {
return DistributedQueryBusConfiguration.DEFAULT.queryThreads(20);
}
@Bean
public LocalQueryDispatchPredicate localQueryDispatchPredicate() {
// Prefer the local handler for every query that this node can handle.
return (query, context) -> true;
}
}
Default configuration
-
Configuration API
-
Spring Boot
Declare dependencies:
<!--somewhere in the POM file-->
<dependency>
<groupId>io.axoniq.framework</groupId>
<artifactId>axon-server-connector</artifactId>
<version>${axoniq.version}</version>
</dependency>
<dependency>
<groupId>org.axonframework</groupId>
<artifactId>axon-messaging</artifactId>
<version>${axon.version}</version>
</dependency>
Configure your application:
import org.axonframework.messaging.core.configuration.MessagingConfigurer;
public class AxonApp {
public static void main(String[] args) {
// Returns a Configurer instance with default components configured.
// `DistributedQueryBus` with `AxonServerQueryBusConnector` is configured as Query Bus by default.
MessagingConfigurer configurer = MessagingConfigurer.create();
}
}
By simply declaring dependency to axon-spring-boot-starter, Axon will automatically configure the AxonServerQueryBusConnector:
<!--somewhere in the POM file-->
<dependency>
<groupId>org.axonframework</groupId>
<artifactId>axon-spring-boot-starter</artifactId>
<version>${axon.version}</version>
</dependency>
|
Excluding the Axon Server Connector
If you exclude the |
SpringCloudQueryBusConnector
The SpringCloudQueryBusConnector distributes queries over HTTP between the nodes reported by a Spring Cloud DiscoveryClient, without an Axon Server in between.
It is part of the same starter as the SpringCloudCommandBusConnector, and shares its discovery, its capabilities endpoint, and its routing ring.
Where a command carries a routing key that pins it to one node, a query asks nothing of where it is handled: any node advertising the query’s name may answer it.
Successive queries of a name therefore rotate over the nodes advertising it, spreading the load rather than sending every query of a name to the same node.
The preferLocalQueryHandler setting still applies, and is handled by the DistributedQueryBus before a query reaches the connector at all.
A query may be answered any number of times, so a node answers with a stream rather than a single reply.
Responses travel back as Server-Sent Events, and appear on the dispatching node’s MessageStream as the answering node produces them.
Bounded buffering
Responses are buffered per query while the application consumes them. A node answering faster than the dispatching application consumes fills that buffer, and the query then fails rather than growing the buffer until memory runs out:
# How many responses to a single query are held before an outpaced query fails.
axon.springcloud.query-buffer-size=1024
# How long a query's response stream may stay open before the container closes it.
axon.springcloud.query-timeout=5m
# How long the responses to a dispatched query are waited for before it is given up on.
axon.springcloud.query-response-timeout=6m
Raise the buffer size where responses legitimately arrive in bursts, and the timeout where a query legitimately takes longer to answer than the default.
query-timeout applies on the answering node and query-response-timeout on the dispatching one.
Keep the latter above the former, so that a node which is merely slow ends the query itself.
The dispatching deadline is a backstop for a node that stops answering without saying so, having been killed or partitioned away.
Its socket reports nothing, so without a deadline the query would stay unanswered indefinitely.
A query whose stream is closed releases the handler’s response stream on the answering node, whether the timeout elapsed or the node that asked stopped reading. A handler is not left producing responses that nothing receives.
Shutting down
A node leaving the cluster stops advertising the queries it handled before it stops answering them, so the other nodes route elsewhere rather than discovering it left by timing out. Queries already dispatched are waited for: shutdown completes once their responses have arrived, and queries dispatched after shutdown began are rejected.
Subscription queries
A subscription query is two activities. First, an update stream is opened on every node advertising the query’s name. Second, one of those nodes is asked for the initial result, which is a query like any other. The two are returned as a single stream: the initial result, and then the updates.
Every node advertising the name is subscribed to, not just the one a plain query would route to. An update is emitted on whichever node’s state changed, and reaches only the subscriptions that node holds a registration for, so a subscriber reaching one node would miss every update emitted on the others. The updates of all of them arrive on the one stream the subscriber reads.
An idle subscription is written to periodically, because a subscription may go a long time without an update and an idle connection is what a load balancer, proxy, or network address translation table reclaims:
# How often a node writes to a subscription it is answering while it has no update to send.
axon.springcloud.subscription-keep-alive-interval=20s
# How long a subscription may hear nothing at all before the subscribing node gives up on it.
axon.springcloud.subscription-inactivity-timeout=60s
Keep the interval comfortably below the idle timeout of whatever sits between the nodes, and the inactivity timeout a few intervals wide. Unlike a plain query, a subscription has no deadline of its own: it lasts as long as the subscriber wants it to, and a node with nothing to report is behaving correctly. What stands in for a deadline is silence, which is how a node that disappeared without its socket reporting so looks.
Ending a subscription
A node that completes a subscription query says so explicitly, and that ends the subscription: there will never be another update to it, so the subscribing node stops waiting on the other nodes as well.
A node that simply stops answering, having shut down or been partitioned away, is a different thing. That ends its own part in the subscription and says nothing about the subscription itself, so the remaining nodes carry on and the subscription ends once the last of them has. A node that fails, rather than ending quietly, fails the whole subscription: carrying on with the rest would leave the subscriber receiving some of the updates and believing it received all of them.
A node that starts handling the query
A node that starts advertising the query’s name while a subscription is active fails that subscription with a SubscriptionQueryMembersChangedException.
The updates that node emitted between advertising the name and being subscribed to are already gone, and no amount of catching up recovers them. Failing is therefore the honest outcome: a subscriber that establishes the subscription query again gets a fresh initial result and a complete stream of updates from that point, where one that carried on would silently be missing whatever it missed. A node that joins without advertising the query’s name emits no updates for it, and leaves the subscription alone.
|
A servlet web application is required
The endpoints nodes reach each other on are Spring Model-View-Controller (MVC) endpoints, and a query’s responses are streamed with an |
Query throughput is claimed against your Axoniq licence, on the same terms as command throughput.
SimpleQueryBus
The SimpleQueryBus is a local, non-distributed query bus implementation that processes queries in the dispatching thread by default. It’s used when Axon Server is not available or when you explicitly configure it.
The SimpleQueryBus provides:
-
Local query handling: All queries are handled within the same JVM.
-
Simple routing: Routes queries to their registered handler based on query name.
-
Subscription query support: Full support for subscription queries within the same JVM.
-
Transaction management: Integrates with a
TransactionManagerfor transactional query handling.
To configure a SimpleQueryBus (instead of an AxonServerQueryBus):
-
Configuration API
-
Spring Boot
import org.axonframework.messaging.core.configuration.MessagingConfigurer;
import org.axonframework.messaging.queryhandling.SimpleQueryBus;
import org.axonframework.messaging.core.unitofwork.UnitOfWorkFactory;
public class AxonConfig {
// omitting other configuration methods...
public void configureQueryBus(MessagingConfigurer configurer) {
configurer.registerQueryBus(
config -> new SimpleQueryBus(config.getComponent(UnitOfWorkFactory.class))
);
}
}
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.axonframework.messaging.queryhandling.QueryBus;
import org.axonframework.messaging.queryhandling.SimpleQueryBus;
import org.axonframework.messaging.core.unitofwork.UnitOfWorkFactory;
@Configuration
public class AxonConfig {
// omitting other configuration methods...
@Bean
public QueryBus queryBus(UnitOfWorkFactory unitOfWorkFactory) {
return new SimpleQueryBus(unitOfWorkFactory);
}
}
Subscription query infrastructure
Subscription queries allow clients to receive an initial result and then continue receiving updates as long as the subscription is active. The subscription query infrastructure consists of several components working together.
QueryUpdateEmitter
The QueryUpdateEmitter is responsible for emitting updates to active subscription queries. The QueryUpdateEmitter is context-aware and must be created from the ProcessingContext:
import org.axonframework.messaging.core.unitofwork.ProcessingContext;
import org.axonframework.messaging.eventhandling.annotation.EventHandler;
import org.axonframework.messaging.queryhandling.QueryUpdateEmitter;
import org.springframework.stereotype.Component;
@Component
public class CardSummaryProjection {
@EventHandler
public void on(CardRedeemedEvent event, ProcessingContext context) {
// Create a context-aware emitter
QueryUpdateEmitter emitter = QueryUpdateEmitter.forContext(context);
// Update the model
CardSummary summary = new CardSummary(event.cardId(), event.amount());
// Emit update to subscription queries
emitter.emit(
FetchCardSummaryQuery.class,
query -> query.cardSummaryId().equals(event.cardId()),
summary
);
}
}
The emitter filters subscription queries based on the provided predicate and emits the update only to matching subscriptions. This allows fine-grained control over which subscribers receive which updates.
|
Automatic emitter creation for annotated methods
When using annotated message handlers (
Explicitly creating the emitter using |
Subscription lifecycle
Subscription queries follow this lifecycle:
-
Subscription creation: Client calls
QueryGateway#subscriptionQuery(Object, Class<R>), which returns aPublisher<R>. -
Initial result: The query is dispatched to a handler, and the initial result is emitted.
-
Update buffering: Updates are buffered until the
Publisheris subscribed to. -
Update streaming: Once subscribed, buffered and new updates are streamed to the client.
-
Completion: The subscription completes when:
-
The client cancels the subscription (by disposing the
Publisher). -
The server calls
QueryUpdateEmitter.complete()to signal no more updates. -
An error occurs, completing the subscription exceptionally.
-
import org.axonframework.messaging.queryhandling.gateway.QueryGateway;
import org.reactivestreams.Publisher;
import reactor.core.Disposable;
import reactor.core.publisher.Flux;
public class QueryDispatcher {
public void dispatchFetchCard(QueryGateway queryGateway) {
String cardId = "...";
// Client-side subscription query
Publisher<CardSummary> results = queryGateway.subscriptionQuery(
new FetchCardSummaryQuery(cardId),
CardSummary.class
);
// Subscribe using Reactor (requires reactor-core dependency)
Disposable subscription = Flux.from(results)
.doOnNext(summary -> System.out.println("Received: " + summary))
.doOnComplete(() -> System.out.println("No more updates"))
.doOnError(error -> System.err.println("Error: " + error))
.subscribe();
// Later: cancel the subscription
subscription.dispose();
}
}
Update buffer
The QueryBus maintains an update buffer for each subscription query. Updates emitted before the client subscribes to the Publisher are stored in this buffer. Once the client subscribes, buffered updates are delivered first, followed by new updates.
The buffer size is configurable when creating a subscription query:
// Default buffer size based on Reactor's Queues.SMALL_BUFFER_SIZE constant
Publisher<CardSummary> results = queryGateway.subscriptionQuery(
new FetchCardSummaryQuery(cardId),
CardSummary.class
);
// Custom buffer size
Publisher<CardSummary> customBufferResults = queryGateway.subscriptionQuery(
new FetchCardSummaryQuery(cardId),
CardSummary.class,
512 // buffer size
);
If the buffer fills up before the client subscribes, attempting to add more updates will complete the subscription exceptionally. Choose an appropriate buffer size based on your expected update rate and subscription delay.