MCPcopy Create free account
hub / github.com/apache/arrow / ReadRowGroups

Method ReadRowGroups

cpp/src/parquet/arrow/reader.cc:1303–1502  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

1301}
1302
1303Result<std::shared_ptr<Table>> FileReaderImpl::ReadRowGroups(
1304 const std::vector<int>& row_groups, const std::vector<int>& column_indices) {
1305 RETURN_NOT_OK(BoundsCheck(row_groups, column_indices));
1306
1307 // PARQUET-1698/PARQUET-1820: pre-buffer row groups/column chunks if enabled
1308 if (reader_properties_.pre_buffer()) {
1309 BEGIN_PARQUET_CATCH_EXCEPTIONS
1310 parquet_reader()->PreBuffer(row_groups, column_indices,
1311 reader_properties_.io_context(),
1312 reader_properties_.cache_options());
1313 END_PARQUET_CATCH_EXCEPTIONS
1314 }
1315
1316 auto fut = DecodeRowGroups(/*self=*/nullptr, row_groups, column_indices,
1317 /*cpu_executor=*/nullptr);
1318 return fut.MoveResult();
1319}
1320
1321Future<std::shared_ptr<Table>> FileReaderImpl::DecodeRowGroups(
1322 std::shared_ptr<FileReaderImpl> self, const std::vector<int>& row_groups,
1323 const std::vector<int>& column_indices, ::arrow::internal::Executor* cpu_executor) {
1324 // `self` is used solely to keep `this` alive in an async context - but we use this
1325 // in a sync context too so use `this` over `self`
1326 std::vector<std::shared_ptr<ColumnReaderImpl>> readers;
1327 std::shared_ptr<::arrow::Schema> result_schema;
1328 RETURN_NOT_OK(GetFieldReaders(column_indices, row_groups, &readers, &result_schema));
1329 // OptionalParallelForAsync requires an executor
1330 if (!cpu_executor) cpu_executor = ::arrow::internal::GetCpuThreadPool();
1331
1332 auto read_column = [row_groups, self, this](size_t i,
1333 std::shared_ptr<ColumnReaderImpl> reader)
1334 -> ::arrow::Result<std::shared_ptr<::arrow::ChunkedArray>> {
1335 std::shared_ptr<::arrow::ChunkedArray> column;
1336 RETURN_NOT_OK(ReadColumn(static_cast<int>(i), row_groups, reader.get(), &column));
1337 return column;
1338 };
1339 auto make_table = [result_schema, row_groups, self,
1340 this](const ::arrow::ChunkedArrayVector& columns)
1341 -> ::arrow::Result<std::shared_ptr<Table>> {
1342 int64_t num_rows = 0;
1343 if (!columns.empty()) {
1344 num_rows = columns[0]->length();
1345 } else {
1346 for (int i : row_groups) {
1347 num_rows += parquet_reader()->metadata()->RowGroup(i)->num_rows();
1348 }
1349 }
1350 auto table = Table::Make(std::move(result_schema), columns, num_rows);
1351 RETURN_NOT_OK(table->Validate());
1352 return table;
1353 };
1354 return ::arrow::internal::OptionalParallelForAsync(reader_properties_.use_threads(),
1355 std::move(readers), read_column,
1356 cpu_executor)
1357 .Then(std::move(make_table));
1358}
1359
1360std::shared_ptr<RowGroupReader> FileReaderImpl::RowGroup(int row_group_index) {

Calls 3

BoundsCheckFunction · 0.85
ARROW_ASSIGN_OR_RAISEFunction · 0.70
OKFunction · 0.50

Tested by

no test coverage detected