This is an automated email from the ASF dual-hosted git repository.

lukasz-antoniak pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/cassandra-analytics.git


The following commit(s) were added to refs/heads/trunk by this push:
     new 2ded78f5 CASSANALYTICS-104: Eliminate redundant filesystem lookups in 
SSTable direct streaming Patch by Lukasz Antoniak; reviewed by Yifan Cai for 
CASSANALYTICS-104
2ded78f5 is described below

commit 2ded78f5329f7b80a18106e4e8e0b43ff34123ba
Author: Lukasz Antoniak <[email protected]>
AuthorDate: Fri Aug 28 17:09:24 2026 +0200

    CASSANALYTICS-104: Eliminate redundant filesystem lookups in SSTable direct 
streaming
    Patch by Lukasz Antoniak; reviewed by Yifan Cai for CASSANALYTICS-104
---
 CHANGES.txt                                        |   1 +
 .../spark/bulkwriter/DirectStreamSession.java      |  90 ++++-----
 .../spark/bulkwriter/SortedSSTableWriter.java      | 211 +++++++++++++++++----
 .../cassandra/spark/bulkwriter/StreamSession.java  |  13 --
 .../cloudstorage/CloudStorageStreamSession.java    |  19 +-
 .../spark/bulkwriter/DirectStreamSessionTest.java  |   2 +-
 .../spark/bulkwriter/SortedSSTableWriterTest.java  |  12 +-
 7 files changed, 225 insertions(+), 123 deletions(-)

diff --git a/CHANGES.txt b/CHANGES.txt
index 0cc8fa70..3992e6e7 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -1,5 +1,6 @@
 0.5.0
 -----
+ * Eliminate redundant filesystem lookups in SSTable direct streaming 
(CASSANALYTICS-104)
  * TokenPartitioner fails to detect range gap in reader (CASSANALYTICS-180)
  * CDC reader stats silently dropped in SidecarCdcBuilder (CASSANALYTICS-191)
  * Add CapturePublishedSchema metric to SidecarCdcStats (CASSANALYTICS-189)
diff --git 
a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/DirectStreamSession.java
 
b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/DirectStreamSession.java
index 7af61709..420acd44 100644
--- 
a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/DirectStreamSession.java
+++ 
b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/DirectStreamSession.java
@@ -21,13 +21,12 @@ package org.apache.cassandra.spark.bulkwriter;
 
 import java.io.IOException;
 import java.math.BigInteger;
-import java.nio.file.DirectoryStream;
 import java.nio.file.Files;
 import java.nio.file.Path;
 import java.util.ArrayList;
+import java.util.Arrays;
 import java.util.HashSet;
 import java.util.List;
-import java.util.Map;
 import java.util.Set;
 import java.util.concurrent.ExecutionException;
 import java.util.concurrent.ExecutorService;
@@ -42,8 +41,6 @@ import org.slf4j.LoggerFactory;
 import org.apache.cassandra.bridge.SSTableDescriptor;
 import org.apache.cassandra.spark.bulkwriter.token.ReplicaAwareFailureHandler;
 import org.apache.cassandra.spark.common.Digest;
-import org.apache.cassandra.spark.common.SSTables;
-import org.apache.cassandra.spark.data.FileType;
 import org.apache.cassandra.util.IntWrapper;
 
 public class DirectStreamSession extends 
StreamSession<TransportContext.DirectDataBulkWriterContext>
@@ -85,20 +82,18 @@ public class DirectStreamSession extends 
StreamSession<TransportContext.DirectDa
                 // 2. validate the sstables
                 // 3. send the sstables to all replicas
                 // 4. remove the sstables once sent
-                Map<Path, Digest> fileDigests = 
sstableWriter.prepareSStablesToSend(writerContext, sstables);
-                // retain only the SSTable data components
+                SortedSSTableWriter.PreparedSSTables preparedSSTables = 
sstableWriter.prepareSStablesToSend(writerContext, sstables);
                 IntWrapper sstableCounter = new IntWrapper();
