|
|
|
|
@@ -34,7 +34,6 @@ import org.apache.kafka.common.serialization.LongDeserializer;
|
|
|
|
|
import org.apache.kafka.common.serialization.StringDeserializer;
|
|
|
|
|
import org.apache.kafka.common.serialization.StringSerializer;
|
|
|
|
|
import org.assertj.core.api.Condition;
|
|
|
|
|
import org.junit.Ignore;
|
|
|
|
|
import org.junit.Rule;
|
|
|
|
|
import org.junit.Test;
|
|
|
|
|
import org.junit.rules.ExpectedException;
|
|
|
|
|
@@ -86,16 +85,16 @@ import static org.junit.Assert.assertTrue;
|
|
|
|
|
* @author Ilayaperumal Gopinathan
|
|
|
|
|
* @author Henryk Konsek
|
|
|
|
|
*/
|
|
|
|
|
public abstract class KafkaBinderTests extends PartitionCapableBinderTests<AbstractKafkaTestBinder, ExtendedConsumerProperties<KafkaConsumerProperties>,
|
|
|
|
|
ExtendedProducerProperties<KafkaProducerProperties>> {
|
|
|
|
|
public abstract class KafkaBinderTests extends
|
|
|
|
|
PartitionCapableBinderTests<AbstractKafkaTestBinder, ExtendedConsumerProperties<KafkaConsumerProperties>, ExtendedProducerProperties<KafkaProducerProperties>> {
|
|
|
|
|
|
|
|
|
|
@Rule
|
|
|
|
|
public ExpectedException expectedProvisioningException = ExpectedException.none();
|
|
|
|
|
|
|
|
|
|
@Override
|
|
|
|
|
protected ExtendedConsumerProperties<KafkaConsumerProperties> createConsumerProperties() {
|
|
|
|
|
final ExtendedConsumerProperties<KafkaConsumerProperties> kafkaConsumerProperties =
|
|
|
|
|
new ExtendedConsumerProperties<>(new KafkaConsumerProperties());
|
|
|
|
|
final ExtendedConsumerProperties<KafkaConsumerProperties> kafkaConsumerProperties = new ExtendedConsumerProperties<>(
|
|
|
|
|
new KafkaConsumerProperties());
|
|
|
|
|
// set the default values that would normally be propagated by Spring Cloud Stream
|
|
|
|
|
kafkaConsumerProperties.setInstanceCount(1);
|
|
|
|
|
kafkaConsumerProperties.setInstanceIndex(0);
|
|
|
|
|
@@ -104,7 +103,8 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
|
|
|
|
|
@Override
|
|
|
|
|
protected ExtendedProducerProperties<KafkaProducerProperties> createProducerProperties() {
|
|
|
|
|
ExtendedProducerProperties<KafkaProducerProperties> producerProperties = new ExtendedProducerProperties<>(new KafkaProducerProperties());
|
|
|
|
|
ExtendedProducerProperties<KafkaProducerProperties> producerProperties = new ExtendedProducerProperties<>(
|
|
|
|
|
new KafkaProducerProperties());
|
|
|
|
|
producerProperties.getExtension().setSync(true);
|
|
|
|
|
return producerProperties;
|
|
|
|
|
}
|
|
|
|
|
@@ -120,10 +120,10 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
protected abstract ZkUtils getZkUtils(KafkaBinderConfigurationProperties kafkaBinderConfigurationProperties);
|
|
|
|
|
|
|
|
|
|
protected abstract void invokeCreateTopic(ZkUtils zkUtils, String topic, int partitions,
|
|
|
|
|
int replicationFactor, Properties topicConfig);
|
|
|
|
|
int replicationFactor, Properties topicConfig);
|
|
|
|
|
|
|
|
|
|
protected abstract int invokePartitionSize(String topic,
|
|
|
|
|
ZkUtils zkUtils);
|
|
|
|
|
ZkUtils zkUtils);
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
@SuppressWarnings("unchecked")
|
|
|
|
|
@@ -307,7 +307,8 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
ExtendedConsumerProperties<KafkaConsumerProperties> dlqConsumerProperties = createConsumerProperties();
|
|
|
|
|
dlqConsumerProperties.setMaxAttempts(1);
|
|
|
|
|
QueueChannel dlqChannel = new QueueChannel();
|
|
|
|
|
Binding<MessageChannel> dlqConsumerBinding = binder.bindConsumer(dlqName, null, dlqChannel, dlqConsumerProperties);
|
|
|
|
|
Binding<MessageChannel> dlqConsumerBinding = binder.bindConsumer(dlqName, null, dlqChannel,
|
|
|
|
|
dlqConsumerProperties);
|
|
|
|
|
|
|
|
|
|
String testMessagePayload = "test." + UUID.randomUUID().toString();
|
|
|
|
|
Message<String> testMessage = MessageBuilder.withPayload(testMessagePayload).build();
|
|
|
|
|
@@ -369,9 +370,9 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
@Test
|
|
|
|
|
@SuppressWarnings("unchecked")
|
|
|
|
|
public void testCompression() throws Exception {
|
|
|
|
|
final KafkaProducerProperties.CompressionType[] codecs = new KafkaProducerProperties.CompressionType[]{
|
|
|
|
|
final KafkaProducerProperties.CompressionType[] codecs = new KafkaProducerProperties.CompressionType[] {
|
|
|
|
|
KafkaProducerProperties.CompressionType.none, KafkaProducerProperties.CompressionType.gzip,
|
|
|
|
|
KafkaProducerProperties.CompressionType.snappy};
|
|
|
|
|
KafkaProducerProperties.CompressionType.snappy };
|
|
|
|
|
byte[] testPayload = new byte[2048];
|
|
|
|
|
Arrays.fill(testPayload, (byte) 65);
|
|
|
|
|
Binder binder = getBinder();
|
|
|
|
|
@@ -541,8 +542,8 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
Binder binder = getBinder(createConfigurationProperties());
|
|
|
|
|
GenericApplicationContext context = new GenericApplicationContext();
|
|
|
|
|
context.refresh();
|
|
|
|
|
//binder.setApplicationContext(context);
|
|
|
|
|
//binder.afterPropertiesSet();
|
|
|
|
|
// binder.setApplicationContext(context);
|
|
|
|
|
// binder.afterPropertiesSet();
|
|
|
|
|
DirectChannel output = new DirectChannel();
|
|
|
|
|
QueueChannel input1 = new QueueChannel();
|
|
|
|
|
|
|
|
|
|
@@ -608,55 +609,6 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
@Ignore("Needs further discussion")
|
|
|
|
|
@SuppressWarnings("unchecked")
|
|
|
|
|
public void testReset() throws Exception {
|
|
|
|
|
Binder binder = getBinder();
|
|
|
|
|
DirectChannel output = new DirectChannel();
|
|
|
|
|
QueueChannel input1 = new QueueChannel();
|
|
|
|
|
|
|
|
|
|
String testTopicName = UUID.randomUUID().toString();
|
|
|
|
|
|
|
|
|
|
Binding<MessageChannel> producerBinding = binder.bindProducer(testTopicName, output,
|
|
|
|
|
createProducerProperties());
|
|
|
|
|
String testPayload1 = "foo1-" + UUID.randomUUID().toString();
|
|
|
|
|
output.send(new GenericMessage<>(testPayload1.getBytes()));
|
|
|
|
|
ExtendedConsumerProperties<KafkaConsumerProperties> properties = createConsumerProperties();
|
|
|
|
|
properties.getExtension().setResetOffsets(true);
|
|
|
|
|
properties.getExtension().setStartOffset(KafkaConsumerProperties.StartOffset.earliest);
|
|
|
|
|
Binding<MessageChannel> consumerBinding = binder.bindConsumer(testTopicName, "startOffsets", input1,
|
|
|
|
|
properties);
|
|
|
|
|
Message<byte[]> receivedMessage1 = (Message<byte[]>) receive(input1);
|
|
|
|
|
assertThat(receivedMessage1).isNotNull();
|
|
|
|
|
assertThat(new String(receivedMessage1.getPayload())).isEqualTo(testPayload1);
|
|
|
|
|
String testPayload2 = "foo2-" + UUID.randomUUID().toString();
|
|
|
|
|
output.send(new GenericMessage<>(testPayload2.getBytes()));
|
|
|
|
|
Message<byte[]> receivedMessage2 = (Message<byte[]>) receive(input1);
|
|
|
|
|
assertThat(receivedMessage2).isNotNull();
|
|
|
|
|
assertThat(new String(receivedMessage2.getPayload())).isEqualTo(testPayload2);
|
|
|
|
|
consumerBinding.unbind();
|
|
|
|
|
|
|
|
|
|
String testPayload3 = "foo3-" + UUID.randomUUID().toString();
|
|
|
|
|
output.send(new GenericMessage<>(testPayload3.getBytes()));
|
|
|
|
|
|
|
|
|
|
ExtendedConsumerProperties<KafkaConsumerProperties> properties2 = createConsumerProperties();
|
|
|
|
|
properties2.getExtension().setResetOffsets(true);
|
|
|
|
|
properties2.getExtension().setStartOffset(KafkaConsumerProperties.StartOffset.earliest);
|
|
|
|
|
consumerBinding = binder.bindConsumer(testTopicName, "startOffsets", input1, properties2);
|
|
|
|
|
Message<byte[]> receivedMessage4 = (Message<byte[]>) receive(input1);
|
|
|
|
|
assertThat(receivedMessage4).isNotNull();
|
|
|
|
|
assertThat(new String(receivedMessage4.getPayload())).isEqualTo(testPayload1);
|
|
|
|
|
Message<byte[]> receivedMessage5 = (Message<byte[]>) receive(input1);
|
|
|
|
|
assertThat(receivedMessage5).isNotNull();
|
|
|
|
|
assertThat(new String(receivedMessage5.getPayload())).isEqualTo(testPayload2);
|
|
|
|
|
Message<byte[]> receivedMessage6 = (Message<byte[]>) receive(input1);
|
|
|
|
|
assertThat(receivedMessage6).isNotNull();
|
|
|
|
|
assertThat(new String(receivedMessage6.getPayload())).isEqualTo(testPayload3);
|
|
|
|
|
consumerBinding.unbind();
|
|
|
|
|
producerBinding.unbind();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Test
|
|
|
|
|
@SuppressWarnings("unchecked")
|
|
|
|
|
public void testResume() throws Exception {
|
|
|
|
|
@@ -742,7 +694,6 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
moduleOutputChannel1.send(message1);
|
|
|
|
|
moduleOutputChannel2.send(message2);
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
Message<?>[] messages = new Message[2];
|
|
|
|
|
messages[0] = receive(moduleInputChannel);
|
|
|
|
|
messages[1] = receive(moduleInputChannel);
|
|
|
|
|
@@ -768,13 +719,15 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
createProducerBindingProperties(createProducerProperties()));
|
|
|
|
|
QueueChannel moduleInputChannel = new QueueChannel();
|
|
|
|
|
|
|
|
|
|
Binding<MessageChannel> producerBinding = binder.bindProducer("testManualAckSucceedsWhenAutoCommitOffsetIsTurnedOff", moduleOutputChannel,
|
|
|
|
|
Binding<MessageChannel> producerBinding = binder.bindProducer(
|
|
|
|
|
"testManualAckSucceedsWhenAutoCommitOffsetIsTurnedOff", moduleOutputChannel,
|
|
|
|
|
createProducerProperties());
|
|
|
|
|
|
|
|
|
|
ExtendedConsumerProperties<KafkaConsumerProperties> consumerProperties = createConsumerProperties();
|
|
|
|
|
consumerProperties.getExtension().setAutoCommitOffset(false);
|
|
|
|
|
|
|
|
|
|
Binding<MessageChannel> consumerBinding = binder.bindConsumer("testManualAckSucceedsWhenAutoCommitOffsetIsTurnedOff", "test", moduleInputChannel,
|
|
|
|
|
Binding<MessageChannel> consumerBinding = binder.bindConsumer(
|
|
|
|
|
"testManualAckSucceedsWhenAutoCommitOffsetIsTurnedOff", "test", moduleInputChannel,
|
|
|
|
|
consumerProperties);
|
|
|
|
|
|
|
|
|
|
String testPayload1 = "foo" + UUID.randomUUID().toString();
|
|
|
|
|
@@ -788,7 +741,8 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
Message<?> receivedMessage = receive(moduleInputChannel);
|
|
|
|
|
assertThat(receivedMessage).isNotNull();
|
|
|
|
|
assertThat(receivedMessage.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT)).isNotNull();
|
|
|
|
|
Acknowledgment acknowledgment = receivedMessage.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT, Acknowledgment.class);
|
|
|
|
|
Acknowledgment acknowledgment = receivedMessage.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT,
|
|
|
|
|
Acknowledgment.class);
|
|
|
|
|
try {
|
|
|
|
|
acknowledgment.acknowledge();
|
|
|
|
|
}
|
|
|
|
|
@@ -810,12 +764,14 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
createProducerBindingProperties(createProducerProperties()));
|
|
|
|
|
QueueChannel moduleInputChannel = new QueueChannel();
|
|
|
|
|
|
|
|
|
|
Binding<MessageChannel> producerBinding = binder.bindProducer("testManualAckIsNotPossibleWhenAutoCommitOffsetIsEnabledOnTheBinder", moduleOutputChannel,
|
|
|
|
|
Binding<MessageChannel> producerBinding = binder.bindProducer(
|
|
|
|
|
"testManualAckIsNotPossibleWhenAutoCommitOffsetIsEnabledOnTheBinder", moduleOutputChannel,
|
|
|
|
|
createProducerProperties());
|
|
|
|
|
|
|
|
|
|
ExtendedConsumerProperties<KafkaConsumerProperties> consumerProperties = createConsumerProperties();
|
|
|
|
|
|
|
|
|
|
Binding<MessageChannel> consumerBinding = binder.bindConsumer("testManualAckIsNotPossibleWhenAutoCommitOffsetIsEnabledOnTheBinder", "test", moduleInputChannel,
|
|
|
|
|
Binding<MessageChannel> consumerBinding = binder.bindConsumer(
|
|
|
|
|
"testManualAckIsNotPossibleWhenAutoCommitOffsetIsEnabledOnTheBinder", "test", moduleInputChannel,
|
|
|
|
|
consumerProperties);
|
|
|
|
|
|
|
|
|
|
String testPayload1 = "foo" + UUID.randomUUID().toString();
|
|
|
|
|
@@ -995,8 +951,8 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
Binding<MessageChannel> outputBinding = binder.bindProducer("partJ.0", output, producerProperties);
|
|
|
|
|
if (usesExplicitRouting()) {
|
|
|
|
|
Object endpoint = extractEndpoint(outputBinding);
|
|
|
|
|
assertThat(getEndpointRouting(endpoint)).
|
|
|
|
|
contains(getExpectedRoutingBaseDestination("partJ.0", "test") + "-' + headers['partition']");
|
|
|
|
|
assertThat(getEndpointRouting(endpoint))
|
|
|
|
|
.contains(getExpectedRoutingBaseDestination("partJ.0", "test") + "-' + headers['partition']");
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
output.send(new GenericMessage<>(2));
|
|
|
|
|
@@ -1044,7 +1000,7 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
QueueChannel input2 = new QueueChannel();
|
|
|
|
|
Binding<MessageChannel> binding2 = binder.bindConsumer("defaultGroup.0", null, input2,
|
|
|
|
|
consumerProperties);
|
|
|
|
|
//Since we don't provide any topic info, let Kafka bind the consumer successfully
|
|
|
|
|
// Since we don't provide any topic info, let Kafka bind the consumer successfully
|
|
|
|
|
Thread.sleep(1000);
|
|
|
|
|
String testPayload1 = "foo-" + UUID.randomUUID().toString();
|
|
|
|
|
output.send(new GenericMessage<>(testPayload1.getBytes()));
|
|
|
|
|
@@ -1063,7 +1019,7 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
output.send(new GenericMessage<>(testPayload2.getBytes()));
|
|
|
|
|
|
|
|
|
|
binding2 = binder.bindConsumer("defaultGroup.0", null, input2, consumerProperties);
|
|
|
|
|
//Since we don't provide any topic info, let Kafka bind the consumer successfully
|
|
|
|
|
// Since we don't provide any topic info, let Kafka bind the consumer successfully
|
|
|
|
|
Thread.sleep(1000);
|
|
|
|
|
String testPayload3 = "foo-" + UUID.randomUUID().toString();
|
|
|
|
|
output.send(new GenericMessage<>(testPayload3.getBytes()));
|
|
|
|
|
@@ -1226,7 +1182,8 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
consumerProperties.setInstanceIndex(2);
|
|
|
|
|
consumerProperties.getExtension().setAutoRebalanceEnabled(false);
|
|
|
|
|
expectedProvisioningException.expect(ProvisioningException.class);
|
|
|
|
|
expectedProvisioningException.expectMessage("The number of expected partitions was: 3, but 1 has been found instead");
|
|
|
|
|
expectedProvisioningException
|
|
|
|
|
.expectMessage("The number of expected partitions was: 3, but 1 has been found instead");
|
|
|
|
|
Binding binding = binder.bindConsumer(testTopicName, "test", output, consumerProperties);
|
|
|
|
|
if (binding != null) {
|
|
|
|
|
binding.unbind();
|
|
|
|
|
@@ -1355,8 +1312,10 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
ExtendedConsumerProperties<KafkaConsumerProperties> consumerProperties = createConsumerProperties();
|
|
|
|
|
binding = binder.bindConsumer(testTopicName, "test", output, consumerProperties);
|
|
|
|
|
DirectFieldAccessor consumerAccessor = new DirectFieldAccessor(getKafkaConsumer(binding));
|
|
|
|
|
assertTrue("Expected StringDeserializer as a custom key deserializer", consumerAccessor.getPropertyValue("keyDeserializer") instanceof StringDeserializer);
|
|
|
|
|
assertTrue("Expected LongDeserializer as a custom value deserializer", consumerAccessor.getPropertyValue("valueDeserializer") instanceof LongDeserializer);
|
|
|
|
|
assertTrue("Expected StringDeserializer as a custom key deserializer",
|
|
|
|
|
consumerAccessor.getPropertyValue("keyDeserializer") instanceof StringDeserializer);
|
|
|
|
|
assertTrue("Expected LongDeserializer as a custom value deserializer",
|
|
|
|
|
consumerAccessor.getPropertyValue("valueDeserializer") instanceof LongDeserializer);
|
|
|
|
|
}
|
|
|
|
|
finally {
|
|
|
|
|
if (binding != null) {
|
|
|
|
|
@@ -1367,11 +1326,15 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
|
|
|
|
|
private KafkaConsumer getKafkaConsumer(Binding binding) {
|
|
|
|
|
DirectFieldAccessor bindingAccessor = new DirectFieldAccessor((DefaultBinding) binding);
|
|
|
|
|
KafkaMessageDrivenChannelAdapter adapter = (KafkaMessageDrivenChannelAdapter) bindingAccessor.getPropertyValue("lifecycle");
|
|
|
|
|
KafkaMessageDrivenChannelAdapter adapter = (KafkaMessageDrivenChannelAdapter) bindingAccessor
|
|
|
|
|
.getPropertyValue("lifecycle");
|
|
|
|
|
DirectFieldAccessor adapterAccessor = new DirectFieldAccessor(adapter);
|
|
|
|
|
ConcurrentMessageListenerContainer messageListenerContainer = (ConcurrentMessageListenerContainer) adapterAccessor.getPropertyValue("messageListenerContainer");
|
|
|
|
|
DirectFieldAccessor containerAccessor = new DirectFieldAccessor((ConcurrentMessageListenerContainer) messageListenerContainer);
|
|
|
|
|
DefaultKafkaConsumerFactory consumerFactory = (DefaultKafkaConsumerFactory) containerAccessor.getPropertyValue("consumerFactory");
|
|
|
|
|
ConcurrentMessageListenerContainer messageListenerContainer = (ConcurrentMessageListenerContainer) adapterAccessor
|
|
|
|
|
.getPropertyValue("messageListenerContainer");
|
|
|
|
|
DirectFieldAccessor containerAccessor = new DirectFieldAccessor(
|
|
|
|
|
(ConcurrentMessageListenerContainer) messageListenerContainer);
|
|
|
|
|
DefaultKafkaConsumerFactory consumerFactory = (DefaultKafkaConsumerFactory) containerAccessor
|
|
|
|
|
.getPropertyValue("consumerFactory");
|
|
|
|
|
return (KafkaConsumer) consumerFactory.createConsumer();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@@ -1397,11 +1360,13 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
QueueChannel moduleInputChannel = new QueueChannel();
|
|
|
|
|
ExtendedProducerProperties<KafkaProducerProperties> producerProperties = createProducerProperties();
|
|
|
|
|
producerProperties.setUseNativeEncoding(true);
|
|
|
|
|
producerProperties.getExtension().getConfiguration().put("value.serializer", "org.apache.kafka.common.serialization.IntegerSerializer");
|
|
|
|
|
producerProperties.getExtension().getConfiguration().put("value.serializer",
|
|
|
|
|
"org.apache.kafka.common.serialization.IntegerSerializer");
|
|
|
|
|
producerBinding = binder.bindProducer(testTopicName, moduleOutputChannel, producerProperties);
|
|
|
|
|
ExtendedConsumerProperties<KafkaConsumerProperties> consumerProperties = createConsumerProperties();
|
|
|
|
|
consumerProperties.getExtension().setAutoRebalanceEnabled(false);
|
|
|
|
|
consumerProperties.getExtension().getConfiguration().put("value.deserializer", "org.apache.kafka.common.serialization.IntegerDeserializer");
|
|
|
|
|
consumerProperties.getExtension().getConfiguration().put("value.deserializer",
|
|
|
|
|
"org.apache.kafka.common.serialization.IntegerDeserializer");
|
|
|
|
|
consumerBinding = binder.bindConsumer(testTopicName, "test", moduleInputChannel, consumerProperties);
|
|
|
|
|
// Let the consumer actually bind to the producer before sending a msg
|
|
|
|
|
binderBindUnbindLatency();
|
|
|
|
|
@@ -1497,9 +1462,9 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
input2.setBeanName("test.input2J");
|
|
|
|
|
Binding<MessageChannel> input2Binding = binder.bindConsumer("partJ.raw.0", "test", input2, consumerProperties);
|
|
|
|
|
|
|
|
|
|
output.send(new GenericMessage<>(new byte[]{(byte) 0}));
|
|
|
|
|
output.send(new GenericMessage<>(new byte[]{(byte) 1}));
|
|
|
|
|
output.send(new GenericMessage<>(new byte[]{(byte) 2}));
|
|
|
|
|
output.send(new GenericMessage<>(new byte[] { (byte) 0 }));
|
|
|
|
|
output.send(new GenericMessage<>(new byte[] { (byte) 1 }));
|
|
|
|
|
output.send(new GenericMessage<>(new byte[] { (byte) 2 }));
|
|
|
|
|
|
|
|
|
|
Message<?> receive0 = receive(input0);
|
|
|
|
|
assertThat(receive0).isNotNull();
|
|
|
|
|
@@ -1538,7 +1503,6 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
catch (UnsupportedOperationException ignored) {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
ExtendedConsumerProperties<KafkaConsumerProperties> consumerProperties = createConsumerProperties();
|
|
|
|
|
consumerProperties.setConcurrency(2);
|
|
|
|
|
consumerProperties.setInstanceIndex(0);
|
|
|
|
|
@@ -1558,13 +1522,13 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
input2.setBeanName("test.input2S");
|
|
|
|
|
Binding<MessageChannel> input2Binding = binder.bindConsumer("part.raw.0", "test", input2, consumerProperties);
|
|
|
|
|
|
|
|
|
|
Message<byte[]> message2 = org.springframework.integration.support.MessageBuilder.withPayload(new byte[]{2})
|
|
|
|
|
Message<byte[]> message2 = org.springframework.integration.support.MessageBuilder.withPayload(new byte[] { 2 })
|
|
|
|
|
.setHeader(IntegrationMessageHeaderAccessor.CORRELATION_ID, "kafkaBinderTestCommonsDelegate")
|
|
|
|
|
.setHeader(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER, 42)
|
|
|
|
|
.setHeader(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE, 43).build();
|
|
|
|
|
output.send(message2);
|
|
|
|
|
output.send(new GenericMessage<>(new byte[]{1}));
|
|
|
|
|
output.send(new GenericMessage<>(new byte[]{0}));
|
|
|
|
|
output.send(new GenericMessage<>(new byte[] { 1 }));
|
|
|
|
|
output.send(new GenericMessage<>(new byte[] { 0 }));
|
|
|
|
|
Message<?> receive0 = receive(input0);
|
|
|
|
|
assertThat(receive0).isNotNull();
|
|
|
|
|
Message<?> receive1 = receive(input1);
|
|
|
|
|
@@ -1593,7 +1557,8 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
consumerProperties.setHeaderMode(HeaderMode.raw);
|
|
|
|
|
Binding<MessageChannel> consumerBinding = binder.bindConsumer("raw.0", "test", moduleInputChannel,
|
|
|
|
|
consumerProperties);
|
|
|
|
|
Message<?> message = org.springframework.integration.support.MessageBuilder.withPayload("testSendAndReceiveWithRawMode".getBytes()).build();
|
|
|
|
|
Message<?> message = org.springframework.integration.support.MessageBuilder
|
|
|
|
|
.withPayload("testSendAndReceiveWithRawMode".getBytes()).build();
|
|
|
|
|
// Let the consumer actually bind to the producer before sending a msg
|
|
|
|
|
binderBindUnbindLatency();
|
|
|
|
|
moduleOutputChannel.send(message);
|
|
|
|
|
@@ -1618,13 +1583,15 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
consumerProperties.setHeaderMode(HeaderMode.raw);
|
|
|
|
|
Binding<MessageChannel> consumerBinding = binder.bindConsumer("raw.string.0", "test", moduleInputChannel,
|
|
|
|
|
consumerProperties);
|
|
|
|
|
Message<?> message = org.springframework.integration.support.MessageBuilder.withPayload("testSendAndReceiveWithRawModeAndStringPayload").build();
|
|
|
|
|
Message<?> message = org.springframework.integration.support.MessageBuilder
|
|
|
|
|
.withPayload("testSendAndReceiveWithRawModeAndStringPayload").build();
|
|
|
|
|
// Let the consumer actually bind to the producer before sending a msg
|
|
|
|
|
binderBindUnbindLatency();
|
|
|
|
|
moduleOutputChannel.send(message);
|
|
|
|
|
Message<?> inbound = receive(moduleInputChannel);
|
|
|
|
|
assertThat(inbound).isNotNull();
|
|
|
|
|
assertThat(new String((byte[]) inbound.getPayload())).isEqualTo("testSendAndReceiveWithRawModeAndStringPayload");
|
|
|
|
|
assertThat(new String((byte[]) inbound.getPayload()))
|
|
|
|
|
.isEqualTo("testSendAndReceiveWithRawModeAndStringPayload");
|
|
|
|
|
producerBinding.unbind();
|
|
|
|
|
consumerBinding.unbind();
|
|
|
|
|
}
|
|
|
|
|
@@ -1656,14 +1623,16 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
Binding<MessageChannel> input3Binding = binder.bindConsumer(barTapName, "tap2", module3InputChannel,
|
|
|
|
|
consumerProperties);
|
|
|
|
|
|
|
|
|
|
Message<?> message = org.springframework.integration.support.MessageBuilder.withPayload("testSendAndReceiveWithExplicitConsumerGroupWithRawMode".getBytes()).build();
|
|
|
|
|
Message<?> message = org.springframework.integration.support.MessageBuilder
|
|
|
|
|
.withPayload("testSendAndReceiveWithExplicitConsumerGroupWithRawMode".getBytes()).build();
|
|
|
|
|
boolean success = false;
|
|
|
|
|
boolean retried = false;
|
|
|
|
|
while (!success) {
|
|
|
|
|
moduleOutputChannel.send(message);
|
|
|
|
|
Message<?> inbound = receive(module1InputChannel);
|
|
|
|
|
assertThat(inbound).isNotNull();
|
|
|
|
|
assertThat(new String((byte[]) inbound.getPayload())).isEqualTo("testSendAndReceiveWithExplicitConsumerGroupWithRawMode");
|
|
|
|
|
assertThat(new String((byte[]) inbound.getPayload()))
|
|
|
|
|
.isEqualTo("testSendAndReceiveWithExplicitConsumerGroupWithRawMode");
|
|
|
|
|
|
|
|
|
|
Message<?> tapped1 = receive(module2InputChannel);
|
|
|
|
|
Message<?> tapped2 = receive(module3InputChannel);
|
|
|
|
|
@@ -1674,12 +1643,15 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests<Abstr
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
|
success = true;
|
|
|
|
|
assertThat(new String((byte[]) tapped1.getPayload())).isEqualTo("testSendAndReceiveWithExplicitConsumerGroupWithRawMode");
|
|
|
|
|
assertThat(new String((byte[]) tapped2.getPayload())).isEqualTo("testSendAndReceiveWithExplicitConsumerGroupWithRawMode");
|
|
|
|
|
assertThat(new String((byte[]) tapped1.getPayload()))
|
|
|
|
|
.isEqualTo("testSendAndReceiveWithExplicitConsumerGroupWithRawMode");
|
|
|
|
|
assertThat(new String((byte[]) tapped2.getPayload()))
|
|
|
|
|
.isEqualTo("testSendAndReceiveWithExplicitConsumerGroupWithRawMode");
|
|
|
|
|
}
|
|
|
|
|
// delete one tap stream is deleted
|
|
|
|
|
input3Binding.unbind();
|
|
|
|
|
Message<?> message2 = org.springframework.integration.support.MessageBuilder.withPayload("bar".getBytes()).build();
|
|
|
|
|
Message<?> message2 = org.springframework.integration.support.MessageBuilder.withPayload("bar".getBytes())
|
|
|
|
|
.build();
|
|
|
|
|
moduleOutputChannel.send(message2);
|
|
|
|
|
|
|
|
|
|
// other tap still receives messages
|
|
|
|
|
|