From 18f4ab7d4a095a794edbf3aa433739838f47e3e7 Mon Sep 17 00:00:00 2001 From: Alan Wang Date: Thu, 10 Sep 2026 10:25:19 -0700 Subject: [PATCH] fix --- .../apache/cassandra/db/commitlog/CommitLogSegment.java | 2 ++ src/java/org/apache/cassandra/hints/HintsWriter.java | 1 + .../cassandra/io/sstable/metadata/MetadataSerializer.java | 7 +++++-- src/java/org/apache/cassandra/io/util/FileUtils.java | 2 ++ .../service/paxos/uncommitted/PaxosBallotTracker.java | 2 ++ .../service/paxos/uncommitted/UncommittedDataFile.java | 2 ++ 6 files changed, 14 insertions(+), 2 deletions(-) diff --git a/src/java/org/apache/cassandra/db/commitlog/CommitLogSegment.java b/src/java/org/apache/cassandra/db/commitlog/CommitLogSegment.java index 98b8f9dc0da9..a8a7a5aba88a 100644 --- a/src/java/org/apache/cassandra/db/commitlog/CommitLogSegment.java +++ b/src/java/org/apache/cassandra/db/commitlog/CommitLogSegment.java @@ -53,6 +53,7 @@ import org.apache.cassandra.schema.TableId; import org.apache.cassandra.schema.TableMetadata; import org.apache.cassandra.utils.IntegerInterval; +import org.apache.cassandra.utils.SyncUtil; import org.apache.cassandra.utils.concurrent.OpOrder; import org.apache.cassandra.utils.concurrent.WaitQueue; @@ -167,6 +168,7 @@ static long getNextId() throw new FSWriteError(e, logFile); } + SyncUtil.trySyncDir(logFile.parent()); this.buffer = createBuffer(); } diff --git a/src/java/org/apache/cassandra/hints/HintsWriter.java b/src/java/org/apache/cassandra/hints/HintsWriter.java index ffe47318e8f6..f25271ef3947 100644 --- a/src/java/org/apache/cassandra/hints/HintsWriter.java +++ b/src/java/org/apache/cassandra/hints/HintsWriter.java @@ -70,6 +70,7 @@ static HintsWriter create(File directory, HintsDescriptor descriptor) throws IOE File file = descriptor.file(directory); FileChannel channel = FileChannel.open(file.toPath(), StandardOpenOption.WRITE, StandardOpenOption.CREATE_NEW); + SyncUtil.trySyncDir(directory); int fd = NativeLibrary.getfd(channel); CRC32 crc = new CRC32(); diff --git a/src/java/org/apache/cassandra/io/sstable/metadata/MetadataSerializer.java b/src/java/org/apache/cassandra/io/sstable/metadata/MetadataSerializer.java index 314694a611e8..8a224a613304 100644 --- a/src/java/org/apache/cassandra/io/sstable/metadata/MetadataSerializer.java +++ b/src/java/org/apache/cassandra/io/sstable/metadata/MetadataSerializer.java @@ -41,10 +41,11 @@ import org.apache.cassandra.io.util.DataInputBuffer; import org.apache.cassandra.io.util.DataOutputBuffer; import org.apache.cassandra.io.util.DataOutputPlus; -import org.apache.cassandra.io.util.DataOutputStreamPlus; import org.apache.cassandra.io.util.File; import org.apache.cassandra.io.util.FileDataInput; +import org.apache.cassandra.io.util.FileOutputStreamPlus; import org.apache.cassandra.io.util.RandomAccessReader; +import org.apache.cassandra.utils.SyncUtil; import org.apache.cassandra.utils.TimeUUID; import static org.apache.cassandra.utils.FBUtilities.updateChecksumInt; @@ -267,10 +268,11 @@ private void mutate(Descriptor descriptor, UnaryOperator transfor public void rewriteSSTableMetadata(Descriptor descriptor, Map currentComponents) throws IOException { File file = descriptor.tmpFileFor(Components.STATS); - try (DataOutputStreamPlus out = file.newOutputStream(File.WriteMode.OVERWRITE)) + try (FileOutputStreamPlus out = file.newOutputStream(File.WriteMode.OVERWRITE)) { serialize(currentComponents, out, descriptor.version); out.flush(); + out.sync(); } catch (IOException e) { @@ -278,5 +280,6 @@ public void rewriteSSTableMetadata(Descriptor descriptor, Map lines, StandardOpenOption ... o writer.newLine(); } + writer.flush(); + if (sync) { SyncUtil.force(fc, true); diff --git a/src/java/org/apache/cassandra/service/paxos/uncommitted/PaxosBallotTracker.java b/src/java/org/apache/cassandra/service/paxos/uncommitted/PaxosBallotTracker.java index d374cdd60d02..e0855bde4851 100644 --- a/src/java/org/apache/cassandra/service/paxos/uncommitted/PaxosBallotTracker.java +++ b/src/java/org/apache/cassandra/service/paxos/uncommitted/PaxosBallotTracker.java @@ -37,6 +37,7 @@ import org.apache.cassandra.service.accord.AccordService; import org.apache.cassandra.service.paxos.Ballot; import org.apache.cassandra.service.paxos.Commit; +import org.apache.cassandra.utils.SyncUtil; import static org.apache.cassandra.io.util.SequentialWriterOption.FINISH_ON_CLOSE; import static org.apache.cassandra.utils.Crc.crc32; @@ -136,6 +137,7 @@ public synchronized void flush() throws IOException writer.writeInt(Integer.reverseBytes((int) crc.getValue())); } file.move(new File(directory, FNAME)); + SyncUtil.trySyncDir(directory); } public synchronized void truncate() diff --git a/src/java/org/apache/cassandra/service/paxos/uncommitted/UncommittedDataFile.java b/src/java/org/apache/cassandra/service/paxos/uncommitted/UncommittedDataFile.java index e440d60a1b33..78853ff3628b 100644 --- a/src/java/org/apache/cassandra/service/paxos/uncommitted/UncommittedDataFile.java +++ b/src/java/org/apache/cassandra/service/paxos/uncommitted/UncommittedDataFile.java @@ -50,6 +50,7 @@ import org.apache.cassandra.utils.AbstractIterator; import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.CloseableIterator; +import org.apache.cassandra.utils.SyncUtil; import org.apache.cassandra.utils.Throwables; public class UncommittedDataFile @@ -252,6 +253,7 @@ UncommittedDataFile finish() { crcFile.move(finalCrc); file.move(finalData); + SyncUtil.trySyncDir(directory); return new UncommittedDataFile(tableId, finalData, finalCrc, generation); } catch (Throwable e)