diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaInboundChannelAdapterSpec.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaInboundChannelAdapterSpec.java index 5bac3b660f..155a832d37 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaInboundChannelAdapterSpec.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaInboundChannelAdapterSpec.java @@ -16,8 +16,6 @@ package org.springframework.integration.kafka.dsl; -import java.lang.reflect.Type; - import org.apache.kafka.clients.consumer.ConsumerRebalanceListener; import org.springframework.integration.dsl.MessageSourceSpec; @@ -70,7 +68,7 @@ public class KafkaInboundChannelAdapterSpec return this; } - public KafkaInboundChannelAdapterSpec payloadType(Type type) { + public KafkaInboundChannelAdapterSpec payloadType(Class type) { this.target.setPayloadType(type); return this; } diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageSource.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageSource.java index 85f1a2ee3f..5bbc521bf7 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageSource.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageSource.java @@ -16,7 +16,6 @@ package org.springframework.integration.kafka.inbound; -import java.lang.reflect.Type; import java.time.Duration; import java.util.ArrayList; import java.util.Arrays; @@ -113,7 +112,7 @@ public class KafkaMessageSource extends AbstractMessageSource impl private RecordMessageConverter messageConverter = new MessagingMessageConverter(); - private Type payloadType; + private Class payloadType; private ConsumerRebalanceListener rebalanceListener; @@ -202,7 +201,7 @@ public class KafkaMessageSource extends AbstractMessageSource impl this.messageConverter = messageConverter; } - protected Type getPayloadType() { + protected Class getPayloadType() { return this.payloadType; } @@ -211,7 +210,7 @@ public class KafkaMessageSource extends AbstractMessageSource impl * Only applies if a type-aware message converter is provided. * @param payloadType the type to convert to. */ - public void setPayloadType(Type payloadType) { + public void setPayloadType(Class payloadType) { this.payloadType = payloadType; } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaInboundChannelAdapterParserTests-context.xml b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaInboundChannelAdapterParserTests-context.xml index ce8c20f00f..72238ed63a 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaInboundChannelAdapterParserTests-context.xml +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaInboundChannelAdapterParserTests-context.xml @@ -16,7 +16,7 @@ client-id="client" group-id="group" message-converter="converter" - payload-type="#{T(java.lang.String)}" + payload-type="java.lang.String" raw-header="true" auto-startup="false" rebalance-listener="rebal">