Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
39 commits
Select commit Hold shift + click to select a range
357704d
Delete dynamic config that were removed by Kafka
0xffff-zhiyan Nov 25, 2025
c3da308
fix
0xffff-zhiyan Nov 26, 2025
fa4dac9
fix
0xffff-zhiyan Dec 16, 2025
c70f511
fix
0xffff-zhiyan Dec 16, 2025
3767e7a
fix
0xffff-zhiyan Dec 22, 2025
12b3d99
filter existingConfigsSnapshot
0xffff-zhiyan Jan 9, 2026
f734b75
fix
0xffff-zhiyan Jan 14, 2026
1ce628f
fix
0xffff-zhiyan Jan 14, 2026
d2b7916
fix
0xffff-zhiyan Jan 14, 2026
cba7df5
fix
0xffff-zhiyan Jan 14, 2026
56a08ca
fix
0xffff-zhiyan Jan 15, 2026
097c9ba
fix
0xffff-zhiyan Jan 21, 2026
12b4193
fix
0xffff-zhiyan Jan 21, 2026
c3a1f16
fix
0xffff-zhiyan Jan 21, 2026
1744a30
fix format
0xffff-zhiyan Jan 21, 2026
4584c49
rename a test
0xffff-zhiyan Jan 23, 2026
d680dd5
fix
0xffff-zhiyan Jan 27, 2026
4d870c0
fix
0xffff-zhiyan Jan 27, 2026
ed2b8a6
fix
0xffff-zhiyan Jan 30, 2026
1ad2841
fix
0xffff-zhiyan Jan 30, 2026
d33536d
fix
0xffff-zhiyan Feb 2, 2026
c3df5b1
fix
0xffff-zhiyan Feb 9, 2026
4f10cc1
fix
0xffff-zhiyan Feb 9, 2026
aabc28d
fix
0xffff-zhiyan Feb 9, 2026
b5bc16d
fix
0xffff-zhiyan Feb 9, 2026
c3b7998
fix format
0xffff-zhiyan Feb 9, 2026
3d78470
fix format
0xffff-zhiyan Feb 9, 2026
fd9a5b3
fix format
0xffff-zhiyan Feb 9, 2026
1031a62
Merge remote-tracking branch 'upstream/trunk' into KAFKA-19851
0xffff-zhiyan Mar 12, 2026
bdf7975
move DefaultSupportedConfigChecker to Java
0xffff-zhiyan Mar 12, 2026
162ed38
fix
0xffff-zhiyan Mar 12, 2026
4c29f85
fix: add logs
0xffff-zhiyan Mar 17, 2026
77fb3a6
fix: added supportedConfigChecker to MetadataBatchLoader
0xffff-zhiyan Mar 17, 2026
280c856
fix: added test cases for handleLoadSnapshot and MetadataBatchLoader
0xffff-zhiyan Mar 18, 2026
0eda9db
fix test
0xffff-zhiyan Mar 18, 2026
ad1dc2e
fix BROKER type configs whitelist
0xffff-zhiyan Mar 18, 2026
a6dc4da
fix DefaultSupportedConfigChecker using Predicate<>
0xffff-zhiyan Mar 24, 2026
24ef0b3
fix
0xffff-zhiyan Mar 24, 2026
6c972f4
add test for MetadataBatchLoader
0xffff-zhiyan Mar 25, 2026
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
1 change: 1 addition & 0 deletions build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -2655,6 +2655,7 @@ project(':shell') {
implementation project(':core')
implementation project(':metadata')
implementation project(':raft')
implementation project(':server')

implementation libs.jose4j // for SASL/OAUTHBEARER JWT validation
implementation libs.jacksonJakartarsJsonProvider
Expand Down
1 change: 1 addition & 0 deletions checkstyle/import-control.xml
Original file line number Diff line number Diff line change
Expand Up @@ -271,6 +271,7 @@
<allow pkg="org.apache.kafka.queue"/>
<allow pkg="org.apache.kafka.raft"/>
<allow pkg="org.apache.kafka.server.common" />
<allow pkg="org.apache.kafka.server.config" />
<allow pkg="org.apache.kafka.server.fault" />
<allow pkg="org.apache.kafka.server.util" />
<allow pkg="org.apache.kafka.shell"/>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -76,7 +76,13 @@ public Optional<TopicMetadata> topicMetadata(String topicName) {

@Override
public CoordinatorMetadataDelta emptyDelta() {
return new KRaftCoordinatorMetadataDelta(new MetadataDelta(metadataImage));
// Note: supportedConfigChecker is not set because CoordinatorMetadataDelta only exposes topic-related methods.
// No ConfigRecord replay happens through this path, so the checker is never invoked.
return new KRaftCoordinatorMetadataDelta(
new MetadataDelta.Builder()
.setImage(metadataImage)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

The returned value is used in production. Should we set supportedConfigChecker to be the production one?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

The emptyDelta() returns a CoordinatorMetadataDelta which only exposes topic-related methods (createdTopicIds, changedTopicIds, deletedTopicIds). The underlying MetadataDelta is not exposed, and no ConfigRecord replay happens through this path, so the SupportedConfigChecker would never be invoked. I think we don't need to set supportedConfigChecker here, WDYT?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Could we add a comment on why supportedConfigChecker is not set?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

added

.build()
);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2229,7 +2229,11 @@ public void testOnMetadataUpdate() {
verify(coordinator0).onLoaded(CoordinatorMetadataImage.EMPTY);

// Publish a new image.
CoordinatorMetadataDelta delta = new KRaftCoordinatorMetadataDelta(new MetadataDelta(MetadataImage.EMPTY));
CoordinatorMetadataDelta delta = new KRaftCoordinatorMetadataDelta(
new MetadataDelta.Builder()
.setImage(MetadataImage.EMPTY)
.build()
);
CoordinatorMetadataImage newImage = CoordinatorMetadataImage.EMPTY;
runtime.onMetadataUpdate(delta, newImage);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -57,7 +57,9 @@ public void testKRaftCoordinatorDelta() {
.addTopic(deletedTopicId, deletedTopicName, 1)
.addTopic(changedTopicId, changedTopicName, 1)
.build();
MetadataDelta delta = new MetadataDelta(image);
MetadataDelta delta = new MetadataDelta.Builder()
.setImage(image)
.build();
delta.replay(new TopicRecord().setTopicId(topicId).setName(topicName));
delta.replay(new TopicRecord().setTopicId(topicId2).setName(topicName2));
delta.replay(new RemoveTopicRecord().setTopicId(deletedTopicId));
Expand Down Expand Up @@ -107,14 +109,18 @@ public void testEqualsAndHashcode() {
Uuid topicId3 = Uuid.randomUuid();
String topicName3 = "test-topic3";

MetadataDelta delta = new MetadataDelta(MetadataImage.EMPTY);
MetadataDelta delta = new MetadataDelta.Builder()
.setImage(MetadataImage.EMPTY)
.build();
delta.replay(new TopicRecord().setTopicId(topicId).setName(topicName));
delta.replay(new TopicRecord().setTopicId(topicId2).setName(topicName2));

KRaftCoordinatorMetadataDelta coordinatorDelta = new KRaftCoordinatorMetadataDelta(delta);
KRaftCoordinatorMetadataDelta coordinatorDeltaCopy = new KRaftCoordinatorMetadataDelta(delta);

MetadataDelta delta2 = new MetadataDelta(MetadataImage.EMPTY);
MetadataDelta delta2 = new MetadataDelta.Builder()
.setImage(MetadataImage.EMPTY)
.build();
delta.replay(new TopicRecord().setTopicId(topicId3).setName(topicName3));
KRaftCoordinatorMetadataDelta coordinatorDelta2 = new KRaftCoordinatorMetadataDelta(delta2);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,9 @@ public MetadataImageBuilder() {
}

public MetadataImageBuilder(MetadataImage image) {
this.delta = new MetadataDelta(image);
this.delta = new MetadataDelta.Builder()
.setImage(image)
.build();
}

public MetadataImageBuilder addTopic(
Expand Down
1 change: 1 addition & 0 deletions core/src/main/scala/kafka/server/ControllerServer.scala
Original file line number Diff line number Diff line change
Expand Up @@ -248,6 +248,7 @@ class ControllerServer(
setCreateTopicPolicy(createTopicPolicy.toJava).
setAlterConfigPolicy(alterConfigPolicy.toJava).
setConfigurationValidator(new ControllerConfigurationValidator(sharedServer.brokerConfig)).
setSupportedConfigChecker(sharedServer.supportedConfigChecker).
setStaticConfig(config.originals).
setBootstrapMetadata(bootstrapMetadata).
setFatalFaultHandler(sharedServer.fatalQuorumControllerFaultHandler).
Expand Down
8 changes: 5 additions & 3 deletions core/src/main/scala/kafka/server/SharedServer.scala
Original file line number Diff line number Diff line change
Expand Up @@ -30,11 +30,11 @@ import org.apache.kafka.image.loader.MetadataLoader
import org.apache.kafka.image.loader.metrics.MetadataLoaderMetrics
import org.apache.kafka.image.publisher.metrics.SnapshotEmitterMetrics
import org.apache.kafka.image.publisher.{SnapshotEmitter, SnapshotGenerator}
import org.apache.kafka.metadata.ListenerInfo
import org.apache.kafka.metadata.MetadataRecordSerde
import org.apache.kafka.metadata.{SupportedConfigChecker, ListenerInfo, MetadataRecordSerde}
import org.apache.kafka.metadata.properties.MetaPropertiesEnsemble
import org.apache.kafka.raft.{Endpoints, ExternalKRaftMetrics}
import org.apache.kafka.server.{ProcessRole, ServerSocketFactory}
import org.apache.kafka.server.config.DefaultSupportedConfigChecker
import org.apache.kafka.server.common.ApiMessageAndVersion
import org.apache.kafka.server.fault.{FaultHandler, LoggingFaultHandler, ProcessTerminatingFaultHandler}
import org.apache.kafka.server.metrics.{BrokerServerMetrics, KafkaYammerMetrics, NodeMetrics}
Expand Down Expand Up @@ -112,6 +112,7 @@ class SharedServer(
private var usedByController: Boolean = false
val brokerConfig = new KafkaConfig(sharedServerConfig.props, false)
val controllerConfig = new KafkaConfig(sharedServerConfig.props, false)
val supportedConfigChecker: SupportedConfigChecker = new DefaultSupportedConfigChecker()

// Factory for creating request handler pools with shared aggregate thread counter
val requestHandlerPoolFactory = new KafkaRequestHandlerPoolFactory()
Expand Down Expand Up @@ -326,7 +327,8 @@ class SharedServer(
setThreadNamePrefix(s"kafka-${sharedServerConfig.nodeId}-").
setFaultHandler(metadataLoaderFaultHandler).
setHighWaterMarkAccessor(() => _raftManager.client.highWatermark()).
setMetrics(metadataLoaderMetrics)
setMetrics(metadataLoaderMetrics).
setSupportedConfigChecker(supportedConfigChecker)
loader = loaderBuilder.build()
snapshotEmitter = new SnapshotEmitter.Builder().
setNodeId(nodeId).
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -89,7 +89,9 @@ class LocalLeaderEndPointTest extends Logging {
alterPartitionManager = alterPartitionManager
)

val delta = new MetadataDelta(MetadataImage.EMPTY)
val delta = new MetadataDelta.Builder()
.setImage(MetadataImage.EMPTY)
.build()
delta.replay(new FeatureLevelRecord()
.setName(MetadataVersion.FEATURE_NAME)
.setFeatureLevel(MetadataVersion.MINIMUM_VERSION.featureLevel())
Expand Down Expand Up @@ -410,7 +412,9 @@ class LocalLeaderEndPointTest extends Logging {
}

private def bumpLeaderEpoch(): Unit = {
val delta = new MetadataDelta(image)
val delta = new MetadataDelta.Builder()
.setImage(image)
.build()
delta.replay(new PartitionChangeRecord()
.setTopicId(topicId)
.setPartitionId(partition)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,9 @@ class DefaultApiVersionManagerTest {
private val brokerFeatures = BrokerFeatures.createDefault(true)
private val metadataCache = {
val cache = new KRaftMetadataCache(1, () => KRaftVersion.LATEST_PRODUCTION)
val delta = new MetadataDelta(MetadataImage.EMPTY)
val delta = new MetadataDelta.Builder()
.setImage(MetadataImage.EMPTY)
.build()
delta.replay(new FeatureLevelRecord()
.setName(MetadataVersion.FEATURE_NAME)
.setFeatureLevel(MetadataVersion.latestProduction().featureLevel())
Expand Down
28 changes: 21 additions & 7 deletions core/src/test/scala/unit/kafka/server/KafkaApisTest.scala
Original file line number Diff line number Diff line change
Expand Up @@ -230,7 +230,9 @@ class KafkaApisTest extends Logging {

def initializeMetadataCacheWithShareGroupsEnabled(enableShareGroups: Boolean = true): MetadataCache = {
val cache = new KRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_1)
val delta = new MetadataDelta(MetadataImage.EMPTY)
val delta = new MetadataDelta.Builder()
.setImage(MetadataImage.EMPTY)
.build()
delta.replay(new FeatureLevelRecord()
.setName(MetadataVersion.FEATURE_NAME)
.setFeatureLevel(MetadataVersion.MINIMUM_VERSION.featureLevel())
Expand Down Expand Up @@ -7027,7 +7029,9 @@ class KafkaApisTest extends Logging {

metadataCache = {
val cache = new KRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_1)
val delta = new MetadataDelta(MetadataImage.EMPTY)
val delta = new MetadataDelta.Builder()
.setImage(MetadataImage.EMPTY)
.build()
delta.replay(new FeatureLevelRecord()
.setName(MetadataVersion.FEATURE_NAME)
.setFeatureLevel(MetadataVersion.MINIMUM_VERSION.featureLevel())
Expand Down Expand Up @@ -7255,7 +7259,9 @@ class KafkaApisTest extends Logging {

metadataCache = {
val cache = new KRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_1)
val delta = new MetadataDelta(MetadataImage.EMPTY)
val delta = new MetadataDelta.Builder()
.setImage(MetadataImage.EMPTY)
.build()
delta.replay(new FeatureLevelRecord()
.setName(MetadataVersion.FEATURE_NAME)
.setFeatureLevel(MetadataVersion.MINIMUM_VERSION.featureLevel())
Expand Down Expand Up @@ -10847,7 +10853,9 @@ class KafkaApisTest extends Logging {
val requestChannelRequest = buildRequest(new ConsumerGroupHeartbeatRequest.Builder(consumerGroupHeartbeatRequest).build())
metadataCache = {
val cache = new KRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_1)
val delta = new MetadataDelta(MetadataImage.EMPTY)
val delta = new MetadataDelta.Builder()
.setImage(MetadataImage.EMPTY)
.build()
delta.replay(new FeatureLevelRecord()
.setName(MetadataVersion.FEATURE_NAME)
.setFeatureLevel(MetadataVersion.MINIMUM_VERSION.featureLevel())
Expand Down Expand Up @@ -10981,7 +10989,9 @@ class KafkaApisTest extends Logging {
val requestChannelRequest = buildRequest(new StreamsGroupHeartbeatRequest.Builder(streamsGroupHeartbeatRequest).build())
metadataCache = {
val cache = new KRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_1)
val delta = new MetadataDelta(MetadataImage.EMPTY)
val delta = new MetadataDelta.Builder()
.setImage(MetadataImage.EMPTY)
.build()
delta.replay(new FeatureLevelRecord()
.setName(MetadataVersion.FEATURE_NAME)
.setFeatureLevel(MetadataVersion.MINIMUM_VERSION.featureLevel())
Expand Down Expand Up @@ -11544,7 +11554,9 @@ class KafkaApisTest extends Logging {
expectedResponse.groups.add(expectedDescribedGroup)
metadataCache = {
val cache = new KRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_1)
val delta = new MetadataDelta(MetadataImage.EMPTY)
val delta = new MetadataDelta.Builder()
.setImage(MetadataImage.EMPTY)
.build()
delta.replay(new FeatureLevelRecord()
.setName(MetadataVersion.FEATURE_NAME)
.setFeatureLevel(MetadataVersion.MINIMUM_VERSION.featureLevel())
Expand Down Expand Up @@ -11707,7 +11719,9 @@ class KafkaApisTest extends Logging {
expectedResponse.groups.add(expectedDescribedGroup)
metadataCache = {
val cache = new KRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_1)
val delta = new MetadataDelta(MetadataImage.EMPTY)
val delta = new MetadataDelta.Builder()
.setImage(MetadataImage.EMPTY)
.build()
delta.replay(new FeatureLevelRecord()
.setName(MetadataVersion.FEATURE_NAME)
.setFeatureLevel(MetadataVersion.MINIMUM_VERSION.featureLevel())
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -215,7 +215,9 @@ class BrokerMetadataPublisherTest {
)

val topicId = Uuid.randomUuid()
var delta = new MetadataDelta(MetadataImage.EMPTY)
var delta = new MetadataDelta.Builder()
.setImage(MetadataImage.EMPTY)
.build()
delta.replay(new TopicRecord()
.setName(Topic.GROUP_METADATA_TOPIC_NAME)
.setTopicId(topicId)
Expand All @@ -232,7 +234,9 @@ class BrokerMetadataPublisherTest {
)
val image = delta.apply(MetadataProvenance.EMPTY)

delta = new MetadataDelta(image)
delta = new MetadataDelta.Builder()
.setImage(image)
.build()
delta.replay(new RemoveTopicRecord()
.setTopicId(topicId)
)
Expand Down Expand Up @@ -340,7 +344,9 @@ class BrokerMetadataPublisherTest {
)

// Share version 1 is getting passed to features delta.
val delta = new MetadataDelta(image)
val delta = new MetadataDelta.Builder()
.setImage(image)
.build()
delta.replay(new FeatureLevelRecord().setName(ShareVersion.FEATURE_NAME).setFeatureLevel(1))

metadataPublisher.onMetadataUpdate(
Expand Down
Loading
Loading