zhjwpku commented on code in PR #880:
URL: https://github.com/apache/iceberg-cpp/pull/880#discussion_r3742495324
##########
src/iceberg/parquet/parquet_schema_util.cc:
##########
@@ -60,6 +60,193 @@ std::optional<int32_t> GetFieldId(const
::parquet::arrow::SchemaField& parquet_f
return FieldIdFromMetadata(parquet_field.field->metadata());
}
+Result<::parquet::LogicalType::EdgeInterpolationAlgorithm> ToParquetAlgorithm(
+ EdgeAlgorithm algorithm) {
+ using ParquetAlgorithm = ::parquet::LogicalType::EdgeInterpolationAlgorithm;
+ switch (algorithm) {
+ case EdgeAlgorithm::kSpherical:
+ return ParquetAlgorithm::SPHERICAL;
+ case EdgeAlgorithm::kVincenty:
+ return ParquetAlgorithm::VINCENTY;
+ case EdgeAlgorithm::kThomas:
+ return ParquetAlgorithm::THOMAS;
+ case EdgeAlgorithm::kAndoyer:
+ return ParquetAlgorithm::ANDOYER;
+ case EdgeAlgorithm::kKarney:
+ return ParquetAlgorithm::KARNEY;
+ }
+ return InvalidArgument("Unknown Iceberg edge algorithm");
+}
+
+Result<::parquet::schema::NodePtr> ToParquetNode(const SchemaField& field);
+
+Result<::parquet::schema::NodePtr> ToParquetPrimitive(
+ const PrimitiveType& primitive, std::string name,
+ ::parquet::Repetition::type repetition, int32_t field_id) {
+ using ParquetType = ::parquet::Type;
+ using TimeUnit = ::parquet::LogicalType::TimeUnit;
+
+ auto logical_type = ::parquet::LogicalType::None();
+ auto physical_type = ParquetType::UNDEFINED;
+ int32_t type_length = -1;
+
+ switch (primitive.type_id()) {
+ case TypeId::kBoolean:
+ physical_type = ParquetType::BOOLEAN;
+ break;
+ case TypeId::kInt:
+ physical_type = ParquetType::INT32;
+ break;
+ case TypeId::kLong:
+ physical_type = ParquetType::INT64;
+ break;
+ case TypeId::kFloat:
+ physical_type = ParquetType::FLOAT;
+ break;
+ case TypeId::kDouble:
+ physical_type = ParquetType::DOUBLE;
+ break;
+ case TypeId::kDecimal: {
+ const auto& decimal = internal::checked_cast<const
DecimalType&>(primitive);
+ if (decimal.precision() <= 9) {
+ physical_type = ParquetType::INT32;
+ } else if (decimal.precision() <= 18) {
+ physical_type = ParquetType::INT64;
+ } else {
+ physical_type = ParquetType::FIXED_LEN_BYTE_ARRAY;
+ type_length = ::arrow::DecimalType::DecimalSize(decimal.precision());
+ }
+ logical_type =
+ ::parquet::LogicalType::Decimal(decimal.precision(),
decimal.scale());
+ } break;
+ case TypeId::kDate:
+ physical_type = ParquetType::INT32;
+ logical_type = ::parquet::LogicalType::Date();
+ break;
+ case TypeId::kTime:
+ physical_type = ParquetType::INT64;
+ logical_type =
+ ::parquet::LogicalType::Time(/*is_adjusted_to_utc=*/false,
TimeUnit::MICROS);
+ break;
+ case TypeId::kTimestamp:
+ physical_type = ParquetType::INT64;
+ logical_type =
::parquet::LogicalType::Timestamp(/*is_adjusted_to_utc=*/false,
+ TimeUnit::MICROS);
+ break;
+ case TypeId::kTimestampTz:
+ physical_type = ParquetType::INT64;
+ logical_type =
::parquet::LogicalType::Timestamp(/*is_adjusted_to_utc=*/true,
+ TimeUnit::MICROS);
+ break;
+ case TypeId::kTimestampNs:
+ physical_type = ParquetType::INT64;
+ logical_type =
::parquet::LogicalType::Timestamp(/*is_adjusted_to_utc=*/false,
+ TimeUnit::NANOS);
+ break;
+ case TypeId::kTimestampTzNs:
+ physical_type = ParquetType::INT64;
+ logical_type =
+ ::parquet::LogicalType::Timestamp(/*is_adjusted_to_utc=*/true,
TimeUnit::NANOS);
+ break;
+ case TypeId::kString:
+ physical_type = ParquetType::BYTE_ARRAY;
+ logical_type = ::parquet::LogicalType::String();
+ break;
+ case TypeId::kUuid:
+ physical_type = ParquetType::FIXED_LEN_BYTE_ARRAY;
+ type_length = 16;
+ logical_type = ::parquet::LogicalType::UUID();
+ break;
+ case TypeId::kFixed: {
+ const auto& fixed = internal::checked_cast<const FixedType&>(primitive);
+ physical_type = ParquetType::FIXED_LEN_BYTE_ARRAY;
+ type_length = fixed.length();
+ } break;
+ case TypeId::kBinary:
+ physical_type = ParquetType::BYTE_ARRAY;
+ break;
+ case TypeId::kGeometry: {
+ const auto& geometry = internal::checked_cast<const
GeometryType&>(primitive);
+ physical_type = ParquetType::BYTE_ARRAY;
+ logical_type =
::parquet::LogicalType::Geometry(std::string(geometry.crs()));
+ } break;
+ case TypeId::kGeography: {
+ const auto& geography = internal::checked_cast<const
GeographyType&>(primitive);
+ ICEBERG_ASSIGN_OR_RAISE(auto algorithm,
ToParquetAlgorithm(geography.algorithm()));
+ physical_type = ParquetType::BYTE_ARRAY;
+ logical_type =
+ ::parquet::LogicalType::Geography(std::string(geography.crs()),
algorithm);
+ } break;
+ case TypeId::kUnknown:
+ if (repetition != ::parquet::Repetition::OPTIONAL) {
+ return InvalidSchema("Iceberg unknown type must be optional");
+ }
+ physical_type = ParquetType::INT32;
+ logical_type = ::parquet::LogicalType::Null();
+ break;
+ case TypeId::kVariant:
+ case TypeId::kStruct:
+ case TypeId::kList:
+ case TypeId::kMap:
+ return NotSupported("Cannot write Iceberg type {} to Parquet",
primitive);
+ }
+
+ return ::parquet::schema::PrimitiveNode::Make(name, repetition,
std::move(logical_type),
+ physical_type, type_length,
field_id);
+}
+
+Result<::parquet::schema::NodePtr> ToParquetNode(const SchemaField& field) {
+ const auto repetition = field.optional() ? ::parquet::Repetition::OPTIONAL
+ : ::parquet::Repetition::REQUIRED;
+ const auto& type = *field.type();
+ if (type.is_primitive()) {
+ return ToParquetPrimitive(internal::checked_cast<const
PrimitiveType&>(type),
+ std::string(field.name()), repetition,
field.field_id());
+ }
+
+ switch (type.type_id()) {
+ case TypeId::kStruct: {
+ const auto& struct_type = internal::checked_cast<const
StructType&>(type);
+ if (struct_type.fields().empty()) {
+ return NotImplemented("Cannot write empty struct '{}' to Parquet",
field.name());
+ }
+ ::parquet::schema::NodeVector children;
+ children.reserve(struct_type.fields().size());
+ for (const auto& child_field : struct_type.fields()) {
+ ICEBERG_ASSIGN_OR_RAISE(auto child, ToParquetNode(child_field));
+ children.push_back(std::move(child));
+ }
+ return ::parquet::schema::GroupNode::Make(
+ std::string(field.name()), repetition, std::move(children),
+ /*logical_type=*/nullptr, field.field_id());
+ }
+ case TypeId::kList: {
+ const auto& list_type = internal::checked_cast<const ListType&>(type);
+ ICEBERG_ASSIGN_OR_RAISE(auto element,
ToParquetNode(list_type.element()));
+ auto repeated = ::parquet::schema::GroupNode::Make(
+ "list", ::parquet::Repetition::REPEATED, {std::move(element)});
+ return ::parquet::schema::GroupNode::Make(
+ std::string(field.name()), repetition, {std::move(repeated)},
+ ::parquet::LogicalType::List(), field.field_id());
+ }
+ case TypeId::kMap: {
+ const auto& map_type = internal::checked_cast<const MapType&>(type);
+ ICEBERG_ASSIGN_OR_RAISE(auto key, ToParquetNode(map_type.key()));
+ ICEBERG_ASSIGN_OR_RAISE(auto value, ToParquetNode(map_type.value()));
+ auto repeated =
+ ::parquet::schema::GroupNode::Make("key_value",
::parquet::Repetition::REPEATED,
+ {std::move(key),
std::move(value)});
+ return ::parquet::schema::GroupNode::Make(
+ std::string(field.name()), repetition, {std::move(repeated)},
+ ::parquet::LogicalType::Map(), field.field_id());
+ }
+ case TypeId::kVariant:
+ return NotSupported("Cannot write Iceberg variant type to Parquet");
+ default:
+ return InvalidSchema("Expected nested Iceberg type, got {}", type);
Review Comment:
nit: `Variant` is not a nested type in spec, maybe we can use unreachable
here.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]