Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -882,17 +882,15 @@ private void initializeAndRestorePhase() {
// Check if the topology has been updated since we last checked, ie via #addNamedTopology or #removeNamedTopology
private void checkForTopologyUpdates() {
if (topologyMetadata.isEmpty() || topologyMetadata.needsUpdate(getName())) {
log.info("StreamThread has detected an update to the topology");

taskManager.handleTopologyUpdates();
log.info("StreamThread has detected an update to the topology, triggering a rebalance to refresh the assignment");
if (topologyMetadata.isEmpty()) {
mainConsumer.unsubscribe();
}
topologyMetadata.maybeNotifyTopologyVersionWaitersAndUpdateThreadsTopologyVersion(getName());

topologyMetadata.maybeWaitForNonEmptyTopology(() -> state);

// We don't need to manually trigger a rebalance to pick up tasks from the new topology, as
// a rebalance will always occur when the metadata is updated after a change in subscription
log.info("Updating consumer subscription following topology update");
subscribeConsumer();
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1129,13 +1129,25 @@ public void updateTaskEndMetadata(final TopicPartition topicPartition, final Lon
* added NamedTopology and create them if so, then close any tasks whose named topology no longer exists
*/
void handleTopologyUpdates() {
tasks.maybeCreateTasksFromNewTopologies();
final Set<String> currentNamedTopologies = topologyMetadata.updateThreadTopologyVersion(Thread.currentThread().getName());

tasks.maybeCreateTasksFromNewTopologies(currentNamedTopologies);
maybeCloseTasksFromRemovedTopologies(currentNamedTopologies);

if (topologyMetadata.isEmpty()) {
log.info("Proactively unsubscribing from all topics due to empty topology");
mainConsumer.unsubscribe();
}

topologyMetadata.maybeNotifyTopologyVersionListeners();
}

void maybeCloseTasksFromRemovedTopologies(final Set<String> currentNamedTopologies) {
try {
final Set<Task> activeTasksToRemove = new HashSet<>();
final Set<Task> standbyTasksToRemove = new HashSet<>();
for (final Task task : tasks.allTasks()) {
if (!topologyMetadata.namedTopologiesView().contains(task.id().topologyName())) {
if (!currentNamedTopologies.contains(task.id().topologyName())) {
if (task.isActive()) {
activeTasksToRemove.add(task);
} else {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -93,8 +93,7 @@ void handleNewAssignmentAndCreateTasks(final Map<TaskId, Set<TopicPartition>> ac
createTasks(activeTasksToCreate, standbyTasksToCreate);
}

void maybeCreateTasksFromNewTopologies() {
final Set<String> currentNamedTopologies = topologyMetadata.namedTopologiesView();
void maybeCreateTasksFromNewTopologies(final Set<String> currentNamedTopologies) {
createTasks(
activeTaskCreator.uncreatedTasksForTopologies(currentNamedTopologies),
standbyTaskCreator.uncreatedTasksForTopologies(currentNamedTopologies)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -84,14 +84,14 @@ public static class TopologyVersion {
public AtomicLong topologyVersion = new AtomicLong(0L); // the local topology version
public ReentrantLock topologyLock = new ReentrantLock();
public Condition topologyCV = topologyLock.newCondition();
public List<TopologyVersionWaiters> activeTopologyWaiters = new LinkedList<>();
public List<TopologyVersionListener> activeTopologyUpdateListeners = new LinkedList<>();
}

public static class TopologyVersionWaiters {
public static class TopologyVersionListener {
final long topologyVersion; // the (minimum) version to wait for these threads to cross
final KafkaFutureImpl<Void> future; // the future waiting on all threads to be updated

public TopologyVersionWaiters(final long topologyVersion, final KafkaFutureImpl<Void> future) {
public TopologyVersionListener(final long topologyVersion, final KafkaFutureImpl<Void> future) {
this.topologyVersion = topologyVersion;
this.future = future;
}
Expand Down Expand Up @@ -162,35 +162,47 @@ public void registerThread(final String threadName) {

public void unregisterThread(final String threadName) {
threadVersions.remove(threadName);
maybeNotifyTopologyVersionWaitersAndUpdateThreadsTopologyVersion(threadName);
maybeNotifyTopologyVersionListeners();
}

public TaskExecutionMetadata taskExecutionMetadata() {
return taskExecutionMetadata;
}

public void maybeNotifyTopologyVersionWaitersAndUpdateThreadsTopologyVersion(final String threadName) {
public Set<String> updateThreadTopologyVersion(final String threadName) {
try {
lock();
final Iterator<TopologyVersionWaiters> iterator = version.activeTopologyWaiters.listIterator();
TopologyVersionWaiters topologyVersionWaiters;
version.topologyLock.lock();
threadVersions.put(threadName, topologyVersion());
return namedTopologiesView();
} finally {
version.topologyLock.unlock();
}
}

public void maybeNotifyTopologyVersionListeners() {
try {
lock();
final long minThreadVersion = getMinimumThreadVersion();
final Iterator<TopologyVersionListener> iterator = version.activeTopologyUpdateListeners.listIterator();
TopologyVersionListener topologyVersionListener;
while (iterator.hasNext()) {
topologyVersionWaiters = iterator.next();
final long topologyVersionWaitersVersion = topologyVersionWaiters.topologyVersion;
if (topologyVersionWaitersVersion <= threadVersions.get(threadName)) {
if (threadVersions.values().stream().allMatch(t -> t >= topologyVersionWaitersVersion)) {
topologyVersionWaiters.future.complete(null);
iterator.remove();
log.info("All threads are now on topology version {}", topologyVersionWaiters.topologyVersion);
}
topologyVersionListener = iterator.next();
final long topologyVersionWaitersVersion = topologyVersionListener.topologyVersion;
if (minThreadVersion >= topologyVersionWaitersVersion) {
topologyVersionListener.future.complete(null);
iterator.remove();
log.info("All threads are now on topology version {}", topologyVersionListener.topologyVersion);
}
}
} finally {
unlock();
}
}

private long getMinimumThreadVersion() {
return threadVersions.values().stream().min(Long::compare).get();
}

public void wakeupThreads() {
try {
lock();
Expand Down Expand Up @@ -225,9 +237,9 @@ public void registerAndBuildNewTopology(final KafkaFutureImpl<Void> future, fina
try {
lock();
buildAndVerifyTopology(newTopologyBuilder);
log.info("New NamedTopology passed validation and will be added {}, old topology version is {}", newTopologyBuilder.topologyName(), version.topologyVersion.get());
log.info("New NamedTopology {} passed validation and will be added, old topology version is {}", newTopologyBuilder.topologyName(), version.topologyVersion.get());
version.topologyVersion.incrementAndGet();
version.activeTopologyWaiters.add(new TopologyVersionWaiters(topologyVersion(), future));
version.activeTopologyUpdateListeners.add(new TopologyVersionListener(topologyVersion(), future));
builders.put(newTopologyBuilder.topologyName(), newTopologyBuilder);
wakeupThreads();
log.info("Added NamedTopology {} and updated topology version to {}", newTopologyBuilder.topologyName(), version.topologyVersion.get());
Expand All @@ -248,7 +260,7 @@ public KafkaFuture<Void> unregisterTopology(final KafkaFutureImpl<Void> removeTo
lock();
log.info("Beginning removal of NamedTopology {}, old topology version is {}", topologyName, version.topologyVersion.get());
version.topologyVersion.incrementAndGet();
version.activeTopologyWaiters.add(new TopologyVersionWaiters(topologyVersion(), removeTopologyFuture));
version.activeTopologyUpdateListeners.add(new TopologyVersionListener(topologyVersion(), removeTopologyFuture));
final InternalTopologyBuilder removedBuilder = builders.remove(topologyName);
removedBuilder.fullSourceTopicNames().forEach(allInputTopics::remove);
removedBuilder.allSourcePatternStrings().forEach(allInputTopics::remove);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -621,38 +621,40 @@ public void shouldAllowMixedCollectionAndPatternSubscriptionWithMultipleNamedTop
@Test
public void shouldAddToEmptyInitialTopologyRemoveResetOffsetsThenAddSameNamedTopology() throws Exception {
CLUSTER.createTopics(SUM_OUTPUT, COUNT_OUTPUT);
// Build up named topology with two stateful subtopologies
final KStream<String, Long> inputStream1 = topology1Builder.stream(INPUT_STREAM_1);
inputStream1.groupByKey().count().toStream().to(COUNT_OUTPUT);
inputStream1.groupByKey().reduce(Long::sum).toStream().to(SUM_OUTPUT);
streams.start();
final NamedTopology namedTopology = topology1Builder.build();
streams.addNamedTopology(namedTopology).all().get();

assertThat(waitUntilMinKeyValueRecordsReceived(consumerConfig, COUNT_OUTPUT, 3), equalTo(COUNT_OUTPUT_DATA));
assertThat(waitUntilMinKeyValueRecordsReceived(consumerConfig, SUM_OUTPUT, 3), equalTo(SUM_OUTPUT_DATA));
streams.removeNamedTopology("topology-1", true).all().get();
streams.cleanUpNamedTopology("topology-1");

CLUSTER.getAllTopicsInCluster().stream().filter(t -> t.contains("changelog")).forEach(t -> {
try {
CLUSTER.deleteTopicAndWait(t);
} catch (final InterruptedException e) {
e.printStackTrace();
}
});
try {
// Build up named topology with two stateful subtopologies
final KStream<String, Long> inputStream1 = topology1Builder.stream(INPUT_STREAM_1);
inputStream1.groupByKey().count().toStream().to(COUNT_OUTPUT);
inputStream1.groupByKey().reduce(Long::sum).toStream().to(SUM_OUTPUT);
streams.start();
final NamedTopology namedTopology = topology1Builder.build();
streams.addNamedTopology(namedTopology).all().get();

final KStream<String, Long> inputStream = topology1BuilderDup.stream(INPUT_STREAM_1);
inputStream.groupByKey().count().toStream().to(COUNT_OUTPUT);
inputStream.groupByKey().reduce(Long::sum).toStream().to(SUM_OUTPUT);
assertThat(waitUntilMinKeyValueRecordsReceived(consumerConfig, COUNT_OUTPUT, 3), equalTo(COUNT_OUTPUT_DATA));
assertThat(waitUntilMinKeyValueRecordsReceived(consumerConfig, SUM_OUTPUT, 3), equalTo(SUM_OUTPUT_DATA));
streams.removeNamedTopology("topology-1", true).all().get();
streams.cleanUpNamedTopology("topology-1");

CLUSTER.getAllTopicsInCluster().stream().filter(t -> t.contains("changelog")).forEach(t -> {
try {
CLUSTER.deleteTopicAndWait(t);
} catch (final InterruptedException e) {
e.printStackTrace();
}
});

final NamedTopology namedTopologyDup = topology1BuilderDup.build();
streams.addNamedTopology(namedTopologyDup).all().get();
final KStream<String, Long> inputStream = topology1BuilderDup.stream(INPUT_STREAM_1);
inputStream.groupByKey().count().toStream().to(COUNT_OUTPUT);
inputStream.groupByKey().reduce(Long::sum).toStream().to(SUM_OUTPUT);

assertThat(waitUntilMinKeyValueRecordsReceived(consumerConfig, COUNT_OUTPUT, 3), equalTo(COUNT_OUTPUT_DATA));
assertThat(waitUntilMinKeyValueRecordsReceived(consumerConfig, SUM_OUTPUT, 3), equalTo(SUM_OUTPUT_DATA));
final NamedTopology namedTopologyDup = topology1BuilderDup.build();
streams.addNamedTopology(namedTopologyDup).all().get();

CLUSTER.deleteTopicsAndWait(SUM_OUTPUT, COUNT_OUTPUT);
assertThat(waitUntilMinKeyValueRecordsReceived(consumerConfig, COUNT_OUTPUT, 3), equalTo(COUNT_OUTPUT_DATA));
assertThat(waitUntilMinKeyValueRecordsReceived(consumerConfig, SUM_OUTPUT, 3), equalTo(SUM_OUTPUT_DATA));
} finally {
CLUSTER.deleteTopicsAndWait(SUM_OUTPUT, COUNT_OUTPUT);
}
}

@Test
Expand Down