Merge branch '2.2.x'
This commit is contained in:
@@ -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]
|
||||
|
||||
@@ -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):
|
||||
|
||||
====
|
||||
|
||||
@@ -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<Object, Object> {
|
||||
|
||||
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<Object, Object> {
|
||||
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<Object, Object> {
|
||||
}
|
||||
}
|
||||
|
||||
private MessageHeaders headers(Headers headers) {
|
||||
Map<String, Object> map = new HashMap<>();
|
||||
for (Header header : headers) {
|
||||
map.put(header.key(), header.value());
|
||||
}
|
||||
return new MessageHeaders(map);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onMessage(ConsumerRecord<Object, Object> data, Acknowledgment acknowledgment) {
|
||||
onMessage(data);
|
||||
|
||||
@@ -45,7 +45,15 @@ class ContractVerifierKafkaStubMessagesInitializer implements KafkaStubMessagesI
|
||||
private Consumer prepareListener(EmbeddedKafkaBroker broker, String destination, KafkaProperties kafkaProperties) {
|
||||
Map<String, Object> 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<String, String> consumerFactory = new DefaultKafkaConsumerFactory<>(
|
||||
consumerProperties);
|
||||
Consumer<String, String> consumer = consumerFactory.createConsumer();
|
||||
|
||||
@@ -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<Message<?>> {
|
||||
|
||||
class Receiver {
|
||||
|
||||
private final Map<String, Consumer> consumers;
|
||||
|
||||
private static final Log log = LogFactory.getLog(Receiver.class);
|
||||
|
||||
private final MessagingMessageConverter messagingMessageConverter = new MessagingMessageConverter();
|
||||
|
||||
private final Map<String, Consumer> consumers;
|
||||
|
||||
Receiver(Map<String, Consumer> consumers) {
|
||||
this.consumers = consumers;
|
||||
}
|
||||
@@ -107,22 +111,44 @@ class Receiver {
|
||||
if (consumer == null) {
|
||||
throw new IllegalStateException("No consumer set up for topic [" + topic + "]");
|
||||
}
|
||||
ConsumerRecord<String, String> 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<String, Object> 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<String, Object> 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("\"")) {
|
||||
|
||||
@@ -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[]
|
||||
}
|
||||
@@ -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:
|
||||
|
||||
@@ -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')
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user