Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,36 @@ public void init(ZeppelinConfiguration zConf, NoteParser noteParser) throws IOEx
this.notebookDir = this.fs.makeQualified(new Path(zConf.getNotebookDir()));
LOGGER.info("Using folder {} to store notebook", notebookDir);
this.fs.tryMkDir(notebookDir);
recoverInterruptedSaves();
}

/**
* A save interrupted between its two renames leaves the .bak and .tmp files of a note but
* no .zpln file, so {@link #list} does not pick it up. Restore those notes before they are
* listed. A single leftover .tmp or .bak file is not restored, because it can't be told apart
* from a leftover of a removed or moved note. It is only logged.
*/
private void recoverInterruptedSaves() {
try {
List<Path> recovered =
fs.recoverInterruptedWrites(notebookDir, ".zpln", this::isValidNote);
if (!recovered.isEmpty()) {
LOGGER.warn("Recovered {} note(s) whose last save was interrupted: {}",
recovered.size(), recovered);
}
} catch (IOException e) {
LOGGER.warn("Fail to recover notes whose last save was interrupted", e);
}
}

private boolean isValidNote(Path noteFile, String content) {
try {
noteParser.fromJson(getNoteId(noteFile.getName()), content);
return true;
} catch (IOException | RuntimeException e) {
LOGGER.warn("Content of the temp file of {} is not a valid note", noteFile);
return false;
}
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,11 +35,16 @@
import java.io.File;
import java.io.IOException;
import java.io.OutputStream;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Paths;
import java.util.HashMap;
import java.util.Map;
import java.util.stream.Stream;

import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;

class FileSystemNotebookRepoTest {

Expand Down Expand Up @@ -151,4 +156,131 @@ void testComplicatedScenarios() throws IOException {
hdfsNotebookRepo.save(note, authInfo);
assertEquals(1, hdfsNotebookRepo.list(authInfo).size());
}

@Test
void testSaveLeavesNoTempOrBackupFile() throws IOException {
Note note = createNote("/title_1", "value_1");
hdfsNotebookRepo.save(note, authInfo);
note.getConfig().put("config_1", "value_2");
hdfsNotebookRepo.save(note, authInfo);

assertEquals("value_2", getConfigValue(note));
try (Stream<java.nio.file.Path> files = Files.list(Paths.get(notebookDir))) {
assertTrue(files.allMatch(f -> f.getFileName().toString().endsWith(".zpln")));
}
}

@Test
void testRecoverFromTempFileWhenSaveStoppedBeforeLastRename() throws IOException {
Note note = createNote("/title_1", "value_1");
hdfsNotebookRepo.save(note, authInfo);
// The original was moved to .bak and the new content fully written to .tmp
note.getConfig().put("config_1", "value_2");
Files.move(noteFile(note, ""), noteFile(note, ".bak"));
writeString(noteFile(note, ".tmp"), note.toJson());

restartRepo();

assertEquals(1, hdfsNotebookRepo.list(authInfo).size());
assertEquals("value_2", getConfigValue(note));
assertFalse(Files.exists(noteFile(note, ".bak")));
}

@Test
void testRecoverFromBackupFileWhenTempFileIsIncomplete() throws IOException {
Note note = createNote("/title_1", "value_1");
hdfsNotebookRepo.save(note, authInfo);
Files.move(noteFile(note, ""), noteFile(note, ".bak"));
String json = note.toJson();
writeString(noteFile(note, ".tmp"), json.substring(0, json.length() / 2));

restartRepo();

assertEquals(1, hdfsNotebookRepo.list(authInfo).size());
assertEquals("value_1", getConfigValue(note));
}

@Test
void testDoNotRecoverFromTempFileOnly() throws IOException {
// A lone .tmp cannot be told apart from a leftover of a deleted note
Note note = createNote("/title_1", "value_1");
writeString(noteFile(note, ".tmp"), note.toJson());

restartRepo();

assertEquals(0, hdfsNotebookRepo.list(authInfo).size());
assertTrue(Files.exists(noteFile(note, ".tmp")));
}

@Test
void testRemovedNoteIsNotRestoredFromLeftoverTmp() throws IOException {
Note note = createNote("/title_1", "value_1");
hdfsNotebookRepo.save(note, authInfo);
writeString(noteFile(note, ".tmp"), note.toJson());
hdfsNotebookRepo.remove(note.getId(), note.getPath(), authInfo);

restartRepo();

assertEquals(0, hdfsNotebookRepo.list(authInfo).size());
}

@Test
void testMovedNoteIsNotDuplicatedFromLeftoverTmp() throws IOException {
Note note = createNote("/title_1", "value_1");
hdfsNotebookRepo.save(note, authInfo);
writeString(noteFile(note, ".tmp"), note.toJson());
hdfsNotebookRepo.move(note.getId(), "/title_1", "/dir/title_2", authInfo);

restartRepo();

long zplnCount;
try (Stream<java.nio.file.Path> files = Files.walk(Paths.get(notebookDir))) {
zplnCount = files.filter(f -> f.toString().endsWith(".zpln")).count();
}
assertEquals(1, zplnCount);
}

@Test
void testKeepExistingNoteWhenTempFileIsLeftOver() throws IOException {
Note note = createNote("/title_1", "value_1");
hdfsNotebookRepo.save(note, authInfo);
writeString(noteFile(note, ".tmp"), "{");

restartRepo();

assertEquals(1, hdfsNotebookRepo.list(authInfo).size());
assertEquals("value_1", getConfigValue(note));
assertTrue(Files.exists(noteFile(note, ".tmp")));
assertFalse(Files.exists(noteFile(note, ".bak")));
}

private Note createNote(String path, String configValue) {
Note note = new Note();
note.setZeppelinConfiguration(zConf);
note.setNoteParser(noteParser);
note.setPath(path);
Map<String, Object> config = new HashMap<>();
config.put("config_1", configValue);
note.setConfig(config);
return note;
}

private java.nio.file.Path noteFile(Note note, String suffix) throws IOException {
return Paths.get(notebookDir,
hdfsNotebookRepo.buildNoteFileName(note.getId(), note.getPath()) + suffix);
}

private void writeString(java.nio.file.Path file, String content) throws IOException {
Files.write(file, content.getBytes(StandardCharsets.UTF_8));
}

private Object getConfigValue(Note note) throws IOException {
return hdfsNotebookRepo.get(note.getId(), note.getPath(), authInfo)
.getConfig().get("config_1");
}

private void restartRepo() throws IOException {
hdfsNotebookRepo = new FileSystemNotebookRepo();
hdfsNotebookRepo.init(zConf, noteParser);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -40,8 +40,10 @@
import java.nio.file.attribute.PosixFilePermissions;
import java.security.PrivilegedExceptionAction;
import java.util.ArrayList;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Set;
import java.util.function.BiPredicate;


/**
Expand All @@ -52,6 +54,8 @@ public class FileSystemStorage {
private static final Logger LOGGER = LoggerFactory.getLogger(FileSystemStorage.class);
private static final String S3A = "s3a";
private static final String FS_DEFAULTFS = "fs.defaultFS";
static final String TMP_SUFFIX = ".tmp";
static final String BACKUP_SUFFIX = ".bak";

// only do UserGroupInformation.loginUserFromKeytab one time, otherwise you will still get
// your ticket expired.
Expand Down Expand Up @@ -239,16 +243,136 @@ public void writeFile(final String content, final Path file, boolean writeTempFi
public Void call() throws IOException {
InputStream in = new ByteArrayInputStream(content.getBytes(
zConf.getString(ZeppelinConfiguration.ConfVars.ZEPPELIN_ENCODING)));
Path tmpFile = new Path(file.toString() + ".tmp");
Path tmpFile = new Path(file.toString() + TMP_SUFFIX);
IOUtils.copyBytes(in, fs.create(tmpFile), hadoopConf);
fs.setPermission(tmpFile, fsPermission);
fs.delete(file, true);
fs.rename(tmpFile, file);
replaceFile(tmpFile, file);
return null;
}
});
}

/**
* Replaces file with tmpFile without deleting the original first. The original is renamed to
* a backup and deleted only after tmpFile is in place, so an interrupted write always leaves a
* complete copy behind.
*/
private void replaceFile(Path tmpFile, Path file) throws IOException {
Path backupFile = new Path(file.toString() + BACKUP_SUFFIX);
boolean hasOriginal = fs.exists(file);
if (hasOriginal) {
// A backup next to an existing file is left over from an earlier write and is older
// than the file itself.
fs.delete(backupFile, false);
if (!fs.rename(file, backupFile)) {
throw new IOException("Fail to back up " + file + " to " + backupFile);
}
}
if (!fs.rename(tmpFile, file)) {
if (hasOriginal) {
if (fs.rename(backupFile, file)) {
// The original is back in place, so the temp file is no longer needed. Leaving it
// would let it be mistaken for an interrupted write later.
fs.delete(tmpFile, false);
} else {
LOGGER.error("Fail to restore {} from {}, please restore it manually",
file, backupFile);
}
}
throw new IOException("Fail to rename " + tmpFile + " to " + file);
}
if (hasOriginal && !fs.delete(backupFile, false)) {
LOGGER.warn("Fail to delete backup file {}", backupFile);
}
}

/**
* Restores files under dir (recursively) whose last {@link #writeFile} was interrupted.
* A file is restored only when it is missing and both its temp and backup files exist:
* from the temp file if isComplete accepts its content, otherwise from the backup file.
* A single leftover file is only logged. Leftover files next to an existing file are kept.
*
* @param dir folder to scan recursively
* @param targetSuffix suffix of the files to restore, e.g. ".zpln"
* @param isComplete given the file to restore and the content of its temp file, tells whether
* that content is complete
* @return the restored files
*/
public List<Path> recoverInterruptedWrites(final Path dir, final String targetSuffix,
final BiPredicate<Path, String> isComplete) throws IOException {
return callHdfsOperation(new HdfsOperation<List<Path>>() {
@Override
public List<Path> call() throws IOException {
List<Path> recovered = new ArrayList<>();
if (!fs.exists(dir)) {
return recovered;
}
Set<Path> missingFiles = new LinkedHashSet<>();
collectMissingFiles(dir, targetSuffix, missingFiles);
for (Path file : missingFiles) {
if (recoverFile(file, isComplete)) {
recovered.add(file);
}
}
return recovered;
}
});
}

private void collectMissingFiles(Path folder, String targetSuffix, Set<Path> missingFiles)
throws IOException {
for (FileStatus status : fs.listStatus(folder)) {
Path path = status.getPath();
if (status.isDirectory()) {
collectMissingFiles(path, targetSuffix, missingFiles);
continue;
}
for (String suffix : new String[] {TMP_SUFFIX, BACKUP_SUFFIX}) {
if (path.getName().endsWith(targetSuffix + suffix)) {
String pathString = path.toString();
Path file = new Path(pathString.substring(0, pathString.length() - suffix.length()));
if (!fs.exists(file)) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Keeping the original until the new file is in place looks like the right direction to me. I ran into one case, though, where this check treats a note that was already deleted or moved as an interrupted save and brings it back.

What happens

recoverInterruptedWrites restores a .zpln whenever it is missing and a complete .zpln.tmp exists. But FileSystemNotebookRepo.remove() and move() only delete/rename the .zpln itself, so a .tmp (or .bak) left next to it stays behind. On the next restart, that leftover is treated as an interrupted save:

  • note removed → it comes back after restart
  • note moved → a second .zpln with the same noteId appears at the old path; list() keys by noteId, so only one of the two is listed

A complete .tmp can be left next to an existing .zpln in two ways with this PR:

  1. Zeppelin stops after the .tmp is fully written but before file is renamed to .bak.
  2. rename(tmpFile, file) fails: replaceFile restores the original from .bak and throws, but does not delete .tmp. This one does not need a crash.

I checked (2) end to end by making rename(*.zpln.tmp → *.zpln) fail once: the save throws, the original .zpln and a complete .tmp remain, and after remove() + restart the note comes back, with the content of the save that had failed.

Repro

These two tests fail on the PR branch when added to FileSystemNotebookRepoTest:

@Test
void testRemovedNoteIsNotRestoredFromLeftoverTmp() throws IOException {
  Note note = createNote("/title_1", "value_1");
  hdfsNotebookRepo.save(note, authInfo);
  writeString(noteFile(note, ".tmp"), note.toJson());
  hdfsNotebookRepo.remove(note.getId(), note.getPath(), authInfo);

  restartRepo();

  assertEquals(0, hdfsNotebookRepo.list(authInfo).size());  // actual: 1
}

@Test
void testMovedNoteIsNotDuplicatedFromLeftoverTmp() throws IOException {
  Note note = createNote("/title_1", "value_1");
  hdfsNotebookRepo.save(note, authInfo);
  writeString(noteFile(note, ".tmp"), note.toJson());
  hdfsNotebookRepo.move(note.getId(), "/title_1", "/dir/title_2", authInfo);

  restartRepo();

  long zplnCount;
  try (Stream<java.nio.file.Path> files = Files.walk(Paths.get(notebookDir))) {
    zplnCount = files.filter(f -> f.toString().endsWith(".zpln")).count();
  }
  assertEquals(1, zplnCount);  // actual: 2
}

Possible fixes (just ideas, happy to go with whatever you prefer)

  • In remove() / move(), also delete (or move along) the sibling .tmp / .bak.
  • In replaceFile, delete tmpFile when the final rename fails and the original has been restored.
  • Only auto-recover when both .bak and .tmp exist without the .zpln, which is exactly the state this PR's writeFile leaves when it is interrupted between the two renames. For the .tmp-only case (notes lost by earlier versions), a complete .tmp can't be told apart from a leftover of a deleted note, so logging it or moving it aside for manual recovery might be safer than restoring it automatically.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the careful review and the repro tests, that was a real gap.

I went with options 2 and 3 in comment:

  • replaceFile now deletes .tmp when the final rename fails and the original has been restored.
  • Recovery only runs when both .bak and .tmp exist without the .zpln, which is the state writeFile leaves when it stops between its two renames. A single leftover .tmp or .bak is logged instead of restored.

I also added your two tests, and they pass now. As you pointed out, this means notes lost by earlier versions (only .tmp left) are no longer restored automatically, so I updated the PR description accordingly.

I didn't change remove() / move() for option 1, to keep this PR focused on the save path. Happy to add it here or in a follow-up if you think it's needed.

missingFiles.add(file);
}
}
}
}
}

private boolean recoverFile(Path file, BiPredicate<Path, String> isComplete)
throws IOException {
Path tmpFile = new Path(file.toString() + TMP_SUFFIX);
Path backupFile = new Path(file.toString() + BACKUP_SUFFIX);
// writeFile leaves both files only when it stops between its two renames. A single leftover
// may belong to a file that was deleted or moved on purpose, so it is not restored.
if (!fs.exists(tmpFile) || !fs.exists(backupFile)) {
LOGGER.warn("Found {} without {}, not restoring it automatically. "
+ "Rename it manually if it should be restored.",
fs.exists(tmpFile) ? tmpFile : backupFile, file);
return false;
}
// The temp file is newer than the backup, but it may be incomplete if the write stopped
// while it was being written.
if (isComplete.test(file, readContent(tmpFile)) && fs.rename(tmpFile, file)) {
fs.delete(backupFile, false);
LOGGER.warn("Recovered {} from {}", file, tmpFile);
return true;
}
if (fs.rename(backupFile, file)) {
LOGGER.warn("Recovered {} from {}", file, backupFile);
return true;
}
LOGGER.error("Fail to recover {}, please check {} and {} manually",
file, tmpFile, backupFile);
return false;
}

private String readContent(Path file) throws IOException {
ByteArrayOutputStream bytes = new ByteArrayOutputStream();
IOUtils.copyBytes(fs.open(file), bytes, hadoopConf);
return bytes.toString(zConf.getString(ZeppelinConfiguration.ConfVars.ZEPPELIN_ENCODING));
}

public void move(Path src, Path dest) throws IOException {
callHdfsOperation(() -> {
fs.rename(src, dest);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,9 @@
import java.net.URISyntaxException;
import java.net.URLDecoder;
import java.nio.charset.StandardCharsets;
import java.nio.file.AtomicMoveNotSupportedException;
import java.nio.file.Files;
import java.nio.file.StandardCopyOption;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
Expand Down Expand Up @@ -175,10 +178,36 @@ public synchronized void save(Note note, AuthenticationInfo subject) throws IOEx
out.close();
}
}
noteJson.moveTo(rootNotebookFileObject.resolveFile(
moveReplacing(noteJson, rootNotebookFileObject.resolveFile(
buildNoteFileName(note), NameScope.DESCENDENT));
}

/**
* Moves src over dest. FileObject.moveTo deletes dest before renaming src, so a crash in
* between loses the note. For local files, replace dest in a single atomic rename instead.
* Other file systems, and local ones without atomic rename, keep using moveTo.
*/
private void moveReplacing(FileObject src, FileObject dest) throws IOException {
if ("file".equals(src.getName().getScheme())) {
try {
Files.move(src.getPath(), dest.getPath(),
StandardCopyOption.ATOMIC_MOVE, StandardCopyOption.REPLACE_EXISTING);
// The move bypassed VFS, so drop the state it cached for these files.
src.refresh();
dest.refresh();
FileObject parent = dest.getParent();
if (parent != null) {
parent.refresh();
}
return;
} catch (AtomicMoveNotSupportedException e) {
LOGGER.warn("Atomic move is not supported for {}, falling back to a non-atomic move",
dest.getName(), e);
}
}
src.moveTo(dest);
}

@Override
public void move(String noteId,
String notePath,
Expand Down
Loading
Loading