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 @@ -114,21 +114,24 @@ public static ApiVersionsResponse createApiVersionsResponse(
NodeApiVersions controllerApiVersions,
ListenerType listenerType,
boolean enableUnstableLastVersion,
boolean zkMigrationEnabled
boolean zkMigrationEnabled,
boolean clientTelemetryEnabled
) {
ApiVersionCollection apiKeys;
if (controllerApiVersions != null) {
apiKeys = intersectForwardableApis(
listenerType,
minRecordVersion,
controllerApiVersions.allSupportedApiVersions(),
enableUnstableLastVersion
enableUnstableLastVersion,
clientTelemetryEnabled
);
} else {
apiKeys = filterApis(
minRecordVersion,
listenerType,
enableUnstableLastVersion
enableUnstableLastVersion,
clientTelemetryEnabled
);
}

Expand Down Expand Up @@ -167,16 +170,21 @@ public static ApiVersionCollection filterApis(
RecordVersion minRecordVersion,
ApiMessageType.ListenerType listenerType
) {
return filterApis(minRecordVersion, listenerType, false);
return filterApis(minRecordVersion, listenerType, false, false);
}

public static ApiVersionCollection filterApis(
RecordVersion minRecordVersion,
ApiMessageType.ListenerType listenerType,
boolean enableUnstableLastVersion
boolean enableUnstableLastVersion,
boolean clientTelemetryEnabled
) {
ApiVersionCollection apiKeys = new ApiVersionCollection();
for (ApiKeys apiKey : ApiKeys.apisForListener(listenerType)) {
// Skip telemetry APIs if client telemetry is disabled.
if ((apiKey == ApiKeys.GET_TELEMETRY_SUBSCRIPTIONS || apiKey == ApiKeys.PUSH_TELEMETRY) && !clientTelemetryEnabled)
continue;

if (apiKey.minRequiredInterBrokerMagic <= minRecordVersion.value) {
apiKey.toApiVersion(enableUnstableLastVersion).ifPresent(apiKeys::add);
}
Expand All @@ -203,13 +211,15 @@ public static ApiVersionCollection collectApis(
* @param minRecordVersion min inter broker magic
* @param activeControllerApiVersions controller ApiVersions
* @param enableUnstableLastVersion whether unstable versions should be advertised or not
* @param clientTelemetryEnabled whether client telemetry is enabled or not
* @return commonly agreed ApiVersion collection
*/
public static ApiVersionCollection intersectForwardableApis(
final ApiMessageType.ListenerType listenerType,
final RecordVersion minRecordVersion,
final Map<ApiKeys, ApiVersion> activeControllerApiVersions,
boolean enableUnstableLastVersion
boolean enableUnstableLastVersion,
boolean clientTelemetryEnabled
) {
ApiVersionCollection apiKeys = new ApiVersionCollection();
for (ApiKeys apiKey : ApiKeys.apisForListener(listenerType)) {
Expand All @@ -220,6 +230,10 @@ public static ApiVersionCollection intersectForwardableApis(
continue;
}

// Skip telemetry APIs if client telemetry is disabled.
if ((apiKey == ApiKeys.GET_TELEMETRY_SUBSCRIPTIONS || apiKey == ApiKeys.PUSH_TELEMETRY) && !clientTelemetryEnabled)
continue;

final ApiVersion finalApiVersion;
if (!apiKey.forwardable) {
finalApiVersion = brokerApiVersion.get();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -103,7 +103,8 @@ public void shouldHaveCommonlyAgreedApiVersionResponseWithControllerOnForwardabl
ApiMessageType.ListenerType.ZK_BROKER,
RecordVersion.current(),
activeControllerApiVersions,
true
true,
false
);

verifyVersions(forwardableAPIKey.id, minVersion, maxVersion, commonResponse);
Expand All @@ -123,7 +124,8 @@ public void shouldCreateApiResponseOnlyWithKeysSupportedByMagicValue() {
null,
ListenerType.ZK_BROKER,
true,
false
false,
true
);
verifyApiKeysForMagic(response, RecordBatch.MAGIC_VALUE_V1);
assertEquals(10, response.throttleTimeMs());
Expand All @@ -144,7 +146,8 @@ public void shouldReturnFeatureKeysWhenMagicIsCurrentValueAndThrottleMsIsDefault
null,
ListenerType.ZK_BROKER,
true,
false
false,
true
);

verifyApiKeysForMagic(response, RecordBatch.MAGIC_VALUE_V1);
Expand Down Expand Up @@ -174,7 +177,8 @@ public void shouldReturnAllKeysWhenMagicIsCurrentValueAndThrottleMsIsDefaultThro
null,
listenerType,
true,
false
false,
true
);
assertEquals(new HashSet<>(ApiKeys.apisForListener(listenerType)), apiKeysInResponse(response));
assertEquals(AbstractResponse.DEFAULT_THROTTLE_TIME, response.throttleTimeMs());
Expand All @@ -183,6 +187,40 @@ public void shouldReturnAllKeysWhenMagicIsCurrentValueAndThrottleMsIsDefaultThro
assertEquals(ApiVersionsResponse.UNKNOWN_FINALIZED_FEATURES_EPOCH, response.data().finalizedFeaturesEpoch());
}

@Test
public void shouldCreateApiResponseWithTelemetryWhenEnabled() {
ApiVersionsResponse response = ApiVersionsResponse.createApiVersionsResponse(
10,
RecordVersion.V1,
Features.emptySupportedFeatures(),
Collections.emptyMap(),
ApiVersionsResponse.UNKNOWN_FINALIZED_FEATURES_EPOCH,
null,
ListenerType.BROKER,
true,
false,
true
);
verifyApiKeysForTelemetry(response, 2);
}

@Test
public void shouldNotCreateApiResponseWithTelemetryWhenDisabled() {
ApiVersionsResponse response = ApiVersionsResponse.createApiVersionsResponse(
10,
RecordVersion.V1,
Features.emptySupportedFeatures(),
Collections.emptyMap(),
ApiVersionsResponse.UNKNOWN_FINALIZED_FEATURES_EPOCH,
null,
ListenerType.BROKER,
true,
false,
false
);
verifyApiKeysForTelemetry(response, 0);
}

@Test
public void testMetadataQuorumApisAreDisabled() {
ApiVersionsResponse response = ApiVersionsResponse.createApiVersionsResponse(
Expand All @@ -194,7 +232,8 @@ public void testMetadataQuorumApisAreDisabled() {
null,
ListenerType.ZK_BROKER,
true,
false
false,
true
);

// Ensure that APIs needed for the KRaft mode are not exposed through ApiVersions until we are ready for them
Expand Down Expand Up @@ -254,6 +293,16 @@ private void verifyApiKeysForMagic(ApiVersionsResponse response, Byte maxMagic)
}
}

private void verifyApiKeysForTelemetry(ApiVersionsResponse response, int expectedCount) {
int count = 0;
for (ApiVersion version : response.data().apiKeys()) {
if (version.apiKey() == ApiKeys.GET_TELEMETRY_SUBSCRIPTIONS.id || version.apiKey() == ApiKeys.PUSH_TELEMETRY.id) {
count++;
}
}
assertEquals(expectedCount, count);
}

private HashSet<ApiKeys> apiKeysInResponse(ApiVersionsResponse apiVersions) {
HashSet<ApiKeys> apiKeys = new HashSet<>();
for (ApiVersion version : apiVersions.data().apiKeys()) {
Expand Down
4 changes: 2 additions & 2 deletions clients/src/test/java/org/apache/kafka/test/TestUtils.java
Original file line number Diff line number Diff line change
Expand Up @@ -600,7 +600,7 @@ public static ApiVersionsResponse defaultApiVersionsResponse(
) {
return createApiVersionsResponse(
throttleTimeMs,
ApiVersionsResponse.filterApis(RecordVersion.current(), listenerType, true),
ApiVersionsResponse.filterApis(RecordVersion.current(), listenerType, true, true),
Features.emptySupportedFeatures(),
false
);
Expand All @@ -613,7 +613,7 @@ public static ApiVersionsResponse defaultApiVersionsResponse(
) {
return createApiVersionsResponse(
throttleTimeMs,
ApiVersionsResponse.filterApis(RecordVersion.current(), listenerType, enableUnstableLastVersion),
ApiVersionsResponse.filterApis(RecordVersion.current(), listenerType, enableUnstableLastVersion, true),
Features.emptySupportedFeatures(),
false
);
Expand Down
18 changes: 14 additions & 4 deletions core/src/main/scala/kafka/server/ApiVersionManager.scala
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import org.apache.kafka.common.feature.SupportedVersionRange
import org.apache.kafka.common.message.ApiMessageType.ListenerType
import org.apache.kafka.common.protocol.ApiKeys
import org.apache.kafka.common.requests.ApiVersionsResponse
import org.apache.kafka.server.ClientMetricsManager
import org.apache.kafka.server.common.Features

import scala.jdk.CollectionConverters._
Expand All @@ -47,15 +48,17 @@ object ApiVersionManager {
config: KafkaConfig,
forwardingManager: Option[ForwardingManager],
supportedFeatures: BrokerFeatures,
metadataCache: MetadataCache
metadataCache: MetadataCache,
clientMetricsManager: Option[ClientMetricsManager]
): ApiVersionManager = {
new DefaultApiVersionManager(
listenerType,
forwardingManager,
supportedFeatures,
metadataCache,
config.unstableApiVersionsEnabled,
config.migrationEnabled
config.migrationEnabled,
clientMetricsManager
)
}
}
Expand Down Expand Up @@ -123,14 +126,16 @@ class SimpleApiVersionManager(
* @param metadataCache the metadata cache, used to get the finalized features and the metadata version
* @param enableUnstableLastVersion whether to enable unstable last version, see [[KafkaConfig.unstableApiVersionsEnabled]]
* @param zkMigrationEnabled whether to enable zk migration, see [[KafkaConfig.migrationEnabled]]
* @param clientMetricsManager the client metrics manager, helps to determine whether client telemetry is enabled
*/
class DefaultApiVersionManager(
val listenerType: ListenerType,
forwardingManager: Option[ForwardingManager],
brokerFeatures: BrokerFeatures,
metadataCache: MetadataCache,
val enableUnstableLastVersion: Boolean,
val zkMigrationEnabled: Boolean = false
val zkMigrationEnabled: Boolean = false,
val clientMetricsManager: Option[ClientMetricsManager] = None
Comment thread
apoorvmittal10 marked this conversation as resolved.
) extends ApiVersionManager {

val enabledApis = ApiKeys.apisForListener(listenerType).asScala
Expand All @@ -139,6 +144,10 @@ class DefaultApiVersionManager(
val supportedFeatures = brokerFeatures.supportedFeatures
val finalizedFeatures = metadataCache.features()
val controllerApiVersions = forwardingManager.flatMap(_.controllerApiVersions)
val clientTelemetryEnabled = clientMetricsManager match {
case Some(manager) => manager.isTelemetryReceiverConfigured
case None => false
}

ApiVersionsResponse.createApiVersionsResponse(
throttleTimeMs,
Expand All @@ -149,7 +158,8 @@ class DefaultApiVersionManager(
controllerApiVersions.orNull,
listenerType,
enableUnstableLastVersion,
zkMigrationEnabled
zkMigrationEnabled,
clientTelemetryEnabled
)
}

Expand Down
7 changes: 3 additions & 4 deletions core/src/main/scala/kafka/server/BrokerServer.scala
Original file line number Diff line number Diff line change
Expand Up @@ -235,14 +235,15 @@ class BrokerServer(
)
clientToControllerChannelManager.start()
forwardingManager = new ForwardingManagerImpl(clientToControllerChannelManager)

clientMetricsManager = new ClientMetricsManager(clientMetricsReceiverPlugin, config.clientTelemetryMaxBytes, time)

val apiVersionManager = ApiVersionManager(
ListenerType.BROKER,
config,
Some(forwardingManager),
brokerFeatures,
metadataCache
metadataCache,
Some(clientMetricsManager)
Comment thread
apoorvmittal10 marked this conversation as resolved.
)

// Create and start the socket server acceptor threads so that the bound port is known.
Expand Down Expand Up @@ -347,8 +348,6 @@ class BrokerServer(
config, Some(clientToControllerChannelManager), None, None,
groupCoordinator, transactionCoordinator)

clientMetricsManager = new ClientMetricsManager(clientMetricsReceiverPlugin, config.clientTelemetryMaxBytes, time)

dynamicConfigHandlers = Map[String, ConfigHandler](
ConfigType.Topic -> new TopicConfigHandler(replicaManager, config, quotaManagers, None),
ConfigType.Broker -> new BrokerConfigHandler(config, quotaManagers),
Expand Down
34 changes: 28 additions & 6 deletions core/src/main/scala/kafka/server/KafkaApis.scala
Original file line number Diff line number Diff line change
Expand Up @@ -3775,16 +3775,38 @@ class KafkaApis(val requestChannel: RequestChannel,

}

// Just a place holder for now.
def handleGetTelemetrySubscriptionsRequest(request: RequestChannel.Request): Unit = {
requestHelper.sendMaybeThrottle(request, request.body[GetTelemetrySubscriptionsRequest].getErrorResponse(Errors.UNSUPPORTED_VERSION.exception))
CompletableFuture.completedFuture[Unit](())
val subscriptionRequest = request.body[GetTelemetrySubscriptionsRequest]

clientMetricsManager match {
case Some(metricsManager) =>
try {
requestHelper.sendMaybeThrottle(request, metricsManager.processGetTelemetrySubscriptionRequest(subscriptionRequest, request.context))
} catch {
case _: Exception =>
requestHelper.sendMaybeThrottle(request, subscriptionRequest.getErrorResponse(Errors.INVALID_REQUEST.exception))
}
case None =>
info("Received get telemetry client request for zookeeper based cluster")
requestHelper.sendMaybeThrottle(request, subscriptionRequest.getErrorResponse(Errors.UNSUPPORTED_VERSION.exception))
}
}

// Just a place holder for now.
def handlePushTelemetryRequest(request: RequestChannel.Request): Unit = {
requestHelper.sendMaybeThrottle(request, request.body[PushTelemetryRequest].getErrorResponse(Errors.UNSUPPORTED_VERSION.exception))
CompletableFuture.completedFuture[Unit](())
val pushTelemetryRequest = request.body[PushTelemetryRequest]

clientMetricsManager match {
case Some(metricsManager) =>
try {
requestHelper.sendMaybeThrottle(request, metricsManager.processPushTelemetryRequest(pushTelemetryRequest, request.context))
} catch {
case _: Exception =>
requestHelper.sendMaybeThrottle(request, pushTelemetryRequest.getErrorResponse(Errors.INVALID_REQUEST.exception))
}
case None =>
info("Received push telemetry client request for zookeeper based cluster")
requestHelper.sendMaybeThrottle(request, pushTelemetryRequest.getErrorResponse(Errors.UNSUPPORTED_VERSION.exception))
}
}

private def updateRecordConversionStats(request: RequestChannel.Request,
Expand Down
3 changes: 2 additions & 1 deletion core/src/main/scala/kafka/server/KafkaServer.scala
Original file line number Diff line number Diff line change
Expand Up @@ -365,7 +365,8 @@ class KafkaServer(
config,
forwardingManager,
brokerFeatures,
metadataCache
metadataCache,
None
)

// Create and start the socket server acceptor threads so that the bound port is known.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -87,7 +87,9 @@ class BrokerApiVersionsCommandTest extends KafkaServerTestHarness {
else s"${apiVersion.minVersion} to ${apiVersion.maxVersion}"
val usableVersion = nodeApiVersions.latestUsableVersion(apiKey)

val line = s"\t${apiKey.name}(${apiKey.id}): $versionRangeStr [usable: $usableVersion]$terminator"
val line =
if (apiKey == ApiKeys.GET_TELEMETRY_SUBSCRIPTIONS || apiKey == ApiKeys.PUSH_TELEMETRY) s"\t${apiKey.name}(${apiKey.id}): UNSUPPORTED$terminator"
else s"\t${apiKey.name}(${apiKey.id}): $versionRangeStr [usable: $usableVersion]$terminator"
assertTrue(lineIter.hasNext)
assertEquals(line, lineIter.next())
} else {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,7 @@ abstract class AbstractApiVersionsRequestTest(cluster: ClusterInstance) {
apiVersionsResponse: ApiVersionsResponse,
listenerName: ListenerName = cluster.clientListener(),
enableUnstableLastVersion: Boolean = false,
clientTelemetryEnabled: Boolean = false,
apiVersion: Short = ApiKeys.API_VERSIONS.latestVersion
): Unit = {
if (cluster.isKRaftTest && apiVersion >= 3) {
Expand All @@ -100,7 +101,8 @@ abstract class AbstractApiVersionsRequestTest(cluster: ClusterInstance) {
ApiMessageType.ListenerType.BROKER,
RecordVersion.current,
NodeApiVersions.create(ApiKeys.controllerApis().asScala.map(ApiVersionsResponse.toApiVersion).asJava).allSupportedApiVersions(),
enableUnstableLastVersion
enableUnstableLastVersion,
clientTelemetryEnabled
)
}

Expand Down
Loading