-                fileDigests.keySet()
-                           .stream()
-                           .filter(p -> 
p.getFileName().toString().endsWith(FileType.DATA.getFileSuffix()))
-                           .forEach(sstable -> {
-                               sstableCounter.value++;
-                               sendSStableToReplicas(sstable);
-                           });
+                preparedSSTables.sstables()
+                                .forEach(preparedSSTable -> {
+                                    sstableCounter.value++;
+                                    sendSStableToReplicas(preparedSSTable);
+                                });
 
                 LOGGER.info("[{}]: Sent newly produced SSTables. sstables={}", 
sessionID, sstableCounter.value);
-                LOGGER.info("[{}]: Removing temporary files after streaming. 
files={}", sessionID, fileDigests);
-                fileDigests.keySet().forEach(path -> {
+                Set<Path> allSSTableFiles = preparedSSTables.files();
+                LOGGER.info("[{}]: Removing temporary files after streaming. 
files={}", sessionID, allSSTableFiles);
+                allSSTableFiles.forEach(path -> {
                     try
                     {
                         Files.deleteIfExists(path);
@@ -153,22 +148,15 @@ public class DirectStreamSession extends 
StreamSession<TransportContext.DirectDa
     @Override
     protected void sendRemainingSSTables()
     {
-        try (DirectoryStream<Path> dataFileStream = 
Files.newDirectoryStream(sstableWriter.getOutDir(), "*Data.db"))
+        try
         {
-            for (Path dataFile : dataFileStream)
-            {
-                if (isFileStreamed(dataFile))
-                {
-                    // the file is already streamed or being streamed; 
skipping it
-                    continue;
-                }
-
-                sendSStableToReplicas(dataFile);
-            }
+            sstableWriter.remainingSSTablesAfterClose()
+                         .sstables()
+                         .forEach(this::sendSStableToReplicas);
 
             LOGGER.info("[{}]: Sent SSTables. sstables={}", sessionID, 
sstableWriter.sstableCount());
         }
-        catch (IOException exception)
+        catch (Exception exception)
         {
             LOGGER.error("[{}]: Unexpected exception while streaming SSTables 
{}",
                          sessionID, sstableWriter.getOutDir());
@@ -182,24 +170,23 @@ public class DirectStreamSession extends 
StreamSession<TransportContext.DirectDa
         }
     }
 
-    private void sendSStableToReplicas(Path dataFile)
+    private void sendSStableToReplicas(SortedSSTableWriter.PreparedSSTable 
preparedSSTable)
     {
         int ssTableIdx = nextSSTableIdx.getAndIncrement();
 
         LOGGER.info("[{}]: Pushing SSTable {} to replicas {}",
-                    sessionID, dataFile,
+                    sessionID, preparedSSTable.dataFile(),
                     
replicas.stream().map(RingInstance::nodeName).collect(Collectors.joining(",")));
-        replicas.removeIf(replica -> !trySendSSTableToOneReplica(dataFile, 
ssTableIdx, replica, sstableWriter.fileDigestMap()));
+        replicas.removeIf(replica -> 
!trySendSSTableToOneReplica(preparedSSTable, ssTableIdx, replica));
     }
 
-    private boolean trySendSSTableToOneReplica(Path dataFile,
+    private boolean 
trySendSSTableToOneReplica(SortedSSTableWriter.PreparedSSTable preparedSSTable,
                                                int ssTableIdx,
-                                               RingInstance replica,
-                                               Map<Path, Digest> fileDigests)
+                                               RingInstance replica)
     {
         try
         {
-            sendSSTableToOneReplica(dataFile, ssTableIdx, replica, 
fileDigests);
+            sendSSTableToOneReplica(preparedSSTable, ssTableIdx, replica);
             return true;
         }
         catch (Exception exception)
@@ -207,32 +194,36 @@ public class DirectStreamSession extends 
StreamSession<TransportContext.DirectDa
             LOGGER.error("[{}]: Failed to stream range {} to instance {}",
                          sessionID, tokenRange, replica.nodeName(), exception);
             writerContext.cluster().refreshClusterInfo();
-            failureHandler.addFailure(this.tokenRange, replica, 
exception.getMessage());
-            errors.add(new StreamError(this.tokenRange, replica, 
exception.getMessage()));
+            // Sometimes error can contain just file name (e.g. when it is 
missing).
+            // Let us return 3 latest stacktrace lines for easier 
troubleshooting.
+            String stackTrace = Arrays.stream(exception.getStackTrace())
+                                      .limit(3)
+                                      .map(StackTraceElement::toString)
+                                      .collect(Collectors.joining("\n"));
+            String errorMessage = exception.getClass().getName() + ": " + 
exception.getMessage()
+                                  + "\n" + stackTrace;
+            failureHandler.addFailure(this.tokenRange, replica, errorMessage);
+            errors.add(new StreamError(this.tokenRange, replica, 
errorMessage));
             clean(replica, sessionID);
             return false;
         }
     }
 
-    private void sendSSTableToOneReplica(Path dataFile,
+    private void sendSSTableToOneReplica(SortedSSTableWriter.PreparedSSTable 
preparedSSTable,
                                          int ssTableIdx,
-                                         RingInstance instance,
-                                         Map<Path, Digest> fileHashes) throws 
IOException
+                                         RingInstance instance) throws 
IOException
     {
-        try (DirectoryStream<Path> componentFileStream = 
Files.newDirectoryStream(dataFile.getParent(),
-                                                                               
   SSTables.getSSTableBaseName(dataFile) + "*"))
+        for (Path componentFile : preparedSSTable.files())
         {
-            for (Path componentFile : componentFileStream)
+            // send data component the last
+            if (componentFile.equals(preparedSSTable.dataFile()))
             {
-                // send data component the last
-                if (componentFile.getFileName().toString().endsWith("Data.db"))
-                {
-                    continue;
-                }
-                sendSSTableComponent(componentFile, ssTableIdx, instance, 
fileHashes.get(componentFile));
+                continue;
             }
-            sendSSTableComponent(dataFile, ssTableIdx, instance, 
fileHashes.get(dataFile));
+            sendSSTableComponent(componentFile, ssTableIdx, instance, 
preparedSSTable.getDigest(componentFile));
         }
+        Preconditions.checkNotNull(preparedSSTable.dataFile(), "Data file not 
present in SSTable: {}", preparedSSTable);
+        sendSSTableComponent(preparedSSTable.dataFile(), ssTableIdx, instance, 
preparedSSTable.getDigest(preparedSSTable.dataFile()));
     }
 
     private void sendSSTableComponent(Path componentFile,
@@ -244,7 +235,6 @@ public class DirectStreamSession extends 
StreamSession<TransportContext.DirectDa
         LOGGER.info("[{}]: Uploading {} to {}: size={} digest={}",
                     sessionID, componentFile, instance.nodeName(), 
Files.size(componentFile), digest);
         directDataTransferApi.uploadSSTableComponent(componentFile, 
ssTableIdx, instance, this.sessionID, digest);
-        recordStreamedFile(componentFile);
     }
 
     private List<CommitResult> commit(DirectStreamResult streamResult) throws 
ExecutionException, InterruptedException
diff --git 
a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/SortedSSTableWriter.java
 
b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/SortedSSTableWriter.java
index 593208dd..73fbc6a2 100644
--- 
a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/SortedSSTableWriter.java
+++ 
b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/SortedSSTableWriter.java
@@ -28,11 +28,15 @@ import java.util.Collection;
 import java.util.Collections;
 import java.util.HashMap;
 import java.util.HashSet;
+import java.util.List;
 import java.util.Map;
 import java.util.Set;
 import java.util.function.Consumer;
+import java.util.stream.Collectors;
 
 import com.google.common.annotations.VisibleForTesting;
+import com.google.common.base.Preconditions;
+import com.google.common.collect.ImmutableList;
 import com.google.common.collect.Range;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -99,6 +103,8 @@ public class SortedSSTableWriter
 
     // Fields protected by synchronization - accessed from both RecordWriter 
thread and executor threads
     private final Map<Path, Digest> overallFileDigests = new HashMap<>();
+    // holds the list of newly created sstables after SSTableWriter.close() 
method was called
+    private PreparedSSTables remainingSSTablesAfterClose = 
PreparedSSTables.EMPTY;
     private boolean isClosed = false;
     private int sstableCount = 0;
     private long bytesWritten = 0;
@@ -209,17 +215,17 @@ public class SortedSSTableWriter
      *
      * @param writerContext the bulk writer context
      * @param sstables the set of SSTable descriptors to prepare
-     * @return a map of file paths to their digests, or an empty map if the 
writer is already closed
+     * @return an object containing list of prepared sstables, or {@link 
PreparedSSTables#EMPTY} if the writer is already closed
      * @throws IOException if an I/O error occurs
      */
-    public synchronized Map<Path, Digest> prepareSStablesToSend(@NotNull 
BulkWriterContext writerContext, Set<SSTableDescriptor> sstables) throws 
IOException
+    public synchronized PreparedSSTables prepareSStablesToSend(@NotNull 
BulkWriterContext writerContext, Set<SSTableDescriptor> sstables) throws 
IOException
     {
         // If the writer is already closed, return empty map
         // The remaining SSTables will be handled by sendRemainingSSTables()
         if (isClosed)
         {
             LOGGER.debug("Writer is already closed, returning empty digest 
map. Remaining SSTables will be handled by sendRemainingSSTables()");
-            return Collections.emptyMap();
+            return PreparedSSTables.EMPTY;
         }
 
         // Filter for SSTables that match the requested descriptors AND 
haven't been hashed yet
@@ -227,8 +233,7 @@ public class SortedSSTableWriter
             SSTableDescriptor baseName = SSTables.getSSTableDescriptor(path);
             return sstables.contains(baseName) && 
!overallFileDigests.containsKey(path);
         };
-        Set<Path> dataFilePaths = new HashSet<>();
-        Map<Path, Digest> fileDigests = new HashMap<>();
+        PreparedSSTables preparedSSTables = new PreparedSSTables();
         // FIXME: CQLSSTableWriter may produce incomplete Filter.db file, 
rebuilding it manually (see CASSANDRA-21423).
         // rebuild Filter.db files before calculating their digest
         rebuildFilterComponents(writerContext, sstableFilter);
@@ -236,21 +241,20 @@ public class SortedSSTableWriter
         {
             for (Path path : stream)
             {
-                if (path.getFileName().toString().endsWith("-" + 
FileType.DATA.getFileSuffix()))
-                {
-                    dataFilePaths.add(path);
-                    sstableCount += 1;
-                }
+                PreparedSSTable preparedSSTable = 
preparedSSTables.getOrPrepareSSTable(path);
 
                 Digest digest = digestAlgorithm.calculateFileDigest(path);
-                fileDigests.put(path, digest);
+                preparedSSTable.addComponent(path, digest);
                 LOGGER.debug("Calculated digest={} for path={}", digest, path);
             }
         }
+        Preconditions.checkState(preparedSSTables.hasAllDataFiles(), "Data 
file not present in some SSTables: {}", preparedSSTables);
+        sstableCount += preparedSSTables.sstablesCount();
+        Map<Path, Digest> fileDigests = preparedSSTables.digests();
         bytesWritten += calculatedTotalSize(fileDigests.keySet());
         overallFileDigests.putAll(fileDigests);
-        validateSSTables(writerContext, getOutDir(), dataFilePaths);
-        return fileDigests;
+        validateSSTables(writerContext, getOutDir(), 
preparedSSTables.dataFiles());
+        return preparedSSTables;
     }
 
     /**
@@ -289,19 +293,28 @@ public class SortedSSTableWriter
         // FIXME: CQLSSTableWriter may produce incomplete Filter.db file, 
rebuilding it manually (see CASSANDRA-21423).
         rebuildFilterComponents(writerContext, unhashedFilter);
 
-        try (DirectoryStream<Path> dataFileStream = 
getDataFileStream(unhashedFilter))
+        PreparedSSTables prepared = new PreparedSSTables();
+
+        try (DirectoryStream<Path> fileStream = 
Files.newDirectoryStream(getOutDir(), unhashedFilter))
         {
-            for (Path dataFile : dataFileStream)
+            for (Path path : fileStream)
             {
                 // NOTE: We calculate file hashes before re-reading so that we 
know what we hashed
                 //       is what we validated. Then we send these along with 
the files and the
                 //       receiving end re-hashes the files to make sure they 
still match.
-                Map<Path, Digest> newFileDigests = 
calculateFileDigestMap(dataFile);
-                overallFileDigests.putAll(newFileDigests);
-                newlyHashedFiles.addAll(newFileDigests.keySet());
-                sstableCount += 1;
+                Digest digest = digestAlgorithm.calculateFileDigest(path);
+                LOGGER.debug("Calculated digest={} for path={}", digest, path);
+
+                overallFileDigests.put(path, digest);
+                newlyHashedFiles.add(path);
+
+                prepared.getOrPrepareSSTable(path)
+                        .addComponent(path, digest);
             }
         }
+        Preconditions.checkState(prepared.hasAllDataFiles(), "Data file not 
present in some SSTables: {}", prepared);
+        sstableCount += prepared.sstablesCount();
+        remainingSSTablesAfterClose = prepared;
         // Only calculate size for newly hashed files, not all files in 
overallFileDigests
         // (previously hashed files may have been deleted by 
prepareSStablesToSend)
         bytesWritten += calculatedTotalSize(newlyHashedFiles);
@@ -388,29 +401,10 @@ public class SortedSSTableWriter
     private DirectoryStream<Path> 
getDataFileStream(DirectoryStream.Filter<Path> filter) throws IOException
     {
         // Combine the data file filter with the provided filter
-        DirectoryStream.Filter<Path> combinedFilter = path -> {
-            String fileName = path.getFileName().toString();
-            return fileName.endsWith("Data.db") && filter.accept(path);
-        };
+        DirectoryStream.Filter<Path> combinedFilter = path -> isDataFile(path) 
&& filter.accept(path);
         return Files.newDirectoryStream(getOutDir(), combinedFilter);
     }
 
-    private Map<Path, Digest> calculateFileDigestMap(Path dataFile) throws 
IOException
-    {
-        Map<Path, Digest> fileHashes = new HashMap<>();
-        try (DirectoryStream<Path> filesToHash =
-             Files.newDirectoryStream(dataFile.getParent(), 
SSTables.getSSTableBaseName(dataFile) + "*"))
-        {
-            for (Path path : filesToHash)
-            {
-                Digest digest = digestAlgorithm.calculateFileDigest(path);
-                fileHashes.put(path, digest);
-                LOGGER.debug("Calculated digest={} for path={}", digest, path);
-            }
-        }
-        return fileHashes;
-    }
-
     private long calculatedTotalSize(Collection<Path> paths) throws IOException
     {
         long totalSize = 0;
@@ -438,4 +432,143 @@ public class SortedSSTableWriter
     {
         return Collections.unmodifiableMap(overallFileDigests);
     }
+
+    public PreparedSSTables remainingSSTablesAfterClose()
+    {
+        return remainingSSTablesAfterClose;
+    }
+
+    /**
+     * Helper class representing list of newly generated sstables.
+     */
+    public static class PreparedSSTables
+    {
+        private static final PreparedSSTables EMPTY = new PreparedSSTables()
+        {
+            @Override
+            protected PreparedSSTable getOrPrepareSSTable(Path path)
+            {
+                // assert that state is never modified
+                throw new IllegalStateException();
+            }
+        };
+
+        private final Map<String, PreparedSSTable> sstables = new HashMap<>(); 
// indexed by base file name
+
+        protected PreparedSSTable getOrPrepareSSTable(Path path)
+        {
+            String baseName = SSTables.getSSTableDescriptor(path).baseFilename;
+            return sstables.computeIfAbsent(baseName, (__) -> new 
PreparedSSTable());
+        }
+
+        /**
+         * @return all files from every generated sstable
+         */
+        public Set<Path> files()
+        {
+            return sstables.values().stream()
+                           .flatMap(c -> c.components.keySet().stream())
+                           .collect(Collectors.toSet());
+        }
+
+        /**
+         * @return data files from all sstables
+         */
+        public Set<Path> dataFiles()
+        {
+            return sstables.values().stream()
+                           .map(c -> c.dataFile)
+                           .collect(Collectors.toSet());
+        }
+
+        /**
+         * @return positive when all produced sstables have {@code *-Data.db} 
file present
+         */
+        public boolean hasAllDataFiles()
+        {
+            return sstables.values().stream()
+                           .allMatch(c -> c.dataFile != null);
+        }
+
+        /**
+         * @return digests for all files from every generated sstable
+         */
+        public Map<Path, Digest> digests()
+        {
+            return sstables.values().stream()
+                           .map(c -> c.components)
+                           .flatMap(c -> c.entrySet().stream())
+                           .collect(Collectors.toMap(Map.Entry::getKey, 
Map.Entry::getValue));
+        }
+
+        public int sstablesCount()
+        {
+            return sstables.size();
+        }
+
+        public Collection<PreparedSSTable> sstables()
+        {
+            return sstables.values();
+        }
+
+        @Override
+        public String toString()
+        {
+            return sstables.values().stream()
+                           .map(PreparedSSTable::toString)
+                           .collect(Collectors.joining(", "));
+        }
+    }
+
+    public static class PreparedSSTable
+    {
+        private Path dataFile;
+        private final Map<Path, Digest> components = new HashMap<>();
+
+        public void addComponent(Path path, Digest digest)
+        {
+            if (isDataFile(path))
+            {
+                dataFile = path;
+            }
+            components.put(path, digest);
+        }
+
+        /**
+         * @return path to data file
+         */
+        public Path dataFile()
+        {
+            return dataFile;
+        }
+
+        /**
+         * @return list of all sstable files
+         */
+        public List<Path> files()
+        {
+            return ImmutableList.copyOf(components.keySet());
+        }
+
+        /**
+         * @param path component path
+         * @return digest value, if it is known for given component
+         */
+        @Nullable
+        public Digest getDigest(Path path)
+        {
+            return components.get(path);
+        }
+
+        @Override
+        public String toString()
+        {
+            return components.keySet().toString();
+        }
+    }
+
+    private static boolean isDataFile(Path path)
+    {
+        return path.getFileName().toString().endsWith("-" + 
FileType.DATA.getFileSuffix());
+    }
 }
diff --git 
a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/StreamSession.java
 
b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/StreamSession.java
index 40cbbdea..fb0cc9d2 100644
--- 
a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/StreamSession.java
+++ 
b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/StreamSession.java
@@ -22,14 +22,12 @@ package org.apache.cassandra.spark.bulkwriter;
 import java.io.File;
 import java.io.IOException;
 import java.math.BigInteger;
-import java.nio.file.Path;
 import java.util.ArrayList;
 import java.util.Collection;
 import java.util.Collections;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
-import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Future;
 import java.util.concurrent.atomic.AtomicReference;
@@ -61,7 +59,6 @@ public abstract class StreamSession<T extends 
TransportContext>
     protected final SortedSSTableWriter sstableWriter;
     protected final ExecutorService executorService;
 
-    private final Set<Path> streamedFiles = ConcurrentHashMap.newKeySet();
     private final AtomicReference<Exception> lastStreamFailure = new 
AtomicReference<>();
     private volatile boolean isStreamFinalized = false;
 
@@ -164,16 +161,6 @@ public abstract class StreamSession<T extends 
TransportContext>
         return lastStreamFailure.compareAndSet(null, streamFailure);
     }
 
-    protected void recordStreamedFile(Path file)
-    {
-        streamedFiles.add(file);
-    }
-
-    protected boolean isFileStreamed(Path file)
-    {
-        return streamedFiles.contains(file);
-    }
-
     @VisibleForTesting
     List<RingInstance> getReplicas()
     {
diff --git 
a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageStreamSession.java
 
b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageStreamSession.java
index b749a359..f3b3a233 100644
--- 
a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageStreamSession.java
+++ 
b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/cloudstorage/CloudStorageStreamSession.java
@@ -23,23 +23,21 @@ import java.io.IOException;
 import java.math.BigInteger;
 import java.nio.file.Path;
 import java.util.HashSet;
-import java.util.Map;
 import java.util.Set;
 import java.util.concurrent.ExecutorService;
-import java.util.stream.Collectors;
 
 import com.google.common.annotations.VisibleForTesting;
 import com.google.common.collect.Range;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import o.a.c.sidecar.client.shaded.client.SidecarInstance;
 import 
o.a.c.sidecar.client.shaded.common.request.data.CreateSliceRequestPayload;
 import 
o.a.c.sidecar.client.shaded.common.response.data.RestoreJobSummaryResponsePayload;
 import org.apache.cassandra.bridge.CassandraBridge;
 import org.apache.cassandra.bridge.CassandraBridgeFactory;
 import org.apache.cassandra.bridge.SSTableDescriptor;
 import org.apache.cassandra.clients.Sidecar;
-import o.a.c.sidecar.client.shaded.client.SidecarInstance;
 import org.apache.cassandra.spark.bulkwriter.BulkWriteValidator;
 import org.apache.cassandra.spark.bulkwriter.BulkWriterContext;
 import org.apache.cassandra.spark.bulkwriter.JobInfo;
@@ -51,8 +49,6 @@ import org.apache.cassandra.spark.bulkwriter.StreamSession;
 import org.apache.cassandra.spark.bulkwriter.TransportContext;
 import 
org.apache.cassandra.spark.bulkwriter.cloudstorage.coordinated.CoordinatedCloudStorageDataTransferApi;
 import org.apache.cassandra.spark.bulkwriter.token.ReplicaAwareFailureHandler;
-import org.apache.cassandra.spark.common.Digest;
-import org.apache.cassandra.spark.common.SSTables;
 import org.apache.cassandra.spark.data.QualifiedTableName;
 import org.apache.cassandra.spark.exception.ConsistencyNotSatisfiedException;
 import org.apache.cassandra.spark.exception.S3ApiCallException;
@@ -119,15 +115,12 @@ public class CloudStorageStreamSession extends 
StreamSession<TransportContext.Cl
         executorService.submit(() -> {
             try
             {
-                Map<Path, Digest> fileDigests = 
sstableWriter.prepareSStablesToSend(writerContext, sstables);
-                sstablesBundler.includeFileDigests(fileDigests);
+                SortedSSTableWriter.PreparedSSTables preparedSSTables = 
sstableWriter.prepareSStablesToSend(writerContext, sstables);
+                sstablesBundler.includeFileDigests(preparedSSTables.digests());
                 // sstablesBundler keeps track of the known files. No need to 
record the streamed files.
-                // group the files by sstable (unique) basename and add to 
bundler
-                fileDigests.keySet()
-                           .stream()
-                           
.collect(Collectors.groupingBy(SSTables::getSSTableBaseName))
-                           .values()
-                           .forEach(sstablesBundler::includeSSTable);
+                // add to bundler files grouped by base sstable name
+                preparedSSTables.sstables()
+                                .forEach(s -> 
sstablesBundler.includeSSTable(s.files()));
 
                 if (!sstablesBundler.hasNext())
                 {
diff --git 
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/DirectStreamSessionTest.java
 
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/DirectStreamSessionTest.java
index f6fe8c5e..de876f67 100644
--- 
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/DirectStreamSessionTest.java
+++ 
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/DirectStreamSessionTest.java
@@ -184,7 +184,7 @@ public class DirectStreamSessionTest
             fut.get();
         })
                 .isInstanceOf(ExecutionException.class)
-                .hasRootCauseInstanceOf(NoSuchFileException.class);
+                .hasMessageContaining(NoSuchFileException.class.getName());
         List<String> actualInstances = 
writerContext.getCleanedInstances().stream()
                                                     
.map(CassandraInstance::nodeName)
                                                     
.collect(Collectors.toList());
diff --git 
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/SortedSSTableWriterTest.java
 
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/SortedSSTableWriterTest.java
index fba73361..95e19e96 100644
--- 
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/SortedSSTableWriterTest.java
+++ 
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/SortedSSTableWriterTest.java
@@ -32,7 +32,6 @@ import java.util.Arrays;
 import java.util.Collections;
 import java.util.HashSet;
 import java.util.List;
-import java.util.Map;
 import java.util.Set;
 import java.util.SortedMap;
 import java.util.TreeMap;
@@ -60,7 +59,6 @@ import org.apache.cassandra.bridge.CassandraVersionFeatures;
 import org.apache.cassandra.bridge.SSTableDescriptor;
 import org.apache.cassandra.spark.bulkwriter.token.ConsistencyLevel;
 import org.apache.cassandra.spark.bulkwriter.token.TokenRangeMapping;
-import org.apache.cassandra.spark.common.Digest;
 import org.apache.cassandra.spark.data.CqlTable;
 import org.apache.cassandra.spark.data.FileSystemSSTable;
 import org.apache.cassandra.spark.data.ReplicationFactor;
@@ -425,13 +423,13 @@ public class SortedSSTableWriterTest
         writer.addRow(BigInteger.valueOf(100), ImmutableMap.of("id", 2, 
"date", 2, "course", "test2", "marks", 200));
 
         // Call prepareSStablesToSend with the existing SSTables
-        Map<Path, Digest> processedFiles = 
writer.prepareSStablesToSend(writerContext, new HashSet<>(existingSSTables));
-        assertThat(processedFiles).as("Should have processed existing 
SSTables").isNotEmpty();
+        SortedSSTableWriter.PreparedSSTables processedFiles = 
writer.prepareSStablesToSend(writerContext, new HashSet<>(existingSSTables));
+        assertThat(processedFiles.sstables()).as("Should have processed 
existing SSTables").isNotEmpty();
 
         long bytesAfterPrepare = writer.bytesWritten();
 
         // Delete the files that were processed (simulating 
DirectStreamSession behavior)
-        for (Path path : processedFiles.keySet())
+        for (Path path : processedFiles.files())
         {
             Files.deleteIfExists(path);
         }
@@ -474,10 +472,10 @@ public class SortedSSTableWriterTest
         int fileDigestCountAfterClose = writer.fileDigestMap().size();
 
         // Try to call prepareSStablesToSend after close - it should return 
empty map
-        Map<Path, Digest> result = writer.prepareSStablesToSend(writerContext, 
new HashSet<>());
+        SortedSSTableWriter.PreparedSSTables result = 
writer.prepareSStablesToSend(writerContext, new HashSet<>());
 
         // Verify it returned an empty map
-        assertThat(result)
+        assertThat(result.sstables())
         .as("prepareSStablesToSend should return empty map when called after 
close")
         .isEmpty();
 


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

Reply via email to