diff --git a/src/Processors/Formats/Impl/Parquet/Reader.cpp b/src/Processors/Formats/Impl/Parquet/Reader.cpp index 24ba5277dd70..01cd1fedc8c5 100644 --- a/src/Processors/Formats/Impl/Parquet/Reader.cpp +++ b/src/Processors/Formats/Impl/Parquet/Reader.cpp @@ -1315,7 +1315,7 @@ void Reader::decodePrimitiveColumn(ColumnChunk & column, const PrimitiveColumnIn size_t end_row_idx = start_row_idx + num_rows; row_subidx += num_rows; - skipToRow(start_row_idx, column, column_info); + skipToRowOrNextPage(start_row_idx, column, column_info); while (true) // loop over pages { @@ -1331,7 +1331,7 @@ void Reader::decodePrimitiveColumn(ColumnChunk & column, const PrimitiveColumnIn /// Advance to next page. chassert(page.value_idx == page.num_values); - skipToRow(page.next_row_idx, column, column_info); + skipToRowOrNextPage(std::nullopt, column, column_info); chassert(page.value_idx == 0); } } @@ -1386,32 +1386,38 @@ void Reader::decodePrimitiveColumn(ColumnChunk & column, const PrimitiveColumnIn } } -void Reader::skipToRow(size_t row_idx, ColumnChunk & column, const PrimitiveColumnInfo & column_info) +void Reader::skipToRowOrNextPage(std::optional row_idx, ColumnChunk & column, const PrimitiveColumnInfo & column_info) { /// True if column.page is initialized and contains the requested row_idx. bool found_page = false; auto & page = column.page; - if (page.initialized && page.value_idx < page.num_values && page.end_row_idx.has_value() && *page.end_row_idx > row_idx) + if (!row_idx.has_value()) + chassert(page.initialized); + + if (row_idx.has_value() && page.initialized && page.value_idx < page.num_values && + page.end_row_idx.has_value() && *page.end_row_idx > *row_idx) /// Fast path: we're just continuing reading the same page as before. found_page = true; if (!found_page && !column.data_pages.empty()) { /// If we have offset index, find the row index there and jump to the correct page. + if (!row_idx.has_value()) + row_idx = column.data_pages[column.data_pages_idx].end_row_idx; while (column.data_pages_idx < column.data_pages.size() && - column.data_pages[column.data_pages_idx].end_row_idx <= row_idx) + column.data_pages[column.data_pages_idx].end_row_idx <= *row_idx) ++column.data_pages_idx; if (column.data_pages_idx == column.data_pages.size()) throw Exception(ErrorCodes::INCORRECT_DATA, "Parquet offset index covers too few rows"); const auto & page_info = column.data_pages[column.data_pages_idx]; size_t first_row_idx = size_t(page_info.meta->first_row_index); - if (first_row_idx > row_idx) + if (first_row_idx > *row_idx) throw Exception(ErrorCodes::LOGICAL_ERROR, "Row passes filters but its page was not selected for reading. This is a bug."); auto data = prefetcher.getRangeData(page_info.prefetch); const char * ptr = data.data(); - if (!initializeDataPage(ptr, ptr + data.size(), first_row_idx, page_info.end_row_idx, row_idx, column, column_info)) + if (!initializeDataPage(ptr, ptr + data.size(), first_row_idx, page_info.end_row_idx, *row_idx, column, column_info)) throw Exception(ErrorCodes::LOGICAL_ERROR, "Page doesn't contain requested row"); found_page = true; } @@ -1419,8 +1425,8 @@ void Reader::skipToRow(size_t row_idx, ColumnChunk & column, const PrimitiveColu while (true) { /// Skip rows inside the page. - if (page.initialized && page.value_idx < page.num_values && - skipRowsInPage(row_idx, page, column, column_info)) + if (row_idx.has_value() && page.initialized && page.value_idx < page.num_values && + skipRowsInPage(*row_idx, page, column, column_info)) return; if (found_page) @@ -1433,8 +1439,10 @@ void Reader::skipToRow(size_t row_idx, ColumnChunk & column, const PrimitiveColu chassert(column.next_page_offset <= all_pages.size()); const char * ptr = all_pages.data() + column.next_page_offset; const char * end = all_pages.data() + all_pages.size(); - initializeDataPage(ptr, end, page.next_row_idx, /*end_row_idx=*/ std::nullopt, row_idx, column, column_info); + initializeDataPage(ptr, end, page.next_row_idx, /*end_row_idx=*/ std::nullopt, row_idx.value_or(page.next_row_idx), column, column_info); column.next_page_offset = ptr - all_pages.data(); + if (!row_idx.has_value()) + return; } } diff --git a/src/Processors/Formats/Impl/Parquet/Reader.h b/src/Processors/Formats/Impl/Parquet/Reader.h index 296cb41419a2..ea0917fbf8bb 100644 --- a/src/Processors/Formats/Impl/Parquet/Reader.h +++ b/src/Processors/Formats/Impl/Parquet/Reader.h @@ -523,7 +523,8 @@ struct Reader void initializePrefetches(); double estimateAverageStringLengthPerRow(const ColumnChunk & column, const RowGroup & row_group) const; void decodeDictionaryPageImpl(const parq::PageHeader & header, std::span data, ColumnChunk & column, const PrimitiveColumnInfo & column_info); - void skipToRow(size_t row_idx, ColumnChunk & column, const PrimitiveColumnInfo & column_info); + /// If row_idx is provided, jump to the start of that row. Otherwise go to the start of next page. + void skipToRowOrNextPage(std::optional row_idx, ColumnChunk & column, const PrimitiveColumnInfo & column_info); std::tuple> decodeAndCheckPageHeader(const char * & data_ptr, const char * data_end) const; bool initializeDataPage(const char * & data_ptr, const char * data_end, size_t next_row_idx, std::optional end_row_idx, size_t target_row_idx, ColumnChunk & column, const PrimitiveColumnInfo & column_info); void decompressPageIfCompressed(PageState & page); diff --git a/tests/queries/0_stateless/00900_long_parquet_load_2.reference b/tests/queries/0_stateless/00900_long_parquet_load_2.reference index e3f9cd89e089..56dc8c5bae18 100644 --- a/tests/queries/0_stateless/00900_long_parquet_load_2.reference +++ b/tests/queries/0_stateless/00900_long_parquet_load_2.reference @@ -111,6 +111,15 @@ 1 7 false 1 1 1 10 1.1 10.1 04/01/09 1 2009-04-01 00:01:00.000000000 1 10118128356831336697 \N \N \N \N \N \N \N \N \N \N \N \N 2 14757089928194662262 +=== Try load data from array_across_pages_1.parquet +0 [1,2] 1 5845551876324623491 + +\N [] 1 5845551876324623491 +=== Try load data from array_across_pages_2.parquet +0 [1,2] 1 5845551876324623491 +1 [3] 1 204286521036002596 + +\N [] 2 6049838397360626087 === Try load data from array_float.parquet 9 idx10 [10.2,8.2] 1 6500220947925738428 0 idx1 [] 1 12173579943307849648 diff --git a/tests/queries/0_stateless/data_parquet/array_across_pages_1.parquet b/tests/queries/0_stateless/data_parquet/array_across_pages_1.parquet new file mode 100644 index 000000000000..303a6819a065 Binary files /dev/null and b/tests/queries/0_stateless/data_parquet/array_across_pages_1.parquet differ diff --git a/tests/queries/0_stateless/data_parquet/array_across_pages_2.parquet b/tests/queries/0_stateless/data_parquet/array_across_pages_2.parquet new file mode 100644 index 000000000000..51740029cc3a Binary files /dev/null and b/tests/queries/0_stateless/data_parquet/array_across_pages_2.parquet differ