This is an automated email from the ASF dual-hosted git repository.

zhouky pushed a commit to branch branch-0.3
in repository https://gitbox.apache.org/repos/asf/incubator-celeborn.git


The following commit(s) were added to refs/heads/branch-0.3 by this push:
     new 9f8b0fbfb [MINOR] Fix some typos
9f8b0fbfb is described below

commit 9f8b0fbfb6e484e8a1266b695be59d4a588e509c
Author: onebox-li <[email protected]>
AuthorDate: Thu Oct 12 20:34:07 2023 +0800

    [MINOR] Fix some typos
    
    Fix some typos
    
    Ditto
    
    No
    
    -
    
    Closes #1983 from onebox-li/fix-typo.
    
    Authored-by: onebox-li <[email protected]>
    Signed-off-by: zky.zhoukeyong <[email protected]>
    (cherry picked from commit a47f6169d82285c9abcbdd2a9166fed6f22fa3ff)
    Signed-off-by: zky.zhoukeyong <[email protected]>
---
 .../celeborn/common/network/TransportContext.java  |  2 +-
 .../common/network/client/TransportClient.java     |  2 +-
 .../common/network/util/TransportConf.java         |  2 +-
 .../org/apache/celeborn/common/CelebornConf.scala  |  4 ++--
 docs/configuration/index.md                        |  2 +-
 docs/configuration/worker.md                       |  2 +-
 .../service/deploy/master/SlotsAllocator.java      | 26 +++++++++++-----------
 .../deploy/master/clustermeta/ha/HARaftServer.java |  2 +-
 .../deploy/worker/storage/ObservedDevice.scala     | 10 ++++-----
 9 files changed, 26 insertions(+), 26 deletions(-)

