diff --git a/pixels-cli/src/main/java/io/pixelsdb/pixels/cli/load/IndexedPixelsConsumer.java b/pixels-cli/src/main/java/io/pixelsdb/pixels/cli/load/IndexedPixelsConsumer.java index 5c7c2b3ad..1c1ff16d0 100644 --- a/pixels-cli/src/main/java/io/pixelsdb/pixels/cli/load/IndexedPixelsConsumer.java +++ b/pixels-cli/src/main/java/io/pixelsdb/pixels/cli/load/IndexedPixelsConsumer.java @@ -144,7 +144,8 @@ protected void processSourceFile(String originalFilePath) throws IOException, Me } } } - originStorage.close(); + // Do not close originStorage: StorageFactory caches one Storage instance per scheme, and the same + // instance is used by flushRemainingData() to write the data files when they share the scheme. } @Override @@ -263,8 +264,6 @@ private void flushRowBatch(PerVirtualNodeWriter bucketWriter) throws IOException queue.clear(); } } - bucketWriter.indexService.flushIndexEntriesOfFile(index.getTableId(), index.getId(), - bucketWriter.currFile.getId(), true, bucketWriter.defaultIndexOption); } private void closePixelsFile(PerVirtualNodeWriter bucketWriter) throws IOException, IndexException @@ -275,6 +274,10 @@ private void closePixelsFile(PerVirtualNodeWriter bucketWriter) throws IOExcepti flushRowBatch(bucketWriter); } + // The main index of a file can only be flushed once, hence it is flushed here instead of in flushRowBatch() + bucketWriter.indexService.flushIndexEntriesOfFile(index.getTableId(), index.getId(), + bucketWriter.currFile.getId(), true, bucketWriter.defaultIndexOption); + closeWriterAndAddFile(bucketWriter.pixelsWriter, bucketWriter.currFile, bucketWriter.currTargetPath, bucketWriter.targetNode); }