Skip to content
Closed
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 @@ -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
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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])
Expand Down