diff --git a/spring-cloud-stream-binder-kafka-streams/pom.xml b/spring-cloud-stream-binder-kafka-streams/pom.xml index a75033024..feb13e631 100644 --- a/spring-cloud-stream-binder-kafka-streams/pom.xml +++ b/spring-cloud-stream-binder-kafka-streams/pom.xml @@ -73,12 +73,7 @@ kafka_2.13 test - - - org.springframework.cloud - spring-cloud-schema-registry-client - test - + org.apache.avro avro diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/serde/MessageConverterDelegateSerde.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/serde/MessageConverterDelegateSerde.java index 1f6b5a5d7..a64e28695 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/serde/MessageConverterDelegateSerde.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/serde/MessageConverterDelegateSerde.java @@ -75,7 +75,9 @@ import org.springframework.util.MimeTypeUtils; * @param type of the object to marshall * @author Soby Chacko * @since 3.0 + * @deprecated in favor of other schema registry providers instead of Spring Cloud Schema Registry. See its motivation above. */ +@Deprecated public class MessageConverterDelegateSerde implements Serde { private static final String VALUE_CLASS_HEADER = "valueClass"; diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/PerRecordAvroContentTypeTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/PerRecordAvroContentTypeTests.java index 00c745a3b..59eee2591 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/PerRecordAvroContentTypeTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/PerRecordAvroContentTypeTests.java @@ -37,8 +37,8 @@ import org.junit.Test; import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.cloud.schema.registry.avro.AvroSchemaMessageConverter; -import org.springframework.cloud.schema.registry.avro.AvroSchemaServiceManagerImpl; +import org.springframework.cloud.function.context.converter.avro.AvroSchemaMessageConverter; +import org.springframework.cloud.function.context.converter.avro.AvroSchemaServiceManagerImpl; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.annotation.Input; import org.springframework.cloud.stream.annotation.StreamListener; @@ -146,7 +146,7 @@ public class PerRecordAvroContentTypeTests { // Convert the byte[] received back to avro object and verify that it is // the same as the one we sent ^^. - AvroSchemaMessageConverter avroSchemaMessageConverter = new AvroSchemaMessageConverter(); + AvroSchemaMessageConverter avroSchemaMessageConverter = new AvroSchemaMessageConverter(new AvroSchemaServiceManagerImpl()); Message receivedMessage = MessageBuilder.withPayload(value) .setHeader("contentType", diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/utils/TestAvroSerializer.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/utils/TestAvroSerializer.java index 6bbf3180e..761636ea9 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/utils/TestAvroSerializer.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/utils/TestAvroSerializer.java @@ -21,8 +21,8 @@ import java.util.Map; import org.apache.kafka.common.serialization.Serializer; -import org.springframework.cloud.schema.registry.avro.AvroSchemaMessageConverter; -import org.springframework.cloud.schema.registry.avro.AvroSchemaServiceManagerImpl; +import org.springframework.cloud.function.context.converter.avro.AvroSchemaMessageConverter; +import org.springframework.cloud.function.context.converter.avro.AvroSchemaServiceManagerImpl; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.support.MessageBuilder; diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/serde/MessageConverterDelegateSerdeTest.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/serde/MessageConverterDelegateSerdeTest.java deleted file mode 100644 index 82c68ff0b..000000000 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/serde/MessageConverterDelegateSerdeTest.java +++ /dev/null @@ -1,74 +0,0 @@ -/* - * Copyright 2018-2019 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.streams.serde; - -import java.util.ArrayList; -import java.util.HashMap; -import java.util.List; -import java.util.Map; -import java.util.Random; -import java.util.UUID; - -import com.example.Sensor; -import com.fasterxml.jackson.databind.ObjectMapper; -import org.junit.Test; - -import org.springframework.cloud.schema.registry.avro.AvroSchemaMessageConverter; -import org.springframework.cloud.schema.registry.avro.AvroSchemaServiceManagerImpl; -import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; -import org.springframework.messaging.converter.MessageConverter; - -import static org.assertj.core.api.Assertions.assertThat; - -/** - * Refer {@link MessageConverterDelegateSerde} for motivations. - * - * @author Soby Chacko - */ -public class MessageConverterDelegateSerdeTest { - - @Test - @SuppressWarnings("unchecked") - public void testCompositeNonNativeSerdeUsingAvroContentType() { - Random random = new Random(); - Sensor sensor = new Sensor(); - sensor.setId(UUID.randomUUID().toString() + "-v1"); - sensor.setAcceleration(random.nextFloat() * 10); - sensor.setVelocity(random.nextFloat() * 100); - sensor.setTemperature(random.nextFloat() * 50); - - List messageConverters = new ArrayList<>(); - messageConverters.add(new AvroSchemaMessageConverter(new AvroSchemaServiceManagerImpl())); - CompositeMessageConverterFactory compositeMessageConverterFactory = new CompositeMessageConverterFactory( - messageConverters, new ObjectMapper()); - MessageConverterDelegateSerde messageConverterDelegateSerde = new MessageConverterDelegateSerde( - compositeMessageConverterFactory.getMessageConverterForAllRegistered()); - - Map configs = new HashMap<>(); - configs.put("valueClass", Sensor.class); - configs.put("contentType", "application/avro"); - messageConverterDelegateSerde.configure(configs, false); - final byte[] serialized = messageConverterDelegateSerde.serializer().serialize(null, - sensor); - - final Object deserialized = messageConverterDelegateSerde.deserializer() - .deserialize(null, serialized); - - assertThat(deserialized).isEqualTo(sensor); - } - -}