Initial commit of KafkaNull changes to SmartCompositeMessageConverter

This commit is contained in:
Oleg Zhurakousky
2022-06-08 12:22:14 +02:00
parent 83901cf16d
commit 808fb28f5c

View File

@@ -19,6 +19,7 @@ package org.springframework.cloud.function.context.config;
import java.lang.reflect.Type; import java.lang.reflect.Type;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Collection; import java.util.Collection;
import java.util.Iterator;
import java.util.List; import java.util.List;
import org.apache.commons.logging.Log; import org.apache.commons.logging.Log;
@@ -34,7 +35,6 @@ import org.springframework.messaging.converter.MessageConverter;
import org.springframework.messaging.converter.SmartMessageConverter; import org.springframework.messaging.converter.SmartMessageConverter;
import org.springframework.messaging.support.MessageBuilder; import org.springframework.messaging.support.MessageBuilder;
import org.springframework.messaging.support.MessageHeaderAccessor; import org.springframework.messaging.support.MessageHeaderAccessor;
import org.springframework.util.CollectionUtils;
import org.springframework.util.MimeType; import org.springframework.util.MimeType;
import org.springframework.util.StringUtils; import org.springframework.util.StringUtils;
@@ -73,40 +73,41 @@ public class SmartCompositeMessageConverter extends CompositeMessageConverter {
return null; return null;
} }
@SuppressWarnings("unchecked")
@Override @Override
@Nullable
public Object fromMessage(Message<?> message, Class<?> targetClass, @Nullable Object conversionHint) { public Object fromMessage(Message<?> message, Class<?> targetClass, @Nullable Object conversionHint) {
for (MessageConverter converter : getConverters()) { if (!(message.getPayload() instanceof byte[]) && targetClass.isInstance(message.getPayload()) && !(message.getPayload() instanceof Collection<?>)) {
if (!(message.getPayload() instanceof byte[]) && targetClass.isInstance(message.getPayload()) && !(message.getPayload() instanceof Collection<?>)) { return message.getPayload();
return message.getPayload(); }
} Object result = null;
if (message.getPayload() instanceof Iterable && conversionHint != null) {
if (message.getPayload() instanceof Iterable && conversionHint != null) { Iterable<Object> iterablePayload = (Iterable<Object>) message.getPayload();
Iterable<Object> iterablePayload = (Iterable) message.getPayload(); Type genericItemType = FunctionTypeUtils.getImmediateGenericType((Type) conversionHint, 0);
Type t = FunctionTypeUtils.getImmediateGenericType((Type) conversionHint, 0); Class<?> genericItemRawType = FunctionTypeUtils.getRawType(genericItemType);
Class rawType = FunctionTypeUtils.getRawType(t); List<Object> resultList = new ArrayList<>();
List<Object> resultList = new ArrayList<>(); for (Object item : iterablePayload) {
for (Object item : iterablePayload) { boolean isConverted = false;
/* if (item.getClass().getName().startsWith("org.springframework.kafka.support.KafkaNull")) {
* Somewhere here we can do KafkaNull check or see below resultList.add(item);
*/ isConverted = true;
Message m = MessageBuilder.withPayload(item).copyHeaders(message.getHeaders()).build(); }
Object result = (converter instanceof SmartMessageConverter & rawType != t ? for (Iterator<MessageConverter> iterator = getConverters().iterator(); iterator.hasNext() && !isConverted;) {
((SmartMessageConverter) converter).fromMessage(m, rawType, t) : Message<?> m = MessageBuilder.withPayload(item).copyHeaders(message.getHeaders()).build(); // TODO Message creating may be expensive
converter.fromMessage(m, rawType)); MessageConverter converter = (MessageConverter) iterator.next();
if (result != null) { Object conversionResult = (converter instanceof SmartMessageConverter & genericItemRawType != genericItemType ?
/* ((SmartMessageConverter) converter).fromMessage(m, genericItemRawType, genericItemType) :
* Or most likely here we can do the KafkaNull check and not add it to the list converter.fromMessage(m, genericItemRawType));
*/ if (conversionResult != null) {
resultList.add(result); resultList.add(conversionResult);
isConverted = true;
} }
} }
if (!CollectionUtils.isEmpty(resultList)) {
return resultList;
}
} }
else { result = resultList;
Object result = (converter instanceof SmartMessageConverter ? }
else {
for (MessageConverter converter : getConverters()) {
result = (converter instanceof SmartMessageConverter ?
((SmartMessageConverter) converter).fromMessage(message, targetClass, conversionHint) : ((SmartMessageConverter) converter).fromMessage(message, targetClass, conversionHint) :
converter.fromMessage(message, targetClass)); converter.fromMessage(message, targetClass));
if (result != null) { if (result != null) {
@@ -114,7 +115,8 @@ public class SmartCompositeMessageConverter extends CompositeMessageConverter {
} }
} }
} }
return null;
return result;
} }
@Override @Override