Skip to content
Merged
Show file tree
Hide file tree
Changes from 7 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 @@ -326,4 +326,14 @@ PartitionsOnReplicaIterator partitionsWithBrokerInIsr(int brokerId) {
boolean hasLeaderships(int brokerId) {
return iterator(brokerId, true).hasNext();
}

int offlinePartitionCount() {
Comment thread
cmccabe marked this conversation as resolved.
PartitionsOnReplicaIterator noLeaderIterator = partitionsWithNoLeader();
int offlinePartitionCount = 0;
while (noLeaderIterator.hasNext()) {
noLeaderIterator.next();
offlinePartitionCount++;
}
return offlinePartitionCount;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -26,4 +26,20 @@ public interface ControllerMetrics {
void updateEventQueueTime(long durationMs);

void updateEventQueueProcessingTime(long durationMs);

void setGlobalTopicsCount(int topicCount);

int globalTopicsCount();

void setGlobalPartitionCount(int partitionCount);

int globalPartitionCount();

void setOfflinePartitionCount(int offlinePartitions);

int offlinePartitionCount();

void setPreferredReplicaImbalanceCount(int replicaImbalances);

int preferredReplicaImbalanceCount();
}
Original file line number Diff line number Diff line change
Expand Up @@ -928,7 +928,7 @@ private QuorumController(LogContext logContext,
this.snapshotGeneratorManager = new SnapshotGeneratorManager(snapshotWriterBuilder);
this.replicationControl = new ReplicationControlManager(snapshotRegistry,
logContext, defaultReplicationFactor, defaultNumPartitions,
configurationControl, clusterControl);
configurationControl, clusterControl, controllerMetrics);
this.logManager = logManager;
this.metaLogListener = new QuorumMetaLogListener();
this.curClaimEpoch = -1L;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,14 +30,35 @@ public final class QuorumControllerMetrics implements ControllerMetrics {
"kafka.controller", "ControllerEventManager", "EventQueueTimeMs", null);
private final static MetricName EVENT_QUEUE_PROCESSING_TIME_MS = new MetricName(
"kafka.controller", "ControllerEventManager", "EventQueueProcessingTimeMs", null);
private final static MetricName GLOBAL_TOPIC_COUNT = new MetricName(
"kafka.controller", "ReplicationControlManager", "GlobalTopicCount", null);
private final static MetricName GLOBAL_PARTITION_COUNT = new MetricName(
"kafka.controller", "ReplicationControlManager", "GlobalPartitionCount", null);
private final static MetricName OFFLINE_PARTITION_COUNT = new MetricName(
"kafka.controller", "ReplicationControlManager", "OfflinePartitionCount", null);
private final static MetricName PREFERRED_REPLICA_IMBALANCE_COUNT = new MetricName(
"kafka.controller", "ReplicationControlManager", "PreferredReplicaImbalanceCount", null);


private volatile boolean active;
private volatile int topics;
private volatile int partitions;
private volatile int offlinePartitions;
private volatile int preferredReplicaImbalances;
private final Gauge<Integer> activeControllerCount;
private final Gauge<Integer> globalPartitionCount;
private final Gauge<Integer> globalTopicCount;
private final Gauge<Integer> offlinePartitionCount;
private final Gauge<Integer> preferredReplicaImbalanceCount;
private final Histogram eventQueueTime;
private final Histogram eventQueueProcessingTime;

public QuorumControllerMetrics(MetricsRegistry registry) {
this.active = false;
this.topics = 0;
this.partitions = 0;
this.offlinePartitions = 0;
this.preferredReplicaImbalances = 0;
this.activeControllerCount = registry.newGauge(ACTIVE_CONTROLLER_COUNT, new Gauge<Integer>() {
@Override
public Integer value() {
Expand All @@ -46,6 +67,30 @@ public Integer value() {
});
this.eventQueueTime = registry.newHistogram(EVENT_QUEUE_TIME_MS, true);
this.eventQueueProcessingTime = registry.newHistogram(EVENT_QUEUE_PROCESSING_TIME_MS, true);
this.globalTopicCount = registry.newGauge(GLOBAL_TOPIC_COUNT, new Gauge<Integer>() {
@Override
public Integer value() {
return topics;
}
});
this.globalPartitionCount = registry.newGauge(GLOBAL_PARTITION_COUNT, new Gauge<Integer>() {
@Override
public Integer value() {
return partitions;
}
});
this.offlinePartitionCount = registry.newGauge(OFFLINE_PARTITION_COUNT, new Gauge<Integer>() {
@Override
public Integer value() {
return offlinePartitions;
}
});
this.preferredReplicaImbalanceCount = registry.newGauge(PREFERRED_REPLICA_IMBALANCE_COUNT, new Gauge<Integer>() {
@Override
public Integer value() {
return preferredReplicaImbalances;
}
});
}

@Override
Expand All @@ -67,4 +112,44 @@ public void updateEventQueueTime(long durationMs) {
public void updateEventQueueProcessingTime(long durationMs) {
eventQueueTime.update(durationMs);
}

@Override
public void setGlobalTopicsCount(int topicCount) {
this.topics = topicCount;
}

@Override
public int globalTopicsCount() {
return this.topics;
}

@Override
public void setGlobalPartitionCount(int partitionCount) {
this.partitions = partitionCount;
}

@Override
public int globalPartitionCount() {
return this.partitions;
}

@Override
public void setOfflinePartitionCount(int offlinePartitions) {
this.offlinePartitions = offlinePartitions;
}

@Override
public int offlinePartitionCount() {
return this.offlinePartitions;
}

@Override
public void setPreferredReplicaImbalanceCount(int replicaImbalances) {
this.preferredReplicaImbalances = replicaImbalances;
}

@Override
public int preferredReplicaImbalanceCount() {
return this.preferredReplicaImbalances;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -279,6 +279,8 @@ public String toString() {
*/
private final int defaultNumPartitions;

private int globalPartitionCount;

/**
* A reference to the controller's configuration control manager.
*/
Expand All @@ -289,6 +291,11 @@ public String toString() {
*/
private final ClusterControlManager clusterControl;

/**
* A reference to the controller's metrics registry.
*/
private final ControllerMetrics controllerMetrics;

/**
* Maps topic names to topic UUIDs.
*/
Expand All @@ -309,13 +316,16 @@ public String toString() {
short defaultReplicationFactor,
int defaultNumPartitions,
ConfigurationControlManager configurationControl,
ClusterControlManager clusterControl) {
ClusterControlManager clusterControl,
ControllerMetrics controllerMetrics) {
this.snapshotRegistry = snapshotRegistry;
this.log = logContext.logger(ReplicationControlManager.class);
this.defaultReplicationFactor = defaultReplicationFactor;
this.defaultNumPartitions = defaultNumPartitions;
this.configurationControl = configurationControl;
this.controllerMetrics = controllerMetrics;
this.clusterControl = clusterControl;
this.globalPartitionCount = 0;
this.topicsByName = new TimelineHashMap<>(snapshotRegistry, 0);
this.topics = new TimelineHashMap<>(snapshotRegistry, 0);
this.brokersToIsrs = new BrokersToIsrs(snapshotRegistry);
Expand All @@ -325,6 +335,7 @@ public void replay(TopicRecord record) {
topicsByName.put(record.name(), record.topicId());
topics.put(record.topicId(),
new TopicControlInfo(record.name(), snapshotRegistry, record.topicId()));
controllerMetrics.setGlobalTopicsCount(topics.size());
log.info("Created topic {} with topic ID {}.", record.name(), record.topicId());
}

Expand All @@ -343,12 +354,16 @@ public void replay(PartitionRecord record) {
topicInfo.parts.put(record.partitionId(), newPartInfo);
brokersToIsrs.update(record.topicId(), record.partitionId(), null,
newPartInfo.isr, NO_LEADER, newPartInfo.leader);
globalPartitionCount++;
} else if (!newPartInfo.equals(prevPartInfo)) {
newPartInfo.maybeLogPartitionChange(log, description, prevPartInfo);
topicInfo.parts.put(record.partitionId(), newPartInfo);
brokersToIsrs.update(record.topicId(), record.partitionId(), prevPartInfo.isr,
newPartInfo.isr, prevPartInfo.leader, newPartInfo.leader);
}
controllerMetrics.setGlobalPartitionCount(globalPartitionCount);
controllerMetrics.setOfflinePartitionCount(brokersToIsrs.offlinePartitionCount());
controllerMetrics.setPreferredReplicaImbalanceCount(preferredReplicaImbalanceCount());
}

public void replay(PartitionChangeRecord record) {
Expand All @@ -370,6 +385,8 @@ public void replay(PartitionChangeRecord record) {
String topicPart = topicInfo.name + "-" + record.partitionId() + " with topic ID " +
record.topicId();
newPartitionInfo.maybeLogPartitionChange(log, topicPart, prevPartitionInfo);
controllerMetrics.setOfflinePartitionCount(brokersToIsrs.offlinePartitionCount());
controllerMetrics.setPreferredReplicaImbalanceCount(preferredReplicaImbalanceCount());
}

public void replay(RemoveTopicRecord record) {
Expand All @@ -389,9 +406,13 @@ public void replay(RemoveTopicRecord record) {
for (int i = 0; i < partition.isr.length; i++) {
brokersToIsrs.removeTopicEntryForBroker(topic.id, partition.isr[i]);
}
globalPartitionCount--;
}
brokersToIsrs.removeTopicEntryForBroker(topic.id, NO_LEADER);

controllerMetrics.setGlobalTopicsCount(topics.size());
controllerMetrics.setGlobalPartitionCount(globalPartitionCount);
controllerMetrics.setOfflinePartitionCount(brokersToIsrs.offlinePartitionCount());
controllerMetrics.setPreferredReplicaImbalanceCount(preferredReplicaImbalanceCount());
log.info("Removed topic {} with ID {}.", topic.name, record.topicId());
}

Expand Down Expand Up @@ -456,6 +477,18 @@ public void replay(RemoveTopicRecord record) {
return ControllerResult.atomicOf(records, data);
}

private int preferredReplicaImbalanceCount() {
int count = 0;
for (TopicControlInfo topic : topics.values()) {
for (PartitionControlInfo part : topic.parts.values()) {
if (part.leader != part.preferredReplica()) {
count++;
}
}
}
return count;
}

private ApiError createTopic(CreatableTopic topic,
List<ApiMessageAndVersion> records,
Map<String, CreatableTopicResult> successes) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,12 +17,19 @@

package org.apache.kafka.controller;


public final class MockControllerMetrics implements ControllerMetrics {
private volatile boolean active;
private volatile int topics;
private volatile int partitions;
private volatile int offlinePartitions;
private volatile int preferredReplicaImbalances;

public MockControllerMetrics() {
this.active = false;
this.topics = 0;
this.partitions = 0;
this.offlinePartitions = 0;
this.preferredReplicaImbalances = 0;
}

@Override
Expand All @@ -44,4 +51,44 @@ public void updateEventQueueTime(long durationMs) {
public void updateEventQueueProcessingTime(long durationMs) {
// nothing to do
}

@Override
public void setGlobalTopicsCount(int topicCount) {
this.topics = topicCount;
}

@Override
public int globalTopicsCount() {
return this.topics;
}

@Override
public void setGlobalPartitionCount(int partitionCount) {
this.partitions = partitionCount;
}

@Override
public int globalPartitionCount() {
return this.partitions;
}

@Override
public void setOfflinePartitionCount(int offlinePartitions) {
this.offlinePartitions = offlinePartitions;
}

@Override
public int offlinePartitionCount() {
return this.offlinePartitions;
}

@Override
public void setPreferredReplicaImbalanceCount(int replicaImbalances) {
this.preferredReplicaImbalances = replicaImbalances;
}

@Override
public int preferredReplicaImbalanceCount() {
return this.preferredReplicaImbalances;
}
}
Loading