diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/DefaultConnectionTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/DefaultConnectionTests.java index 51e06fcbe1..016a6cca4b 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/DefaultConnectionTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/DefaultConnectionTests.java @@ -49,7 +49,8 @@ public class DefaultConnectionTests extends AbstractBrokerTests { @Test public void testFetchPartitionMetadata() throws Exception { createTopic(TEST_TOPIC, 1, 1, 1); - Connection brokerConnection = new DefaultConnection(getKafkaRule().getBrokerAddresses().get(0), "client", 64*1024, 3000, 1, 100); + Connection brokerConnection = + new DefaultConnection(getKafkaRule().getBrokerAddresses().get(0), "client", 64*1024, 3000, 1, 10000); Partition partition = new Partition(TEST_TOPIC, 0); Result result = brokerConnection.fetchInitialOffset(-1, partition); Assert.assertEquals(0, result.getErrors().size()); @@ -62,7 +63,8 @@ public class DefaultConnectionTests extends AbstractBrokerTests { createTopic(TEST_TOPIC, 1, 1, 1); Producer producer = createStringProducer(0); producer.send( createMessages(10, TEST_TOPIC)); - Connection brokerConnection = new DefaultConnection(getKafkaRule().getBrokerAddresses().get(0), "client", 64*1024, 3000, 1, 100); + Connection brokerConnection = + new DefaultConnection(getKafkaRule().getBrokerAddresses().get(0), "client", 64*1024, 3000, 1, 10000); Partition partition = new Partition(TEST_TOPIC, 0); FetchRequest fetchRequest = new FetchRequest(partition, 0L, 1000); Result result = brokerConnection.fetch(fetchRequest); @@ -84,7 +86,8 @@ public class DefaultConnectionTests extends AbstractBrokerTests { createTopic(TEST_TOPIC, 1, 1, 1); Producer producer = createStringProducer(1); producer.send( createMessages(10, TEST_TOPIC)); - Connection brokerConnection = new DefaultConnection(getKafkaRule().getBrokerAddresses().get(0), "client", 64*1024, 3000, 1, 100); + Connection brokerConnection = + new DefaultConnection(getKafkaRule().getBrokerAddresses().get(0), "client", 64*1024, 3000, 1, 10000); Partition partition = new Partition(TEST_TOPIC, 0); FetchRequest fetchRequest = new FetchRequest(partition, 0L, 1000); Result result = brokerConnection.fetch(fetchRequest); @@ -106,7 +109,8 @@ public class DefaultConnectionTests extends AbstractBrokerTests { createTopic(TEST_TOPIC, 1, 1, 1); Producer producer = createStringProducer(2); producer.send( createMessages(10, TEST_TOPIC)); - Connection brokerConnection = new DefaultConnection(getKafkaRule().getBrokerAddresses().get(0), "client", 64*1024, 3000, 1, 100); + Connection brokerConnection = + new DefaultConnection(getKafkaRule().getBrokerAddresses().get(0), "client", 64*1024, 3000, 1, 10000); Partition partition = new Partition(TEST_TOPIC, 0); FetchRequest fetchRequest = new FetchRequest(partition, 0L, 1000); Result result = brokerConnection.fetch(fetchRequest);