Manage log.dir in the EmbeddedKafkaBroker

Resolves https://github.com/spring-projects/spring-kafka/issues/194

Create the temporary directory in EKB instead of the broker to avoid
`NoSuchFileException`s during shutdown.

**cherry-pick to 2.4.x, 2.3.x, 2.2.x**

# Conflicts:
#	spring-kafka-test/src/main/java/org/springframework/kafka/test/EmbeddedKafkaBroker.java
This commit is contained in:
Gary Russell
2020-05-01 10:52:27 -04:00
committed by Artem Bilan
parent 65dba07634
commit a21fdebdd7

View File

@@ -16,8 +16,9 @@
package org.springframework.kafka.test;
import static org.assertj.core.api.Assertions.assertThat;
import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.file.Files;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Arrays;
@@ -29,6 +30,7 @@ import java.util.List;
import java.util.Map;
import java.util.Properties;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.stream.Collectors;
@@ -212,6 +214,7 @@ public class EmbeddedKafkaBroker implements InitializingBean, DisposableBean {
this.zookeeperClient = new ZkClient(this.zkConnect, zkSessionTimeout, zkConnectionTimeout,
ZKStringSerializer$.MODULE$);
this.kafkaServers.clear();
boolean userLogDir = this.brokerProperties.get(KafkaConfig.LogDirProp()) != null && this.count == 1;
for (int i = 0; i < this.count; i++) {
Properties brokerConfigProperties = createBrokerProperties(i);
brokerConfigProperties.setProperty(KafkaConfig.ReplicaSocketTimeoutMsProp(), "1000");
@@ -222,6 +225,9 @@ public class EmbeddedKafkaBroker implements InitializingBean, DisposableBean {
if (this.brokerProperties != null) {
this.brokerProperties.forEach(brokerConfigProperties::put);
}
if (!userLogDir) {
logDir(brokerConfigProperties);
}
KafkaServer server = TestUtils.createServer(new KafkaConfig(brokerConfigProperties), Time.SYSTEM);
this.kafkaServers.add(server);
if (this.kafkaPorts[i] == 0) {
@@ -237,6 +243,16 @@ public class EmbeddedKafkaBroker implements InitializingBean, DisposableBean {
System.setProperty(SPRING_EMBEDDED_ZOOKEEPER_CONNECT, getZookeeperConnectionString());
}
private void logDir(Properties brokerConfigProperties) {
try {
brokerConfigProperties.put(KafkaConfig.LogDirProp(),
Files.createTempDirectory("spring.kafka." + UUID.randomUUID()).toString());
}
catch (IOException e) {
throw new UncheckedIOException(e);
}
}
private void overrideExitMethods() {
String exitMsg = "Exit.%s(%d, %s) called";
Exit.setExitProcedure((statusCode, message) -> {
@@ -478,11 +494,12 @@ public class EmbeddedKafkaBroker implements InitializingBean, DisposableBean {
* @param topics the topics.
*/
public void consumeFromEmbeddedTopics(Consumer<?, ?> consumer, String... topics) {
HashSet<String> diff = new HashSet<>(Arrays.asList(topics));
diff.removeAll(new HashSet<>(this.topics));
assertThat(this.topics)
.as("topic(s):'" + diff + "' are not in embedded topic list")
.containsAll(new HashSet<>(Arrays.asList(topics)));
List<String> notEmbedded = Arrays.stream(topics)
.filter(topic -> !this.topics.contains(topic))
.collect(Collectors.toList());
if (notEmbedded.size() > 0) {
throw new IllegalStateException("topic(s):'" + notEmbedded + "' are not in embedded topic list");
}
final AtomicBoolean assigned = new AtomicBoolean();
consumer.subscribe(Arrays.asList(topics), new ConsumerRebalanceListener() {
@@ -516,9 +533,9 @@ public class EmbeddedKafkaBroker implements InitializingBean, DisposableBean {
}
consumer.seekToBeginning(records.partitions());
}
assertThat(assigned.get())
.as("Failed to be assigned partitions from the embedded topics")
.isTrue();
if (!assigned.get()) {
throw new IllegalStateException("Failed to be assigned partitions from the embedded topics");
}
logger.debug("Subscription Initiated");
}