From 9a5c5a03f11868a655bc14bba6aea75aacaa56cf Mon Sep 17 00:00:00 2001 From: bono007 Date: Thu, 4 Feb 2021 06:44:10 -0600 Subject: [PATCH] Adds support for Kafka message keys in contracts (#1585) * Adds support for Kafka message keys in contracts * Updated doc to include note on Kafka key SerDes * Removes unintentional changes from previous commit. Fixes gh-1267 --- README.adoc | 1 + .../asciidoc/_project-features-messaging.adoc | 9 +- .../kafka/StubRunnerKafkaRouter.java | 37 ++------ ...tVerifierKafkaStubMessagesInitializer.java | 26 ++--- .../messaging/kafka/KafkaStubMessages.java | 94 +++++++++---------- .../kafka/KafkaStubRunnerSpec.groovy | 91 ++++++------------ .../src/test/resources/application.yml | 7 +- .../test/resources/stubs/bookReturned3.groovy | 22 +++++ 8 files changed, 133 insertions(+), 154 deletions(-) create mode 100644 tests/spring-cloud-contract-stub-runner-kafka/src/test/resources/stubs/bookReturned3.groovy diff --git a/README.adoc b/README.adoc index 4412bca9da..84644f2332 100644 --- a/README.adoc +++ b/README.adoc @@ -19,6 +19,7 @@ If you prefer to learn about the project by doing some tutorials, you can check workshops under https://cloud-samples.spring.io/spring-cloud-contract-samples/workshops.html[this link]. + == Project page You can read more about Spring Cloud Contract by going to https://spring.io/projects/spring-cloud-contract[the project page] diff --git a/docs/src/main/asciidoc/_project-features-messaging.adoc b/docs/src/main/asciidoc/_project-features-messaging.adoc index cadc70db7c..3231525e82 100644 --- a/docs/src/main/asciidoc/_project-features-messaging.adoc +++ b/docs/src/main/asciidoc/_project-features-messaging.adoc @@ -1201,18 +1201,21 @@ spring: kafka: bootstrap-servers: ${spring.embedded.kafka.brokers} producer: + value-serializer: org.springframework.kafka.support.serializer.JsonSerializer properties: - "value.serializer": "org.springframework.kafka.support.serializer.JsonSerializer" "spring.json.trusted.packages": "*" consumer: + value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer properties: - "value.deserializer": "org.springframework.kafka.support.serializer.JsonDeserializer" - "value.serializer": "org.springframework.kafka.support.serializer.JsonSerializer" "spring.json.trusted.packages": "*" group-id: groupId ---- ==== +NOTE: If your application uses non-integer record keys you will need to set the `spring.kafka.producer.key-serializer` +and `spring.kafka.consumer.key-deserializer` properties accordingly because the Kafka de/serialization expects non-null +record keys to be of integer type. + Now consider the following contracts (we number them 1 and 2): ==== diff --git a/spring-cloud-contract-stub-runner/src/main/java/org/springframework/cloud/contract/stubrunner/messaging/kafka/StubRunnerKafkaRouter.java b/spring-cloud-contract-stub-runner/src/main/java/org/springframework/cloud/contract/stubrunner/messaging/kafka/StubRunnerKafkaRouter.java index 2f972bfd2c..907085a5fe 100644 --- a/spring-cloud-contract-stub-runner/src/main/java/org/springframework/cloud/contract/stubrunner/messaging/kafka/StubRunnerKafkaRouter.java +++ b/spring-cloud-contract-stub-runner/src/main/java/org/springframework/cloud/contract/stubrunner/messaging/kafka/StubRunnerKafkaRouter.java @@ -16,25 +16,20 @@ package org.springframework.cloud.contract.stubrunner.messaging.kafka; -import java.util.HashMap; import java.util.List; -import java.util.Map; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; -import org.apache.kafka.common.header.Header; -import org.apache.kafka.common.header.Headers; import org.springframework.beans.factory.BeanFactory; import org.springframework.cloud.contract.spec.Contract; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.listener.MessageListener; import org.springframework.kafka.support.Acknowledgment; +import org.springframework.kafka.support.converter.MessagingMessageConverter; import org.springframework.messaging.Message; -import org.springframework.messaging.MessageHeaders; -import org.springframework.messaging.support.MessageBuilder; /** * @author Marcin Grzejszczak @@ -43,6 +38,8 @@ class StubRunnerKafkaRouter implements MessageListener { private static final Log log = LogFactory.getLog(StubRunnerKafkaRouter.class); + private final MessagingMessageConverter messageConverter = new MessagingMessageConverter(); + private final StubRunnerKafkaMessageSelector selector; private final BeanFactory beanFactory; @@ -69,19 +66,15 @@ class StubRunnerKafkaRouter implements MessageListener { if (log.isDebugEnabled()) { log.debug("Received message [" + data + "]"); } - Message message = MessageBuilder.createMessage(data.value(), - headers(data.headers())); + Message message = messageConverter.toMessage(data, null, null, null); Contract dsl = this.selector.matchingContract(message); - if (dsl != null && dsl.getOutputMessage() != null - && dsl.getOutputMessage().getSentTo() != null) { + if (dsl != null && dsl.getOutputMessage() != null && dsl.getOutputMessage().getSentTo() != null) { String destination = dsl.getOutputMessage().getSentTo().getClientValue(); if (log.isDebugEnabled()) { - log.debug( - "Found a matching contract with an output message. Will send it to the [" - + destination + "] destination"); + log.debug("Found a matching contract with an output message. Will send it to the [" + destination + + "] destination"); } - Message transform = new StubRunnerKafkaTransformer(this.contracts) - .transform(dsl); + Message transform = new StubRunnerKafkaTransformer(this.contracts).transform(dsl); String defaultTopic = kafkaTemplate().getDefaultTopic(); try { kafkaTemplate().setDefaultTopic(destination); @@ -93,17 +86,8 @@ class StubRunnerKafkaRouter implements MessageListener { } } - private MessageHeaders headers(Headers headers) { - Map map = new HashMap<>(); - for (Header header : headers) { - map.put(header.key(), header.value()); - } - return new MessageHeaders(map); - } - @Override - public void onMessage(ConsumerRecord data, - Acknowledgment acknowledgment) { + public void onMessage(ConsumerRecord data, Acknowledgment acknowledgment) { onMessage(data); } @@ -113,8 +97,7 @@ class StubRunnerKafkaRouter implements MessageListener { } @Override - public void onMessage(ConsumerRecord data, - Acknowledgment acknowledgment, Consumer consumer) { + public void onMessage(ConsumerRecord data, Acknowledgment acknowledgment, Consumer consumer) { onMessage(data); } diff --git a/spring-cloud-contract-verifier/src/main/java/org/springframework/cloud/contract/verifier/messaging/kafka/ContractVerifierKafkaStubMessagesInitializer.java b/spring-cloud-contract-verifier/src/main/java/org/springframework/cloud/contract/verifier/messaging/kafka/ContractVerifierKafkaStubMessagesInitializer.java index a240b8461c..8c9f2988ee 100644 --- a/spring-cloud-contract-verifier/src/main/java/org/springframework/cloud/contract/verifier/messaging/kafka/ContractVerifierKafkaStubMessagesInitializer.java +++ b/spring-cloud-contract-verifier/src/main/java/org/springframework/cloud/contract/verifier/messaging/kafka/ContractVerifierKafkaStubMessagesInitializer.java @@ -29,15 +29,12 @@ import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.test.EmbeddedKafkaBroker; import org.springframework.kafka.test.utils.KafkaTestUtils; -class ContractVerifierKafkaStubMessagesInitializer - implements KafkaStubMessagesInitializer { +class ContractVerifierKafkaStubMessagesInitializer implements KafkaStubMessagesInitializer { - private static final Log log = LogFactory - .getLog(ContractVerifierKafkaStubMessagesInitializer.class); + private static final Log log = LogFactory.getLog(ContractVerifierKafkaStubMessagesInitializer.class); @Override - public Map initialize(EmbeddedKafkaBroker broker, - KafkaProperties kafkaProperties) { + public Map initialize(EmbeddedKafkaBroker broker, KafkaProperties kafkaProperties) { Map map = new HashMap<>(); for (String topic : broker.getTopics()) { map.put(topic, prepareListener(broker, topic, kafkaProperties)); @@ -45,11 +42,18 @@ class ContractVerifierKafkaStubMessagesInitializer return map; } - private Consumer prepareListener(EmbeddedKafkaBroker broker, String destination, - KafkaProperties kafkaProperties) { - Map consumerProperties = KafkaTestUtils.consumerProps( - kafkaProperties.getConsumer().getGroupId(), "false", broker); - consumerProperties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + private Consumer prepareListener(EmbeddedKafkaBroker broker, String destination, KafkaProperties kafkaProperties) { + Map consumerProperties = KafkaTestUtils + .consumerProps(kafkaProperties.getConsumer().getGroupId(), "false", broker); + + // Respect custom key/value deserializers and any additional props under + // 'spring.kafka.consumer.properties' + consumerProperties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, + kafkaProperties.getConsumer().getKeyDeserializer()); + consumerProperties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, + kafkaProperties.getConsumer().getValueDeserializer()); + consumerProperties.putAll(kafkaProperties.getConsumer().getProperties()); + DefaultKafkaConsumerFactory consumerFactory = new DefaultKafkaConsumerFactory<>( consumerProperties); Consumer consumer = consumerFactory.createConsumer(); diff --git a/spring-cloud-contract-verifier/src/main/java/org/springframework/cloud/contract/verifier/messaging/kafka/KafkaStubMessages.java b/spring-cloud-contract-verifier/src/main/java/org/springframework/cloud/contract/verifier/messaging/kafka/KafkaStubMessages.java index e4cd390394..ec69ea4857 100644 --- a/spring-cloud-contract-verifier/src/main/java/org/springframework/cloud/contract/verifier/messaging/kafka/KafkaStubMessages.java +++ b/spring-cloud-contract-verifier/src/main/java/org/springframework/cloud/contract/verifier/messaging/kafka/KafkaStubMessages.java @@ -33,6 +33,8 @@ import org.apache.kafka.common.header.Headers; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.cloud.contract.verifier.messaging.MessageVerifier; import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.support.KafkaHeaders; +import org.springframework.kafka.support.converter.MessagingMessageConverter; import org.springframework.kafka.test.EmbeddedKafkaBroker; import org.springframework.kafka.test.utils.KafkaTestUtils; import org.springframework.messaging.Message; @@ -47,11 +49,10 @@ class KafkaStubMessages implements MessageVerifier> { private final Receiver receiver; - KafkaStubMessages(KafkaTemplate kafkaTemplate, EmbeddedKafkaBroker broker, - KafkaProperties kafkaProperties, KafkaStubMessagesInitializer initializer) { + KafkaStubMessages(KafkaTemplate kafkaTemplate, EmbeddedKafkaBroker broker, KafkaProperties kafkaProperties, + KafkaStubMessagesInitializer initializer) { this.kafkaTemplate = kafkaTemplate; - Map topicToConsumer = initializer.initialize(broker, - kafkaProperties); + Map topicToConsumer = initializer.initialize(broker, kafkaProperties); this.receiver = new Receiver(topicToConsumer); } @@ -61,8 +62,7 @@ class KafkaStubMessages implements MessageVerifier> { try { this.kafkaTemplate.setDefaultTopic(destination); if (log.isDebugEnabled()) { - log.debug("Will send a message [" + message + "] to destination [" - + destination + "]"); + log.debug("Will send a message [" + message + "] to destination [" + destination + "]"); } this.kafkaTemplate.send(message).get(5, TimeUnit.SECONDS); this.kafkaTemplate.flush(); @@ -87,8 +87,7 @@ class KafkaStubMessages implements MessageVerifier> { @Override public void send(Object payload, Map headers, String destination) { - Message message = MessageBuilder.createMessage(payload, - new MessageHeaders(headers)); + Message message = MessageBuilder.createMessage(payload, new MessageHeaders(headers)); send(message, destination); } @@ -96,10 +95,12 @@ class KafkaStubMessages implements MessageVerifier> { class Receiver { - private final Map consumers; - private static final Log log = LogFactory.getLog(Receiver.class); + private final MessagingMessageConverter messagingMessageConverter = new MessagingMessageConverter(); + + private final Map consumers; + Receiver(Map consumers) { this.consumers = consumers; } @@ -107,25 +108,46 @@ class Receiver { Message receive(String topic, long timeout, TimeUnit timeUnit) { Consumer consumer = this.consumers.get(topic); if (consumer == null) { - throw new IllegalStateException( - "No consumer set up for topic [" + topic + "]"); + throw new IllegalStateException("No consumer set up for topic [" + topic + "]"); } - ConsumerRecord record = KafkaTestUtils.getSingleRecord(consumer, - topic, timeUnit.toMillis(timeout)); + ConsumerRecord record = KafkaTestUtils.getSingleRecord(consumer, topic, timeUnit.toMillis(timeout)); if (log.isDebugEnabled()) { log.debug("Got a single record for destination [" + topic + "]"); } - return new Record(record).toMessage(); + return toMessage(consumer, record); } -} + Message toMessage(Consumer consumer, ConsumerRecord record) { + Map headersMap = toMap(record.headers()); -class Record { + // Leverage spring-kafka to add the headers + messagingMessageConverter.commonHeaders(null, consumer, headersMap, record.key(), record.topic(), + record.partition(), record.offset(), + record.timestampType() != null ? record.timestampType().name() : null, record.timestamp()); + // commonHeaders() maps the record key under 'kafka_receivedMessageKey' - put + // under 'kafka_messageKey' as well to satisfy both client/server usages as there + // is not currently a way to set a header name based on client/server + headersMap.put(KafkaHeaders.MESSAGE_KEY, record.key()); - private final ConsumerRecord record; - - Record(ConsumerRecord record) { - this.record = record; + // TODO explore using MessagingMessageConverter to do all of the conversion + // (ideally delete this entire method) + Object textPayload = record.value(); + // sometimes it's a message sometimes just payload + if (textPayload instanceof String && ((String) textPayload).contains("payload") + && ((String) textPayload).contains("headers")) { + try { + Object object = new JSONParser(JSONParser.DEFAULT_PERMISSIVE_MODE).parse((String) textPayload); + JSONObject jo = (JSONObject) object; + String payload = (String) jo.get("payload"); + JSONObject headersInJson = (JSONObject) jo.get("headers"); + headersMap.putAll(headersInJson); + return MessageBuilder.createMessage(unquoted(payload), new MessageHeaders(headersMap)); + } + catch (ParseException ex) { + throw new IllegalStateException(ex); + } + } + return MessageBuilder.createMessage(unquoted(textPayload), new MessageHeaders(headersMap)); } private Map toMap(Headers headers) { @@ -136,36 +158,10 @@ class Record { return map; } - Message toMessage() { - Object textPayload = record.value(); - // sometimes it's a message sometimes just payload - MessageHeaders headers = new MessageHeaders(toMap(record.headers())); - if (textPayload instanceof String && ((String) textPayload).contains("payload") - && ((String) textPayload).contains("headers")) { - try { - Object object = new JSONParser(JSONParser.DEFAULT_PERMISSIVE_MODE) - .parse((String) textPayload); - JSONObject jo = (JSONObject) object; - String payload = (String) jo.get("payload"); - JSONObject headersInJson = (JSONObject) jo.get("headers"); - Map newHeaders = new HashMap(headers); - newHeaders.putAll(headersInJson); - return MessageBuilder.createMessage(unquoted(payload), - new MessageHeaders(newHeaders)); - } - catch (ParseException ex) { - throw new IllegalStateException(ex); - } - } - return MessageBuilder.createMessage(unquoted(textPayload), headers); - } - private Object unquoted(Object value) { - String textPayload = value instanceof byte[] ? new String((byte[]) value) - : value.toString(); + String textPayload = value instanceof byte[] ? new String((byte[]) value) : value.toString(); if (textPayload.startsWith("\"") && textPayload.endsWith("\"")) { - return textPayload.substring(1, textPayload.length() - 1).replace("\\\"", - "\""); + return textPayload.substring(1, textPayload.length() - 1).replace("\\\"", "\""); } return textPayload; } diff --git a/tests/spring-cloud-contract-stub-runner-kafka/src/test/groovy/org/springframework/cloud/contract/stubrunner/messaging/kafka/KafkaStubRunnerSpec.groovy b/tests/spring-cloud-contract-stub-runner-kafka/src/test/groovy/org/springframework/cloud/contract/stubrunner/messaging/kafka/KafkaStubRunnerSpec.groovy index 8f8994311e..d6310784c3 100644 --- a/tests/spring-cloud-contract-stub-runner-kafka/src/test/groovy/org/springframework/cloud/contract/stubrunner/messaging/kafka/KafkaStubRunnerSpec.groovy +++ b/tests/spring-cloud-contract-stub-runner-kafka/src/test/groovy/org/springframework/cloud/contract/stubrunner/messaging/kafka/KafkaStubRunnerSpec.groovy @@ -54,7 +54,7 @@ import org.springframework.test.context.ContextConfiguration @SpringBootTest(properties = ["debug=true"]) @AutoConfigureStubRunner @IgnoreIf({ os.windows }) -@EmbeddedKafka(topics = ["input", "output", "delete"]) +@EmbeddedKafka(topics = ["input", "input2", "output", "delete"]) @Commons class KafkaStubRunnerSpec extends Specification { @@ -104,6 +104,30 @@ class KafkaStubRunnerSpec extends Specification { } } + def 'should propagate the Kafka record key via message headers'() { + expect: + await.eventually { + log.info("Sending the message") + // tag::client_send[] + Message message = MessageBuilder.createMessage(new BookReturned('bar'), new MessageHeaders([kafka_messageKey: "bar5150",])) + kafkaTemplate.setDefaultTopic('input2') + kafkaTemplate.send(message) + // end::client_send[] + log.info("Message sent") + log.info("Receiving the message") + // tag::client_receive[] + Message receivedMessage = receiveFromOutput() + // end::client_receive[] + log.info("Message received [" + receivedMessage + "]") + // tag::client_receive_message[] + assert receivedMessage != null + assert assertThatBodyContainsBookName(receivedMessage.getPayload(), 'bar') + assert receivedMessage.getHeaders().get('BOOK-NAME') == 'bar' + assert receivedMessage.getHeaders().get("kafka_receivedMessageKey") == 'bar5150' + // end::client_receive_message[] + } + } + def 'should trigger a message by label'() { expect: await.eventually { @@ -198,11 +222,15 @@ class KafkaStubRunnerSpec extends Specification { } private boolean assertThatBodyContainsBookNameFoo(Object payload) { + return assertThatBodyContainsBookName(payload, 'foo') + } + + private boolean assertThatBodyContainsBookName(Object payload, String expectedValue) { log.info("Got payload [" + payload + "]") String objectAsString = payload instanceof String ? payload : JsonOutput.toJson(payload) def json = new JsonSlurper().parseText(objectAsString) - return json.bookName == 'foo' + return json.bookName == expectedValue } @Configuration @@ -243,63 +271,4 @@ class KafkaStubRunnerSpec extends Specification { return this.output } } - - Contract dsl = - // tag::sample_dsl[] - Contract.make { - label 'return_book_1' - input { - triggeredBy('bookReturnedTriggered()') - } - outputMessage { - sentTo('output') - body('''{ "bookName" : "foo" }''') - headers { - header('BOOK-NAME', 'foo') - } - } - } - // end::sample_dsl[] - - Contract dsl2 = - // tag::sample_dsl_2[] - Contract.make { - label 'return_book_2' - input { - messageFrom('input') - messageBody([ - bookName: 'foo' - ]) - messageHeaders { - header('sample', 'header') - } - } - outputMessage { - sentTo('output') - body([ - bookName: 'foo' - ]) - headers { - header('BOOK-NAME', 'foo') - } - } - } - // end::sample_dsl_2[] - - Contract dsl3 = - // tag::sample_dsl_3[] - Contract.make { - label 'delete_book' - input { - messageFrom('delete') - messageBody([ - bookName: 'foo' - ]) - messageHeaders { - header('sample', 'header') - } - assertThat('bookWasDeleted()') - } - } - // end::sample_dsl_3[] } \ No newline at end of file diff --git a/tests/spring-cloud-contract-stub-runner-kafka/src/test/resources/application.yml b/tests/spring-cloud-contract-stub-runner-kafka/src/test/resources/application.yml index d1a88f5869..5f8b6211aa 100644 --- a/tests/spring-cloud-contract-stub-runner-kafka/src/test/resources/application.yml +++ b/tests/spring-cloud-contract-stub-runner-kafka/src/test/resources/application.yml @@ -6,13 +6,14 @@ spring: kafka: bootstrap-servers: ${spring.embedded.kafka.brokers} producer: + key-serializer: org.apache.kafka.common.serialization.StringSerializer + value-serializer: org.springframework.kafka.support.serializer.JsonSerializer properties: - "value.serializer": "org.springframework.kafka.support.serializer.JsonSerializer" "spring.json.trusted.packages": "*" consumer: + key-deserializer: org.apache.kafka.common.serialization.StringDeserializer + value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer properties: - "value.deserializer": "org.springframework.kafka.support.serializer.JsonDeserializer" - "value.serializer": "org.springframework.kafka.support.serializer.JsonSerializer" "spring.json.trusted.packages": "*" group-id: ${random.value} server: diff --git a/tests/spring-cloud-contract-stub-runner-kafka/src/test/resources/stubs/bookReturned3.groovy b/tests/spring-cloud-contract-stub-runner-kafka/src/test/resources/stubs/bookReturned3.groovy new file mode 100644 index 0000000000..1b27bd1a34 --- /dev/null +++ b/tests/spring-cloud-contract-stub-runner-kafka/src/test/resources/stubs/bookReturned3.groovy @@ -0,0 +1,22 @@ +org.springframework.cloud.contract.spec.Contract.make { + label 'return_book_3' + input { + messageFrom('input2') + messageBody([ + bookName: 'bar' + ]) + messageHeaders { + header('kafka_receivedMessageKey', 'bar5150') + } + } + outputMessage { + sentTo('output') + body([ + bookName: 'bar' + ]) + headers { + header('BOOK-NAME', 'bar') + header('kafka_messageKey', 'bar5150') + } + } +} \ No newline at end of file