Skip to content

Commit 8ae7ff9

Browse files
committed
Addressing PR review
Signed-off-by: Soby Chacko <[email protected]>
1 parent 69f4e3b commit 8ae7ff9

File tree

7 files changed

+192
-319
lines changed

7 files changed

+192
-319
lines changed

spring-kafka/src/main/java/org/springframework/kafka/config/AbstractShareKafkaListenerContainerFactory.java

Lines changed: 0 additions & 214 deletions
This file was deleted.

spring-kafka/src/main/java/org/springframework/kafka/config/MethodKafkaListenerEndpoint.java

Lines changed: 1 addition & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -31,13 +31,11 @@
3131
import org.springframework.expression.BeanResolver;
3232
import org.springframework.kafka.listener.KafkaListenerErrorHandler;
3333
import org.springframework.kafka.listener.MessageListenerContainer;
34-
import org.springframework.kafka.listener.ShareKafkaMessageListenerContainer;
3534
import org.springframework.kafka.listener.adapter.BatchMessagingMessageListenerAdapter;
3635
import org.springframework.kafka.listener.adapter.BatchToRecordAdapter;
3736
import org.springframework.kafka.listener.adapter.HandlerAdapter;
3837
import org.springframework.kafka.listener.adapter.MessagingMessageListenerAdapter;
3938
import org.springframework.kafka.listener.adapter.RecordMessagingMessageListenerAdapter;
40-
import org.springframework.kafka.listener.adapter.ShareRecordMessagingMessageListenerAdapter;
4139
import org.springframework.kafka.support.JavaUtils;
4240
import org.springframework.kafka.support.converter.BatchMessageConverter;
4341
import org.springframework.kafka.support.converter.MessageConverter;
@@ -179,12 +177,7 @@ protected MessagingMessageListenerAdapter<K, V> createMessageListener(MessageLis
179177
"Could not create message listener - MessageHandlerMethodFactory not set");
180178

181179
final MessagingMessageListenerAdapter<K, V> messageListener;
182-
if (container instanceof ShareKafkaMessageListenerContainer<?, ?>) {
183-
messageListener = createShareMessageListenerInstance(messageConverter);
184-
}
185-
else {
186-
messageListener = createMessageListenerInstance(messageConverter);
187-
}
180+
messageListener = createMessageListenerInstance(messageConverter);
188181
messageListener.setHandlerMethod(configureListenerAdapter(messageListener));
189182
JavaUtils.INSTANCE
190183
.acceptIfNotNull(getReplyTopic(), replyTopic -> {
@@ -210,31 +203,6 @@ protected HandlerAdapter configureListenerAdapter(MessagingMessageListenerAdapte
210203
return new HandlerAdapter(invocableHandlerMethod);
211204
}
212205

213-
/**
214-
* Create an empty {@link MessagingMessageListenerAdapter} instance.
215-
* @param messageConverter the converter (may be null).
216-
* @return the {@link MessagingMessageListenerAdapter} instance.
217-
*/
218-
protected MessagingMessageListenerAdapter<K, V> createShareMessageListenerInstance(
219-
@Nullable MessageConverter messageConverter) {
220-
221-
MessagingMessageListenerAdapter<K, V> listener;
222-
ShareRecordMessagingMessageListenerAdapter<K, V> messageListener = new ShareRecordMessagingMessageListenerAdapter<>(
223-
this.bean, this.method, this.errorHandler);
224-
if (messageConverter instanceof RecordMessageConverter recordMessageConverter) {
225-
messageListener.setMessageConverter(recordMessageConverter);
226-
}
227-
listener = messageListener;
228-
if (this.messagingConverter != null) {
229-
listener.setMessagingConverter(this.messagingConverter);
230-
}
231-
BeanResolver resolver = getBeanResolver();
232-
if (resolver != null) {
233-
listener.setBeanResolver(resolver);
234-
}
235-
return listener;
236-
}
237-
238206
protected MessagingMessageListenerAdapter<K, V> createMessageListenerInstance(
239207
@Nullable MessageConverter messageConverter) {
240208

0 commit comments

Comments
 (0)