From 24b52809edd7dd41488ed427b6580f707ce16485 Mon Sep 17 00:00:00 2001 From: Walliee Date: Wed, 1 May 2019 23:12:25 -0400 Subject: [PATCH] Switch back to use DefaultKafkaHeaderMapper * update spring-kafka to v2.2.5 * use DefaultKafkaHeaderMapper instead of BinderHeaderMapper * delete BinderHeaderMapper resolves #652 --- pom.xml | 2 +- .../binder/kafka/BinderHeaderMapper.java | 387 ------------------ .../kafka/KafkaMessageChannelBinder.java | 7 +- 3 files changed, 5 insertions(+), 391 deletions(-) delete mode 100644 spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/BinderHeaderMapper.java diff --git a/pom.xml b/pom.xml index f46181f16..80be01920 100644 --- a/pom.xml +++ b/pom.xml @@ -12,7 +12,7 @@ 1.8 - 2.2.2.RELEASE + 2.2.5.RELEASE 3.1.0.RELEASE 2.0.0 3.0.0.BUILD-SNAPSHOT diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/BinderHeaderMapper.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/BinderHeaderMapper.java deleted file mode 100644 index 3b4a5d748..000000000 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/BinderHeaderMapper.java +++ /dev/null @@ -1,387 +0,0 @@ -/* - * Copyright 2017-2018 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.springframework.cloud.stream.binder.kafka; - -import java.io.IOException; -import java.nio.charset.StandardCharsets; -import java.util.Arrays; -import java.util.HashMap; -import java.util.Iterator; -import java.util.LinkedHashSet; -import java.util.List; -import java.util.Map; -import java.util.Set; - -import com.fasterxml.jackson.core.JsonProcessingException; -import com.fasterxml.jackson.databind.DeserializationContext; -import com.fasterxml.jackson.databind.JsonNode; -import com.fasterxml.jackson.databind.ObjectMapper; -import com.fasterxml.jackson.databind.deser.std.StdNodeBasedDeserializer; -import com.fasterxml.jackson.databind.module.SimpleModule; -import com.fasterxml.jackson.databind.type.TypeFactory; -import org.apache.kafka.common.header.Header; -import org.apache.kafka.common.header.Headers; -import org.apache.kafka.common.header.internals.RecordHeader; - -import org.springframework.kafka.support.AbstractKafkaHeaderMapper; -import org.springframework.lang.Nullable; -import org.springframework.messaging.MessageHeaders; -import org.springframework.util.Assert; -import org.springframework.util.ClassUtils; -import org.springframework.util.MimeType; - -/** - * Default header mapper for Apache Kafka. Most headers in {@link KafkaHeaders} are not - * mapped on outbound messages. The exceptions are correlation and reply headers for - * request/reply messaging. Header types are added to a special header - * {@link #JSON_TYPES}. - * - * @author Gary Russell - * @author Artem Bilan - * @since 2.0 - * @deprecated will be removed in the next point release after 2.1.0. See issue - * https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/509 - */ -public class BinderHeaderMapper extends AbstractKafkaHeaderMapper { - - private static final List DEFAULT_TRUSTED_PACKAGES = Arrays - .asList("java.util", "java.lang", "org.springframework.util"); - - private static final List DEFAULT_TO_STRING_CLASSES = Arrays.asList( - "org.springframework.util.MimeType", "org.springframework.http.MediaType"); - - /** - * Header name for java types of other headers. - */ - public static final String JSON_TYPES = "spring_json_header_types"; - - private final ObjectMapper objectMapper; - - private final Set trustedPackages = new LinkedHashSet<>( - DEFAULT_TRUSTED_PACKAGES); - - private final Set toStringClasses = new LinkedHashSet<>( - DEFAULT_TO_STRING_CLASSES); - - /** - * Construct an instance with the default object mapper and default header patterns - * for outbound headers; all inbound headers are mapped. The default pattern list is - * {@code "!id", "!timestamp" and "*"}. In addition, most of the headers in - * {@link KafkaHeaders} are never mapped as headers since they represent data in - * consumer/producer records. - * @see #BinderHeaderMapper(ObjectMapper) - */ - public BinderHeaderMapper() { - this(new ObjectMapper()); - } - - /** - * Construct an instance with the provided object mapper and default header patterns - * for outbound headers; all inbound headers are mapped. The patterns are applied in - * order, stopping on the first match (positive or negative). Patterns are negated by - * preceding them with "!". The default pattern list is - * {@code "!id", "!timestamp" and "*"}. In addition, most of the headers in - * {@link KafkaHeaders} are never mapped as headers since they represent data in - * consumer/producer records. - * @param objectMapper the object mapper. - * @see org.springframework.util.PatternMatchUtils#simpleMatch(String, String) - */ - public BinderHeaderMapper(ObjectMapper objectMapper) { - this(objectMapper, "!" + MessageHeaders.ID, "!" + MessageHeaders.TIMESTAMP, "*"); - } - - /** - * Construct an instance with a default object mapper and the provided header patterns - * for outbound headers; all inbound headers are mapped. The patterns are applied in - * order, stopping on the first match (positive or negative). Patterns are negated by - * preceding them with "!". The patterns will replace the default patterns; you - * generally should not map the {@code "id" and "timestamp"} headers. Note: most of - * the headers in {@link KafkaHeaders} are ever mapped as headers since they represent - * data in consumer/producer records. - * @param patterns the patterns. - * @see org.springframework.util.PatternMatchUtils#simpleMatch(String, String) - */ - public BinderHeaderMapper(String... patterns) { - this(new ObjectMapper(), patterns); - } - - /** - * Construct an instance with the provided object mapper and the provided header - * patterns for outbound headers; all inbound headers are mapped. The patterns are - * applied in order, stopping on the first match (positive or negative). Patterns are - * negated by preceding them with "!". The patterns will replace the default patterns; - * you generally should not map the {@code "id" and "timestamp"} headers. Note: most - * of the headers in {@link KafkaHeaders} are never mapped as headers since they - * represent data in consumer/producer records. - * @param objectMapper the object mapper. - * @param patterns the patterns. - * @see org.springframework.util.PatternMatchUtils#simpleMatch(String, String) - */ - public BinderHeaderMapper(ObjectMapper objectMapper, String... patterns) { - super(patterns); - Assert.notNull(objectMapper, "'objectMapper' must not be null"); - Assert.noNullElements(patterns, "'patterns' must not have null elements"); - this.objectMapper = objectMapper; - this.objectMapper.registerModule(new SimpleModule() - .addDeserializer(MimeType.class, new MimeTypeJsonDeserializer())); - } - - /** - * Return the object mapper. - * @return the mapper. - */ - protected ObjectMapper getObjectMapper() { - return this.objectMapper; - } - - /** - * Provide direct access to the trusted packages set for subclasses. - * @return the trusted packages. - * @since 2.2 - */ - protected Set getTrustedPackages() { - return this.trustedPackages; - } - - /** - * Provide direct access to the toString() classes by subclasses. - * @return the toString() classes. - * @since 2.2 - */ - protected Set getToStringClasses() { - return this.toStringClasses; - } - - /** - * Add packages to the trusted packages list (default {@code java.util, java.lang}) - * used when constructing objects from JSON. If any of the supplied packages is - * {@code "*"}, all packages are trusted. If a class for a non-trusted package is - * encountered, the header is returned to the application with value of type - * {@link NonTrustedHeaderType}. - * @param trustedPackages the packages to trust. - */ - public void addTrustedPackages(String... trustedPackages) { - if (trustedPackages != null) { - for (String whiteList : trustedPackages) { - if ("*".equals(whiteList)) { - this.trustedPackages.clear(); - break; - } - else { - this.trustedPackages.add(whiteList); - } - } - } - } - - /** - * Add class names that the outbound mapper should perform toString() operations on - * before mapping. - * @param classNames the class names. - * @since 2.2 - */ - public void addToStringClasses(String... classNames) { - this.toStringClasses.addAll(Arrays.asList(classNames)); - } - - @Override - public void fromHeaders(MessageHeaders headers, Headers target) { - final Map jsonHeaders = new HashMap<>(); - headers.forEach((k, v) -> { - if (matches(k, v)) { - if (v instanceof byte[]) { - target.add(new RecordHeader(k, (byte[]) v)); - } - else { - try { - Object value = v; - String className = v.getClass().getName(); - if (this.toStringClasses.contains(className)) { - value = v.toString(); - className = "java.lang.String"; - } - target.add(new RecordHeader(k, - getObjectMapper().writeValueAsBytes(value))); - jsonHeaders.put(k, className); - } - catch (Exception e) { - if (logger.isDebugEnabled()) { - logger.debug("Could not map " + k + " with type " - + v.getClass().getName()); - } - } - } - } - }); - if (jsonHeaders.size() > 0) { - try { - target.add(new RecordHeader(JSON_TYPES, - getObjectMapper().writeValueAsBytes(jsonHeaders))); - } - catch (IllegalStateException | JsonProcessingException e) { - logger.error("Could not add json types header", e); - } - } - } - - @Override - public void toHeaders(Headers source, final Map headers) { - final Map jsonTypes = decodeJsonTypes(source); - source.forEach(h -> { - if (!(h.key().equals(JSON_TYPES))) { - if (jsonTypes != null && jsonTypes.containsKey(h.key())) { - Class type = Object.class; - String requestedType = jsonTypes.get(h.key()); - boolean trusted = false; - try { - trusted = trusted(requestedType); - if (trusted) { - type = ClassUtils.forName(requestedType, null); - } - } - catch (Exception e) { - logger.error("Could not load class for header: " + h.key(), e); - } - if (trusted) { - try { - headers.put(h.key(), - getObjectMapper().readValue(h.value(), type)); - } - catch (IOException e) { - logger.error("Could not decode json type: " - + new String(h.value()) + " for key: " + h.key(), e); - headers.put(h.key(), h.value()); - } - } - else { - headers.put(h.key(), - new NonTrustedHeaderType(h.value(), requestedType)); - } - } - else { - headers.put(h.key(), h.value()); - } - } - }); - } - - @SuppressWarnings("unchecked") - @Nullable - private Map decodeJsonTypes(Headers source) { - Map types = null; - Iterator
iterator = source.iterator(); - while (iterator.hasNext()) { - Header next = iterator.next(); - if (next.key().equals(JSON_TYPES)) { - try { - types = getObjectMapper().readValue(next.value(), Map.class); - } - catch (IOException e) { - logger.error( - "Could not decode json types: " + new String(next.value()), - e); - } - break; - } - } - return types; - } - - protected boolean trusted(String requestedType) { - if (!this.trustedPackages.isEmpty()) { - int lastDot = requestedType.lastIndexOf('.'); - if (lastDot < 0) { - return false; - } - String packageName = requestedType.substring(0, lastDot); - for (String trustedPackage : this.trustedPackages) { - if (packageName.equals(trustedPackage) - || packageName.startsWith(trustedPackage + ".")) { - return true; - } - } - return false; - } - return true; - } - - /** - * The {@link StdNodeBasedDeserializer} extension for {@link MimeType} - * deserialization. It is presented here for backward compatibility when older - * producers send {@link MimeType} headers as serialization version. - */ - private class MimeTypeJsonDeserializer extends StdNodeBasedDeserializer { - - private static final long serialVersionUID = 1L; - - MimeTypeJsonDeserializer() { - super(MimeType.class); - } - - @Override - public MimeType convert(JsonNode root, DeserializationContext ctxt) - throws IOException { - JsonNode type = root.get("type"); - JsonNode subType = root.get("subtype"); - JsonNode parameters = root.get("parameters"); - Map params = BinderHeaderMapper.this.objectMapper - .readValue(parameters.traverse(), TypeFactory.defaultInstance() - .constructMapType(HashMap.class, String.class, String.class)); - return new MimeType(type.asText(), subType.asText(), params); - } - - } - - /** - * Represents a header that could not be decoded due to an untrusted type. - */ - public static class NonTrustedHeaderType { - - private final byte[] headerValue; - - private final String untrustedType; - - NonTrustedHeaderType(byte[] headerValue, String untrustedType) { // NOSONAR - this.headerValue = headerValue; // NOSONAR - this.untrustedType = untrustedType; - } - - public byte[] getHeaderValue() { - return this.headerValue; // NOSONAR - } - - public String getUntrustedType() { - return this.untrustedType; - } - - @Override - public String toString() { - try { - return "NonTrustedHeaderType [headerValue=" - + new String(this.headerValue, StandardCharsets.UTF_8) - + ", untrustedType=" + this.untrustedType + "]"; - } - catch (Exception e) { - return "NonTrustedHeaderType [headerValue=" - + Arrays.toString(this.headerValue) + ", untrustedType=" - + this.untrustedType + "]"; - } - } - - } - -} diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index 288907aab..b70966455 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -96,6 +96,7 @@ import org.springframework.kafka.listener.AbstractMessageListenerContainer; import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; import org.springframework.kafka.listener.ConsumerAwareRebalanceListener; import org.springframework.kafka.listener.ContainerProperties; +import org.springframework.kafka.support.DefaultKafkaHeaderMapper; import org.springframework.kafka.support.KafkaHeaderMapper; import org.springframework.kafka.support.KafkaHeaders; import org.springframework.kafka.support.ProducerListener; @@ -376,11 +377,11 @@ public class KafkaMessageChannelBinder extends if (!patterns.contains("!" + MessageHeaders.ID)) { patterns.add(0, "!" + MessageHeaders.ID); } - mapper = new BinderHeaderMapper( + mapper = new DefaultKafkaHeaderMapper( patterns.toArray(new String[patterns.size()])); } else { - mapper = new BinderHeaderMapper(); + mapper = new DefaultKafkaHeaderMapper(); } } handler.setHeaderMapper(mapper); @@ -825,7 +826,7 @@ public class KafkaMessageChannelBinder extends KafkaHeaderMapper.class); } if (mapper == null) { - BinderHeaderMapper headerMapper = new BinderHeaderMapper() { + DefaultKafkaHeaderMapper headerMapper = new DefaultKafkaHeaderMapper() { @Override public void toHeaders(Headers source, Map headers) {