Annotation-driven Listener Endpoints
The easiest way to receive a message asynchronously is to use the annotated listener endpoint infrastructure. In a nutshell, it lets you expose a method of a managed bean as a Valkey listener endpoint. The following example shows how to use it:
@Componentpublic class MyService {
@ValkeyListener(topic = "my-channel") public void processOrder(String data) { ... }}The idea of the preceding example is that, whenever a message is received through a channel my-channel, the processOrder method is invoked accordingly (in this case, with the content of the Valkey Pub/Sub message, similar to what the MessageListenerAdapter provides).
The annotated endpoint infrastructure creates a message listener behind the scenes for each annotated method and registers it to ValkeyMessageListenerContainer.
Endpoints are not registered against the application context but can be easily located for management purposes by using the ValkeyListenerEndpointRegistry bean.
Note: All three drivers (Valkey GLIDE, Lettuce, and Jedis) support annotation-driven listener endpoints.
Enable Listener Endpoint Annotations
Section titled “Enable Listener Endpoint Annotations”To enable support for @ValkeyListener annotations, you can add @EnableValkeyListeners to one of your @Configuration classes, as the following example shows:
@Configuration@EnableValkeyListenerspublic class ValkeyConfiguration {
@Bean public ValkeyMessageListenerContainer valkeyMessageListenerContainer(ValkeyConnectionFactory connectionFactory) {
ValkeyMessageListenerContainer factory = new ValkeyMessageListenerContainer(); factory.setConnectionFactory(connectionFactory); return factory; }}@Configuration@EnableValkeyListenersclass ValkeyConfiguration {
@Bean fun valkeyMessageListenerContainer(connectionFactory: ValkeyConnectionFactory) = ValkeyMessageListenerContainer().apply { setConnectionFactory(connectionFactory); }}You can customize the listener registrar by implementing the ValkeyListenerConfigurer interface.
See the javadoc of classes that implement ValkeyListenerConfigurer for details and examples.
Programmatic Endpoint Registration
Section titled “Programmatic Endpoint Registration”ValkeyListenerEndpoint provides a model of a Valkey endpoint and is responsible for configuring the container for that model.
The infrastructure lets you programmatically configure endpoints in addition to the ones that are detected by the @ValkeyListener annotation.
The following example shows how to do so:
@Configuration@EnableValkeyListenerspublic class AppConfig implements ValkeyListenerConfigurer {
@Override public void configureValkeyListeners(ValkeyListenerEndpointRegistrar registrar) { SimpleValkeyListenerEndpoint endpoint = new SimpleValkeyListenerEndpoint(); endpoint.setId("myValkeyEndpoint"); endpoint.setTopic("my-channel"); endpoint.setMessageListener((message, pattern) -> { // processing }); registrar.registerEndpoint(endpoint); }}In the preceding example, we used SimpleValkeyListenerEndpoint, which provides the actual MessageListener to invoke.
However, you could also build your own endpoint variant to describe a custom invocation mechanism.
Note that you could skip the use of @ValkeyListener altogether and programmatically register only your endpoints through ValkeyListenerConfigurer.
Annotated Endpoint Method Signature
Section titled “Annotated Endpoint Method Signature”So far, we have been injecting a simple String in our endpoint, but it can actually have a very flexible method signature.
In the following example, we rewrite it to inject the Order with a header:
@Componentpublic class MyService {
@ValkeyListener("my-order-channels*") public void processOrder(Order order, @Header("pattern") Topic pattern) { ... }}The main elements you can inject in Valkey listener endpoints are as follows:
- The
org.springframework.messaging.Messagethat represents the incoming Valkey message. Note that this message holds headers (as defined byPubSubHeaders). @Header-annotated method arguments to extract a specific header value. Since Valkey Pub/Sub messages do consist only of the body, header values such as the topic or pattern that matched the message are synthetically added as headers.- A
@Headers-annotated argument that must also be assignable tojava.util.Mapfor getting access to all headers. - A non-annotated element that is not one of the supported types is considered to be the payload. You can make that explicit by annotating the parameter with
@Payload. You can also turn on validation by adding an extra@Valid.
The ability to inject Spring’s Message abstraction is particularly useful to benefit from all the information stored in the transport-specific message without relying on transport-specific API.
The following example shows how to do so:
@ValkeyListener("my-channel")public void processOrder(Message<byte[]> order) { ... }Handling of method arguments is provided by DefaultMessageHandlerMethodFactory, which you can further customize to support additional method arguments.
Annotated-based endpoints support flexible message conversion, which is provided by ValkeyMessageConverters along with a default set of converters for String, byte array, and JSON (if a supported library is present on the classpath).
If you wish to customize the message conversion, you can do so by implementing ValkeyListenerConfigurer and overriding the configureMessageConverters method, as the following example shows:
@Configuration@EnableValkeyListenerspublic class AppConfig implements ValkeyListenerConfigurer {
@Override public void configureMessageConverters(ValkeyMessageConverters.Builder builder) { builder.withStringConverter(StandardCharsets.UTF_8) .addCustomConverter(new JdkSerializationValkeySerializer()); }}ValkeyMessageConverters registers by default the following converters:
StringMessageConverterByteArrayMessageConverter- JSON converters (if a supported library like Jackson, Gson, JSON-B, or Kotlin Serialization is present on the classpath)
@ValkeyListener(consumes) is useful to indicate the desired content type when using different message converters to select the appropriate converter.
The next example shows converter selection for the JdkSerializationValkeySerializer converter assuming a registration as per the previous example:
@ValkeyListener(topic = "my-channel", consumes = JdkSerializerMessageConverter.APPLICATION_JAVA_SERIALIZED_OBJECT_VALUE)public void processOrder(Person person) { ... }If you require more control over the method argument resolution, you can also configure a custom MessageHandlerMethodFactory through ValkeyListenerConfigurer.
You can customize the conversion and validation support there as well.
For instance, if we want to make sure our Order is valid before processing it, we can annotate the payload with @Valid and configure the necessary validator, as the following example shows:
@Configuration@EnableValkeyListenerspublic class AppConfig implements ValkeyListenerConfigurer {
@Override public void configureValkeyListeners(ValkeyListenerEndpointRegistrar registrar) { registrar.setMessageHandlerMethodFactory(myValkeyHandlerMethodFactory()); }
@Bean public DefaultMessageHandlerMethodFactory myValkeyHandlerMethodFactory() { DefaultMessageHandlerMethodFactory factory = new DefaultMessageHandlerMethodFactory(); factory.setValidator(myValidator()); return factory; }}