Skip to content
Open
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
3 changes: 3 additions & 0 deletions src/iceberg/file_reader.h
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,9 @@ class ICEBERG_EXPORT ReaderProperties : public ConfigBase<ReaderProperties> {

/// \brief The batch size to read.
inline static Entry<int64_t> kBatchSize{"read.batch-size", 4096};
/// \brief Read list columns as Arrow large_list (64-bit offsets) instead of list.
/// Default: false (use 32-bit offset list).
inline static Entry<bool> kArrowUseLargeList{"read.arrow.use-large-list", false};
/// \brief Skip GenericDatum in Avro reader for better performance.
/// When true, decode directly from Avro to Arrow without GenericDatum intermediate.
/// Default: true (skip GenericDatum for better performance).
Expand Down
104 changes: 98 additions & 6 deletions src/iceberg/parquet/parquet_reader.cc
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@

#include "iceberg/parquet/parquet_reader.h"

#include <algorithm>
#include <numeric>

#include <arrow/c/bridge.h>
Expand All @@ -41,6 +42,7 @@
#include "iceberg/result.h"
#include "iceberg/schema_internal.h"
#include "iceberg/schema_util.h"
#include "iceberg/util/checked_cast.h"
#include "iceberg/util/macros.h"

namespace iceberg::parquet {
Expand Down Expand Up @@ -84,6 +86,78 @@ class EmptyRecordBatchReader : public ::arrow::RecordBatchReader {
}
};

// forward declaration to unblock cycle dependence.
std::shared_ptr<::arrow::Field> UseLargeListField(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

add comment:

Suggested change
std::shared_ptr<::arrow::Field> UseLargeListField(
// forward declaration to unblock cycle dependence.
std::shared_ptr<::arrow::Field> UseLargeListField(

const std::shared_ptr<::arrow::Field>& field);

// Rebuild a data type with all nested list types replaced by large_list.
std::shared_ptr<::arrow::DataType> UseLargeListType(
const std::shared_ptr<::arrow::DataType>& type) {
switch (type->id()) {
case ::arrow::Type::LIST: {
const auto& list_type = internal::checked_cast<const ::arrow::ListType&>(*type);
return ::arrow::large_list(UseLargeListField(list_type.value_field()));
}
case ::arrow::Type::STRUCT: {
::arrow::FieldVector fields;
fields.reserve(type->num_fields());
for (const auto& field : type->fields()) {
fields.push_back(UseLargeListField(field));
}
return ::arrow::struct_(std::move(fields));
}
case ::arrow::Type::MAP: {
const auto& map_type = internal::checked_cast<const ::arrow::MapType&>(*type);
return std::make_shared<::arrow::MapType>(UseLargeListField(map_type.key_field()),
UseLargeListField(map_type.item_field()),
map_type.keys_sorted());
}
default:
return type;
}
}

std::shared_ptr<::arrow::Field> UseLargeListField(
const std::shared_ptr<::arrow::Field>& field) {
return field->WithType(UseLargeListType(field->type()));
}

// Rewrite all fields in a field vector to use large_list instead of list.
::arrow::FieldVector UseLargeListFields(const ::arrow::FieldVector& fields) {
::arrow::FieldVector rewritten;
rewritten.reserve(fields.size());
for (const auto& field : fields) {
rewritten.push_back(UseLargeListField(field));
}
return rewritten;
}

// Returns true if the type contains a large_list, at any level of nesting.
bool ContainsLargeList(const ::arrow::DataType& type) {
if (type.id() == ::arrow::Type::LARGE_LIST) {
return true;
}
return std::ranges::any_of(
type.fields(), [](const auto& field) { return ContainsLargeList(*field->type()); });
}

// Returns true if the reader produces large_list arrays.
//
// Arrow honors the requested large_list type only when it derives the Arrow schema from
// the Parquet schema. A file that carries serialized ARROW:schema metadata keeps its
// original list type instead, so whether large lists are produced can only be told from
// the schema of the reader.
bool ProducesLargeList(const ::arrow::RecordBatchReader& reader) {
const auto& schema = reader.schema();
if (schema == nullptr) {
// an empty reader produces no arrays to be described
return false;
}
return std::ranges::any_of(schema->fields(), [](const auto& field) {
return ContainsLargeList(*field->type());
});
}

} // namespace

// A stateful context to keep track of the reading progress.
Expand Down Expand Up @@ -118,6 +192,10 @@ class ParquetReader::Impl {
arrow_reader_properties.set_batch_size(
options.properties.Get(ReaderProperties::kBatchSize));
arrow_reader_properties.set_arrow_extensions_enabled(true);
use_large_list_ = options.properties.Get(ReaderProperties::kArrowUseLargeList);
if (use_large_list_) {
arrow_reader_properties.set_list_type(::arrow::Type::LARGE_LIST);
}

// Open the Parquet file reader
ICEBERG_ASSIGN_OR_RAISE(input_stream_, OpenInputStream(options));
Expand Down Expand Up @@ -212,12 +290,6 @@ class ParquetReader::Impl {
Status InitReadContext() {
context_ = std::make_unique<ReadContext>();

// Build the output Arrow schema
ArrowSchema arrow_schema;
ICEBERG_RETURN_UNEXPECTED(ToArrowSchema(*read_schema_, &arrow_schema));
ICEBERG_ARROW_ASSIGN_OR_RETURN(context_->output_arrow_schema_,
::arrow::ImportSchema(&arrow_schema));

// Row group pruning based on the split
// TODO(gangwu): add row group filtering based on zone map, bloom filter, etc.
std::vector<int> row_group_indices;
Expand Down Expand Up @@ -250,6 +322,24 @@ class ParquetReader::Impl {
reader_->GetRecordBatchReader(row_group_indices, column_indices));
}

// Build the output Arrow schema from the projected Iceberg schema. This schema is the
// target of ProjectRecordBatch, so it must describe the projected schema rather than
// the schema of the file.
ArrowSchema arrow_schema;
ICEBERG_RETURN_UNEXPECTED(ToArrowSchema(*read_schema_, &arrow_schema));
ICEBERG_ARROW_ASSIGN_OR_RETURN(context_->output_arrow_schema_,
::arrow::ImportSchema(&arrow_schema));

if (use_large_list_ && ProducesLargeList(*context_->record_batch_reader_)) {
// Align the output schema with the large_list arrays produced by the Parquet
// reader. Note that Arrow ignores the requested list type when the file carries
// serialized ARROW:schema metadata, in which case the reader keeps producing plain
// list arrays and the output schema must keep describing them as such.
context_->output_arrow_schema_ =
::arrow::schema(UseLargeListFields(context_->output_arrow_schema_->fields()),
context_->output_arrow_schema_->metadata());
}

return {};
}

Expand All @@ -258,6 +348,8 @@ class ParquetReader::Impl {
::arrow::MemoryPool* pool_ = ::arrow::default_memory_pool();
// The split to read from the Parquet file.
std::optional<Split> split_;
// Whether to read list columns as large_list (64-bit offsets).
bool use_large_list_ = false;
// Schema to read from the Parquet file.
std::shared_ptr<::iceberg::Schema> read_schema_;
// The projection result to apply to the read schema.
Expand Down
3 changes: 2 additions & 1 deletion src/iceberg/parquet/parquet_schema_util.cc
Original file line number Diff line number Diff line change
Expand Up @@ -251,7 +251,8 @@ Status ValidateParquetSchemaEvolution(
}
break;
case TypeId::kList:
if (arrow_type->id() == ::arrow::Type::LIST) {
if (arrow_type->id() == ::arrow::Type::LIST ||
arrow_type->id() == ::arrow::Type::LARGE_LIST) {
return {};
}
break;
Expand Down
46 changes: 45 additions & 1 deletion src/iceberg/test/parquet_schema_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -111,12 +111,14 @@ ::parquet::schema::NodePtr MakeMapNode(const std::string& name,

// Helper to create SchemaManifest from Parquet schema
::parquet::arrow::SchemaManifest MakeSchemaManifest(
const ::parquet::schema::NodePtr& parquet_schema) {
const ::parquet::schema::NodePtr& parquet_schema,
::arrow::Type::type list_type = ::arrow::Type::LIST) {
auto parquet_schema_descriptor = std::make_shared<::parquet::SchemaDescriptor>();
parquet_schema_descriptor->Init(parquet_schema);

auto properties = ::parquet::default_arrow_reader_properties();
properties.set_arrow_extensions_enabled(true);
properties.set_list_type(list_type);

::parquet::arrow::SchemaManifest manifest;
auto status = ::parquet::arrow::SchemaManifest::Make(parquet_schema_descriptor.get(),
Expand Down Expand Up @@ -340,6 +342,16 @@ TEST(ParquetSchemaProjectionTest, ValidateSchemaEvolutionAllowsNullPhysicalType)
ASSERT_THAT(status, IsOk());
}

TEST(ParquetSchemaProjectionTest, ValidateSchemaEvolutionAllowsLargeList) {
::parquet::arrow::SchemaField parquet_field;
parquet_field.field = ::arrow::field("numbers", ::arrow::large_list(::arrow::int32()));

ListType expected_type(
SchemaField::MakeOptional(/*field_id=*/101, "element", iceberg::int32()));
auto status = ValidateParquetSchemaEvolution(expected_type, parquet_field);
ASSERT_THAT(status, IsOk());
}

TEST(ParquetSchemaProjectionTest, ProjectNullPhysicalFieldsAsNull) {
Schema expected_schema({
SchemaField::MakeOptional(/*field_id=*/1, "age", iceberg::int32()),
Expand Down Expand Up @@ -624,6 +636,38 @@ TEST(ParquetSchemaProjectionTest, ProjectListType) {
ASSERT_EQ(SelectedColumnIndices(projection), std::vector<int32_t>({0, 1}));
}

TEST(ParquetSchemaProjectionTest, ProjectLargeListType) {
Schema expected_schema({
SchemaField::MakeOptional(
/*field_id=*/2, "numbers",
std::make_shared<ListType>(SchemaField::MakeOptional(
/*field_id=*/101, "element", iceberg::int32()))),
});

auto parquet_schema = MakeGroupNode(
"iceberg_schema",
{
MakeListNode("numbers", MakeInt32Node("element", /*field_id=*/101),
/*field_id=*/2),
});

auto schema_manifest = MakeSchemaManifest(parquet_schema, ::arrow::Type::LARGE_LIST);
ASSERT_EQ(schema_manifest.schema_fields[0].field->type()->id(),
::arrow::Type::LARGE_LIST);

auto projection_result = Project(expected_schema, schema_manifest);
ASSERT_THAT(projection_result, IsOk());

const auto& projection = *projection_result;
ASSERT_EQ(projection.fields.size(), 1);
ASSERT_PROJECTED_FIELD(projection.fields[0], 0);

ASSERT_EQ(projection.fields[0].children.size(), 1);
ASSERT_PROJECTED_FIELD(projection.fields[0].children[0], 0);

ASSERT_EQ(SelectedColumnIndices(projection), std::vector<int32_t>({0}));
}

TEST(ParquetSchemaProjectionTest, ProjectMapType) {
Schema expected_schema({
SchemaField::MakeOptional(
Expand Down
Loading
Loading