Skip to content
Open
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
8 changes: 6 additions & 2 deletions core/src/main/scala/kafka/server/DynamicBrokerConfig.scala
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ import org.apache.kafka.common.utils.internals.BufferSupplier
import org.apache.kafka.common.utils.Utils
import org.apache.kafka.common.utils.internals.ConfigUtils
import org.apache.kafka.network.SocketServer
import org.apache.kafka.raft.KafkaRaftClient
import org.apache.kafka.raft.{KRaftConfigs, KafkaRaftClient}
import org.apache.kafka.server.{DynamicThreadPool, ProcessRole}
import org.apache.kafka.server.common.{ApiMessageAndVersion, DirectoryEventHandler}
import org.apache.kafka.server.config.{BrokerReconfigurable => JBrokerReconfigurable, DynamicConfig, DynamicProducerStateManagerConfig, ServerConfigs, ServerLogConfigs, DynamicBrokerConfig => JDynamicBrokerConfig}
Expand Down Expand Up @@ -741,7 +741,11 @@ class DynamicMetricsReporters(brokerId: Int, config: KafkaConfig, metrics: Metri

class DynamicMetricReporterState(brokerId: Int, config: KafkaConfig, metrics: Metrics, clusterId: String) {
private[server] val dynamicConfig = config.dynamicConfig
private val propsOverride = Map[String, AnyRef](ServerConfigs.BROKER_ID_CONFIG -> brokerId.toString)
// broker.id will no longer be passed from 5.0 (KIP-1232).
private val propsOverride = Map[String, AnyRef](
ServerConfigs.BROKER_ID_CONFIG -> brokerId.toString,
KRaftConfigs.NODE_ID_CONFIG -> brokerId.toString
)
private[server] val currentReporters = mutable.Map[String, MetricsReporter]()
createReporters(config, clusterId, metricsReporterClasses(dynamicConfig.currentKafkaConfig.values()).asJava,
Collections.emptyMap[String, Object])
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -699,6 +699,40 @@ class DynamicBrokerConfigTest {
assertTrue(m.currentReporters.isEmpty)
}

@Test
def testMetricReportersAreConfiguredWithBothIdConfigs(): Unit = {
val nodeId = 0
val reporterName = classOf[TestConfigCapturingReporter].getName
// createBrokerConfig only sets node.id, so broker.id is only visible through the synonym mechanism
val origProps = TestUtils.createBrokerConfig(nodeId)
origProps.put(MetricConfigs.METRIC_REPORTER_CLASSES_CONFIG, reporterName)

val config = KafkaConfig(origProps)
config.dynamicConfig.initialize(None)
val m = new DynamicMetricsReporters(nodeId, config, mock(classOf[Metrics]), "clusterId")
config.dynamicConfig.addReconfigurable(m)

def updateReporters(reporterNames: String): Unit = {
val props = new Properties()
props.put(MetricConfigs.METRIC_REPORTER_CLASSES_CONFIG, reporterNames)
config.dynamicConfig.updateDefaultConfig(props)
}

def assertBothIdsArePassedAsStrings(): Unit = {
val configs = m.currentReporters(reporterName).asInstanceOf[TestConfigCapturingReporter].configs
assertEquals(nodeId.toString, configs.get(KRaftConfigs.NODE_ID_CONFIG))
assertEquals(nodeId.toString, configs.get(ServerConfigs.BROKER_ID_CONFIG))
}

assertBothIdsArePassedAsStrings()

// a dynamic update recreates the reporter with the parsed config values, where the ids are integers
updateReporters("")
assertTrue(m.currentReporters.isEmpty)
updateReporters(reporterName)
assertBothIdsArePassedAsStrings()
}

@Test
def testDynamicLogLocalRetentionMsConfig(): Unit = {
val props = TestUtils.createBrokerConfig(0, port = 8181)
Expand Down Expand Up @@ -1336,6 +1370,16 @@ class TestDynamicThreadPool extends BrokerReconfigurable {
}
}

class TestConfigCapturingReporter extends MetricsReporter {
var configs: util.Map[String, _] = _

override def configure(configs: util.Map[String, _]): Unit = this.configs = configs
override def init(metrics: util.List[KafkaMetric]): Unit = {}
override def metricChange(metric: KafkaMetric): Unit = {}
override def metricRemoval(metric: KafkaMetric): Unit = {}
override def close(): Unit = {}
}

class TestExporterOnly extends MetricsReporter with ClientTelemetryExporterProvider {
override def configure(configs: util.Map[String, _]): Unit = {}
override def init(metrics: util.List[KafkaMetric]): Unit = {}
Expand Down
Loading