Repository: spark
Updated Branches:
  refs/heads/master 9f6435763 -> 8f5c827b0


[SPARK-5344][WebUI] HistoryServer cannot recognize that inprogress file was 
renamed to completed file

`FsHistoryProvider` tries to update application status but if `checkForLogs` is 
called before `.inprogress` file is renamed to completed file, the file is not 
recognized as completed.

Author: Kousuke Saruta <[email protected]>

Closes #4132 from sarutak/SPARK-5344 and squashes the following commits:

9658008 [Kousuke Saruta] Merge branch 'master' of git://git.apache.org/spark 
into SPARK-5344
d2c72b6 [Kousuke Saruta] Fixed update issue of FsHistoryProvider


Project: http://git-wip-us.apache.org/repos/asf/spark/repo
Commit: http://git-wip-us.apache.org/repos/asf/spark/commit/8f5c827b
Tree: http://git-wip-us.apache.org/repos/asf/spark/tree/8f5c827b
Diff: http://git-wip-us.apache.org/repos/asf/spark/diff/8f5c827b

Branch: refs/heads/master
Commit: 8f5c827b01026bf45fc774ed7387f11a941abea8
Parents: 9f64357
Author: Kousuke Saruta <[email protected]>
Authored: Sun Jan 25 15:34:20 2015 -0800
Committer: Andrew Or <[email protected]>
Committed: Sun Jan 25 15:34:20 2015 -0800

----------------------------------------------------------------------
 .../deploy/history/FsHistoryProvider.scala      |  4 +++-
 .../deploy/history/FsHistoryProviderSuite.scala | 23 ++++++++++++++++++++
 2 files changed, 26 insertions(+), 1 deletion(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/spark/blob/8f5c827b/core/src/main/scala/org/apache/spark/deploy/history/FsHistoryProvider.scala
----------------------------------------------------------------------
diff --git 
a/core/src/main/scala/org/apache/spark/deploy/history/FsHistoryProvider.scala 
b/core/src/main/scala/org/apache/spark/deploy/history/FsHistoryProvider.scala
index 2b084a2..0ae45f4 100644
--- 
a/core/src/main/scala/org/apache/spark/deploy/history/FsHistoryProvider.scala
+++ 
b/core/src/main/scala/org/apache/spark/deploy/history/FsHistoryProvider.scala
@@ -203,7 +203,9 @@ private[history] class FsHistoryProvider(conf: SparkConf) 
extends ApplicationHis
       if (!logInfos.isEmpty) {
         val newApps = new mutable.LinkedHashMap[String, 
FsApplicationHistoryInfo]()
         def addIfAbsent(info: FsApplicationHistoryInfo) = {
-          if (!newApps.contains(info.id)) {
+          if (!newApps.contains(info.id) ||
+              
newApps(info.id).logPath.endsWith(EventLoggingListener.IN_PROGRESS) &&
+              !info.logPath.endsWith(EventLoggingListener.IN_PROGRESS)) {
             newApps += (info.id -> info)
           }
         }

http://git-wip-us.apache.org/repos/asf/spark/blob/8f5c827b/core/src/test/scala/org/apache/spark/deploy/history/FsHistoryProviderSuite.scala
----------------------------------------------------------------------
diff --git 
a/core/src/test/scala/org/apache/spark/deploy/history/FsHistoryProviderSuite.scala
 
b/core/src/test/scala/org/apache/spark/deploy/history/FsHistoryProviderSuite.scala
index 8379883..3fbc1a2 100644
--- 
a/core/src/test/scala/org/apache/spark/deploy/history/FsHistoryProviderSuite.scala
+++ 
b/core/src/test/scala/org/apache/spark/deploy/history/FsHistoryProviderSuite.scala
@@ -167,6 +167,29 @@ class FsHistoryProviderSuite extends FunSuite with 
BeforeAndAfter with Matchers
     list.size should be (1)
   }
 
+  test("history file is renamed from inprogress to completed") {
+    val conf = new SparkConf()
+      .set("spark.history.fs.logDirectory", testDir.getAbsolutePath())
+      .set("spark.testing", "true")
+    val provider = new FsHistoryProvider(conf)
+
+    val logFile1 = new File(testDir, "app1" + EventLoggingListener.IN_PROGRESS)
+    writeFile(logFile1, true, None,
+      SparkListenerApplicationStart("app1", Some("app1"), 1L, "test"),
+      SparkListenerApplicationEnd(2L)
+    )
+    provider.checkForLogs()
+    val appListBeforeRename = provider.getListing()
+    appListBeforeRename.size should be (1)
+    appListBeforeRename.head.logPath should 
endWith(EventLoggingListener.IN_PROGRESS)
+
+    logFile1.renameTo(new File(testDir, "app1"))
+    provider.checkForLogs()
+    val appListAfterRename = provider.getListing()
+    appListAfterRename.size should be (1)
+    appListAfterRename.head.logPath should not 
endWith(EventLoggingListener.IN_PROGRESS)
+  }
+
   private def writeFile(file: File, isNewFormat: Boolean, codec: 
Option[CompressionCodec],
     events: SparkListenerEvent*) = {
     val out =


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to