Skip to content
Merged
Show file tree
Hide file tree
Changes from 28 commits
Commits
Show all changes
30 commits
Select commit Hold shift + click to select a range
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 @@ -20,7 +20,7 @@
import org.apache.kafka.clients.consumer.internals.ConsumerDelegate;
import org.apache.kafka.clients.consumer.internals.ConsumerDelegateCreator;
import org.apache.kafka.clients.consumer.internals.ConsumerMetadata;
import org.apache.kafka.clients.consumer.internals.KafkaConsumerMetrics;
import org.apache.kafka.clients.consumer.internals.metrics.KafkaConsumerMetrics;
import org.apache.kafka.clients.consumer.internals.SubscriptionState;
import org.apache.kafka.common.KafkaException;
import org.apache.kafka.common.Metric;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,8 @@
import org.apache.kafka.clients.consumer.internals.events.TopicMetadataApplicationEvent;
import org.apache.kafka.clients.consumer.internals.events.UnsubscribeApplicationEvent;
import org.apache.kafka.clients.consumer.internals.events.ValidatePositionsApplicationEvent;
import org.apache.kafka.clients.consumer.internals.metrics.RebalanceCallbackMetrics;
import org.apache.kafka.clients.consumer.internals.metrics.KafkaConsumerMetrics;
import org.apache.kafka.clients.consumer.internals.metrics.RebalanceCallbackMetricsManager;
import org.apache.kafka.common.Cluster;
import org.apache.kafka.common.IsolationLevel;
import org.apache.kafka.common.KafkaException;
Expand Down Expand Up @@ -396,7 +397,7 @@ private void process(final ConsumerRebalanceListenerCallbackNeededEvent event) {
logContext,
subscriptions,
time,
new RebalanceCallbackMetrics(metrics)
new RebalanceCallbackMetricsManager(metrics)
);
this.backgroundEventProcessor = new BackgroundEventProcessor(
logContext,
Expand Down Expand Up @@ -540,7 +541,7 @@ private void process(final ConsumerRebalanceListenerCallbackNeededEvent event) {
logContext,
subscriptions,
time,
new RebalanceCallbackMetrics(metrics)
new RebalanceCallbackMetricsManager(metrics)
);
ApiVersions apiVersions = new ApiVersions();
Supplier<NetworkClientDelegate> networkClientDelegateSupplier = () -> new NetworkClientDelegate(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.clients.consumer.RetriableCommitFailedException;
import org.apache.kafka.clients.consumer.internals.metrics.OffsetCommitMetricsManager;
import org.apache.kafka.common.KafkaException;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.errors.DisconnectException;
Expand All @@ -32,11 +33,6 @@
import org.apache.kafka.common.message.OffsetCommitRequestData;
import org.apache.kafka.common.message.OffsetCommitResponseData;
import org.apache.kafka.common.metrics.Metrics;
import org.apache.kafka.common.metrics.Sensor;
import org.apache.kafka.common.metrics.stats.Avg;
import org.apache.kafka.common.metrics.stats.Max;
import org.apache.kafka.common.metrics.stats.Meter;
import org.apache.kafka.common.metrics.stats.WindowedCount;
import org.apache.kafka.common.protocol.Errors;
import org.apache.kafka.common.record.RecordBatch;
import org.apache.kafka.common.requests.AbstractRequest;
Expand Down Expand Up @@ -66,8 +62,6 @@
import java.util.function.BiConsumer;
import java.util.stream.Collectors;

import static org.apache.kafka.clients.consumer.internals.ConsumerUtils.CONSUMER_METRIC_GROUP_PREFIX;
import static org.apache.kafka.clients.consumer.internals.ConsumerUtils.COORDINATOR_METRICS_SUFFIX;
import static org.apache.kafka.clients.consumer.internals.ConsumerUtils.THROW_ON_FETCH_STABLE_OFFSET_UNSUPPORTED;
import static org.apache.kafka.clients.consumer.internals.NetworkClientDelegate.PollResult.EMPTY;
import static org.apache.kafka.common.protocol.Errors.COORDINATOR_LOAD_IN_PROGRESS;
Expand All @@ -88,7 +82,7 @@ public class CommitRequestManager implements RequestManager, MemberStateListener
private final boolean throwOnFetchStableOffsetUnsupported;
final PendingRequests pendingRequests;
private boolean closing = false;
private Sensor commitSensor;
private OffsetCommitMetricsManager metricsManager;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

final?


/**
* Latest member ID and epoch received via the {@link #onMemberEpochUpdated(Optional, Optional)},
Expand Down Expand Up @@ -119,7 +113,6 @@ public CommitRequestManager(
config.getLong(ConsumerConfig.RETRY_BACKOFF_MS_CONFIG),
config.getLong(ConsumerConfig.RETRY_BACKOFF_MAX_MS_CONFIG),
OptionalDouble.empty(),
CONSUMER_METRIC_GROUP_PREFIX,
metrics);
}

Expand All @@ -136,7 +129,6 @@ public CommitRequestManager(
final long retryBackoffMs,
final long retryBackoffMaxMs,
final OptionalDouble jitter,
final String metricGroupPrefix,
final Metrics metrics) {
Objects.requireNonNull(coordinatorRequestManager, "Coordinator is needed upon committing offsets");
this.logContext = logContext;
Expand All @@ -158,8 +150,8 @@ public CommitRequestManager(
this.jitter = jitter;
this.throwOnFetchStableOffsetUnsupported = config.getBoolean(THROW_ON_FETCH_STABLE_OFFSET_UNSUPPORTED);
this.memberInfo = new MemberInfo();
this.metricsManager = new OffsetCommitMetricsManager(metrics);
this.offsetCommitCallbackInvoker = offsetCommitCallbackInvoker;
this.commitSensor = addCommitSensor(metrics, metricGroupPrefix);
}

/**
Expand Down Expand Up @@ -400,23 +392,6 @@ public NetworkClientDelegate.PollResult drainPendingOffsetCommitRequests() {
return new NetworkClientDelegate.PollResult(Long.MAX_VALUE, requests);
}

private Sensor addCommitSensor(Metrics metrics, String metricGrpPrefix) {
String metricGrpName = metricGrpPrefix + COORDINATOR_METRICS_SUFFIX;
Sensor sensor = metrics.sensor("commit-latency");
sensor.add(metrics.metricName("commit-latency-avg",
metricGrpName,
"The average time taken for a commit request"), new Avg());
sensor.add(metrics.metricName("commit-latency-max",
metricGrpName,
"The max time taken for a commit request"), new Max());
sensor.add(new Meter(new WindowedCount(),
metrics.metricName("commit-rate", metricGrpName,
"The number of commit calls per second"),
metrics.metricName("commit-total", metricGrpName,
"The total number of commit calls")));
return sensor;
}

private class OffsetCommitRequestState extends RetriableRequestState {
private final Map<TopicPartition, OffsetAndMetadata> offsets;
private final String groupId;
Expand Down Expand Up @@ -511,8 +486,8 @@ public NetworkClientDelegate.UnsentRequest toUnsentRequest() {
* - fail the future with a non-recoverable KafkaException for all unexpected errors (even if retriable)
*/
@Override
void onResponse(final ClientResponse response) {
commitSensor.record(response.requestLatencyMs());
public void onResponse(final ClientResponse response) {
metricsManager.recordRequestLatency(response.requestLatencyMs());
long currentTimeMs = response.receivedTimeMs();
OffsetCommitResponse commitResponse = (OffsetCommitResponse) response.responseBody();
Set<String> unauthorizedTopics = new HashSet<>();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@
import org.apache.kafka.clients.consumer.OffsetCommitCallback;
import org.apache.kafka.clients.consumer.RetriableCommitFailedException;
import org.apache.kafka.clients.consumer.internals.Utils.TopicPartitionComparator;
import org.apache.kafka.clients.consumer.internals.metrics.RebalanceCallbackMetrics;
import org.apache.kafka.clients.consumer.internals.metrics.RebalanceCallbackMetricsManager;
import org.apache.kafka.common.Cluster;
import org.apache.kafka.common.KafkaException;
import org.apache.kafka.common.Node;
Expand Down Expand Up @@ -235,7 +235,7 @@ public ConsumerCoordinator(GroupRebalanceConfig rebalanceConfig,
logContext,
subscriptions,
time,
new RebalanceCallbackMetrics(metrics, metricGrpPrefix)
new RebalanceCallbackMetricsManager(metrics, metricGrpPrefix)
);
this.metadata.requestUpdate(true);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
package org.apache.kafka.clients.consumer.internals;

import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.internals.metrics.KafkaConsumerMetrics;
import org.apache.kafka.common.metrics.Metrics;
import org.apache.kafka.common.utils.Timer;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@
package org.apache.kafka.clients.consumer.internals;

import org.apache.kafka.clients.consumer.ConsumerRebalanceListener;
import org.apache.kafka.clients.consumer.internals.metrics.RebalanceCallbackMetrics;
import org.apache.kafka.clients.consumer.internals.metrics.RebalanceCallbackMetricsManager;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.errors.InterruptException;
import org.apache.kafka.common.errors.WakeupException;
Expand All @@ -41,16 +41,16 @@ public class ConsumerRebalanceListenerInvoker {
private final Logger log;
private final SubscriptionState subscriptions;
private final Time time;
private final RebalanceCallbackMetrics metrics;
private final RebalanceCallbackMetricsManager metricsManager;

ConsumerRebalanceListenerInvoker(LogContext logContext,
SubscriptionState subscriptions,
Time time,
RebalanceCallbackMetrics metrics) {
RebalanceCallbackMetricsManager metricsManager) {
this.log = logContext.logger(getClass());
this.subscriptions = subscriptions;
this.time = time;
this.metrics = metrics;
this.metricsManager = metricsManager;
}

public Exception invokePartitionsAssigned(final SortedSet<TopicPartition> assignedPartitions) {
Expand All @@ -62,7 +62,7 @@ public Exception invokePartitionsAssigned(final SortedSet<TopicPartition> assign
try {
final long startMs = time.milliseconds();
listener.get().onPartitionsAssigned(assignedPartitions);
metrics.assignCallbackSensor.record(time.milliseconds() - startMs);
metricsManager.recordPartitionsAssignedLatency(time.milliseconds() - startMs);
} catch (WakeupException | InterruptException e) {
throw e;
} catch (Exception e) {
Expand All @@ -88,7 +88,7 @@ public Exception invokePartitionsRevoked(final SortedSet<TopicPartition> revoked
try {
final long startMs = time.milliseconds();
listener.get().onPartitionsRevoked(revokedPartitions);
metrics.revokeCallbackSensor.record(time.milliseconds() - startMs);
metricsManager.recordPartitionsRevokedLatency(time.milliseconds() - startMs);
} catch (WakeupException | InterruptException e) {
throw e;
} catch (Exception e) {
Expand All @@ -114,7 +114,7 @@ public Exception invokePartitionsLost(final SortedSet<TopicPartition> lostPartit
try {
final long startMs = time.milliseconds();
listener.get().onPartitionsLost(lostPartitions);
metrics.loseCallbackSensor.record(time.milliseconds() - startMs);
metricsManager.recordPartitionsLostLatency(time.milliseconds() - startMs);
} catch (WakeupException | InterruptException e) {
throw e;
} catch (Exception e) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,7 @@ public final class ConsumerUtils {
public static final String CONSUMER_JMX_PREFIX = "kafka.consumer";
public static final String CONSUMER_METRIC_GROUP_PREFIX = "consumer";
public static final String COORDINATOR_METRICS_SUFFIX = "-coordinator-metrics";
public static final String CONSUMER_METRICS_SUFFIX = "-metrics";

/**
* A fixed, large enough value will suffice for max.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,10 +23,12 @@
import org.apache.kafka.clients.consumer.internals.events.BackgroundEventHandler;
import org.apache.kafka.clients.consumer.internals.events.ErrorBackgroundEvent;
import org.apache.kafka.clients.consumer.internals.events.GroupMetadataUpdateEvent;
import org.apache.kafka.clients.consumer.internals.metrics.HeartbeatMetricsManager;
import org.apache.kafka.common.Uuid;
import org.apache.kafka.common.errors.GroupAuthorizationException;
import org.apache.kafka.common.errors.RetriableException;
import org.apache.kafka.common.message.ConsumerGroupHeartbeatRequestData;
import org.apache.kafka.common.metrics.Metrics;
import org.apache.kafka.common.protocol.Errors;
import org.apache.kafka.common.requests.ConsumerGroupHeartbeatRequest;
import org.apache.kafka.common.requests.ConsumerGroupHeartbeatResponse;
Expand Down Expand Up @@ -108,6 +110,7 @@ public class HeartbeatRequestManager implements RequestManager {
* sending heartbeat until the next poll.
*/
private final Timer pollTimer;
private final HeartbeatMetricsManager metricsManager;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: I'd either group this with the field below or separate this from pollTimer with a newline. Looks like the javadoc is referring to both fields like this.


private GroupMetadataUpdateEvent previousGroupMetadataUpdateEvent = null;

Expand All @@ -118,7 +121,8 @@ public HeartbeatRequestManager(
final CoordinatorRequestManager coordinatorRequestManager,
final SubscriptionState subscriptions,
final MembershipManager membershipManager,
final BackgroundEventHandler backgroundEventHandler) {
final BackgroundEventHandler backgroundEventHandler,
final Metrics metrics) {
this.coordinatorRequestManager = coordinatorRequestManager;
this.logger = logContext.logger(getClass());
this.membershipManager = membershipManager;
Expand All @@ -130,6 +134,7 @@ public HeartbeatRequestManager(
this.heartbeatRequestState = new HeartbeatRequestState(logContext, time, 0, retryBackoffMs,
retryBackoffMaxMs, maxPollIntervalMs);
this.pollTimer = time.timer(maxPollIntervalMs);
this.metricsManager = new HeartbeatMetricsManager(metrics);
}

// Visible for testing
Expand All @@ -141,7 +146,8 @@ public HeartbeatRequestManager(
final MembershipManager membershipManager,
final HeartbeatState heartbeatState,
final HeartbeatRequestState heartbeatRequestState,
final BackgroundEventHandler backgroundEventHandler) {
final BackgroundEventHandler backgroundEventHandler,
final Metrics metrics) {
this.logger = logContext.logger(this.getClass());
this.maxPollIntervalMs = config.getInt(CommonClientConfigs.MAX_POLL_INTERVAL_MS_CONFIG);
this.coordinatorRequestManager = coordinatorRequestManager;
Expand All @@ -150,6 +156,7 @@ public HeartbeatRequestManager(
this.membershipManager = membershipManager;
this.backgroundEventHandler = backgroundEventHandler;
this.pollTimer = timer;
this.metricsManager = new HeartbeatMetricsManager(metrics);
}

/**
Expand Down Expand Up @@ -245,6 +252,7 @@ private NetworkClientDelegate.UnsentRequest makeHeartbeatRequest(final long curr
NetworkClientDelegate.UnsentRequest request = makeHeartbeatRequest(ignoreResponse);
heartbeatRequestState.onSendAttempt(currentTimeMs);
membershipManager.onHeartbeatRequestSent();
metricsManager.recordHeartbeatSentMs(currentTimeMs);
return request;
}

Expand All @@ -257,6 +265,7 @@ private NetworkClientDelegate.UnsentRequest makeHeartbeatRequest(final boolean i
else
return request.whenComplete((response, exception) -> {
if (response != null) {
metricsManager.recordRequestLatency(response.requestLatencyMs());
onResponse((ConsumerGroupHeartbeatResponse) response.responseBody(), request.handler().completionTimeMs());
} else {
onFailure(exception, request.handler().completionTimeMs());
Expand All @@ -267,6 +276,7 @@ private NetworkClientDelegate.UnsentRequest makeHeartbeatRequest(final boolean i
private NetworkClientDelegate.UnsentRequest logResponse(final NetworkClientDelegate.UnsentRequest request) {
return request.whenComplete((response, exception) -> {
if (response != null) {
metricsManager.recordRequestLatency(response.requestLatencyMs());
Errors error =
Errors.forCode(((ConsumerGroupHeartbeatResponse) response.responseBody()).data().errorCode());
if (error == Errors.NONE)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@
import org.apache.kafka.clients.consumer.OffsetAndTimestamp;
import org.apache.kafka.clients.consumer.OffsetCommitCallback;
import org.apache.kafka.clients.consumer.OffsetResetStrategy;
import org.apache.kafka.clients.consumer.internals.metrics.KafkaConsumerMetrics;
import org.apache.kafka.common.Cluster;
import org.apache.kafka.common.IsolationLevel;
import org.apache.kafka.common.KafkaException;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -194,7 +194,8 @@ protected RequestManagers create() {
coordinator,
subscriptions,
membershipManager,
backgroundEventHandler);
backgroundEventHandler,
metrics);
}

return new RequestManagers(
Expand Down

This file was deleted.

Loading