GH-120: Update to 10.0.0.0 Kafka

Resolves #120
This commit is contained in:
Fouad HAMDI
2016-06-14 15:22:11 +02:00
committed by Gary Russell
parent bae41a1778
commit dba407bd33
9 changed files with 48 additions and 43 deletions

View File

@@ -80,7 +80,7 @@ subprojects { subproject ->
hamcrestVersion = '1.3'
jacksonVersion = '2.6.7'
junitVersion = '4.12'
kafkaVersion = '0.9.0.1'
kafkaVersion = '0.10.0.0'
log4jVersion = '1.2.17'
mockitoVersion = '1.10.19'
scalaVersion = '2.11'

View File

@@ -0,0 +1,9 @@
log4j.rootCategory=WARN, stdout
log4j.appender.stdout=org.apache.log4j.ConsoleAppender
log4j.appender.stdout.layout=org.apache.log4j.PatternLayout
log4j.appender.stdout.layout.ConversionPattern=%d{HH:mm:ss.SSS} %-5p [%t][%c] %m%n
log4j.category.org.springframework.kafka=WARN
log4j.category.org.apache.kafka.clients=WARN
log4j.category.org.apache.kafka.common.network.Selector=ERROR
log4j.category.kafka.server.ReplicaFetcherThread=ERROR

View File

@@ -3,4 +3,4 @@ distributionBase=GRADLE_USER_HOME
distributionPath=wrapper/dists
zipStoreBase=GRADLE_USER_HOME
zipStorePath=wrapper/dists
distributionUrl=https\://services.gradle.org/distributions/gradle-2.14-bin.zip
distributionUrl=https\://services.gradle.org/distributions/gradle-2.14-all.zip

View File

@@ -35,9 +35,11 @@ import org.I0Itec.zkclient.ZkClient;
import org.I0Itec.zkclient.exception.ZkInterruptedException;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRebalanceListener;
import org.apache.kafka.common.Node;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.protocol.Errors;
import org.apache.kafka.common.protocol.SecurityProtocol;
import org.apache.kafka.common.requests.MetadataResponse;
import org.junit.rules.ExternalResource;
import org.springframework.kafka.test.core.BrokerAddress;
@@ -49,9 +51,6 @@ import org.springframework.retry.support.RetryTemplate;
import kafka.admin.AdminUtils;
import kafka.admin.AdminUtils$;
import kafka.api.PartitionMetadata;
import kafka.api.TopicMetadata;
import kafka.cluster.BrokerEndPoint;
import kafka.server.KafkaConfig;
import kafka.server.KafkaServer;
import kafka.server.NotRunning;
@@ -147,7 +146,8 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule {
true, randomPort,
scala.Option.<SecurityProtocol>apply(null),
scala.Option.<File>apply(null),
true, false, 0, false, 0, false, 0);
scala.Option.<Properties>apply(null),
true, false, 0, false, 0, false, 0, scala.Option.<String>apply(null));
brokerConfigProperties.setProperty("replica.socket.timeout.ms", "1000");
brokerConfigProperties.setProperty("controller.socket.timeout.ms", "1000");
brokerConfigProperties.setProperty("offsets.topic.replication.factor", "1");
@@ -157,7 +157,7 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule {
ZkUtils zkUtils = new ZkUtils(getZkClient(), null, false);
Properties props = new Properties();
for (String topic : this.topics) {
AdminUtils.createTopic(zkUtils, topic, this.partitionsPerTopic, this.count, props);
AdminUtils.createTopic(zkUtils, topic, this.partitionsPerTopic, this.count, props, null);
}
System.setProperty(SPRING_EMBEDDED_KAFKA_BROKERS, getBrokersAsString());
}
@@ -176,7 +176,7 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule {
// do nothing
}
try {
CoreUtils.rm(kafkaServer.config().logDirs());
CoreUtils.delete(kafkaServer.config().logDirs());
}
catch (Exception e) {
// do nothing
@@ -266,16 +266,14 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule {
canExit = true;
ZkUtils zkUtils = new ZkUtils(getZkClient(), null, false);
Map<String, Properties> topicProperties = AdminUtils$.MODULE$.fetchAllTopicConfigs(zkUtils);
Set<TopicMetadata> topicMetadatas =
Set<MetadataResponse.TopicMetadata> topicMetadatas =
AdminUtils$.MODULE$.fetchTopicMetadataFromZk(topicProperties.keySet(), zkUtils);
for (TopicMetadata topicMetadata : JavaConversions.asJavaCollection(topicMetadatas)) {
if (Errors.forCode(topicMetadata.errorCode()).exception() == null) {
for (PartitionMetadata partitionMetadata :
JavaConversions.asJavaCollection(topicMetadata.partitionsMetadata())) {
Collection<BrokerEndPoint> inSyncReplicas =
JavaConversions.asJavaCollection(partitionMetadata.isr());
for (BrokerEndPoint broker : inSyncReplicas) {
if (broker.id() == index) {
for (MetadataResponse.TopicMetadata topicMetadata : JavaConversions.asJavaCollection(topicMetadatas)) {
if (Errors.forCode(topicMetadata.error().code()).exception() == null) {
for (MetadataResponse.PartitionMetadata partitionMetadata : topicMetadata.partitionMetadata()) {
Collection<Node> inSyncReplicas = partitionMetadata.isr();
for (Node node : inSyncReplicas) {
if (node.id() == index) {
canExit = false;
}
}
@@ -331,14 +329,13 @@ public class KafkaEmbedded extends ExternalResource implements KafkaRule {
}
canExit = true;
ZkUtils zkUtils = new ZkUtils(getZkClient(), null, false);
TopicMetadata topicMetadata = AdminUtils$.MODULE$.fetchTopicMetadataFromZk(topic, zkUtils);
if (Errors.forCode(topicMetadata.errorCode()).exception() == null) {
for (PartitionMetadata partitionMetadata :
JavaConversions.asJavaCollection(topicMetadata.partitionsMetadata())) {
Collection<BrokerEndPoint> isr = JavaConversions.asJavaCollection(partitionMetadata.isr());
MetadataResponse.TopicMetadata topicMetadata = AdminUtils$.MODULE$.fetchTopicMetadataFromZk(topic, zkUtils);
if (Errors.forCode(topicMetadata.error().code()).exception() == null) {
for (MetadataResponse.PartitionMetadata partitionMetadata : topicMetadata.partitionMetadata()) {
Collection<Node> isr = partitionMetadata.isr();
boolean containsIndex = false;
for (BrokerEndPoint broker : isr) {
if (broker.id() == brokerId) {
for (Node node : isr) {
if (node.id() == brokerId) {
containsIndex = true;
}
}

View File

@@ -80,7 +80,7 @@ public final class KafkaTestUtils {
props.put(ConsumerConfig.GROUP_ID_CONFIG, group);
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, autoCommit);
props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, "100");
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "15000");
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 15000);
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, IntegerDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
return props;

View File

@@ -429,8 +429,7 @@ public class KafkaMessageListenerContainer<K, V> extends AbstractMessageListener
if (this.assignedPartitions != null) {
// avoid group management rebalance due to a slow
// consumer
this.consumer.pause(this.assignedPartitions
.toArray(new TopicPartition[this.assignedPartitions.size()]));
this.consumer.pause(this.assignedPartitions);
this.paused = true;
this.unsent = records;
}
@@ -514,8 +513,7 @@ public class KafkaMessageListenerContainer<K, V> extends AbstractMessageListener
private ConsumerRecords<K, V> checkPause(ConsumerRecords<K, V> unsent) {
if (this.paused && this.recordsToProcess.size() < this.containerProperties.getQueueDepth()) {
// Listener has caught up.
this.consumer.resume(
this.assignedPartitions.toArray(new TopicPartition[this.assignedPartitions.size()]));
this.consumer.resume(this.assignedPartitions);
this.paused = false;
if (unsent != null) {
try {
@@ -693,7 +691,7 @@ public class KafkaMessageListenerContainer<K, V> extends AbstractMessageListener
if (offset < 0) {
if (!metadata.relativeToCurrent) {
this.consumer.seekToEnd(topicPartition);
this.consumer.seekToEnd(Arrays.asList(topicPartition));
}
newOffset = Math.max(0, this.consumer.position(topicPartition) + offset);
}
@@ -708,7 +706,7 @@ public class KafkaMessageListenerContainer<K, V> extends AbstractMessageListener
}
}
catch (Exception e) {
logger.error("Failed to set initial offset for " + topicPartition
this.logger.error("Failed to set initial offset for " + topicPartition
+ " at " + newOffset + ". Position is " + this.consumer.position(topicPartition), e);
}
}

View File

@@ -102,7 +102,7 @@ public class TopicPartitionInitialOffset {
}
public boolean isRelativeToCurrent() {
return relativeToCurrent;
return this.relativeToCurrent;
}
@Override

View File

@@ -330,7 +330,7 @@ public class EnableKafkaIntegrationTests {
@Bean
public Map<String, Object> consumerConfigs() {
Map<String, Object> consumerProps = KafkaTestUtils.consumerProps("testAnnot", "true", embeddedKafka);
Map<String, Object> consumerProps = KafkaTestUtils.consumerProps("testAnnot", "false", embeddedKafka);
consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
return consumerProps;
}

View File

@@ -19,6 +19,7 @@ package org.springframework.kafka.listener;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.BDDMockito.willAnswer;
import static org.mockito.Matchers.any;
import static org.mockito.Matchers.anyObject;
import static org.mockito.Mockito.atLeastOnce;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.times;
@@ -142,8 +143,8 @@ public class KafkaMessageListenerContainerTests {
template.flush();
assertThat(latch.await(60, TimeUnit.SECONDS)).isTrue();
assertThat(bitSet.cardinality()).isEqualTo(6);
verify(consumer, atLeastOnce()).pause(any(TopicPartition.class), any(TopicPartition.class));
verify(consumer, atLeastOnce()).resume(any(TopicPartition.class), any(TopicPartition.class));
verify(consumer, atLeastOnce()).pause(anyObject());
verify(consumer, atLeastOnce()).resume(anyObject());
container.stop();
logger.info("Stop " + this.testName.getMethodName());
}
@@ -205,8 +206,8 @@ public class KafkaMessageListenerContainerTests {
template.flush();
assertThat(latch.await(60, TimeUnit.SECONDS)).isTrue();
assertThat(bitSet.cardinality()).isEqualTo(6);
verify(consumer, atLeastOnce()).pause(any(TopicPartition.class), any(TopicPartition.class));
verify(consumer, atLeastOnce()).resume(any(TopicPartition.class), any(TopicPartition.class));
verify(consumer, atLeastOnce()).pause(anyObject());
verify(consumer, atLeastOnce()).resume(anyObject());
container.stop();
logger.info("Stop " + this.testName.getMethodName() + ackMode);
}
@@ -264,8 +265,8 @@ public class KafkaMessageListenerContainerTests {
// Verify that commitSync is called when paused
assertThat(latch.await(60, TimeUnit.SECONDS)).isTrue();
verify(consumer, atLeastOnce()).pause(any(TopicPartition.class), any(TopicPartition.class));
verify(consumer, atLeastOnce()).resume(any(TopicPartition.class), any(TopicPartition.class));
verify(consumer, atLeastOnce()).pause(anyObject());
verify(consumer, atLeastOnce()).resume(anyObject());
container.stop();
}
@@ -370,8 +371,8 @@ public class KafkaMessageListenerContainerTests {
template.flush();
assertThat(latch.await(60, TimeUnit.SECONDS)).isTrue();
assertThat(bitSet.cardinality()).isEqualTo(6);
verify(consumer, atLeastOnce()).pause(any(TopicPartition.class), any(TopicPartition.class));
verify(consumer, atLeastOnce()).resume(any(TopicPartition.class), any(TopicPartition.class));
verify(consumer, atLeastOnce()).pause(anyObject());
verify(consumer, atLeastOnce()).resume(anyObject());
container.stop();
logger.info("Stop " + this.testName.getMethodName());
}
@@ -439,8 +440,8 @@ public class KafkaMessageListenerContainerTests {
template.flush();
assertThat(latch.await(60, TimeUnit.SECONDS)).isTrue();
assertThat(bitSet.cardinality()).isEqualTo(6);
verify(consumer, atLeastOnce()).pause(any(TopicPartition.class), any(TopicPartition.class));
verify(consumer, atLeastOnce()).resume(any(TopicPartition.class), any(TopicPartition.class));
verify(consumer, atLeastOnce()).pause(anyObject());
verify(consumer, atLeastOnce()).resume(anyObject());
container.stop();
logger.info("Stop " + this.testName.getMethodName());
}