Valdiate pojo using @Valid in sping cloud streams

Viewed 292

How can one enable validation using @Valid inside the following kafka consumer code ? I am using Spring Cloud Stream (Kafka Stream binder implementation), and there after my implemention is using functional model for example.

@Bean
public Consumer<KStream<String, @Valid Pojo>> process() {
    return messages -> messages.foreach((k, v) -> process(v));
}

I tried the following but it didn't work....

@Bean
public DefaultMessageHandlerMethodFactory configureMessageHandlerMethodFactory(
        DefaultMessageHandlerMethodFactory messageHandlerMethodFactory,
        LocalValidatorFactoryBean validatorFactoryBean) {       
    messageHandlerMethodFactory.setValidator(validatorFactoryBean);
    return messageHandlerMethodFactory;
}

This is simple in spring-kafka by implementing KafkaListenerConfigurer and setting LocalValidatorFactoryBean on KafkaListenerEndpointRegistrar

public class KafkaConfiguration implements KafkaListenerConfigurer {

@Override
public void configureKafkaListeners(KafkaListenerEndpointRegistrar registrar) {
  registrar.setValidator(validatorFactoryBean);
}
.....
2 Answers

This is not supported in the functional model at the moment. Even for a non-functional scenario, this is non-trivial for types like KStream. The KafkaListenerConfigurer you mentioned above is for regular Kafka Support with a message channel binder. Your best options for Kafka Streams binder are either using some custom validation in the function itself before continuing with the processing or introducing a schema registry and then perform a schema validation before passing the record to the function.

You can follow the recommendation to create a bean that respects the functional interface of java, that is, it has only a public method, for example:

@Validated
@Component
public class Processor implements Consumer<KStream<String, Pojo>> {

    @Override
    public void accept(final @Valid @NotNull KStream<String, Pojo> stream) {
          stream.foreach((k, v) -> process(v));
    }

    private void process(final Pojo v) {
    }

}

So that generates an execution: javax.validation.ConstraintDeclarationException: HV000151: A method overriding another method must not reset the parameter constraint configuration

It is not possible to overwrite the parameters of the accept method of the consumer functional interface so just remove the interface and leave the component like this:

@Validated
@Component
public class Processor {

     public void accept(final @Valid @NotNull KStream<String, Pojo> stream) {
          stream.foreach((k, v) -> process(v));
     }

     private void process(final Pojo v) {
     }

 }

The problem is that the spring cloud function will not recognize the bean for not extending one of the functional classes.

the workaround I got was:

@RequiredArgsConstructor
public abstract class ValidatedEventListener<T> implements Consumer<T> {

    private final Validator validator;

    @Override
    public void accept(final T t) {
        validate(t);
        listen(t);
    }

    public abstract void listen(final T t);

    public void validate(final Object event) {
        var violations = validator.validate(event);
        if (!violations.isEmpty()) throw new ConstraintViolationException(violations);
    }

}
Related