Baunsgaard commented on code in PR #2535:
URL: https://github.com/apache/systemds/pull/2535#discussion_r3560354354


##########
src/main/java/org/apache/sysds/runtime/io/DeltaKernelUtils.java:
##########
@@ -162,6 +174,199 @@ public static int countSelected(int size, boolean[] 
selected) {
                return n;
        }
 
+       // ------------------------------------------
+       // direct parquet decode of Delta data files
+       // ------------------------------------------
+
+       /** Physical-schema metadata key carrying the parquet field id (column 
mapping mode {@code id}). */
+       private static final String PARQUET_FIELD_ID_KEY = "parquet.field.id";
+
+       /**
+        * Whether data files can be decoded directly into pre-allocated output 
columns: the physical read schema must be a
+        * positional 1:1 image of the logical schema, i.e. no partition 
columns (not stored in the data files, spliced back
+        * in by the kernel) and no kernel metadata columns such as {@code 
row_index} (only requested for deletion-vector
+        * reads). Deletion vectors themselves are excluded separately via the 
exact-row-count check.
+        *
+        * @param logicalSchema  the table's logical schema
+        * @param physicalSchema the physical read schema from the scan state
+        * @return true if data files can be decoded directly
+        */
+       public static boolean supportsDirectDecode(StructType logicalSchema, 
StructType physicalSchema) {
+               if(physicalSchema.length() != logicalSchema.length())
+                       return false;
+               for(int c = 0; c < physicalSchema.length(); c++)
+                       if(physicalSchema.at(c).isMetadataColumn())
+                               return false;
+               return true;
+       }
+
+       /** @param scanFileRow a scan-file row @return the fully-qualified path 
of its data file */
+       public static String dataFilePath(Row scanFileRow) {
+               return 
InternalScanFileUtils.getAddFileStatus(scanFileRow).getPath();
+       }
+
+       /**
+        * Decode one Delta data file into pre-allocated typed column arrays at 
the given absolute row offset, through
+        * parquet-mr's column API ({@link ColumnReadStoreImpl}/{@link 
ColumnReader}) with no kernel engine or intermediate
+        * batch vectors in the path. Columns are resolved by parquet field id 
first (column mapping mode {@code id}) and
+        * physical name second; columns absent from the file (schema 
evolution) keep the array defaults (0 for numerics,
+        * null for strings), matching the kernel-path null semantics.
+        *
+        * @param filePath       fully-qualified path of the parquet data file
+        * @param physicalSchema physical read schema (positionally 1:1 with 
the output columns)
+        * @param readCodes      per-column type codes (see the {@code T_*} 
constants)
+        * @param dest           pre-allocated per-column backing arrays
+        * @param destOff        absolute row offset of this file's first row
+        * @param limit          exclusive upper row bound of this file's slice
+        * @param tablePath      table path for error messages
+        * @return the number of rows decoded
+        * @throws IOException on read failure
+        */
+       public static int decodeDataFileInto(String filePath, StructType 
physicalSchema, int[] readCodes, Object[] dest,
+               int destOff, int limit, String tablePath) throws IOException {
+               final Configuration conf = 
ConfigurationManager.getCachedJobConf();
+               final int ncol = physicalSchema.length();
+               int off = destOff;
+               try(ParquetFileReader reader = 
ParquetFileReader.open(HadoopInputFile.fromPath(new Path(filePath), conf))) {
+                       MessageType parquetSchema = 
reader.getFooter().getFileMetaData().getSchema();
+                       String createdBy = 
reader.getFooter().getFileMetaData().getCreatedBy();

Review Comment:
   Avoid calling reader.getFooter().getFileMetadata() twice, just put another 
variable.



##########
src/main/java/org/apache/sysds/runtime/io/DeltaKernelUtils.java:
##########
@@ -162,6 +174,199 @@ public static int countSelected(int size, boolean[] 
selected) {
                return n;
        }
 
+       // ------------------------------------------
+       // direct parquet decode of Delta data files
+       // ------------------------------------------
+
+       /** Physical-schema metadata key carrying the parquet field id (column 
mapping mode {@code id}). */
+       private static final String PARQUET_FIELD_ID_KEY = "parquet.field.id";
+
+       /**
+        * Whether data files can be decoded directly into pre-allocated output 
columns: the physical read schema must be a
+        * positional 1:1 image of the logical schema, i.e. no partition 
columns (not stored in the data files, spliced back
+        * in by the kernel) and no kernel metadata columns such as {@code 
row_index} (only requested for deletion-vector
+        * reads). Deletion vectors themselves are excluded separately via the 
exact-row-count check.
+        *
+        * @param logicalSchema  the table's logical schema
+        * @param physicalSchema the physical read schema from the scan state
+        * @return true if data files can be decoded directly
+        */
+       public static boolean supportsDirectDecode(StructType logicalSchema, 
StructType physicalSchema) {
+               if(physicalSchema.length() != logicalSchema.length())
+                       return false;
+               for(int c = 0; c < physicalSchema.length(); c++)
+                       if(physicalSchema.at(c).isMetadataColumn())
+                               return false;
+               return true;
+       }
+
+       /** @param scanFileRow a scan-file row @return the fully-qualified path 
of its data file */
+       public static String dataFilePath(Row scanFileRow) {
+               return 
InternalScanFileUtils.getAddFileStatus(scanFileRow).getPath();
+       }
+
+       /**
+        * Decode one Delta data file into pre-allocated typed column arrays at 
the given absolute row offset, through
+        * parquet-mr's column API ({@link ColumnReadStoreImpl}/{@link 
ColumnReader}) with no kernel engine or intermediate
+        * batch vectors in the path. Columns are resolved by parquet field id 
first (column mapping mode {@code id}) and
+        * physical name second; columns absent from the file (schema 
evolution) keep the array defaults (0 for numerics,
+        * null for strings), matching the kernel-path null semantics.
+        *
+        * @param filePath       fully-qualified path of the parquet data file
+        * @param physicalSchema physical read schema (positionally 1:1 with 
the output columns)
+        * @param readCodes      per-column type codes (see the {@code T_*} 
constants)
+        * @param dest           pre-allocated per-column backing arrays
+        * @param destOff        absolute row offset of this file's first row
+        * @param limit          exclusive upper row bound of this file's slice
+        * @param tablePath      table path for error messages
+        * @return the number of rows decoded
+        * @throws IOException on read failure
+        */
+       public static int decodeDataFileInto(String filePath, StructType 
physicalSchema, int[] readCodes, Object[] dest,
+               int destOff, int limit, String tablePath) throws IOException {
+               final Configuration conf = 
ConfigurationManager.getCachedJobConf();
+               final int ncol = physicalSchema.length();
+               int off = destOff;
+               try(ParquetFileReader reader = 
ParquetFileReader.open(HadoopInputFile.fromPath(new Path(filePath), conf))) {
+                       MessageType parquetSchema = 
reader.getFooter().getFileMetaData().getSchema();
+                       String createdBy = 
reader.getFooter().getFileMetaData().getCreatedBy();
+                       String[] colNames = 
resolveParquetColumns(physicalSchema, parquetSchema);
+                       GroupConverter root = 
dummyConverter(parquetSchema.getFieldCount());
+                       PageReadStore pages;
+                       while((pages = reader.readNextRowGroup()) != null) {
+                               int nrow = (int) pages.getRowCount();
+                               // guard before decoding: writing past the 
limit would overflow into the
+                               // next file's slice (or off the array) in the 
pre-allocated output

Review Comment:
   This comment is redundant, the code bellow is sufficient



##########
src/main/java/org/apache/sysds/runtime/io/DeltaKernelUtils.java:
##########
@@ -162,6 +174,199 @@ public static int countSelected(int size, boolean[] 
selected) {
                return n;
        }
 
+       // ------------------------------------------
+       // direct parquet decode of Delta data files
+       // ------------------------------------------
+
+       /** Physical-schema metadata key carrying the parquet field id (column 
mapping mode {@code id}). */
+       private static final String PARQUET_FIELD_ID_KEY = "parquet.field.id";
+
+       /**
+        * Whether data files can be decoded directly into pre-allocated output 
columns: the physical read schema must be a
+        * positional 1:1 image of the logical schema, i.e. no partition 
columns (not stored in the data files, spliced back
+        * in by the kernel) and no kernel metadata columns such as {@code 
row_index} (only requested for deletion-vector
+        * reads). Deletion vectors themselves are excluded separately via the 
exact-row-count check.
+        *
+        * @param logicalSchema  the table's logical schema
+        * @param physicalSchema the physical read schema from the scan state
+        * @return true if data files can be decoded directly
+        */
+       public static boolean supportsDirectDecode(StructType logicalSchema, 
StructType physicalSchema) {
+               if(physicalSchema.length() != logicalSchema.length())
+                       return false;
+               for(int c = 0; c < physicalSchema.length(); c++)
+                       if(physicalSchema.at(c).isMetadataColumn())
+                               return false;
+               return true;
+       }
+
+       /** @param scanFileRow a scan-file row @return the fully-qualified path 
of its data file */
+       public static String dataFilePath(Row scanFileRow) {
+               return 
InternalScanFileUtils.getAddFileStatus(scanFileRow).getPath();
+       }
+
+       /**
+        * Decode one Delta data file into pre-allocated typed column arrays at 
the given absolute row offset, through
+        * parquet-mr's column API ({@link ColumnReadStoreImpl}/{@link 
ColumnReader}) with no kernel engine or intermediate
+        * batch vectors in the path. Columns are resolved by parquet field id 
first (column mapping mode {@code id}) and
+        * physical name second; columns absent from the file (schema 
evolution) keep the array defaults (0 for numerics,
+        * null for strings), matching the kernel-path null semantics.
+        *
+        * @param filePath       fully-qualified path of the parquet data file
+        * @param physicalSchema physical read schema (positionally 1:1 with 
the output columns)
+        * @param readCodes      per-column type codes (see the {@code T_*} 
constants)
+        * @param dest           pre-allocated per-column backing arrays
+        * @param destOff        absolute row offset of this file's first row
+        * @param limit          exclusive upper row bound of this file's slice
+        * @param tablePath      table path for error messages
+        * @return the number of rows decoded
+        * @throws IOException on read failure
+        */
+       public static int decodeDataFileInto(String filePath, StructType 
physicalSchema, int[] readCodes, Object[] dest,
+               int destOff, int limit, String tablePath) throws IOException {
+               final Configuration conf = 
ConfigurationManager.getCachedJobConf();
+               final int ncol = physicalSchema.length();
+               int off = destOff;
+               try(ParquetFileReader reader = 
ParquetFileReader.open(HadoopInputFile.fromPath(new Path(filePath), conf))) {
+                       MessageType parquetSchema = 
reader.getFooter().getFileMetaData().getSchema();
+                       String createdBy = 
reader.getFooter().getFileMetaData().getCreatedBy();
+                       String[] colNames = 
resolveParquetColumns(physicalSchema, parquetSchema);
+                       GroupConverter root = 
dummyConverter(parquetSchema.getFieldCount());
+                       PageReadStore pages;
+                       while((pages = reader.readNextRowGroup()) != null) {
+                               int nrow = (int) pages.getRowCount();
+                               // guard before decoding: writing past the 
limit would overflow into the
+                               // next file's slice (or off the array) in the 
pre-allocated output
+                               if(off + nrow > limit)
+                                       throw new DMLRuntimeException("Delta 
file produced more rows than its "
+                                               + "numRecords statistic; 
refusing direct read of " + tablePath);
+                               ColumnReadStoreImpl store = new 
ColumnReadStoreImpl(pages, root, parquetSchema, createdBy);
+                               for(int c = 0; c < ncol; c++) {
+                                       if(colNames[c] == null)
+                                               continue; // column absent from 
this data file -> keep defaults (nulls)
+                                       ColumnDescriptor desc = 
parquetSchema.getColumnDescription(new String[] {colNames[c]});
+                                       
decodeColumnInto(store.getColumnReader(desc), desc.getMaxDefinitionLevel(), 
nrow, readCodes[c],
+                                               dest[c], off);
+                               }
+                               off += nrow;
+                       }
+               }
+               return off - destOff;
+       }
+
+       /**
+        * Resolve each physical-schema column to the parquet column name of 
the given file: by parquet field id when the
+        * schema carries one, by name otherwise, or null when the file does 
not contain the column at all.
+        */
+       private static String[] resolveParquetColumns(StructType schema, 
MessageType parquetSchema) {
+               Map<Integer, String> idToName = new HashMap<>();
+               Map<String, String> names = new HashMap<>();
+               for(int i = 0; i < parquetSchema.getFieldCount(); i++) {
+                       org.apache.parquet.schema.Type t = 
parquetSchema.getType(i);
+                       names.put(t.getName(), t.getName());
+                       if(t.getId() != null)
+                               idToName.put(t.getId().intValue(), t.getName());
+               }
+               String[] resolved = new String[schema.length()];
+               for(int c = 0; c < schema.length(); c++) {
+                       Object fid = 
schema.at(c).getMetadata().get(PARQUET_FIELD_ID_KEY);
+                       String byId = (fid instanceof Number) ? 
idToName.get(((Number) fid).intValue()) : null;
+                       resolved[c] = (byId != null) ? byId : 
names.get(schema.at(c).getName());
+               }
+               return resolved;
+       }
+
+       /** No-op converter tree; the column API requires one, but values are 
pulled via the typed getters. */
+       private static GroupConverter dummyConverter(int nFields) {
+               final PrimitiveConverter[] leaves = new 
PrimitiveConverter[nFields];
+               for(int i = 0; i < nFields; i++)
+                       leaves[i] = new PrimitiveConverter() {
+                       };
+               return new GroupConverter() {
+                       @Override
+                       public Converter getConverter(int fieldIndex) {
+                               return leaves[fieldIndex];
+                       }
+
+                       @Override
+                       public void start() {
+                       }
+
+                       @Override
+                       public void end() {
+                       }
+               };
+       }
+
+       /**
+        * Decode one parquet column of the current row group into a 
pre-allocated typed array at the given offset. Null
+        * cells (definition level below max) keep the array default (0 for 
numerics, null for strings).
+        */
+       private static void decodeColumnInto(ColumnReader creader, int maxDef, 
int nrow, int readCode, Object dest,
+               int off) {
+               switch(readCode) {
+                       case T_DOUBLE: {
+                               double[] a = (double[]) dest;
+                               for(int r = 0; r < nrow; r++) {
+                                       if(creader.getCurrentDefinitionLevel() 
== maxDef)

Review Comment:
   What is this getCurrentDefinitionLevel actually doing ? 
   
   Is it nessesarry for all the rows to call this function ?



##########
src/main/java/org/apache/sysds/runtime/io/DeltaKernelUtils.java:
##########
@@ -162,6 +174,199 @@ public static int countSelected(int size, boolean[] 
selected) {
                return n;
        }
 
+       // ------------------------------------------
+       // direct parquet decode of Delta data files
+       // ------------------------------------------
+
+       /** Physical-schema metadata key carrying the parquet field id (column 
mapping mode {@code id}). */
+       private static final String PARQUET_FIELD_ID_KEY = "parquet.field.id";
+
+       /**
+        * Whether data files can be decoded directly into pre-allocated output 
columns: the physical read schema must be a
+        * positional 1:1 image of the logical schema, i.e. no partition 
columns (not stored in the data files, spliced back
+        * in by the kernel) and no kernel metadata columns such as {@code 
row_index} (only requested for deletion-vector
+        * reads). Deletion vectors themselves are excluded separately via the 
exact-row-count check.
+        *
+        * @param logicalSchema  the table's logical schema
+        * @param physicalSchema the physical read schema from the scan state
+        * @return true if data files can be decoded directly
+        */
+       public static boolean supportsDirectDecode(StructType logicalSchema, 
StructType physicalSchema) {
+               if(physicalSchema.length() != logicalSchema.length())
+                       return false;
+               for(int c = 0; c < physicalSchema.length(); c++)
+                       if(physicalSchema.at(c).isMetadataColumn())
+                               return false;
+               return true;
+       }
+
+       /** @param scanFileRow a scan-file row @return the fully-qualified path 
of its data file */
+       public static String dataFilePath(Row scanFileRow) {
+               return 
InternalScanFileUtils.getAddFileStatus(scanFileRow).getPath();
+       }
+
+       /**
+        * Decode one Delta data file into pre-allocated typed column arrays at 
the given absolute row offset, through
+        * parquet-mr's column API ({@link ColumnReadStoreImpl}/{@link 
ColumnReader}) with no kernel engine or intermediate
+        * batch vectors in the path. Columns are resolved by parquet field id 
first (column mapping mode {@code id}) and
+        * physical name second; columns absent from the file (schema 
evolution) keep the array defaults (0 for numerics,
+        * null for strings), matching the kernel-path null semantics.
+        *
+        * @param filePath       fully-qualified path of the parquet data file
+        * @param physicalSchema physical read schema (positionally 1:1 with 
the output columns)
+        * @param readCodes      per-column type codes (see the {@code T_*} 
constants)
+        * @param dest           pre-allocated per-column backing arrays
+        * @param destOff        absolute row offset of this file's first row
+        * @param limit          exclusive upper row bound of this file's slice
+        * @param tablePath      table path for error messages
+        * @return the number of rows decoded
+        * @throws IOException on read failure
+        */
+       public static int decodeDataFileInto(String filePath, StructType 
physicalSchema, int[] readCodes, Object[] dest,
+               int destOff, int limit, String tablePath) throws IOException {
+               final Configuration conf = 
ConfigurationManager.getCachedJobConf();
+               final int ncol = physicalSchema.length();
+               int off = destOff;
+               try(ParquetFileReader reader = 
ParquetFileReader.open(HadoopInputFile.fromPath(new Path(filePath), conf))) {
+                       MessageType parquetSchema = 
reader.getFooter().getFileMetaData().getSchema();
+                       String createdBy = 
reader.getFooter().getFileMetaData().getCreatedBy();
+                       String[] colNames = 
resolveParquetColumns(physicalSchema, parquetSchema);
+                       GroupConverter root = 
dummyConverter(parquetSchema.getFieldCount());
+                       PageReadStore pages;
+                       while((pages = reader.readNextRowGroup()) != null) {
+                               int nrow = (int) pages.getRowCount();
+                               // guard before decoding: writing past the 
limit would overflow into the
+                               // next file's slice (or off the array) in the 
pre-allocated output
+                               if(off + nrow > limit)
+                                       throw new DMLRuntimeException("Delta 
file produced more rows than its "
+                                               + "numRecords statistic; 
refusing direct read of " + tablePath);
+                               ColumnReadStoreImpl store = new 
ColumnReadStoreImpl(pages, root, parquetSchema, createdBy);
+                               for(int c = 0; c < ncol; c++) {
+                                       if(colNames[c] == null)
+                                               continue; // column absent from 
this data file -> keep defaults (nulls)
+                                       ColumnDescriptor desc = 
parquetSchema.getColumnDescription(new String[] {colNames[c]});
+                                       
decodeColumnInto(store.getColumnReader(desc), desc.getMaxDefinitionLevel(), 
nrow, readCodes[c],
+                                               dest[c], off);
+                               }
+                               off += nrow;
+                       }
+               }
+               return off - destOff;
+       }
+
+       /**
+        * Resolve each physical-schema column to the parquet column name of 
the given file: by parquet field id when the
+        * schema carries one, by name otherwise, or null when the file does 
not contain the column at all.
+        */
+       private static String[] resolveParquetColumns(StructType schema, 
MessageType parquetSchema) {
+               Map<Integer, String> idToName = new HashMap<>();
+               Map<String, String> names = new HashMap<>();
+               for(int i = 0; i < parquetSchema.getFieldCount(); i++) {
+                       org.apache.parquet.schema.Type t = 
parquetSchema.getType(i);
+                       names.put(t.getName(), t.getName());
+                       if(t.getId() != null)
+                               idToName.put(t.getId().intValue(), t.getName());
+               }
+               String[] resolved = new String[schema.length()];
+               for(int c = 0; c < schema.length(); c++) {
+                       Object fid = 
schema.at(c).getMetadata().get(PARQUET_FIELD_ID_KEY);
+                       String byId = (fid instanceof Number) ? 
idToName.get(((Number) fid).intValue()) : null;
+                       resolved[c] = (byId != null) ? byId : 
names.get(schema.at(c).getName());
+               }
+               return resolved;
+       }
+
+       /** No-op converter tree; the column API requires one, but values are 
pulled via the typed getters. */
+       private static GroupConverter dummyConverter(int nFields) {
+               final PrimitiveConverter[] leaves = new 
PrimitiveConverter[nFields];
+               for(int i = 0; i < nFields; i++)
+                       leaves[i] = new PrimitiveConverter() {
+                       };
+               return new GroupConverter() {
+                       @Override
+                       public Converter getConverter(int fieldIndex) {
+                               return leaves[fieldIndex];
+                       }
+
+                       @Override
+                       public void start() {
+                       }
+
+                       @Override
+                       public void end() {
+                       }
+               };
+       }
+
+       /**
+        * Decode one parquet column of the current row group into a 
pre-allocated typed array at the given offset. Null
+        * cells (definition level below max) keep the array default (0 for 
numerics, null for strings).
+        */
+       private static void decodeColumnInto(ColumnReader creader, int maxDef, 
int nrow, int readCode, Object dest,
+               int off) {
+               switch(readCode) {
+                       case T_DOUBLE: {
+                               double[] a = (double[]) dest;
+                               for(int r = 0; r < nrow; r++) {
+                                       if(creader.getCurrentDefinitionLevel() 
== maxDef)
+                                               a[off + r] = 
creader.getDouble();

Review Comment:
   we can move the off + r out to the forloop, this can be done for all the 
cases.
   
   ```
   int end = nrow + off;
   for(int r = off; r < end; r++) {
   }
   ```



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to