diff --git a/core/src/test/scala/unit/kafka/server/ReplicaFetcherThreadTest.scala b/core/src/test/scala/unit/kafka/server/ReplicaFetcherThreadTest.scala index da7a8d0fee764..ae294683a0d12 100644 --- a/core/src/test/scala/unit/kafka/server/ReplicaFetcherThreadTest.scala +++ b/core/src/test/scala/unit/kafka/server/ReplicaFetcherThreadTest.scala @@ -21,22 +21,22 @@ import kafka.log.{LogManager, UnifiedLog} import kafka.server.AbstractFetcherThread.ResultWithPartitions import kafka.server.QuotaFactory.UNBOUNDED_QUOTA import kafka.server.epoch.util.MockBlockingSender -import kafka.server.metadata.ZkMetadataCache +import kafka.server.metadata.KRaftMetadataCache import kafka.utils.TestUtils import org.apache.kafka.clients.FetchSessionHandler import org.apache.kafka.common.compress.Compression import org.apache.kafka.common.{TopicIdPartition, TopicPartition, Uuid} -import org.apache.kafka.common.message.{FetchResponseData, UpdateMetadataRequestData} +import org.apache.kafka.common.message.{FetchResponseData} import org.apache.kafka.common.message.OffsetForLeaderEpochRequestData.OffsetForLeaderPartition import org.apache.kafka.common.message.OffsetForLeaderEpochResponseData.EpochEndOffset import org.apache.kafka.common.protocol.{ApiKeys, Errors} import org.apache.kafka.common.record.{CompressionType, MemoryRecords, RecordBatch, RecordValidationStats, SimpleRecord} import org.apache.kafka.common.requests.OffsetsForLeaderEpochResponse.{UNDEFINED_EPOCH, UNDEFINED_EPOCH_OFFSET} -import org.apache.kafka.common.requests.{FetchRequest, FetchResponse, UpdateMetadataRequest} +import org.apache.kafka.common.requests.{FetchRequest, FetchResponse} import org.apache.kafka.common.utils.{LogContext, Time} -import org.apache.kafka.server.BrokerFeatures import org.apache.kafka.server.config.ReplicationConfigs import org.apache.kafka.server.common.{MetadataVersion, OffsetAndEpoch} +import org.apache.kafka.server.common.KRaftVersion import org.apache.kafka.server.network.BrokerEndPoint import org.apache.kafka.storage.internals.log.LogAppendInfo import org.apache.kafka.storage.log.metrics.BrokerTopicStats @@ -67,27 +67,7 @@ class ReplicaFetcherThreadTest { private val brokerEndPoint = new BrokerEndPoint(0, "localhost", 1000) private val failedPartitions = new FailedPartitions - - private val partitionStates = List( - new UpdateMetadataRequestData.UpdateMetadataPartitionState() - .setTopicName("topic1") - .setPartitionIndex(0) - .setControllerEpoch(0) - .setLeader(0) - .setLeaderEpoch(0), - new UpdateMetadataRequestData.UpdateMetadataPartitionState() - .setTopicName("topic2") - .setPartitionIndex(0) - .setControllerEpoch(0) - .setLeader(0) - .setLeaderEpoch(0), - ).asJava - - private val updateMetadataRequest = new UpdateMetadataRequest.Builder(ApiKeys.UPDATE_METADATA.latestVersion(), - 0, 0, 0, partitionStates, Collections.emptyList(), topicIds.asJava).build() - // TODO: support raft code? - private var metadataCache = new ZkMetadataCache(0, MetadataVersion.latestTesting(), BrokerFeatures.createEmpty()) - metadataCache.updateMetadata(0, updateMetadataRequest) + private var metadataCache = new KRaftMetadataCache(0, () => KRaftVersion.LATEST_PRODUCTION) private def initialFetchState(topicId: Option[Uuid], fetchOffset: Long, leaderEpoch: Int = 1): InitialFetchState = { InitialFetchState(topicId = topicId, leader = new BrokerEndPoint(0, "localhost", 9092), @@ -205,8 +185,7 @@ class ReplicaFetcherThreadTest { props.setProperty(ReplicationConfigs.INTER_BROKER_PROTOCOL_VERSION_CONFIG, ibp.version) val config = KafkaConfig.fromProps(props) - metadataCache = new ZkMetadataCache(0, ibp, BrokerFeatures.createEmpty()) - metadataCache.updateMetadata(0, updateMetadataRequest) + metadataCache = new KRaftMetadataCache(0, () => KRaftVersion.LATEST_PRODUCTION) //Setup all dependencies val logManager: LogManager = mock(classOf[LogManager])