Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 18 additions & 10 deletions src/Processors/Formats/Impl/Parquet/Reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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
{
Expand All @@ -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);
}
}
Expand Down Expand Up @@ -1386,41 +1386,47 @@ 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<size_t> 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;
}

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)
Expand All @@ -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;
}
}

Expand Down
3 changes: 2 additions & 1 deletion src/Processors/Formats/Impl/Parquet/Reader.h
Original file line number Diff line number Diff line change
Expand Up @@ -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<const char> 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<size_t> row_idx, ColumnChunk & column, const PrimitiveColumnInfo & column_info);
std::tuple<parq::PageHeader, std::span<const char>> 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<size_t> end_row_idx, size_t target_row_idx, ColumnChunk & column, const PrimitiveColumnInfo & column_info);
void decompressPageIfCompressed(PageState & page);
Expand Down
9 changes: 9 additions & 0 deletions tests/queries/0_stateless/00900_long_parquet_load_2.reference
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Binary file not shown.
Binary file not shown.
Loading