diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java index d12086dd1b..37e8db5b80 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java @@ -59,6 +59,7 @@ class InternalParquetRecordWriter { private long recordCount = 0; private long recordCountForNextMemCheck = MINIMUM_RECORD_COUNT_FOR_CHECK; + private long lastRowGroupEndPos = 0; private ColumnWriteStore columnStore; private ColumnChunkPageWriteStore pageStore; @@ -122,6 +123,13 @@ public void write(T value) throws IOException, InterruptedException { checkBlockSizeReached(); } + /** + * @return the total size of data written to the file and buffered in memory + */ + public long getDataSize() { + return lastRowGroupEndPos + columnStore.getBufferedSize(); + } + private void checkBlockSizeReached() throws IOException { if (recordCount >= recordCountForNextMemCheck) { // checking the memory size is relatively expensive, so let's not do it for every record. long memSize = columnStore.getBufferedSize(); @@ -133,6 +141,7 @@ private void checkBlockSizeReached() throws IOException { flushRowGroupToStore(); initStore(); recordCountForNextMemCheck = min(max(MINIMUM_RECORD_COUNT_FOR_CHECK, recordCount / 2), MAXIMUM_RECORD_COUNT_FOR_CHECK); + this.lastRowGroupEndPos = parquetFileWriter.getPos(); } else { recordCountForNextMemCheck = min( max(MINIMUM_RECORD_COUNT_FOR_CHECK, (recordCount + (long)(nextRowGroupSize / ((float)recordSize))) / 2), // will check halfway diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetWriter.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetWriter.java index 70abdaca79..e3b7953760 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetWriter.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetWriter.java @@ -300,6 +300,13 @@ public void close() throws IOException { } } + /** + * @return the total size of data written to the file and buffered in memory + */ + public long getDataSize() { + return writer.getDataSize(); + } + /** * An abstract builder class for ParquetWriter instances. *