Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
34 commits
Select commit Hold shift + click to select a range
96de66e
add interface and tests
jolshan Apr 8, 2024
8f7bc87
Fix 2.12 compilation issue
jolshan Apr 9, 2024
ccb7a50
satisfy scala 2.12
jolshan Apr 12, 2024
b4e2e21
Few fixes to allow the cluster to start up.
jolshan Apr 17, 2024
a31a71e
Merge branch 'trunk' of github.com:apache/kafka into kafka-16308
jolshan Apr 17, 2024
92b7326
supporting supported features
jolshan Apr 18, 2024
0b0b9ab
Some cleanups.
jolshan Apr 25, 2024
70163cb
simplify interface
jolshan May 10, 2024
940653a
Further iterate on interface
jolshan May 11, 2024
50ee633
More fixes
jolshan May 13, 2024
7f162b8
note cleanups
jolshan May 13, 2024
db0b9a0
Merge branch 'trunk' of github.com:apache/kafka into kafka-16308
jolshan May 13, 2024
830a637
fix compilation error
jolshan May 13, 2024
0dcada4
fix default for MV
jolshan May 15, 2024
3988422
updates to the simpler comments
jolshan May 16, 2024
4e76ba4
Renaming and refactoring
jolshan May 20, 2024
1ca5112
Refactor method
jolshan May 20, 2024
f55811f
compilation fixes
jolshan May 20, 2024
435808a
Followups
jolshan May 21, 2024
13a8840
Address feedback
jolshan May 22, 2024
5fc9fea
more comments
jolshan May 22, 2024
4e1989a
Add 0 version back
jolshan May 22, 2024
1a43f18
Re-remove version 0
jolshan May 22, 2024
d86ff02
fix supported versions, set up production features
jolshan May 22, 2024
1d3fab8
Fix when there is no version default flag
jolshan May 23, 2024
e325a88
Fixes
jolshan May 23, 2024
b56707b
More adjustments
jolshan May 23, 2024
2416e95
Simplify production ready
jolshan May 24, 2024
1a474ef
Remove unneeded value
jolshan May 24, 2024
a28ece3
Fix conditional
jolshan May 24, 2024
5f1ac97
Merge branch 'trunk' of github.com:apache/kafka into kafka-16308
jolshan May 28, 2024
24e14ed
Addressing Jun's comments
jolshan May 29, 2024
00b63ad
Fix bug
jolshan May 29, 2024
0704b61
Cleanups
jolshan May 29, 2024
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 @@ -26,7 +26,7 @@
/**
* Represents an immutable basic version range using 2 attributes: min and max, each of type short.
* The min and max attributes need to satisfy 2 rules:
* - they are each expected to be >= 0, as we only consider positive version values to be valid.
* - they are each expected to be >= 0, as we only consider non-negative version values to be valid.
* - max should be >= min.
*
* The class also provides API to convert the version range to a map.
Expand Down
13 changes: 7 additions & 6 deletions core/src/main/scala/kafka/server/ApiVersionManager.scala
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ 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 org.apache.kafka.server.common.FinalizedFeatures

import scala.collection.mutable
import scala.jdk.CollectionConverters._
Expand All @@ -40,7 +40,7 @@ trait ApiVersionManager {
}
def newRequestMetrics: RequestChannel.Metrics = new network.RequestChannel.Metrics(enabledApis)

def features: Features
def features: FinalizedFeatures
}

object ApiVersionManager {
Expand Down Expand Up @@ -73,21 +73,22 @@ object ApiVersionManager {
* @param brokerFeatures the broker features
* @param enableUnstableLastVersion whether to enable unstable last version, see [[KafkaConfig.unstableApiVersionsEnabled]]
* @param zkMigrationEnabled whether to enable zk migration, see [[KafkaConfig.migrationEnabled]]
* @param featuresProvider a provider to the finalized features supported
*/
class SimpleApiVersionManager(
val listenerType: ListenerType,
val enabledApis: collection.Set[ApiKeys],
brokerFeatures: org.apache.kafka.common.feature.Features[SupportedVersionRange],
val enableUnstableLastVersion: Boolean,
val zkMigrationEnabled: Boolean,
val featuresProvider: () => Features
val featuresProvider: () => FinalizedFeatures
Comment thread
jolshan marked this conversation as resolved.
) extends ApiVersionManager {

def this(
listenerType: ListenerType,
enableUnstableLastVersion: Boolean,
zkMigrationEnabled: Boolean,
featuresProvider: () => Features
featuresProvider: () => FinalizedFeatures
) = {
this(
listenerType,
Expand All @@ -113,7 +114,7 @@ class SimpleApiVersionManager(
)
}

override def features: Features = featuresProvider.apply()
override def features: FinalizedFeatures = featuresProvider.apply()
}

/**
Expand Down Expand Up @@ -164,5 +165,5 @@ class DefaultApiVersionManager(
)
}

override def features: Features = metadataCache.features()
override def features: FinalizedFeatures = metadataCache.features()
}
12 changes: 8 additions & 4 deletions core/src/main/scala/kafka/server/BrokerFeatures.scala
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ package kafka.server

import kafka.utils.Logging
import org.apache.kafka.common.feature.{Features, SupportedVersionRange}
import org.apache.kafka.server.common.Features.PRODUCTION_FEATURES
import org.apache.kafka.server.common.MetadataVersion

import java.util
Expand Down Expand Up @@ -75,16 +76,19 @@ object BrokerFeatures extends Logging {
}

def defaultSupportedFeatures(unstableMetadataVersionsEnabled: Boolean): Features[SupportedVersionRange] = {

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.

Should we add a unit test for this change?

@jolshan jolshan May 16, 2024

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

I can look into adding something. The problem is I don't add any new production features in this change. I would add one in the next PR.

Features.supportedFeatures(
java.util.Collections.singletonMap(MetadataVersion.FEATURE_NAME,
val features = new util.HashMap[String, SupportedVersionRange]()
features.put(MetadataVersion.FEATURE_NAME,
new SupportedVersionRange(
MetadataVersion.MINIMUM_KRAFT_VERSION.featureLevel(),
if (unstableMetadataVersionsEnabled) {
MetadataVersion.latestTesting.featureLevel
} else {
MetadataVersion.latestProduction.featureLevel
}
)))
}))
PRODUCTION_FEATURES.forEach { feature =>
features.put(feature.featureName, new SupportedVersionRange(0, feature.latestProduction()))
}
Features.supportedFeatures(features)
}

def createEmpty(): BrokerFeatures = {
Expand Down
4 changes: 2 additions & 2 deletions core/src/main/scala/kafka/server/MetadataCache.scala
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ import org.apache.kafka.admin.BrokerMetadata
import org.apache.kafka.common.message.{MetadataResponseData, UpdateMetadataRequestData}
import org.apache.kafka.common.network.ListenerName
import org.apache.kafka.common.{Cluster, Node, TopicPartition, Uuid}
import org.apache.kafka.server.common.{Features, MetadataVersion}
import org.apache.kafka.server.common.{FinalizedFeatures, MetadataVersion}

import java.util
import scala.collection._
Expand Down Expand Up @@ -109,7 +109,7 @@ trait MetadataCache {

def getRandomAliveBrokerId: Option[Int]

def features(): Features
def features(): FinalizedFeatures
}

object MetadataCache {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ import org.apache.kafka.common.protocol.Errors
import org.apache.kafka.common.requests.MetadataResponse
import org.apache.kafka.image.MetadataImage
import org.apache.kafka.metadata.{BrokerRegistration, PartitionRegistration, Replicas}
import org.apache.kafka.server.common.{Features, MetadataVersion}
import org.apache.kafka.server.common.{FinalizedFeatures, MetadataVersion}

import java.util
import java.util.concurrent.ThreadLocalRandom
Expand Down Expand Up @@ -539,9 +539,9 @@ class KRaftMetadataCache(val brokerId: Int) extends MetadataCache with Logging w

override def metadataVersion(): MetadataVersion = _currentImage.features().metadataVersion()

override def features(): Features = {
override def features(): FinalizedFeatures = {
val image = _currentImage
new Features(image.features().metadataVersion(),
new FinalizedFeatures(image.features().metadataVersion(),
image.features().finalizedVersions(),
image.highestOffsetAndEpoch().offset,
true)
Expand Down
14 changes: 7 additions & 7 deletions core/src/main/scala/kafka/server/metadata/ZkMetadataCache.scala
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ import org.apache.kafka.common.network.ListenerName
import org.apache.kafka.common.protocol.Errors
import org.apache.kafka.common.requests.{AbstractControlRequest, ApiVersionsResponse, MetadataResponse, UpdateMetadataRequest}
import org.apache.kafka.common.security.auth.SecurityProtocol
import org.apache.kafka.server.common.{Features, MetadataVersion}
import org.apache.kafka.server.common.{FinalizedFeatures, MetadataVersion}

import java.util.concurrent.{ThreadLocalRandom, TimeUnit}
import scala.concurrent.TimeoutException
Expand All @@ -53,7 +53,7 @@ class FeatureCacheUpdateException(message: String) extends RuntimeException(mess
trait ZkFinalizedFeatureCache {
def waitUntilFeatureEpochOrThrow(minExpectedEpoch: Long, timeoutMs: Long): Unit

def getFeatureOption: Option[Features]
def getFeatureOption: Option[FinalizedFeatures]
}

case class MetadataSnapshot(partitionStates: mutable.AnyRefMap[String, mutable.LongMap[UpdateMetadataPartitionState]],
Expand Down Expand Up @@ -177,7 +177,7 @@ class ZkMetadataCache(
private val stateChangeLogger = new StateChangeLogger(brokerId, inControllerContext = false, None)

// Features are updated via ZK notification (see FinalizedFeatureChangeListener)
@volatile private var _features: Option[Features] = Option.empty
@volatile private var _features: Option[FinalizedFeatures] = Option.empty
private val featureLock = new ReentrantLock()
private val featureCond = featureLock.newCondition()

Expand Down Expand Up @@ -617,9 +617,9 @@ class ZkMetadataCache(

override def metadataVersion(): MetadataVersion = metadataVersion

override def features(): Features = _features match {
override def features(): FinalizedFeatures = _features match {
case Some(features) => features
case None => new Features(metadataVersion,
case None => new FinalizedFeatures(metadataVersion,
Collections.emptyMap(),
ApiVersionsResponse.UNKNOWN_FINALIZED_FEATURES_EPOCH,
false)
Expand All @@ -639,7 +639,7 @@ class ZkMetadataCache(
* not modified.
*/
def updateFeaturesOrThrow(latestFeatures: Map[String, Short], latestEpoch: Long): Unit = {
val latest = new Features(metadataVersion,
val latest = new FinalizedFeatures(metadataVersion,
latestFeatures.map(kv => (kv._1, kv._2.asInstanceOf[java.lang.Short])).asJava,
latestEpoch,
false)
Expand Down Expand Up @@ -711,5 +711,5 @@ class ZkMetadataCache(
}
}

override def getFeatureOption: Option[Features] = _features
override def getFeatureOption: Option[FinalizedFeatures] = _features
}
126 changes: 108 additions & 18 deletions core/src/main/scala/kafka/tools/StorageTool.scala
Original file line number Diff line number Diff line change
Expand Up @@ -28,18 +28,18 @@ import net.sourceforge.argparse4j.inf.Namespace
import org.apache.kafka.common.Uuid
import org.apache.kafka.common.utils.Utils
import org.apache.kafka.metadata.bootstrap.{BootstrapDirectory, BootstrapMetadata}
import org.apache.kafka.server.common.{ApiMessageAndVersion, MetadataVersion}
import org.apache.kafka.server.common.{ApiMessageAndVersion, Features, MetadataVersion}
import org.apache.kafka.common.metadata.FeatureLevelRecord
import org.apache.kafka.common.metadata.UserScramCredentialRecord
import org.apache.kafka.common.security.scram.internals.ScramMechanism
import org.apache.kafka.common.security.scram.internals.ScramFormatter
import org.apache.kafka.server.config.ReplicationConfigs
import org.apache.kafka.metadata.properties.MetaPropertiesEnsemble.VerificationFlag
import org.apache.kafka.metadata.properties.{MetaProperties, MetaPropertiesEnsemble, MetaPropertiesVersion, PropertiesUtils}
import org.apache.kafka.server.common.FeatureVersion

import java.util
import java.util.Base64
import java.util.Optional
import java.util.{Base64, Collections, Optional}
import scala.collection.mutable
import scala.jdk.CollectionConverters._
import scala.collection.mutable.ArrayBuffer
Expand All @@ -60,24 +60,30 @@ object StorageTool extends Logging {
case "format" =>
val directories = configToLogDirectories(config.get)
val clusterId = namespace.getString("cluster_id")
val metadataVersion = getMetadataVersion(namespace,
Option(config.get.originals.get(ReplicationConfigs.INTER_BROKER_PROTOCOL_VERSION_CONFIG)).map(_.toString))
if (!metadataVersion.isKRaftSupported) {
throw new TerseFailure(s"Must specify a valid KRaft metadata.version of at least ${MetadataVersion.IBP_3_0_IV0}.")
}
if (!metadataVersion.isProduction) {
if (config.get.unstableMetadataVersionsEnabled) {
System.out.println(s"WARNING: using pre-production metadata.version $metadataVersion.")
} else {
throw new TerseFailure(s"The metadata.version $metadataVersion is not ready for production use yet.")
}
}
val metaProperties = new MetaProperties.Builder().
setVersion(MetaPropertiesVersion.V1).
setClusterId(clusterId).
setNodeId(config.get.nodeId).
build()
val metadataRecords : ArrayBuffer[ApiMessageAndVersion] = ArrayBuffer()
val specifiedFeatures: util.List[String] = namespace.getList("feature")
val releaseVersionFlagSpecified = namespace.getString("release_version") != null
if (releaseVersionFlagSpecified && specifiedFeatures != null) {
throw new TerseFailure("Both --release-version and --feature were set. Only one of the two flags can be set.")
}
val featureNamesAndLevelsMap = featureNamesAndLevels(Option(specifiedFeatures).getOrElse(Collections.emptyList).asScala.toList)
val metadataVersion = getMetadataVersion(namespace, featureNamesAndLevelsMap,
Option(config.get.originals.get(ReplicationConfigs.INTER_BROKER_PROTOCOL_VERSION_CONFIG)).map(_.toString))
validateMetadataVersion(metadataVersion, config)
// Get all other features, validate, and create records for them
// Use latest default for features if --release-version is not specified
generateFeatureRecords(
metadataRecords,
metadataVersion,
featureNamesAndLevelsMap,
Features.PRODUCTION_FEATURES.asScala.toList,
releaseVersionFlagSpecified
)
getUserScramCredentialRecords(namespace).foreach(userScramCredentialRecords => {
if (!metadataVersion.isScramSupported) {
throw new TerseFailure(s"SCRAM is only supported in metadata.version ${MetadataVersion.IBP_3_5_IV2} or later.")
Expand All @@ -86,6 +92,7 @@ object StorageTool extends Logging {
metadataRecords.append(new ApiMessageAndVersion(record, 0.toShort))
}
})

val bootstrapMetadata = buildBootstrapMetadata(metadataVersion, Some(metadataRecords), "format command")
val ignoreFormatted = namespace.getBoolean("ignore_formatted")
if (!configToSelfManagedMode(config.get)) {
Expand All @@ -109,6 +116,52 @@ object StorageTool extends Logging {
}
}

private def validateMetadataVersion(metadataVersion: MetadataVersion, config: Option[KafkaConfig]): Unit = {
if (!metadataVersion.isKRaftSupported) {
throw new TerseFailure(s"Must specify a valid KRaft metadata.version of at least ${MetadataVersion.IBP_3_0_IV0}.")
}
if (!metadataVersion.isProduction) {
if (config.get.unstableMetadataVersionsEnabled) {
System.out.println(s"WARNING: using pre-production metadata.version $metadataVersion.")
} else {
throw new TerseFailure(s"The metadata.version $metadataVersion is not ready for production use yet.")
}
}
}

private[tools] def generateFeatureRecords(metadataRecords: ArrayBuffer[ApiMessageAndVersion],
metadataVersion: MetadataVersion,
specifiedFeatures: Map[String, java.lang.Short],
allFeatures: List[Features],
releaseVersionSpecified: Boolean): Unit = {
// If we are using --release-version, the default is based on the metadata version.
val metadataVersionForDefault = if (releaseVersionSpecified) metadataVersion else MetadataVersion.LATEST_PRODUCTION

val allNonZeroFeaturesAndLevels: ArrayBuffer[FeatureVersion] = mutable.ArrayBuffer[FeatureVersion]()

allFeatures.foreach { feature =>
val level: java.lang.Short = specifiedFeatures.getOrElse(feature.featureName, feature.defaultValue(metadataVersionForDefault))
// Only set feature records for levels greater than 0. 0 is assumed if there is no record. Throw an error if level < 0.
if (level != 0) {
allNonZeroFeaturesAndLevels.append(feature.fromFeatureLevel(level))
}
}
val featuresMap = Features.featureImplsToMap(allNonZeroFeaturesAndLevels.asJava)
featuresMap.put(MetadataVersion.FEATURE_NAME, metadataVersion.featureLevel)

try {
for (feature <- allNonZeroFeaturesAndLevels) {
// In order to validate, we need all feature versions set.
Features.validateVersion(feature, featuresMap)
metadataRecords.append(new ApiMessageAndVersion(new FeatureLevelRecord().
setName(feature.featureName).
setFeatureLevel(feature.featureLevel), 0.toShort))
}
} catch {
case e: Throwable => throw new TerseFailure(e.getMessage)
}
}

def parseArguments(args: Array[String]): Namespace = {
val parser = ArgumentParsers.
newArgumentParser("kafka-storage", /* defaultHelp */ true, /* prefixChars */ "-", /* fromFilePrefix */ "@").
Expand Down Expand Up @@ -141,6 +194,9 @@ object StorageTool extends Logging {
formatParser.addArgument("--release-version", "-r").
action(store()).
help(s"A KRaft release version to use for the initial metadata.version. The minimum is ${MetadataVersion.IBP_3_0_IV0}, the default is ${MetadataVersion.LATEST_PRODUCTION}")
formatParser.addArgument("--feature", "-f").
help("A feature upgrade we should perform, in feature=level format. For example: `metadata.version=5`.").
Comment thread
jolshan marked this conversation as resolved.
action(append());

parser.parseArgsOrFail(args)
}
Expand All @@ -156,16 +212,27 @@ object StorageTool extends Logging {

def getMetadataVersion(
namespace: Namespace,
featureNamesAndLevelsMap: Map[String, java.lang.Short],
defaultVersionString: Option[String]
): MetadataVersion = {
val defaultValue = defaultVersionString match {
case Some(versionString) => MetadataVersion.fromVersionString(versionString)
case None => MetadataVersion.LATEST_PRODUCTION
}

Option(namespace.getString("release_version"))
.map(ver => MetadataVersion.fromVersionString(ver))
.getOrElse(defaultValue)
val releaseVersionTag = Option(namespace.getString("release_version"))
val featureTag = featureNamesAndLevelsMap.get(MetadataVersion.FEATURE_NAME)

(releaseVersionTag, featureTag) match {
case (Some(_), Some(_)) => // We should throw an error before we hit this case, but include for completeness
throw new IllegalArgumentException("Both --release_version and --feature were set. Only one of the two flags can be set.")
Comment thread
jolshan marked this conversation as resolved.
case (Some(version), None) =>
MetadataVersion.fromVersionString(version)
case (None, Some(level)) =>
MetadataVersion.fromFeatureLevel(level)
case (None, None) =>
defaultValue
}
}

private def getUserScramCredentialRecord(
Expand Down Expand Up @@ -469,4 +536,27 @@ object StorageTool extends Logging {
}
0
}

private def parseNameAndLevel(input: String): (String, java.lang.Short) = {
val equalsIndex = input.indexOf("=")
if (equalsIndex < 0)
throw new RuntimeException("Can't parse feature=level string " + input + ": equals sign not found.")
val name = input.substring(0, equalsIndex).trim
val levelString = input.substring(equalsIndex + 1).trim
try {
levelString.toShort
} catch {
case _: Throwable =>
throw new RuntimeException("Can't parse feature=level string " + input + ": " + "unable to parse " + levelString + " as a short.")
}
(name, levelString.toShort)
}

def featureNamesAndLevels(features: List[String]): Map[String, java.lang.Short] = {
features.map { (feature: String) =>
// Ensure the feature exists
val nameAndLevel = parseNameAndLevel(feature)
(nameAndLevel._1, nameAndLevel._2)
}.toMap
}
}
Loading