diff --git 
a/common/src/main/java/org/apache/celeborn/common/network/TransportContext.java 
b/common/src/main/java/org/apache/celeborn/common/network/TransportContext.java
index 761304377..50bb5cac1 100644
--- 
a/common/src/main/java/org/apache/celeborn/common/network/TransportContext.java
+++ 
b/common/src/main/java/org/apache/celeborn/common/network/TransportContext.java
@@ -152,7 +152,7 @@ public class TransportContext {
         conf.connectionTimeoutMs(),
         closeIdleConnections,
         enableHeartbeat,
-        conf.clientHearbeatInterval());
+        conf.clientHeartbeatInterval());
   }
 
   public TransportConf getConf() {
diff --git 
a/common/src/main/java/org/apache/celeborn/common/network/client/TransportClient.java
 
b/common/src/main/java/org/apache/celeborn/common/network/client/TransportClient.java
index 697ca2d26..d28546200 100644
--- 
a/common/src/main/java/org/apache/celeborn/common/network/client/TransportClient.java
+++ 
b/common/src/main/java/org/apache/celeborn/common/network/client/TransportClient.java
@@ -299,7 +299,7 @@ public class TransportClient implements Closeable {
   @Override
   public String toString() {
     return Objects.toStringHelper(this)
-        .add("remoteAdress", channel.remoteAddress())
+        .add("remoteAddress", channel.remoteAddress())
         .add("isActive", isActive())
         .toString();
   }
diff --git 
a/common/src/main/java/org/apache/celeborn/common/network/util/TransportConf.java
 
b/common/src/main/java/org/apache/celeborn/common/network/util/TransportConf.java
index aff2dfa71..468b70685 100644
--- 
a/common/src/main/java/org/apache/celeborn/common/network/util/TransportConf.java
+++ 
b/common/src/main/java/org/apache/celeborn/common/network/util/TransportConf.java
@@ -150,7 +150,7 @@ public class TransportConf {
     return celebornConf.fetchDataTimeoutCheckInterval(module);
   }
 
-  public long clientHearbeatInterval() {
+  public long clientHeartbeatInterval() {
     return celebornConf.clientHeartbeatInterval(module);
   }
 }
diff --git 
a/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala 
b/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
index 7f016c4ca..60a5f7f58 100644
--- a/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
+++ b/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
@@ -610,7 +610,7 @@ class CelebornConf(loadDefaults: Boolean) extends Cloneable 
with Logging with Se
   def haMasterRatisRpcTimeoutMin: Long = get(HA_MASTER_RATIS_RPC_TIMEOUT_MIN)
   def haMasterRatisRpcTimeoutMax: Long = get(HA_MASTER_RATIS_RPC_TIMEOUT_MAX)
   def haMasterRatisFirstElectionTimeoutMin: Long = 
get(HA_MASTER_RATIS_FIRSTELECTION_TIMEOUT_MIN)
-  def haMasterRatisFristElectionTimeoutMax: Long = 
get(HA_MASTER_RATIS_FIRSTELECTION_TIMEOUT_MAX)
+  def haMasterRatisFirstElectionTimeoutMax: Long = 
get(HA_MASTER_RATIS_FIRSTELECTION_TIMEOUT_MAX)
   def haMasterRatisNotificationNoLeaderTimeout: Long =
     get(HA_MASTER_RATIS_NOTIFICATION_NO_LEADER_TIMEOUT)
   def haMasterRatisRpcSlownessTimeout: Long = 
get(HA_MASTER_RATIS_RPC_SLOWNESS_TIMEOUT)
@@ -2570,7 +2570,7 @@ object CelebornConf extends Logging {
     buildConf("celeborn.worker.decommission.checkInterval")
       .categories("worker")
       .doc(
-        "The wait interval of checking whether all the shuffle expired during 
worker decomission")
+        "The wait interval of checking whether all the shuffle expired during 
worker decommission")
       .version("0.4.0")
       .timeConf(TimeUnit.MILLISECONDS)
       .createWithDefaultString("30s")
diff --git a/docs/configuration/index.md b/docs/configuration/index.md
index fbda2f5ab..c62e2dbb1 100644
--- a/docs/configuration/index.md
+++ b/docs/configuration/index.md
@@ -39,7 +39,7 @@ off-heap-memory = bufferSize * estimatedTasks * 2 + network 
memory
 For example, if a Celeborn worker has 10 storage directories or disks and the 
buffer size is set to 256 KiB.
 The necessary off-heap memory is 10 GiB.
 
-Network memory will be consumed when netty reads from a TPC channel, there 
will need some extra
+Network memory will be consumed when netty reads from a TCP channel, there 
will need some extra
 memory. Empirically, Celeborn worker off-heap memory should be set to 
`(numDirs  * bufferSize * 1.2)`.
 
 ## All Configurations
diff --git a/docs/configuration/worker.md b/docs/configuration/worker.md
index 6badddfb7..3cc87e5c4 100644
--- a/docs/configuration/worker.md
+++ b/docs/configuration/worker.md
@@ -34,7 +34,7 @@ license: |
 | celeborn.worker.congestionControl.low.watermark | &lt;undefined&gt; | Will 
stop congest users if the total pending bytes of disk buffer is lower than this 
configuration | 0.3.0 | 
 | celeborn.worker.congestionControl.sample.time.window | 10s | The worker 
holds a time sliding list to calculate users' produce/consume rate | 0.3.0 | 
 | celeborn.worker.congestionControl.user.inactive.interval | 10min | How long 
will consider this user is inactive if it doesn't send data | 0.3.0 | 
-| celeborn.worker.decommission.checkInterval | 30s | The wait interval of 
checking whether all the shuffle expired during worker decomission | 0.4.0 | 
+| celeborn.worker.decommission.checkInterval | 30s | The wait interval of 
checking whether all the shuffle expired during worker decommission | 0.4.0 | 
 | celeborn.worker.decommission.forceExitTimeout | 6h | The wait time of 
waiting for all the shuffle expire during worker decommission. | 0.4.0 | 
 | celeborn.worker.directMemoryRatioForMemoryShuffleStorage | 0.0 | Max ratio 
of direct memory to store shuffle data | 0.2.0 | 
 | celeborn.worker.directMemoryRatioForReadBuffer | 0.1 | Max ratio of direct 
memory for read buffer | 0.2.0 | 
diff --git 
a/master/src/main/java/org/apache/celeborn/service/deploy/master/SlotsAllocator.java
 
b/master/src/main/java/org/apache/celeborn/service/deploy/master/SlotsAllocator.java
index cac3bdf13..7d9ffab37 100644
--- 
a/master/src/main/java/org/apache/celeborn/service/deploy/master/SlotsAllocator.java
+++ 
b/master/src/main/java/org/apache/celeborn/service/deploy/master/SlotsAllocator.java
@@ -203,39 +203,39 @@ public class SlotsAllocator {
     Iterator<Integer> iter = partitionIdList.iterator();
     outer:
     while (iter.hasNext()) {
-      int nextPrimayInd = primaryIndex;
+      int nextPrimaryInd = primaryIndex;
 
       int partitionId = iter.next();
       StorageInfo storageInfo = new StorageInfo();
       if (restrictions != null) {
-        while (!haveUsableSlots(restrictions, workers, nextPrimayInd)) {
-          nextPrimayInd = (nextPrimayInd + 1) % workers.size();
-          if (nextPrimayInd == primaryIndex) {
+        while (!haveUsableSlots(restrictions, workers, nextPrimaryInd)) {
+          nextPrimaryInd = (nextPrimaryInd + 1) % workers.size();
+          if (nextPrimaryInd == primaryIndex) {
             break outer;
           }
         }
         storageInfo =
-            getStorageInfo(workers, nextPrimayInd, restrictions, 
workerDiskIndexForPrimary);
+            getStorageInfo(workers, nextPrimaryInd, restrictions, 
workerDiskIndexForPrimary);
       }
       PartitionLocation primaryPartition =
-          createLocation(partitionId, workers.get(nextPrimayInd), null, 
storageInfo, true);
+          createLocation(partitionId, workers.get(nextPrimaryInd), null, 
storageInfo, true);
 
       if (shouldReplicate) {
-        int nextReplicaInd = (nextPrimayInd + 1) % workers.size();
+        int nextReplicaInd = (nextPrimaryInd + 1) % workers.size();
         if (restrictions != null) {
           while (!haveUsableSlots(restrictions, workers, nextReplicaInd)
-              || !satisfyRackAware(shouldRackAware, workers, nextPrimayInd, 
nextReplicaInd)) {
+              || !satisfyRackAware(shouldRackAware, workers, nextPrimaryInd, 
nextReplicaInd)) {
             nextReplicaInd = (nextReplicaInd + 1) % workers.size();
-            if (nextReplicaInd == nextPrimayInd) {
+            if (nextReplicaInd == nextPrimaryInd) {
               break outer;
             }
           }
           storageInfo =
               getStorageInfo(workers, nextReplicaInd, restrictions, 
workerDiskIndexForReplica);
         } else if (shouldRackAware) {
-          while (!satisfyRackAware(true, workers, nextPrimayInd, 
nextReplicaInd)) {
+          while (!satisfyRackAware(true, workers, nextPrimaryInd, 
nextReplicaInd)) {
             nextReplicaInd = (nextReplicaInd + 1) % workers.size();
-            if (nextReplicaInd == nextPrimayInd) {
+            if (nextReplicaInd == nextPrimaryInd) {
               break outer;
             }
           }
@@ -253,9 +253,9 @@ public class SlotsAllocator {
 
       Tuple2<List<PartitionLocation>, List<PartitionLocation>> locations =
           slots.computeIfAbsent(
-              workers.get(nextPrimayInd), v -> new Tuple2<>(new ArrayList<>(), 
new ArrayList<>()));
+              workers.get(nextPrimaryInd), v -> new Tuple2<>(new 
ArrayList<>(), new ArrayList<>()));
       locations._1.add(primaryPartition);
-      primaryIndex = (nextPrimayInd + 1) % workers.size();
+      primaryIndex = (nextPrimaryInd + 1) % workers.size();
       iter.remove();
     }
     return partitionIdList;
diff --git 
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/HARaftServer.java
 
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/HARaftServer.java
index 73149bc01..c5a0e1c78 100644
--- 
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/HARaftServer.java
+++ 
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/HARaftServer.java
@@ -320,7 +320,7 @@ public class HARaftServer {
     TimeDuration firstElectionTimeoutMin =
         TimeDuration.valueOf(conf.haMasterRatisFirstElectionTimeoutMin(), 
TimeUnit.SECONDS);
     TimeDuration firstElectionTimeoutMax =
-        TimeDuration.valueOf(conf.haMasterRatisFristElectionTimeoutMax(), 
TimeUnit.SECONDS);
+        TimeDuration.valueOf(conf.haMasterRatisFirstElectionTimeoutMax(), 
TimeUnit.SECONDS);
     RaftServerConfigKeys.Rpc.setFirstElectionTimeoutMin(properties, 
firstElectionTimeoutMin);
     RaftServerConfigKeys.Rpc.setFirstElectionTimeoutMax(properties, 
firstElectionTimeoutMax);
 
diff --git 
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/ObservedDevice.scala
 
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/ObservedDevice.scala
index b1d47b491..9196ac49a 100644
--- 
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/ObservedDevice.scala
+++ 
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/ObservedDevice.scala
@@ -124,13 +124,13 @@ class ObservedDevice(val deviceInfo: DeviceInfo, conf: 
CelebornConf, workerSourc
       false
     } else {
       var statsSource: Source = null
-      var infligtSource: Source = null
+      var inflightSource: Source = null
 
       try {
         statsSource = Source.fromFile(statFile)
-        infligtSource = Source.fromFile(inFlightFile)
+        inflightSource = Source.fromFile(inFlightFile)
         val stats = statsSource.getLines().next().trim.split("[ \t]+", -1)
-        val inflight = infligtSource.getLines().next().trim.split("[ \t]+", -1)
+        val inflight = inflightSource.getLines().next().trim.split("[ \t]+", 
-1)
         val readComplete = stats(0).toLong
         val writeComplete = stats(4).toLong
         val readInflight = inflight(0).toLong
@@ -172,8 +172,8 @@ class ObservedDevice(val deviceInfo: DeviceInfo, conf: 
CelebornConf, workerSourc
         if (statsSource != null) {
           statsSource.close()
         }
-        if (infligtSource != null) {
-          infligtSource.close()
+        if (inflightSource != null) {
+          inflightSource.close()
         }
       }
     }

Reply via email to