diff --git a/checkstyle/import-control-share.xml b/checkstyle/import-control-share.xml index 89422b705456f..8abc1b9f7b779 100644 --- a/checkstyle/import-control-share.xml +++ b/checkstyle/import-control-share.xml @@ -46,6 +46,7 @@ + diff --git a/core/src/main/java/kafka/server/share/DelayedShareFetch.java b/core/src/main/java/kafka/server/share/DelayedShareFetch.java index 0d25d850ff71a..4cb9ce0cf4241 100644 --- a/core/src/main/java/kafka/server/share/DelayedShareFetch.java +++ b/core/src/main/java/kafka/server/share/DelayedShareFetch.java @@ -25,6 +25,7 @@ import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.requests.FetchRequest; import org.apache.kafka.server.purgatory.DelayedOperation; +import org.apache.kafka.server.share.fetch.DelayedShareFetchGroupKey; import org.apache.kafka.server.share.fetch.ShareFetchData; import org.apache.kafka.server.storage.log.FetchIsolation; import org.apache.kafka.server.storage.log.FetchPartitionData; diff --git a/core/src/main/java/kafka/server/share/SharePartition.java b/core/src/main/java/kafka/server/share/SharePartition.java index 6d828697af928..740bf697de405 100644 --- a/core/src/main/java/kafka/server/share/SharePartition.java +++ b/core/src/main/java/kafka/server/share/SharePartition.java @@ -36,6 +36,8 @@ import org.apache.kafka.common.utils.Time; import org.apache.kafka.coordinator.group.GroupConfigManager; import org.apache.kafka.server.share.acknowledge.ShareAcknowledgementBatch; +import org.apache.kafka.server.share.fetch.DelayedShareFetchGroupKey; +import org.apache.kafka.server.share.fetch.DelayedShareFetchKey; import org.apache.kafka.server.share.fetch.ShareAcquiredRecords; import org.apache.kafka.server.share.persister.GroupTopicPartitionData; import org.apache.kafka.server.share.persister.PartitionAllData; diff --git a/core/src/main/java/kafka/server/share/SharePartitionManager.java b/core/src/main/java/kafka/server/share/SharePartitionManager.java index 804af3e1c87dc..4288dd55703d7 100644 --- a/core/src/main/java/kafka/server/share/SharePartitionManager.java +++ b/core/src/main/java/kafka/server/share/SharePartitionManager.java @@ -46,6 +46,9 @@ import org.apache.kafka.server.share.context.FinalContext; import org.apache.kafka.server.share.context.ShareFetchContext; import org.apache.kafka.server.share.context.ShareSessionContext; +import org.apache.kafka.server.share.fetch.DelayedShareFetchGroupKey; +import org.apache.kafka.server.share.fetch.DelayedShareFetchKey; +import org.apache.kafka.server.share.fetch.DelayedShareFetchPartitionKey; import org.apache.kafka.server.share.fetch.ShareFetchData; import org.apache.kafka.server.share.persister.Persister; import org.apache.kafka.server.share.session.ShareSession; diff --git a/core/src/main/scala/kafka/cluster/Partition.scala b/core/src/main/scala/kafka/cluster/Partition.scala index 636a64c7f01cd..e432ead8edb27 100755 --- a/core/src/main/scala/kafka/cluster/Partition.scala +++ b/core/src/main/scala/kafka/cluster/Partition.scala @@ -25,7 +25,7 @@ import kafka.log._ import kafka.log.remote.RemoteLogManager import kafka.server._ import kafka.server.metadata.{KRaftMetadataCache, ZkMetadataCache} -import kafka.server.share.{DelayedShareFetch, DelayedShareFetchPartitionKey} +import kafka.server.share.DelayedShareFetch import kafka.utils.CoreUtils.{inReadLock, inWriteLock} import kafka.utils._ import kafka.zookeeper.ZooKeeperClientException @@ -46,6 +46,7 @@ import org.apache.kafka.server.common.{MetadataVersion, RequestLocal} import org.apache.kafka.storage.internals.log.{AppendOrigin, FetchDataInfo, LeaderHwChange, LogAppendInfo, LogOffsetMetadata, LogOffsetSnapshot, LogOffsetsListener, LogReadInfo, LogStartOffsetIncrementReason, VerificationGuard} import org.apache.kafka.server.metrics.KafkaMetricsGroup import org.apache.kafka.server.purgatory.{DelayedOperationPurgatory, TopicPartitionOperationKey} +import org.apache.kafka.server.share.fetch.DelayedShareFetchPartitionKey import org.apache.kafka.server.storage.log.{FetchIsolation, FetchParams} import org.apache.kafka.storage.internals.checkpoint.OffsetCheckpoints import org.slf4j.event.Level diff --git a/core/src/main/scala/kafka/server/ReplicaManager.scala b/core/src/main/scala/kafka/server/ReplicaManager.scala index 530dbe5b53f67..d6792bdb1d00a 100644 --- a/core/src/main/scala/kafka/server/ReplicaManager.scala +++ b/core/src/main/scala/kafka/server/ReplicaManager.scala @@ -25,7 +25,7 @@ import kafka.server.HostedPartition.Online import kafka.server.QuotaFactory.QuotaManagers import kafka.server.ReplicaManager.{AtMinIsrPartitionCountMetricName, FailedIsrUpdatesPerSecMetricName, IsrExpandsPerSecMetricName, IsrShrinksPerSecMetricName, LeaderCountMetricName, OfflineReplicaCountMetricName, PartitionCountMetricName, PartitionsWithLateTransactionsCountMetricName, ProducerIdCountMetricName, ReassigningPartitionsMetricName, UnderMinIsrPartitionCountMetricName, UnderReplicatedPartitionsMetricName, createLogReadResult, isListOffsetsTimestampUnsupported} import kafka.server.metadata.ZkMetadataCache -import kafka.server.share.{DelayedShareFetch, DelayedShareFetchKey, DelayedShareFetchPartitionKey} +import kafka.server.share.DelayedShareFetch import kafka.utils._ import kafka.zk.KafkaZkClient import org.apache.kafka.common.errors._ @@ -61,6 +61,7 @@ import org.apache.kafka.server.common.MetadataVersion._ import org.apache.kafka.server.metrics.KafkaMetricsGroup import org.apache.kafka.server.network.BrokerEndPoint import org.apache.kafka.server.purgatory.{DelayedOperationKey, DelayedOperationPurgatory, TopicPartitionOperationKey} +import org.apache.kafka.server.share.fetch.{DelayedShareFetchKey, DelayedShareFetchPartitionKey} import org.apache.kafka.server.storage.log.{FetchParams, FetchPartitionData} import org.apache.kafka.server.util.{Scheduler, ShutdownableThread} import org.apache.kafka.storage.internals.checkpoint.{LazyOffsetCheckpoints, OffsetCheckpointFile, OffsetCheckpoints} diff --git a/core/src/test/java/kafka/server/share/DelayedShareFetchTest.java b/core/src/test/java/kafka/server/share/DelayedShareFetchTest.java index 77d6db80158d4..0c7b488f18020 100644 --- a/core/src/test/java/kafka/server/share/DelayedShareFetchTest.java +++ b/core/src/test/java/kafka/server/share/DelayedShareFetchTest.java @@ -29,6 +29,7 @@ import org.apache.kafka.common.requests.FetchRequest; import org.apache.kafka.server.purgatory.DelayedOperationKey; import org.apache.kafka.server.purgatory.DelayedOperationPurgatory; +import org.apache.kafka.server.share.fetch.DelayedShareFetchGroupKey; import org.apache.kafka.server.share.fetch.ShareAcquiredRecords; import org.apache.kafka.server.share.fetch.ShareFetchData; import org.apache.kafka.server.storage.log.FetchIsolation; @@ -72,7 +73,6 @@ import static org.mockito.Mockito.times; import static org.mockito.Mockito.when; - public class DelayedShareFetchTest { private static final int MAX_WAIT_MS = 5000; private static final int MAX_FETCH_RECORDS = 100; diff --git a/core/src/test/java/kafka/server/share/SharePartitionManagerTest.java b/core/src/test/java/kafka/server/share/SharePartitionManagerTest.java index 4a9283d9dd3ad..67c2a6cce778c 100644 --- a/core/src/test/java/kafka/server/share/SharePartitionManagerTest.java +++ b/core/src/test/java/kafka/server/share/SharePartitionManagerTest.java @@ -60,6 +60,8 @@ import org.apache.kafka.server.share.context.FinalContext; import org.apache.kafka.server.share.context.ShareFetchContext; import org.apache.kafka.server.share.context.ShareSessionContext; +import org.apache.kafka.server.share.fetch.DelayedShareFetchGroupKey; +import org.apache.kafka.server.share.fetch.DelayedShareFetchKey; import org.apache.kafka.server.share.fetch.ShareAcquiredRecords; import org.apache.kafka.server.share.fetch.ShareFetchData; import org.apache.kafka.server.share.persister.NoOpShareStatePersister; diff --git a/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala b/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala index 59e9cdc6e6b0a..f03e253777ea5 100644 --- a/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala +++ b/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala @@ -46,7 +46,7 @@ import java.nio.ByteBuffer import java.util.Optional import java.util.concurrent.{ConcurrentHashMap, CountDownLatch, Semaphore} import kafka.server.metadata.{KRaftMetadataCache, ZkMetadataCache} -import kafka.server.share.{DelayedShareFetch, DelayedShareFetchPartitionKey} +import kafka.server.share.DelayedShareFetch import org.apache.kafka.clients.ClientResponse import org.apache.kafka.common.compress.Compression import org.apache.kafka.common.config.TopicConfig @@ -59,6 +59,7 @@ import org.apache.kafka.server.common.{ControllerRequestCompletionHandler, Metad import org.apache.kafka.server.common.MetadataVersion.IBP_2_6_IV0 import org.apache.kafka.server.metrics.KafkaYammerMetrics import org.apache.kafka.server.purgatory.{DelayedOperationPurgatory, TopicPartitionOperationKey} +import org.apache.kafka.server.share.fetch.DelayedShareFetchPartitionKey import org.apache.kafka.server.storage.log.{FetchIsolation, FetchParams} import org.apache.kafka.server.util.{KafkaScheduler, MockTime} import org.apache.kafka.storage.internals.checkpoint.OffsetCheckpoints diff --git a/core/src/main/java/kafka/server/share/DelayedShareFetchGroupKey.java b/share/src/main/java/org/apache/kafka/server/share/fetch/DelayedShareFetchGroupKey.java similarity index 94% rename from core/src/main/java/kafka/server/share/DelayedShareFetchGroupKey.java rename to share/src/main/java/org/apache/kafka/server/share/fetch/DelayedShareFetchGroupKey.java index 3847518298599..0fe1a4774f501 100644 --- a/core/src/main/java/kafka/server/share/DelayedShareFetchGroupKey.java +++ b/share/src/main/java/org/apache/kafka/server/share/fetch/DelayedShareFetchGroupKey.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package kafka.server.share; +package org.apache.kafka.server.share.fetch; import org.apache.kafka.common.Uuid; @@ -28,7 +28,7 @@ public class DelayedShareFetchGroupKey implements DelayedShareFetchKey { private final Uuid topicId; private final int partition; - DelayedShareFetchGroupKey(String groupId, Uuid topicId, int partition) { + public DelayedShareFetchGroupKey(String groupId, Uuid topicId, int partition) { this.groupId = groupId; this.topicId = topicId; this.partition = partition; diff --git a/core/src/main/java/kafka/server/share/DelayedShareFetchKey.java b/share/src/main/java/org/apache/kafka/server/share/fetch/DelayedShareFetchKey.java similarity index 95% rename from core/src/main/java/kafka/server/share/DelayedShareFetchKey.java rename to share/src/main/java/org/apache/kafka/server/share/fetch/DelayedShareFetchKey.java index 7979514a83091..c7b8f05fa18cd 100644 --- a/core/src/main/java/kafka/server/share/DelayedShareFetchKey.java +++ b/share/src/main/java/org/apache/kafka/server/share/fetch/DelayedShareFetchKey.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package kafka.server.share; +package org.apache.kafka.server.share.fetch; import org.apache.kafka.server.purgatory.DelayedOperationKey; diff --git a/core/src/main/java/kafka/server/share/DelayedShareFetchPartitionKey.java b/share/src/main/java/org/apache/kafka/server/share/fetch/DelayedShareFetchPartitionKey.java similarity index 97% rename from core/src/main/java/kafka/server/share/DelayedShareFetchPartitionKey.java rename to share/src/main/java/org/apache/kafka/server/share/fetch/DelayedShareFetchPartitionKey.java index af7f92aa76e7f..584613cde17e4 100644 --- a/core/src/main/java/kafka/server/share/DelayedShareFetchPartitionKey.java +++ b/share/src/main/java/org/apache/kafka/server/share/fetch/DelayedShareFetchPartitionKey.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package kafka.server.share; +package org.apache.kafka.server.share.fetch; import org.apache.kafka.common.Uuid; diff --git a/core/src/test/java/kafka/server/share/DelayedShareFetchKeyTest.java b/share/src/test/java/org/apache/kafka/server/share/fetch/DelayedShareFetchKeyTest.java similarity index 98% rename from core/src/test/java/kafka/server/share/DelayedShareFetchKeyTest.java rename to share/src/test/java/org/apache/kafka/server/share/fetch/DelayedShareFetchKeyTest.java index bca879d105cf3..a23d7455dd075 100644 --- a/core/src/test/java/kafka/server/share/DelayedShareFetchKeyTest.java +++ b/share/src/test/java/org/apache/kafka/server/share/fetch/DelayedShareFetchKeyTest.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package kafka.server.share; +package org.apache.kafka.server.share.fetch; import org.apache.kafka.common.TopicIdPartition; import org.apache.kafka.common.TopicPartition;