This is an automated email from the ASF dual-hosted git repository.
rexxiong pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/incubator-celeborn.git
The following commit(s) were added to refs/heads/main by this push:
new 5043f46aa [CELEBORN-863][FOLLOWUP] Fix persisted committed file infos
lost
5043f46aa is described below
commit 5043f46aa737a64248230dd5301661657c0541fc
Author: mingji <[email protected]>
AuthorDate: Thu Mar 14 13:51:10 2024 +0800
[CELEBORN-863][FOLLOWUP] Fix persisted committed file infos lost
### What changes were proposed in this pull request?
To fix a bug that might cause persisted committed file info lost.
### Why are the changes needed?
A worker starts will clean its persisted committed file info and won't put
back if this worker restart again, the committed file infos will lost.
### Does this PR introduce _any_ user-facing change?
NO.
### How was this patch tested?
GA.
Closes #2390 from FMX/b863-1.
Authored-by: mingji <[email protected]>
Signed-off-by: Shuang <[email protected]>
---
.../org/apache/celeborn/common/util/PbSerDeUtils.scala | 4 ++--
.../org/apache/celeborn/common/util/PbSerDeUtilsTest.scala | 2 +-
.../service/deploy/worker/storage/StorageManager.scala | 14 +++++++-------
3 files changed, 10 insertions(+), 10 deletions(-)
diff --git
a/common/src/main/scala/org/apache/celeborn/common/util/PbSerDeUtils.scala
b/common/src/main/scala/org/apache/celeborn/common/util/PbSerDeUtils.scala
index 87d18f28a..1ca5759d2 100644
--- a/common/src/main/scala/org/apache/celeborn/common/util/PbSerDeUtils.scala
+++ b/common/src/main/scala/org/apache/celeborn/common/util/PbSerDeUtils.scala
@@ -149,10 +149,10 @@ object PbSerDeUtils {
fileInfoMap
}
- def toPbFileInfoMap(fileInfoMap: ConcurrentHashMap[String, FileInfo]):
Array[Byte] = {
+ def toPbFileInfoMap(fileInfoMap: ConcurrentHashMap[String, DiskFileInfo]):
Array[Byte] = {
val pbFileInfoMap = JavaUtils.newConcurrentHashMap[String, PbFileInfo]()
fileInfoMap.entrySet().asScala.foreach { entry =>
- pbFileInfoMap.put(entry.getKey,
toPbFileInfo(entry.getValue.asInstanceOf[DiskFileInfo]))
+ pbFileInfoMap.put(entry.getKey, toPbFileInfo(entry.getValue))
}
PbFileInfoMap.newBuilder.putAllValues(pbFileInfoMap).build.toByteArray
}
diff --git
a/common/src/test/scala/org/apache/celeborn/common/util/PbSerDeUtilsTest.scala
b/common/src/test/scala/org/apache/celeborn/common/util/PbSerDeUtilsTest.scala
index 0d8063d63..3f4c7b2f4 100644
---
a/common/src/test/scala/org/apache/celeborn/common/util/PbSerDeUtilsTest.scala
+++
b/common/src/test/scala/org/apache/celeborn/common/util/PbSerDeUtilsTest.scala
@@ -69,7 +69,7 @@ class PbSerDeUtilsTest extends CelebornFunSuite {
new ReduceFileMeta(chunkOffsets2),
file2.getAbsolutePath,
6000L)
- val fileInfoMap = JavaUtils.newConcurrentHashMap[String, FileInfo]()
+ val fileInfoMap = JavaUtils.newConcurrentHashMap[String, DiskFileInfo]()
fileInfoMap.put("file1", fileInfo1)
fileInfoMap.put("file2", fileInfo2)
val cache = JavaUtils.newConcurrentHashMap[String, UserIdentifier]()
diff --git
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/StorageManager.scala
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/StorageManager.scala
index ea84272bf..194e13378 100644
---
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/StorageManager.scala
+++
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/StorageManager.scala
@@ -37,7 +37,7 @@ import org.apache.celeborn.common.CelebornConf
import org.apache.celeborn.common.exception.CelebornException
import org.apache.celeborn.common.identity.UserIdentifier
import org.apache.celeborn.common.internal.Logging
-import org.apache.celeborn.common.meta.{DeviceInfo, DiskFileInfo, DiskInfo,
DiskStatus, FileInfo, MapFileMeta, ReduceFileMeta, TimeWindow}
+import org.apache.celeborn.common.meta.{DeviceInfo, DiskFileInfo, DiskInfo,
DiskStatus, MapFileMeta, ReduceFileMeta, TimeWindow}
import org.apache.celeborn.common.metrics.source.{AbstractSource,
ThreadPoolSource}
import org.apache.celeborn.common.network.util.{NettyUtils, TransportConf}
import org.apache.celeborn.common.protocol.{PartitionLocation,
PartitionSplitMode, PartitionType, StorageInfo}
@@ -203,7 +203,8 @@ final private[worker] class StorageManager(conf:
CelebornConf, workerSource: Abs
conf.workerGracefulShutdownSaveCommittedFileInfoSync
private val saveCommittedFileInfoInterval =
conf.workerGracefulShutdownSaveCommittedFileInfoInterval
- private var committedFileInfos: ConcurrentHashMap[String,
ConcurrentHashMap[String, FileInfo]] = _
+ private val committedFileInfos =
+ JavaUtils.newConcurrentHashMap[String, ConcurrentHashMap[String,
DiskFileInfo]]()
// ShuffleClient can fetch data from a restarted worker only
// when the worker's fetching port is stable.
val workerGracefulShutdown = conf.workerGracefulShutdown
@@ -219,8 +220,6 @@ final private[worker] class StorageManager(conf:
CelebornConf, workerSource: Abs
logError("Init level DB failed:", e)
this.db = null
}
- committedFileInfos =
- JavaUtils.newConcurrentHashMap[String, ConcurrentHashMap[String,
FileInfo]]()
saveCommittedFileInfosExecutor =
ThreadUtils.newDaemonSingleThreadScheduledExecutor(
"worker-storage-manager-committed-fileinfo-saver")
@@ -266,6 +265,7 @@ final private[worker] class StorageManager(conf:
CelebornConf, workerSource: Abs
val files = PbSerDeUtils.fromPbFileInfoMap(entry.getValue, cache)
logDebug(s"Reload DB: $shuffleKey -> $files")
diskFileInfos.put(shuffleKey, files)
+ committedFileInfos.put(shuffleKey, files)
db.delete(entry.getKey)
} catch {
case exception: Exception =>
@@ -301,8 +301,8 @@ final private[worker] class StorageManager(conf:
CelebornConf, workerSource: Abs
private def getNextIndex() = counter.getAndUpdate(counterOperator)
private val newMapFunc =
- new java.util.function.Function[String, ConcurrentHashMap[String,
FileInfo]]() {
- override def apply(key: String): ConcurrentHashMap[String, FileInfo] =
+ new java.util.function.Function[String, ConcurrentHashMap[String,
DiskFileInfo]]() {
+ override def apply(key: String): ConcurrentHashMap[String, DiskFileInfo]
=
JavaUtils.newConcurrentHashMap()
}
@@ -778,7 +778,7 @@ final private[worker] class StorageManager(conf:
CelebornConf, workerSource: Abs
def notifyFileInfoCommitted(
shuffleKey: String,
fileName: String,
- fileInfo: FileInfo): Unit = {
+ fileInfo: DiskFileInfo): Unit = {
committedFileInfos.computeIfAbsent(shuffleKey, newMapFunc).put(fileName,
fileInfo)
}