diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/SerdeResolverUtils.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/SerdeResolverUtils.java index ac6a4ae3a8..4294846ca3 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/SerdeResolverUtils.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/SerdeResolverUtils.java @@ -36,10 +36,13 @@ import org.springframework.core.ResolvableType; import org.springframework.kafka.support.serializer.JacksonJsonSerde; +import tools.jackson.databind.json.JsonMapper; + /** * Utility class that contains various methods to help resolve {@link Serde Serdes}. * * @author Chris Bono + * @author adityaanikam * @since 4.0 */ abstract class SerdeResolverUtils { @@ -102,7 +105,8 @@ static Serde resolveForType(ConfigurableApplicationContext context, Resolvabl // Use JsonSerde if type is not exactly Object if (!genericRawClazz.isAssignableFrom((Object.class))) { - return new JacksonJsonSerde<>(genericRawClazz); + JsonMapper jsonMapper = context.getBeanProvider(JsonMapper.class).getIfUnique(); + return jsonMapper != null ? new JacksonJsonSerde<>(genericRawClazz, jsonMapper) : new JacksonJsonSerde<>(genericRawClazz); } // Finally, just resort to using the fallback diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/SerdeResolverUtilsTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/SerdeResolverUtilsTests.java index e81f46670c..f962e33bc3 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/SerdeResolverUtilsTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/SerdeResolverUtilsTests.java @@ -16,6 +16,7 @@ package org.springframework.cloud.stream.binder.kafka.streams; +import java.lang.reflect.Field; import java.nio.ByteBuffer; import java.util.Date; import java.util.UUID; @@ -40,6 +41,9 @@ import org.springframework.core.ParameterizedTypeReference; import org.springframework.core.ResolvableType; import org.springframework.kafka.support.serializer.JacksonJsonSerde; +import org.springframework.kafka.support.serializer.JacksonJsonSerializer; + +import tools.jackson.databind.json.JsonMapper; import static org.assertj.core.api.Assertions.assertThat; import static org.junit.jupiter.params.provider.Arguments.arguments; @@ -49,6 +53,7 @@ * Unit tests for {@link SerdeResolverUtils}. * * @author Chris Bono + * @author adityaanikam */ @SuppressWarnings({ "rawtypes", "unchecked" }) class SerdeResolverUtilsTests { @@ -157,6 +162,38 @@ void returnsFallbackSerdeForJavaLangObject() { .isNull()); } } + + @Nested + class WithJsonMapperBean { + + @Test + void usesConfiguredJsonMapperBeanWhenPresent() throws Exception { + new ApplicationContextRunner() + .withPropertyValues("spring.cloud.function.ineligible-definitions: sendToDlqAndContinue") + .withUserConfiguration(SerdeResolverJsonMapperTestApp.class) + .run((context) -> { + Serde serde = SerdeResolverUtils.resolveForType(context, ResolvableType.forClass(Foo.class), null); + assertThat(serde).isInstanceOf(JacksonJsonSerde.class); + + JsonMapper expectedMapper = context.getBean(JsonMapper.class); + JacksonJsonSerializer serializer = (JacksonJsonSerializer) serde.serializer(); + Field mapperField = JacksonJsonSerializer.class.getDeclaredField("jsonMapper"); + mapperField.setAccessible(true); + assertThat(mapperField.get(serializer)).isSameAs(expectedMapper); + }); + } + + @EnableAutoConfiguration + static class SerdeResolverJsonMapperTestApp { + + @Bean + public JsonMapper customJsonMapper() { + return JsonMapper.builder().build(); + } + + } + + } } }