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 | <undefined> | 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()
}
}
}