diff --git a/README.adoc b/README.adoc index c236be76fc..8b55dc2868 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 1b5a8e0029..305db71386 100644 --- a/docs/src/main/asciidoc/_project-features-messaging.adoc +++ b/docs/src/main/asciidoc/_project-features-messaging.adoc @@ -1222,18 +1222,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 42f137d42d..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,7 +66,7 @@ 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) { String destination = dsl.getOutputMessage().getSentTo().getClientValue(); @@ -89,14 +86,6 @@ 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) { 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 2cfbfee8f2..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 @@ -45,7 +45,15 @@ class ContractVerifierKafkaStubMessagesInitializer implements KafkaStubMessagesI 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"); + + // 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 2a0784cb30..489cdad2be 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 @@ -34,6 +34,8 @@ import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.cloud.contract.verifier.converter.YamlContract; 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; @@ -94,10 +96,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,22 +111,44 @@ class Receiver { if (consumer == null) { 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) { @@ -133,28 +159,6 @@ 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(); if (textPayload.startsWith("\"") && textPayload.endsWith("\"")) { 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