GH-2374 Add initial support for function-based error handling
Resolves #2374
This commit is contained in:
@@ -20,6 +20,7 @@ import java.io.IOException;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import com.fasterxml.jackson.core.JsonGenerator;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
@@ -32,6 +33,9 @@ import org.springframework.beans.factory.BeanFactoryAware;
|
||||
import org.springframework.beans.factory.DisposableBean;
|
||||
import org.springframework.beans.factory.InitializingBean;
|
||||
import org.springframework.beans.factory.support.DefaultSingletonBeanRegistry;
|
||||
import org.springframework.cloud.function.context.FunctionCatalog;
|
||||
import org.springframework.cloud.stream.config.BindingProperties;
|
||||
import org.springframework.cloud.stream.config.BindingServiceProperties;
|
||||
import org.springframework.cloud.stream.config.ConsumerEndpointCustomizer;
|
||||
import org.springframework.cloud.stream.config.ListenerContainerCustomizer;
|
||||
import org.springframework.cloud.stream.config.MessageSourceCustomizer;
|
||||
@@ -66,9 +70,11 @@ import org.springframework.messaging.MessageHandler;
|
||||
import org.springframework.messaging.MessageHeaders;
|
||||
import org.springframework.messaging.SubscribableChannel;
|
||||
import org.springframework.messaging.support.ChannelInterceptor;
|
||||
import org.springframework.messaging.support.ErrorMessage;
|
||||
import org.springframework.retry.RecoveryCallback;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.CollectionUtils;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
/**
|
||||
* {@link AbstractBinder} that serves as base class for {@link MessageChannel} binders.
|
||||
@@ -687,6 +693,25 @@ public abstract class AbstractMessageChannelBinder<C extends ConsumerProperties,
|
||||
return registerErrorInfrastructure(destination, group, consumerProperties, false);
|
||||
}
|
||||
|
||||
private void subscribeFunctionErrorHandler(String errorChannelName, String bindingName) {
|
||||
if (!StringUtils.hasText(bindingName)) {
|
||||
return;
|
||||
}
|
||||
BindingServiceProperties bsp = getApplicationContext().getBean(BindingServiceProperties.class);
|
||||
BindingProperties bp = bsp.getBindingProperties(bindingName);
|
||||
if (StringUtils.hasText(bp.getErrorHandlerDefinition())) {
|
||||
FunctionCatalog catalog = getApplicationContext().getBean(FunctionCatalog.class);
|
||||
Consumer<ErrorMessage> errorHandler = catalog.lookup(Consumer.class, bp.getErrorHandlerDefinition());
|
||||
if (errorHandler == null) {
|
||||
logger.warn("Failed to retrieve error handling function with definition: " + bp.getErrorHandlerDefinition() + ", for binding: " + bindingName);
|
||||
}
|
||||
else {
|
||||
SubscribableChannel functionErrorChannel = getApplicationContext().getBean(errorChannelName, SubscribableChannel.class);
|
||||
functionErrorChannel.subscribe(errorMessage -> errorHandler.accept((ErrorMessage) errorMessage));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Build an errorChannelRecoverer that writes to a pub/sub channel for the destination
|
||||
* when an exception is thrown to a consumer.
|
||||
@@ -720,7 +745,10 @@ public abstract class AbstractMessageChannelBinder<C extends ConsumerProperties,
|
||||
|
||||
((GenericApplicationContext) getApplicationContext()).registerBean(
|
||||
errorChannelName, SubscribableChannel.class, () -> errorChannel);
|
||||
|
||||
this.subscribeFunctionErrorHandler(errorChannelName, consumerProperties.getBindingName());
|
||||
}
|
||||
|
||||
ErrorMessageSendingRecoverer recoverer;
|
||||
if (errorMessageStrategy == null) {
|
||||
recoverer = new ErrorMessageSendingRecoverer(errorChannel);
|
||||
|
||||
@@ -46,6 +46,11 @@ public class BindingProperties {
|
||||
|
||||
private static final String COMMA = ",";
|
||||
|
||||
/**
|
||||
* Function that corresponds to the error handler of the underlying function.
|
||||
*/
|
||||
private String errorHandlerDefinition;
|
||||
|
||||
/**
|
||||
* The physical name at the broker that the binder binds to.
|
||||
*/
|
||||
@@ -157,4 +162,12 @@ public class BindingProperties {
|
||||
return "BindingProperties{" + sb.toString() + "}";
|
||||
}
|
||||
|
||||
public String getErrorHandlerDefinition() {
|
||||
return errorHandlerDefinition;
|
||||
}
|
||||
|
||||
public void setErrorHandlerDefinition(String errorHandlerDefinition) {
|
||||
this.errorHandlerDefinition = errorHandlerDefinition;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -20,26 +20,17 @@ import java.util.function.Consumer;
|
||||
import java.util.function.Function;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.mockito.Mockito;
|
||||
|
||||
import org.springframework.boot.SpringApplication;
|
||||
import org.springframework.boot.WebApplicationType;
|
||||
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
|
||||
import org.springframework.boot.builder.SpringApplicationBuilder;
|
||||
import org.springframework.cloud.stream.binder.test.InputDestination;
|
||||
import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration;
|
||||
import org.springframework.context.ApplicationContext;
|
||||
import org.springframework.context.ConfigurableApplicationContext;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.integration.annotation.ServiceActivator;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.ArgumentMatchers.eq;
|
||||
import static org.mockito.ArgumentMatchers.isNull;
|
||||
|
||||
/**
|
||||
* @author Marius Bogoevici
|
||||
@@ -47,25 +38,6 @@ import static org.mockito.ArgumentMatchers.isNull;
|
||||
*/
|
||||
public class ErrorBindingTests {
|
||||
|
||||
@SuppressWarnings({ "rawtypes", "unchecked" })
|
||||
@Test
|
||||
void testErrorChannelNotBoundByDefault() {
|
||||
ConfigurableApplicationContext applicationContext = SpringApplication.run(
|
||||
TestProcessor.class, "--server.port=0",
|
||||
"--spring.cloud.stream.default-binder=mock",
|
||||
"--spring.jmx.enabled=false");
|
||||
BinderFactory binderFactory = applicationContext.getBean(BinderFactory.class);
|
||||
|
||||
Binder binder = binderFactory.getBinder(null, MessageChannel.class);
|
||||
|
||||
Mockito.verify(binder).bindConsumer(eq("processor-in-0"), isNull(),
|
||||
any(MessageChannel.class), any(ConsumerProperties.class));
|
||||
Mockito.verify(binder).bindProducer(eq("processor-out-0"), any(MessageChannel.class),
|
||||
any(ProducerProperties.class));
|
||||
Mockito.verifyNoMoreInteractions(binder);
|
||||
applicationContext.close();
|
||||
}
|
||||
|
||||
@Test
|
||||
void testConfigurationWithDefaultErrorHandler() {
|
||||
ApplicationContext context = new SpringApplicationBuilder(
|
||||
@@ -73,6 +45,8 @@ public class ErrorBindingTests {
|
||||
ErrorBindingTests.ErrorConfigurationDefault.class))
|
||||
.web(WebApplicationType.NONE)
|
||||
.run("--spring.cloud.stream.bindings.handle-in-0.consumer.max-attempts=1",
|
||||
"--spring.cloud.function.definition=handle",
|
||||
"--spring.cloud.stream.default.error-handler-definition=errorHandler",
|
||||
"--spring.jmx.enabled=false");
|
||||
|
||||
InputDestination source = context.getBean(InputDestination.class);
|
||||
@@ -82,16 +56,18 @@ public class ErrorBindingTests {
|
||||
|
||||
ErrorConfigurationDefault errorConfiguration = context
|
||||
.getBean(ErrorConfigurationDefault.class);
|
||||
assertThat(errorConfiguration.counter == 3);
|
||||
assertThat(errorConfiguration.counter).isEqualTo(6);
|
||||
}
|
||||
|
||||
@Test
|
||||
void testConfigurationWithCustomErrorHandler() {
|
||||
void testConfigurationWithBindingSpecificErrorHandler() {
|
||||
ApplicationContext context = new SpringApplicationBuilder(
|
||||
TestChannelBinderConfiguration.getCompleteConfiguration(
|
||||
ErrorBindingTests.ErrorConfigurationWithCustomErrorHandler.class))
|
||||
.web(WebApplicationType.NONE)
|
||||
.run("--spring.cloud.stream.bindings.handle-in-0.consumer.max-attempts=1",
|
||||
"--spring.cloud.function.definition=handle",
|
||||
"--spring.cloud.stream.bindings.handle-in-0.error-handler-definition=errorHandler",
|
||||
"--spring.jmx.enabled=false");
|
||||
|
||||
InputDestination source = context.getBean(InputDestination.class);
|
||||
@@ -101,7 +77,7 @@ public class ErrorBindingTests {
|
||||
|
||||
ErrorConfigurationWithCustomErrorHandler errorConfiguration = context
|
||||
.getBean(ErrorConfigurationWithCustomErrorHandler.class);
|
||||
assertThat(errorConfiguration.counter == 6);
|
||||
assertThat(errorConfiguration.counter).isEqualTo(6);
|
||||
}
|
||||
|
||||
@EnableAutoConfiguration
|
||||
@@ -126,6 +102,13 @@ public class ErrorBindingTests {
|
||||
};
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Consumer<Object> errorHandler() {
|
||||
return v -> {
|
||||
this.counter++;
|
||||
};
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@EnableAutoConfiguration
|
||||
@@ -141,11 +124,12 @@ public class ErrorBindingTests {
|
||||
};
|
||||
}
|
||||
|
||||
@ServiceActivator(inputChannel = "input.anonymous.errors")
|
||||
public void error(Message<?> message) {
|
||||
this.counter++;
|
||||
@Bean
|
||||
public Consumer<Object> errorHandler() {
|
||||
return v -> {
|
||||
this.counter++;
|
||||
};
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user