diff --git a/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteBufferStreamReader.java b/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteBufferStreamReader.java index 0bea1452d..632fa8793 100644 --- a/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteBufferStreamReader.java +++ b/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteBufferStreamReader.java @@ -31,9 +31,9 @@ import org.apache.celeborn.common.network.util.NettyUtils; import org.apache.celeborn.common.protocol.PbReadAddCredit; import org.apache.celeborn.plugin.flink.buffer.CreditListener; import org.apache.celeborn.plugin.flink.buffer.TransferBufferPool; +import org.apache.celeborn.plugin.flink.client.CelebornBufferStream; +import org.apache.celeborn.plugin.flink.client.FlinkShuffleClientImpl; import org.apache.celeborn.plugin.flink.protocol.ReadData; -import org.apache.celeborn.plugin.flink.readclient.CelebornBufferStream; -import org.apache.celeborn.plugin.flink.readclient.FlinkShuffleClientImpl; public class RemoteBufferStreamReader extends CreditListener { private static Logger logger = LoggerFactory.getLogger(RemoteBufferStreamReader.class); diff --git a/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleInputGateDelegation.java b/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleInputGateDelegation.java index 95ad20495..6162d0bbc 100644 --- a/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleInputGateDelegation.java +++ b/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleInputGateDelegation.java @@ -54,7 +54,7 @@ import org.apache.celeborn.common.exception.PartitionUnRetryAbleException; import org.apache.celeborn.common.identity.UserIdentifier; import org.apache.celeborn.plugin.flink.buffer.BufferPacker; import org.apache.celeborn.plugin.flink.buffer.TransferBufferPool; -import org.apache.celeborn.plugin.flink.readclient.FlinkShuffleClientImpl; +import org.apache.celeborn.plugin.flink.client.FlinkShuffleClientImpl; import org.apache.celeborn.plugin.flink.utils.BufferUtils; public class RemoteShuffleInputGateDelegation { diff --git a/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleOutputGate.java b/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleOutputGate.java index f695af14d..dc131f971 100644 --- a/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleOutputGate.java +++ b/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/RemoteShuffleOutputGate.java @@ -35,7 +35,7 @@ import org.apache.celeborn.common.identity.UserIdentifier; import org.apache.celeborn.common.protocol.PartitionLocation; import org.apache.celeborn.plugin.flink.buffer.BufferHeader; import org.apache.celeborn.plugin.flink.buffer.BufferPacker; -import org.apache.celeborn.plugin.flink.readclient.FlinkShuffleClientImpl; +import org.apache.celeborn.plugin.flink.client.FlinkShuffleClientImpl; import org.apache.celeborn.plugin.flink.utils.BufferUtils; import org.apache.celeborn.plugin.flink.utils.Utils; diff --git a/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/readclient/CelebornBufferStream.java b/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/client/CelebornBufferStream.java similarity index 99% rename from client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/readclient/CelebornBufferStream.java rename to client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/client/CelebornBufferStream.java index 061a918c2..b63757d2d 100644 --- a/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/readclient/CelebornBufferStream.java +++ b/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/client/CelebornBufferStream.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.celeborn.plugin.flink.readclient; +package org.apache.celeborn.plugin.flink.client; import java.io.IOException; import java.nio.ByteBuffer; diff --git a/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/readclient/FlinkShuffleClientImpl.java b/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/client/FlinkShuffleClientImpl.java similarity index 99% rename from client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/readclient/FlinkShuffleClientImpl.java rename to client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/client/FlinkShuffleClientImpl.java index 5602d1aac..bdf64691b 100644 --- a/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/readclient/FlinkShuffleClientImpl.java +++ b/client-flink/common/src/main/java/org/apache/celeborn/plugin/flink/client/FlinkShuffleClientImpl.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.celeborn.plugin.flink.readclient; +package org.apache.celeborn.plugin.flink.client; import java.io.IOException; import java.nio.ByteBuffer; diff --git a/client-flink/common/src/test/java/org/apache/celeborn/plugin/flink/FlinkShuffleClientImplSuiteJ.java b/client-flink/common/src/test/java/org/apache/celeborn/plugin/flink/FlinkShuffleClientImplSuiteJ.java index d0e0dee28..40f2d5e07 100644 --- a/client-flink/common/src/test/java/org/apache/celeborn/plugin/flink/FlinkShuffleClientImplSuiteJ.java +++ b/client-flink/common/src/test/java/org/apache/celeborn/plugin/flink/FlinkShuffleClientImplSuiteJ.java @@ -37,7 +37,7 @@ import org.apache.celeborn.common.network.client.TransportClient; import org.apache.celeborn.common.network.client.TransportClientFactory; import org.apache.celeborn.common.protocol.PartitionLocation; import org.apache.celeborn.common.protocol.message.StatusCode; -import org.apache.celeborn.plugin.flink.readclient.FlinkShuffleClientImpl; +import org.apache.celeborn.plugin.flink.client.FlinkShuffleClientImpl; public class FlinkShuffleClientImplSuiteJ { static int BufferSize = 64; diff --git a/client-flink/common/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleOutputGateSuiteJ.java b/client-flink/common/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleOutputGateSuiteJ.java index e5cf37c04..33c31b2e4 100644 --- a/client-flink/common/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleOutputGateSuiteJ.java +++ b/client-flink/common/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleOutputGateSuiteJ.java @@ -32,7 +32,7 @@ import org.junit.Before; import org.junit.Test; import org.apache.celeborn.common.protocol.PartitionLocation; -import org.apache.celeborn.plugin.flink.readclient.FlinkShuffleClientImpl; +import org.apache.celeborn.plugin.flink.client.FlinkShuffleClientImpl; public class RemoteShuffleOutputGateSuiteJ { private final RemoteShuffleOutputGate remoteShuffleOutputGate = diff --git a/client-flink/flink-1.16/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java b/client-flink/flink-1.16/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java index 023f0b29e..14c2a7b59 100644 --- a/client-flink/flink-1.16/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java +++ b/client-flink/flink-1.16/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java @@ -63,7 +63,7 @@ import org.junit.Test; import org.apache.celeborn.common.CelebornConf; import org.apache.celeborn.plugin.flink.buffer.BufferPacker; import org.apache.celeborn.plugin.flink.buffer.DataBuffer; -import org.apache.celeborn.plugin.flink.readclient.FlinkShuffleClientImpl; +import org.apache.celeborn.plugin.flink.client.FlinkShuffleClientImpl; import org.apache.celeborn.plugin.flink.utils.BufferUtils; public class RemoteShuffleResultPartitionSuiteJ { diff --git a/client-flink/flink-1.17/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java b/client-flink/flink-1.17/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java index 2147c6888..c02e1ad8b 100644 --- a/client-flink/flink-1.17/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java +++ b/client-flink/flink-1.17/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java @@ -63,7 +63,7 @@ import org.junit.Test; import org.apache.celeborn.common.CelebornConf; import org.apache.celeborn.plugin.flink.buffer.BufferPacker; import org.apache.celeborn.plugin.flink.buffer.DataBuffer; -import org.apache.celeborn.plugin.flink.readclient.FlinkShuffleClientImpl; +import org.apache.celeborn.plugin.flink.client.FlinkShuffleClientImpl; import org.apache.celeborn.plugin.flink.utils.BufferUtils; public class RemoteShuffleResultPartitionSuiteJ { diff --git a/client-flink/flink-1.18/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java b/client-flink/flink-1.18/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java index 2147c6888..c02e1ad8b 100644 --- a/client-flink/flink-1.18/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java +++ b/client-flink/flink-1.18/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java @@ -63,7 +63,7 @@ import org.junit.Test; import org.apache.celeborn.common.CelebornConf; import org.apache.celeborn.plugin.flink.buffer.BufferPacker; import org.apache.celeborn.plugin.flink.buffer.DataBuffer; -import org.apache.celeborn.plugin.flink.readclient.FlinkShuffleClientImpl; +import org.apache.celeborn.plugin.flink.client.FlinkShuffleClientImpl; import org.apache.celeborn.plugin.flink.utils.BufferUtils; public class RemoteShuffleResultPartitionSuiteJ { diff --git a/client-flink/flink-1.19/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java b/client-flink/flink-1.19/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java index 19718f42c..fece38532 100644 --- a/client-flink/flink-1.19/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java +++ b/client-flink/flink-1.19/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java @@ -63,7 +63,7 @@ import org.junit.Test; import org.apache.celeborn.common.CelebornConf; import org.apache.celeborn.plugin.flink.buffer.BufferPacker; import org.apache.celeborn.plugin.flink.buffer.DataBuffer; -import org.apache.celeborn.plugin.flink.readclient.FlinkShuffleClientImpl; +import org.apache.celeborn.plugin.flink.client.FlinkShuffleClientImpl; import org.apache.celeborn.plugin.flink.utils.BufferUtils; public class RemoteShuffleResultPartitionSuiteJ { diff --git a/client-flink/flink-1.20/src/main/java/org/apache/celeborn/plugin/flink/tiered/CelebornChannelBufferReader.java b/client-flink/flink-1.20/src/main/java/org/apache/celeborn/plugin/flink/tiered/CelebornChannelBufferReader.java index a387cf9fb..534a2d880 100644 --- a/client-flink/flink-1.20/src/main/java/org/apache/celeborn/plugin/flink/tiered/CelebornChannelBufferReader.java +++ b/client-flink/flink-1.20/src/main/java/org/apache/celeborn/plugin/flink/tiered/CelebornChannelBufferReader.java @@ -42,9 +42,9 @@ import org.apache.celeborn.common.protocol.PbNotifyRequiredSegment; import org.apache.celeborn.common.protocol.PbReadAddCredit; import org.apache.celeborn.common.util.JavaUtils; import org.apache.celeborn.plugin.flink.ShuffleResourceDescriptor; +import org.apache.celeborn.plugin.flink.client.CelebornBufferStream; +import org.apache.celeborn.plugin.flink.client.FlinkShuffleClientImpl; import org.apache.celeborn.plugin.flink.protocol.SubPartitionReadData; -import org.apache.celeborn.plugin.flink.readclient.CelebornBufferStream; -import org.apache.celeborn.plugin.flink.readclient.FlinkShuffleClientImpl; /** * Wrap the {@link CelebornBufferStream}, utilize in flink hybrid shuffle integration strategy now. diff --git a/client-flink/flink-1.20/src/main/java/org/apache/celeborn/plugin/flink/tiered/CelebornTierConsumerAgent.java b/client-flink/flink-1.20/src/main/java/org/apache/celeborn/plugin/flink/tiered/CelebornTierConsumerAgent.java index 8d06ba77c..0c5c454ee 100644 --- a/client-flink/flink-1.20/src/main/java/org/apache/celeborn/plugin/flink/tiered/CelebornTierConsumerAgent.java +++ b/client-flink/flink-1.20/src/main/java/org/apache/celeborn/plugin/flink/tiered/CelebornTierConsumerAgent.java @@ -63,7 +63,7 @@ import org.apache.celeborn.common.identity.UserIdentifier; import org.apache.celeborn.plugin.flink.RemoteShuffleResource; import org.apache.celeborn.plugin.flink.ShuffleResourceDescriptor; import org.apache.celeborn.plugin.flink.buffer.ReceivedNoHeaderBufferPacker; -import org.apache.celeborn.plugin.flink.readclient.FlinkShuffleClientImpl; +import org.apache.celeborn.plugin.flink.client.FlinkShuffleClientImpl; public class CelebornTierConsumerAgent implements TierConsumerAgent { diff --git a/client-flink/flink-1.20/src/main/java/org/apache/celeborn/plugin/flink/tiered/CelebornTierProducerAgent.java b/client-flink/flink-1.20/src/main/java/org/apache/celeborn/plugin/flink/tiered/CelebornTierProducerAgent.java index fc2c14982..8cf9edebd 100644 --- a/client-flink/flink-1.20/src/main/java/org/apache/celeborn/plugin/flink/tiered/CelebornTierProducerAgent.java +++ b/client-flink/flink-1.20/src/main/java/org/apache/celeborn/plugin/flink/tiered/CelebornTierProducerAgent.java @@ -54,7 +54,7 @@ import org.apache.celeborn.common.protocol.PartitionLocation; import org.apache.celeborn.plugin.flink.buffer.BufferHeader; import org.apache.celeborn.plugin.flink.buffer.BufferPacker; import org.apache.celeborn.plugin.flink.buffer.ReceivedNoHeaderBufferPacker; -import org.apache.celeborn.plugin.flink.readclient.FlinkShuffleClientImpl; +import org.apache.celeborn.plugin.flink.client.FlinkShuffleClientImpl; import org.apache.celeborn.plugin.flink.utils.BufferUtils; import org.apache.celeborn.plugin.flink.utils.Utils; diff --git a/client-flink/flink-1.20/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java b/client-flink/flink-1.20/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java index 0f17a156b..99f502d87 100644 --- a/client-flink/flink-1.20/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java +++ b/client-flink/flink-1.20/src/test/java/org/apache/celeborn/plugin/flink/RemoteShuffleResultPartitionSuiteJ.java @@ -64,7 +64,7 @@ import org.junit.Test; import org.apache.celeborn.common.CelebornConf; import org.apache.celeborn.plugin.flink.buffer.BufferPacker; import org.apache.celeborn.plugin.flink.buffer.DataBuffer; -import org.apache.celeborn.plugin.flink.readclient.FlinkShuffleClientImpl; +import org.apache.celeborn.plugin.flink.client.FlinkShuffleClientImpl; import org.apache.celeborn.plugin.flink.utils.BufferUtils; public class RemoteShuffleResultPartitionSuiteJ { diff --git a/tests/flink-it/src/test/scala/org/apache/celeborn/tests/flink/HeartbeatTest.scala b/tests/flink-it/src/test/scala/org/apache/celeborn/tests/flink/HeartbeatTest.scala index dbd7e543f..2c1008871 100644 --- a/tests/flink-it/src/test/scala/org/apache/celeborn/tests/flink/HeartbeatTest.scala +++ b/tests/flink-it/src/test/scala/org/apache/celeborn/tests/flink/HeartbeatTest.scala @@ -22,7 +22,7 @@ import org.scalatest.funsuite.AnyFunSuite import org.apache.celeborn.common.identity.UserIdentifier import org.apache.celeborn.common.internal.Logging -import org.apache.celeborn.plugin.flink.readclient.FlinkShuffleClientImpl +import org.apache.celeborn.plugin.flink.client.FlinkShuffleClientImpl import org.apache.celeborn.service.deploy.{HeartbeatFeature, MiniClusterFeature} class HeartbeatTest extends AnyFunSuite with Logging with MiniClusterFeature with HeartbeatFeature