diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Constant.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/Constant.h index 2e72d799606b..a614e5173d24 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/Constant.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/Constant.h @@ -49,6 +49,7 @@ DEFINE_ICEBERG_FIELD(logicalType); /// this field has a camelCase name DEFINE_ICEBERG_FIELD(transform); DEFINE_ICEBERG_FIELD(direction); +DEFINE_ICEBERG_FIELD(unknown); DEFINE_ICEBERG_FIELD(uuid); DEFINE_ICEBERG_FIELD(value); DEFINE_ICEBERG_FIELD(manifest_length); diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/DataFileStatistics.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/DataFileStatistics.cpp index 4e6879496faf..268ff3e98a91 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/DataFileStatistics.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/DataFileStatistics.cpp @@ -58,8 +58,19 @@ void DataFileStatistics::update(const Chunk & chunk) } } +void DataFileStatistics::excludeColumns(std::vector excluded_) +{ + chassert(excluded_.size() == field_ids.size()); + excluded = std::move(excluded_); +} + void DataFileStatistics::merge(const DataFileStatistics & other) { + if (excluded.size() < other.excluded.size()) + excluded.resize(other.excluded.size(), false); + for (size_t i = 0; i < other.excluded.size(); ++i) + excluded[i] = excluded[i] || other.excluded[i]; + if (other.column_sizes.empty()) return; @@ -94,7 +105,8 @@ std::vector> DataFileStatistics::getColumnSizes() cons std::vector> result; for (size_t i = 0; i < column_sizes.size(); ++i) { - result.push_back({field_ids[i], column_sizes[i]}); + if (!isExcluded(i)) + result.push_back({field_ids[i], column_sizes[i]}); } return result; } @@ -104,7 +116,8 @@ std::vector> DataFileStatistics::getNullCounts() const std::vector> result; for (size_t i = 0; i < null_counts.size(); ++i) { - result.push_back({field_ids[i], null_counts[i]}); + if (!isExcluded(i)) + result.push_back({field_ids[i], null_counts[i]}); } return result; } @@ -115,7 +128,8 @@ std::vector> DataFileStatistics::getLowerBounds() const std::vector> result; for (size_t i = 0; i < ranges.size(); ++i) { - result.push_back({field_ids[i], ranges[i].left}); + if (!isExcluded(i)) + result.push_back({field_ids[i], ranges[i].left}); } return result; } @@ -125,7 +139,8 @@ std::vector> DataFileStatistics::getUpperBounds() const std::vector> result; for (size_t i = 0; i < ranges.size(); ++i) { - result.push_back({field_ids[i], ranges[i].right}); + if (!isExcluded(i)) + result.push_back({field_ids[i], ranges[i].right}); } return result; } diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/DataFileStatistics.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/DataFileStatistics.h index 12278b7fcd56..6c90b826979b 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/DataFileStatistics.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/DataFileStatistics.h @@ -27,6 +27,10 @@ class DataFileStatistics void update(const Chunk & chunk); void merge(const DataFileStatistics & other); + /// Marks schema positions whose column is not stored in the data file (e.g. a column of only the + /// Iceberg `unknown` type): the getters omit them, so no statistics describe a column that was not written. + void excludeColumns(std::vector excluded_); + std::vector> getColumnSizes() const; std::vector> getNullCounts() const; std::vector> getLowerBounds() const; @@ -35,8 +39,10 @@ class DataFileStatistics const std::vector & getFieldIds() const { return field_ids; } private: static Range uniteRanges(const Range & left, const Range & right); + bool isExcluded(size_t i) const { return i < excluded.size() && excluded[i]; } std::vector field_ids; + std::vector excluded; std::vector column_sizes; std::vector null_counts; std::vector ranges; diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp index 7c232a26b8fe..ed273fd010e7 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/MetadataGenerator.cpp @@ -726,6 +726,7 @@ void MetadataGenerator::generateAddColumnMetadata(const String & column_name, Da { if (!type->isNullable()) throw Exception(ErrorCodes::BAD_ARGUMENTS, "Iceberg spec doesn't allow to add non-nullable columns"); + Iceberg::checkUnknownTypeAllowed(column_name, type, metadata_object->getValue(Iceberg::f_format_version)); const auto next_schema_id = getNextSchemaId(metadata_object); auto current_schema = deepCopy(getCurrentSchema()); @@ -756,6 +757,7 @@ void MetadataGenerator::generateAddColumnMetadata(const String & column_name, Da bool MetadataGenerator::generateModifyColumnMetadata(const String & column_name, DataTypePtr type, ContextPtr context) { + Iceberg::checkUnknownTypeAllowed(column_name, type, metadata_object->getValue(Iceberg::f_format_version)); auto current_schema = getCurrentSchema(); auto last_column_id = metadata_object->getValue(Iceberg::f_last_column_id); diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp index 3414bf04e4c5..71674f8d6eef 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp @@ -1,5 +1,15 @@ #include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include #include #include #include @@ -10,6 +20,147 @@ namespace DB { +namespace ErrorCodes +{ + extern const int LOGICAL_ERROR; + extern const int NOT_IMPLEMENTED; +} + +namespace Iceberg +{ + +DataTypePtr stripNothing(const DataTypePtr & type) +{ + if (isNothing(type)) + return nullptr; + + if (const auto * nullable_type = typeid_cast(type.get())) + { + const auto & nested = nullable_type->getNestedType(); + auto stripped_nested = stripNothing(nested); + if (!stripped_nested) + return nullptr; + return stripped_nested == nested ? type : makeNullable(stripped_nested); + } + + if (const auto * tuple_type = typeid_cast(type.get())) + { + const auto & elements = tuple_type->getElements(); + DataTypes kept_elements; + Strings kept_names; + bool changed = false; + for (size_t i = 0; i < elements.size(); ++i) + { + auto stripped_element = stripNothing(elements[i]); + if (!stripped_element) + { + changed = true; + continue; + } + changed |= stripped_element != elements[i]; + kept_elements.push_back(stripped_element); + kept_names.push_back(tuple_type->getNameByPosition(i + 1)); + } + if (!changed) + return type; + if (kept_elements.empty()) + return nullptr; + if (tuple_type->hasExplicitNames()) + return std::make_shared(kept_elements, kept_names); + return std::make_shared(kept_elements); + } + + if (const auto * array_type = typeid_cast(type.get())) + { + const auto & nested = array_type->getNestedType(); + auto stripped_nested = stripNothing(nested); + if (!stripped_nested) + return nullptr; + return stripped_nested == nested ? type : std::make_shared(stripped_nested); + } + + if (const auto * map_type = typeid_cast(type.get())) + { + const auto & key = map_type->getKeyType(); + const auto & value = map_type->getValueType(); + auto stripped_key = stripNothing(key); + auto stripped_value = stripNothing(value); + if (!stripped_key || !stripped_value) + return nullptr; + if (stripped_key == key && stripped_value == value) + return type; + return std::make_shared(stripped_key, stripped_value); + } + + return type; +} + +namespace +{ + +ColumnPtr stripNothingColumnImpl(const ColumnPtr & column, const DataTypePtr & type) +{ + auto stripped_type = stripNothing(type); + if (!stripped_type) + throw Exception(ErrorCodes::LOGICAL_ERROR, "Column of type {} has nothing left after stripping Nothing", type->getName()); + if (stripped_type == type) + return column; + + auto full_column = removeSpecialRepresentations(column->convertToFullColumnIfConst()); + + if (const auto * nullable_type = typeid_cast(type.get())) + { + const auto & nullable_column = assert_cast(*full_column); + return ColumnNullable::create( + stripNothingColumnImpl(nullable_column.getNestedColumnPtr(), nullable_type->getNestedType()), + nullable_column.getNullMapColumnPtr()); + } + + if (const auto * tuple_type = typeid_cast(type.get())) + { + const auto & tuple_column = assert_cast(*full_column); + const auto & elements = tuple_type->getElements(); + Columns kept_columns; + for (size_t i = 0; i < elements.size(); ++i) + { + if (stripNothing(elements[i])) + kept_columns.push_back(stripNothingColumnImpl(tuple_column.getColumnPtr(i), elements[i])); + } + return ColumnTuple::create(kept_columns); + } + + if (const auto * array_type = typeid_cast(type.get())) + { + const auto & array_column = assert_cast(*full_column); + return ColumnArray::create( + stripNothingColumnImpl(array_column.getDataPtr(), array_type->getNestedType()), + array_column.getOffsetsPtr()); + } + + if (const auto * map_type = typeid_cast(type.get())) + { + const auto & map_column = assert_cast(*full_column); + const auto & key_value = map_column.getNestedData(); + return ColumnMap::create( + stripNothingColumnImpl(key_value.getColumnPtr(0), map_type->getKeyType()), + stripNothingColumnImpl(key_value.getColumnPtr(1), map_type->getValueType()), + map_column.getNestedColumn().getOffsetsPtr()); + } + + throw Exception(ErrorCodes::LOGICAL_ERROR, "Unexpected type {} while stripping Nothing from a column", type->getName()); +} + +} + +ColumnPtr stripNothingColumn(const ColumnPtr & column, const DataTypePtr & original_type, const DataTypePtr & stripped_type) +{ + auto result = stripNothingColumnImpl(column, original_type); + chassert(stripped_type && stripped_type->equals(*stripNothing(original_type))); + return result; +} + +} + #if USE_AVRO MultipleFileWriter::MultipleFileWriter( @@ -39,6 +190,32 @@ MultipleFileWriter::MultipleFileWriter( , new_file_path_callback(std::move(new_file_path_callback_)) { column_mapper->setStorageColumnEncoding(Iceberg::IcebergSchemaProcessor::traverseSchema(schema_)); + + written_column_types.reserve(sample_block->columns()); + for (const auto & column : *sample_block) + { + written_column_types.push_back(Iceberg::stripNothing(column.type)); + has_nothing_leaves |= written_column_types.back() != column.type; + } + + if (!has_nothing_leaves) + { + filtered_sample_block = sample_block; + } + else + { + aggregate_stats.excludeColumns(getStatisticsExcludedColumns()); + + Block filtered; + for (size_t i = 0; i < sample_block->columns(); ++i) + { + if (!written_column_types[i]) + continue; + const auto & column = sample_block->getByPosition(i); + filtered.insert({written_column_types[i]->createColumn(), written_column_types[i], column.name}); + } + filtered_sample_block = std::make_shared(std::move(filtered)); + } } void MultipleFileWriter::startNewFile() @@ -49,6 +226,8 @@ void MultipleFileWriter::startNewFile() } current_file_stats = std::make_shared(schema); + if (has_nothing_leaves) + current_file_stats->excludeColumns(getStatisticsExcludedColumns()); current_file_num_rows = 0; current_file_num_bytes = 0; auto metadata_path = filename_generator.generateDataFileName(); @@ -69,16 +248,74 @@ void MultipleFileWriter::startNewFile() } FormatFilterInfoPtr format_filter_info = std::make_shared(nullptr, context, column_mapper, nullptr, nullptr); output_format = FormatFactory::instance().getOutputFormatParallelIfPossible( - write_format, *buffer, *sample_block, context, format_settings, format_filter_info); + write_format, *buffer, *filtered_sample_block, context, format_settings, format_filter_info); +} + +std::vector MultipleFileWriter::getStatisticsExcludedColumns() const +{ + std::vector excluded(written_column_types.size()); + for (size_t i = 0; i < written_column_types.size(); ++i) + excluded[i] = !written_column_types[i]; + return excluded; +} + +Columns MultipleFileWriter::filterColumns(const Columns & columns) const +{ + Columns filtered_columns; + filtered_columns.reserve(filtered_sample_block->columns()); + for (size_t i = 0; i < columns.size(); ++i) + { + const auto & original_type = sample_block->getByPosition(i).type; + const auto & written_type = written_column_types[i]; + if (written_type == original_type) + { + filtered_columns.push_back(columns[i]); + } + else if (written_type) + { + filtered_columns.push_back(Iceberg::stripNothingColumn(columns[i], original_type, written_type)); + } + else + { + for (size_t row = 0; row < columns[i]->size(); ++row) + { + if (!columns[i]->isDefaultAt(row)) + throw Exception( + ErrorCodes::NOT_IMPLEMENTED, + "Cannot write column '{}' of type {} into an Iceberg data file: row {} holds a non-default value, " + "but the column contains an Iceberg `unknown` element that no data file format can store, so the " + "value (for example, list elements or map keys) would be lost. Only NULL, empty lists and empty " + "maps can be inserted into such a column", + sample_block->getByPosition(i).name, original_type->getName(), row); + } + } + } + return filtered_columns; } void MultipleFileWriter::consume(const Chunk & chunk) { + /// Validate before starting a file, so a rejected chunk leaves no empty data file behind. + if (filtered_sample_block->columns() == 0 && sample_block->columns() > 0) + throw Exception( + ErrorCodes::NOT_IMPLEMENTED, + "Cannot write an Iceberg data file: every column contains only the Iceberg `unknown` type, " + "which no data file format can store"); + + std::optional filtered_columns; + if (has_nothing_leaves) + filtered_columns = filterColumns(chunk.getColumns()); + if (!current_file_num_rows || *current_file_num_rows >= max_data_file_num_rows || *current_file_num_bytes >= max_data_file_num_bytes) { startNewFile(); } - output_format->write(sample_block->cloneWithColumns(chunk.getColumns())); + + if (filtered_columns) + output_format->write(filtered_sample_block->cloneWithColumns(std::move(*filtered_columns))); + else + output_format->write(sample_block->cloneWithColumns(chunk.getColumns())); + output_format->flush(); *current_file_num_rows += chunk.getNumRows(); *current_file_num_bytes += chunk.bytes(); diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.h index 973b7e4932f3..00712850182a 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.h @@ -10,6 +10,23 @@ namespace DB { +namespace Iceberg +{ + +/// Iceberg `unknown` maps to `Nullable(Nothing)` and can be nested inside structs, lists and maps. +/// No serialisation format (Parquet, ORC, Avro) can represent `Nothing`, so only those leaves are +/// removed before writing; the reader fills a missing struct field with NULL. +/// Returns `type` itself if it has no `Nothing` leaf, the type without its `Nothing` leaves otherwise, +/// and `nullptr` if nothing serialisable remains (`unknown` itself, a struct of only `unknown` fields, +/// or a list or map whose element, key or value strips to nothing). +DataTypePtr stripNothing(const DataTypePtr & type); + +/// Removes from `column` of type `original_type` the parts that `stripNothing` removes from the type. +/// `stripped_type` must be the non-null result of `stripNothing(original_type)`. +ColumnPtr stripNothingColumn(const ColumnPtr & column, const DataTypePtr & original_type, const DataTypePtr & stripped_type); + +} + #if USE_AVRO class MultipleFileWriter @@ -68,6 +85,11 @@ class MultipleFileWriter std::vector getDataFileEntries() const; private: + /// Strips `Nothing` leaves from `columns` and drops fully `Nothing` columns, throwing if a dropped one holds a non-default row. + Columns filterColumns(const Columns & columns) const; + /// Per `sample_block` column: whether it is left out of the data file, so statistics must not describe it. + std::vector getStatisticsExcludedColumns() const; + UInt64 max_data_file_num_rows; UInt64 max_data_file_num_bytes; Poco::JSON::Array::Ptr schema; @@ -92,6 +114,14 @@ class MultipleFileWriter std::optional format_settings; const String& write_format; SharedHeader sample_block; + /// Per `sample_block` column: the type passed to the format writer, i.e. the result of + /// `Iceberg::stripNothing`. It is the original type when the column has no Iceberg `unknown` + /// leaf, and nullptr when nothing serialisable remains: such a column is left out of the data + /// file and may only hold default values, because anything else would be lost. + DataTypes written_column_types; + bool has_nothing_leaves = false; + /// `sample_block` with `Nothing` leaves stripped and fully `Nothing` columns removed; used by the format writer. + SharedHeader filtered_sample_block; UInt64 total_bytes = 0; std::function new_file_path_callback; }; diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/SchemaProcessor.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/SchemaProcessor.cpp index f6c989b13ad3..ece1e4dc2176 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/SchemaProcessor.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/SchemaProcessor.cpp @@ -29,6 +29,7 @@ #include #include #include +#include #include #include #include @@ -454,6 +455,8 @@ DataTypePtr IcebergSchemaProcessor::getSimpleType(const String & type_name_arg, } if (type_name == f_uuid) return std::make_shared(); + if (type_name == f_unknown) + return std::make_shared(); if (type_name.starts_with("fixed[") && type_name.ends_with(']')) { diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp index d378e21f03a5..d99129d1f1e7 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp @@ -534,6 +534,17 @@ static size_t icebergDecimalRequiredBytes(UInt32 precision) return bytes; } +void checkUnknownTypeAllowed(const String & column_name, const DataTypePtr & type, Int64 format_version) +{ + if (format_version >= 3 || Iceberg::stripNothing(type) == type) + return; + throw Exception( + ErrorCodes::BAD_ARGUMENTS, + "Column '{}' of type {} maps to the Iceberg `unknown` type, which requires format version 3, " + "but the table uses format version {} (set `iceberg_format_version = 3`)", + column_name, type->getName(), format_version); +} + /// Returns type and required std::pair getIcebergType(DataTypePtr type, Int32 & iter) { @@ -574,6 +585,8 @@ std::pair getIcebergType(DataTypePtr type, Int32 & ite return {"string", true}; case TypeIndex::UUID: return {"uuid", true}; + case TypeIndex::Nothing: + return {"unknown", false}; case TypeIndex::Decimal32: case TypeIndex::Decimal64: case TypeIndex::Decimal128: @@ -1058,6 +1071,9 @@ std::pair createEmptyMetadataFile( ContextPtr context, UInt64 format_version) { + for (const auto & column : columns) + checkUnknownTypeAllowed(column.name, column.type, static_cast(format_version)); + std::unordered_map column_name_to_source_id; static Poco::UUIDGenerator uuid_generator; diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.h index 4d403a785dc1..12dd43472b75 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.h @@ -93,6 +93,11 @@ Poco::JSON::Object::Ptr getMetadataJSONObject( std::pair getIcebergType(DataTypePtr type, Int32 & iter); + +/// Throws if `type` contains `Nothing` at any depth (the Iceberg `unknown` type) and +/// `format_version` is below 3, the first format version that defines `unknown`. +void checkUnknownTypeAllowed(const String & column_name, const DataTypePtr & type, Int64 format_version); + Poco::Dynamic::Var getAvroType(DataTypePtr type, Int32 field_id); Poco::Dynamic::Var getAvroLogicalType(DataTypePtr type); diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_metadata_generator.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_metadata_generator.cpp index 1f2c5bec0cb3..c98827c495c0 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_metadata_generator.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_metadata_generator.cpp @@ -6,12 +6,14 @@ #include #include +#include #include #include #include #include #include #include +#include #include #include @@ -491,4 +493,37 @@ TEST(IcebergMetadataGenerator, ModifyColumnWideningRecordsTheNewTypeInANewSchema EXPECT_EQ(stored_type.extract(), "long"); } + +TEST(IcebergMetadataGenerator, GetIcebergTypeNothingProducesUnknown) +{ + Int32 iter = 0; + auto [iceberg_type, required] = Iceberg::getIcebergType(std::make_shared(), iter); + ASSERT_TRUE(iceberg_type.isString()); + EXPECT_EQ(iceberg_type.extract(), "unknown"); + EXPECT_FALSE(required); +} + + +TEST(IcebergMetadataGenerator, GetIcebergTypeNullableNothingProducesUnknown) +{ + Int32 iter = 0; + auto [iceberg_type, required] = Iceberg::getIcebergType(makeNullable(std::make_shared()), iter); + ASSERT_TRUE(iceberg_type.isString()); + EXPECT_EQ(iceberg_type.extract(), "unknown"); + EXPECT_FALSE(required); +} + + +TEST(IcebergMetadataGenerator, AddColumnUnknownTypeRecordsUnknownInSchema) +{ + auto metadata = makeMetadataWithGap(); + MetadataGenerator gen(metadata); + + gen.generateAddColumnMetadata("placeholder", makeNullable(std::make_shared())); + + auto stored_type = findCurrentFieldType(metadata, "placeholder"); + ASSERT_TRUE(stored_type.isString()); + EXPECT_EQ(stored_type.extract(), "unknown"); +} + #endif diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_schema_processor.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_schema_processor.cpp index e13a421eda1f..212898ac0f4e 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_schema_processor.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_schema_processor.cpp @@ -121,6 +121,19 @@ TEST(IcebergSchemaProcessor, GetSimpleTypeDecimal) EXPECT_EQ(type->getName(), "Decimal(10, 2)"); } +TEST(IcebergSchemaProcessor, GetSimpleTypeUnknown) +{ + auto type = IcebergSchemaProcessor::getSimpleType("unknown", getContext().context); + EXPECT_EQ(type->getName(), "Nothing"); +} + +TEST(IcebergSchemaProcessor, UnknownFieldInSchemaProducesNullableNothing) +{ + auto schema = parseSchema(R"json({"schema-id":0,"fields":[{"id":1,"name":"placeholder","required":false,"type":"unknown"}]})json"); + IcebergSchemaProcessor processor(getContext().context); + EXPECT_NO_THROW(processor.addIcebergTableSchema(schema, getContext().context)); +} + TEST(IcebergSchemaProcessor, GetSimpleTypeUnknownThrows) { EXPECT_THROW(IcebergSchemaProcessor::getSimpleType("unknown_type", getContext().context), DB::Exception); diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_strip_nothing.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_strip_nothing.cpp new file mode 100644 index 000000000000..42a1617b8d4f --- /dev/null +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_strip_nothing.cpp @@ -0,0 +1,200 @@ +#include + +#include +#include +#include +#include +#include +#include +#include + +using namespace DB; + +namespace +{ + +DataTypePtr type(const String & name) +{ + return DataTypeFactory::instance().get(name); +} + +String strippedName(const String & name) +{ + auto stripped = Iceberg::stripNothing(type(name)); + return stripped ? stripped->getName() : "nullptr"; +} + +} + +TEST(IcebergStripNothing, TypeWithoutNothingIsReturnedAsIs) +{ + for (const auto * name : {"Int64", "Nullable(String)", "Tuple(a Int64, b Nullable(String))", "Array(Int32)", "Map(String, Int64)"}) + { + auto original = type(name); + EXPECT_EQ(Iceberg::stripNothing(original), original) << name; + } +} + +TEST(IcebergStripNothing, NothingLeavesAreRemoved) +{ + EXPECT_EQ(strippedName("Nothing"), "nullptr"); + EXPECT_EQ(strippedName("Nullable(Nothing)"), "nullptr"); + EXPECT_EQ(strippedName("Tuple(a Nullable(Int64), u Nullable(Nothing))"), "Tuple(a Nullable(Int64))"); + EXPECT_EQ(strippedName("Tuple(u Nullable(Nothing), v Nullable(Nothing))"), "nullptr"); + EXPECT_EQ( + strippedName("Tuple(a Int64, s Tuple(b Nullable(String), u Nullable(Nothing)), t Tuple(u Nullable(Nothing)))"), + "Tuple(a Int64, s Tuple(b Nullable(String)))"); + EXPECT_EQ(strippedName("Array(Tuple(a Nullable(Int64), u Nullable(Nothing)))"), "Array(Tuple(a Nullable(Int64)))"); + EXPECT_EQ(strippedName("Map(String, Tuple(a Nullable(Int64), u Nullable(Nothing)))"), "Map(String, Tuple(a Nullable(Int64)))"); +} + +TEST(IcebergStripNothing, ContainerOfOnlyNothingStripsToNothing) +{ + EXPECT_EQ(strippedName("Array(Nullable(Nothing))"), "nullptr"); + EXPECT_EQ(strippedName("Array(Tuple(u Nullable(Nothing)))"), "nullptr"); + EXPECT_EQ(strippedName("Map(String, Nullable(Nothing))"), "nullptr"); +} + +TEST(IcebergStripNothing, TupleColumnKeepsSiblingValues) +{ + auto original = type("Tuple(a Nullable(Int64), u Nullable(Nothing))"); + auto stripped = Iceberg::stripNothing(original); + + auto column = original->createColumn(); + column->insert(Tuple{Field(Int64(42)), Field()}); + column->insert(Tuple{Field(), Field()}); + + auto result = Iceberg::stripNothingColumn(std::move(column), original, stripped); + EXPECT_EQ(result->getName(), stripped->createColumn()->getName()); + ASSERT_EQ(result->size(), 2); + EXPECT_EQ((*result)[0], Field(Tuple{Field(Int64(42))})); + EXPECT_EQ((*result)[1], Field(Tuple{Field()})); +} + +TEST(IcebergStripNothing, ConstTupleColumnIsMaterialized) +{ + auto original = type("Tuple(a Nullable(Int64), u Nullable(Nothing))"); + auto stripped = Iceberg::stripNothing(original); + + auto single = original->createColumn(); + single->insert(Tuple{Field(Int64(7)), Field()}); + ColumnPtr column = ColumnConst::create(std::move(single), 3); + + auto result = Iceberg::stripNothingColumn(column, original, stripped); + EXPECT_EQ(result->getName(), stripped->createColumn()->getName()); + ASSERT_EQ(result->size(), 3); + for (size_t row = 0; row < 3; ++row) + EXPECT_EQ((*result)[row], Field(Tuple{Field(Int64(7))})); +} + +TEST(IcebergStripNothing, ArrayAndMapColumnsKeepOffsetsAndKeys) +{ + { + auto original = type("Array(Tuple(a Nullable(Int64), u Nullable(Nothing)))"); + auto stripped = Iceberg::stripNothing(original); + + auto column = original->createColumn(); + column->insert(Array{Field(Tuple{Field(Int64(1)), Field()}), Field(Tuple{Field(Int64(2)), Field()})}); + column->insert(Array{}); + + auto result = Iceberg::stripNothingColumn(std::move(column), original, stripped); + EXPECT_EQ(result->getName(), stripped->createColumn()->getName()); + ASSERT_EQ(result->size(), 2); + EXPECT_EQ((*result)[0], Field(Array{Field(Tuple{Field(Int64(1))}), Field(Tuple{Field(Int64(2))})})); + EXPECT_EQ((*result)[1], Field(Array{})); + } + { + auto original = type("Map(String, Tuple(a Nullable(Int64), u Nullable(Nothing)))"); + auto stripped = Iceberg::stripNothing(original); + + auto column = original->createColumn(); + column->insert(Map{Field(Tuple{Field("k"), Field(Tuple{Field(Int64(5)), Field()})})}); + + auto result = Iceberg::stripNothingColumn(std::move(column), original, stripped); + EXPECT_EQ(result->getName(), stripped->createColumn()->getName()); + ASSERT_EQ(result->size(), 1); + EXPECT_EQ((*result)[0], Field(Map{Field(Tuple{Field("k"), Field(Tuple{Field(Int64(5))})})})); + } +} + +#if USE_AVRO + +namespace +{ + +/// Schema fields with ids 1, 2, 3 for (id Int64, name Nullable(String), u Nullable(Nothing)). +Poco::JSON::Array::Ptr statisticsSchema() +{ + Poco::JSON::Array::Ptr schema = new Poco::JSON::Array; + for (Int32 id = 1; id <= 3; ++id) + { + Poco::JSON::Object::Ptr field = new Poco::JSON::Object; + field->set(Iceberg::f_id, id); + schema->add(field); + } + return schema; +} + +Chunk statisticsChunk() +{ + auto id = type("Int64")->createColumn(); + id->insert(Field(Int64(10))); + id->insert(Field(Int64(11))); + + auto name = type("Nullable(String)")->createColumn(); + name->insert(Field("x")); + name->insert(Field("y")); + + auto unknown = type("Nullable(Nothing)")->createColumn(); + unknown->insertDefault(); + unknown->insertDefault(); + + Columns columns; + columns.push_back(std::move(id)); + columns.push_back(std::move(name)); + columns.push_back(std::move(unknown)); + return Chunk(std::move(columns), 2); +} + +template +std::vector fieldIdsOf(const std::vector> & stats) +{ + std::vector ids; + for (const auto & [id, _] : stats) + ids.push_back(id); + return ids; +} + +} + +TEST(IcebergStripNothing, ExcludedColumnsHaveNoStatistics) +{ + DataFileStatistics stats(statisticsSchema()); + stats.excludeColumns({false, false, true}); + stats.update(statisticsChunk()); + + const std::vector written = {1, 2}; + EXPECT_EQ(fieldIdsOf(stats.getColumnSizes()), written); + EXPECT_EQ(fieldIdsOf(stats.getNullCounts()), written); + EXPECT_EQ(fieldIdsOf(stats.getLowerBounds()), written); + EXPECT_EQ(fieldIdsOf(stats.getUpperBounds()), written); + EXPECT_EQ(stats.getLowerBounds()[0].second, Field(Int64(10))); + EXPECT_EQ(stats.getUpperBounds()[0].second, Field(Int64(11))); + + DataFileStatistics merged(statisticsSchema()); + merged.merge(stats); + EXPECT_EQ(fieldIdsOf(merged.getColumnSizes()), written); + EXPECT_EQ(fieldIdsOf(merged.getLowerBounds()), written); +} + +TEST(IcebergStripNothing, StatisticsWithoutExclusionCoverEveryColumn) +{ + DataFileStatistics stats(statisticsSchema()); + stats.update(statisticsChunk()); + + const std::vector all = {1, 2, 3}; + EXPECT_EQ(fieldIdsOf(stats.getColumnSizes()), all); + EXPECT_EQ(fieldIdsOf(stats.getLowerBounds()), all); +} + +#endif diff --git a/tests/integration/test_storage_iceberg_no_spark/test_unknown_type.py b/tests/integration/test_storage_iceberg_no_spark/test_unknown_type.py new file mode 100644 index 000000000000..0e076f526c2b --- /dev/null +++ b/tests/integration/test_storage_iceberg_no_spark/test_unknown_type.py @@ -0,0 +1,561 @@ +#!/usr/bin/env python3 + +""" +Integration test for Iceberg v3 `unknown` primitive type. + +Creates an Iceberg v2 table via ClickHouse, then patches the metadata JSON +to v3 and injects an `unknown`-typed field. Verifies that ClickHouse reads +the table correctly: +- The unknown column is mapped to Nullable(Nothing) +- All values in the unknown column are NULL +- Non-unknown columns read correctly +""" + +import json +import re + +from helpers.iceberg_utils import ( + create_iceberg_table, + get_creation_expression, + get_uuid_str, +) + + +def _metadata_dir(table_name): + return ( + f"/var/lib/clickhouse/user_files/iceberg_data/default/{table_name}/metadata" + ) + + +def _read_latest_metadata(instance, table_name): + metadata_dir = _metadata_dir(table_name) + latest = instance.exec_in_container( + ["bash", "-c", f"ls -v {metadata_dir}/v*.metadata.json | tail -1"] + ).strip() + raw = instance.exec_in_container(["cat", latest]) + return json.loads(raw), latest + + +def _write_next_metadata(instance, table_name, meta, prev_path): + metadata_dir = _metadata_dir(table_name) + version_match = re.search(r"/v(\d+)[^/]*\.metadata\.json$", prev_path) + new_version = int(version_match.group(1)) + 1 + new_path = f"{metadata_dir}/v{new_version}.metadata.json" + new_content = json.dumps(meta, indent=4) + instance.exec_in_container( + ["bash", "-c", f"cat > {new_path} << 'JSONEOF'\n{new_content}\nJSONEOF"] + ) + + +def test_unknown_type_read(started_cluster_iceberg_no_spark): + """Create a v2 Iceberg table, patch metadata to v3 with an unknown column, + and verify ClickHouse reads it correctly.""" + instance = started_cluster_iceberg_no_spark.instances["node1"] + table_name = "test_unknown_type_" + get_uuid_str() + + # 1. Create a normal v2 table with two columns. + create_iceberg_table( + "local", + instance, + table_name, + started_cluster_iceberg_no_spark, + "(id Int64, name Nullable(String))", + format_version=2, + ) + instance.query( + f"INSERT INTO {table_name} VALUES (1, 'alice'), (2, 'bob'), (3, 'charlie')" + ) + + # 2. Read the latest metadata JSON. + meta, prev_path = _read_latest_metadata(instance, table_name) + + # 3. Patch: bump format-version to 3 and add an unknown-typed field + # via proper schema evolution (new schema-id, not in-place mutation). + meta["format-version"] = 3 + + # Find the current schema to copy its fields. + current_schema_id = meta.get("current-schema-id", 0) + current_schema = None + for schema in meta.get("schemas", []): + if schema.get("schema-id", 0) == current_schema_id: + current_schema = schema + break + assert current_schema is not None, "current schema not found in metadata" + + # Determine the next field id and next schema id. + last_column_id = meta.get("last-column-id", 0) + new_field_id = last_column_id + 1 + new_schema_id = max(s.get("schema-id", 0) for s in meta.get("schemas", [])) + 1 + + # Create a new schema that includes the unknown field. + new_schema = { + "type": "struct", + "schema-id": new_schema_id, + "fields": current_schema["fields"] + + [ + { + "id": new_field_id, + "name": "placeholder", + "required": False, + "type": "unknown", + } + ], + } + meta.setdefault("schemas", []).append(new_schema) + meta["current-schema-id"] = new_schema_id + meta["last-column-id"] = new_field_id + + # 4. Write the patched metadata as the next version. + _write_next_metadata(instance, table_name, meta, prev_path) + + # 5. Drop and recreate the ClickHouse table so it re-reads metadata. + instance.query(f"DROP TABLE IF EXISTS {table_name}") + create_iceberg_table( + "local", + instance, + table_name, + started_cluster_iceberg_no_spark, + ) + + # 6. Verify row count. + assert instance.query(f"SELECT count() FROM {table_name}").strip() == "3" + + # 7. Verify the unknown column type is Nullable(Nothing). + describe = instance.query(f"DESCRIBE TABLE {table_name}") + lines = [line.split("\t") for line in describe.strip().split("\n")] + placeholder_row = [row for row in lines if row[0] == "placeholder"] + assert len(placeholder_row) == 1, ( + f"Expected one 'placeholder' column, got: {lines}" + ) + assert placeholder_row[0][1] == "Nullable(Nothing)" + + # 8. Verify all values in the unknown column are NULL. + result = instance.query(f"SELECT placeholder FROM {table_name}").strip() + assert result == "\\N\n\\N\n\\N" + + # 9. Verify non-unknown columns read correctly alongside the unknown column. + result = instance.query( + f"SELECT id, name, placeholder FROM {table_name} ORDER BY id" + ).strip() + expected = "1\talice\t\\N\n2\tbob\t\\N\n3\tcharlie\t\\N" + assert result == expected + + +def _patch_table_with_unknown_column(instance, table_name): + """Read the latest metadata, add an unknown-typed column via schema + evolution (v2 -> v3), and write the new metadata version. Returns the + name of the new column.""" + meta, prev_path = _read_latest_metadata(instance, table_name) + meta["format-version"] = 3 + + current_schema_id = meta.get("current-schema-id", 0) + current_schema = None + for schema in meta.get("schemas", []): + if schema.get("schema-id", 0) == current_schema_id: + current_schema = schema + break + assert current_schema is not None + + last_column_id = meta.get("last-column-id", 0) + new_field_id = last_column_id + 1 + new_schema_id = max(s.get("schema-id", 0) for s in meta.get("schemas", [])) + 1 + + new_schema = { + "type": "struct", + "schema-id": new_schema_id, + "fields": current_schema["fields"] + + [ + { + "id": new_field_id, + "name": "placeholder", + "required": False, + "type": "unknown", + } + ], + } + meta.setdefault("schemas", []).append(new_schema) + meta["current-schema-id"] = new_schema_id + meta["last-column-id"] = new_field_id + + _write_next_metadata(instance, table_name, meta, prev_path) + + +def test_unknown_type_write(started_cluster_iceberg_no_spark): + """Verify that inserting into a table with an unknown-typed column + succeeds: the unknown column is silently dropped from the data file + (it has no physical representation) but still reads back as NULLs.""" + instance = started_cluster_iceberg_no_spark.instances["node1"] + table_name = "test_unknown_type_write_" + get_uuid_str() + + # 1. Create a v2 table, insert initial data, and patch to v3 with unknown column. + create_iceberg_table( + "local", + instance, + table_name, + started_cluster_iceberg_no_spark, + "(id Int64, name Nullable(String))", + format_version=2, + ) + instance.query( + f"INSERT INTO {table_name} VALUES (1, 'alice'), (2, 'bob')" + ) + _patch_table_with_unknown_column(instance, table_name) + + # 2. Re-create the ClickHouse table to pick up the new schema. + instance.query(f"DROP TABLE IF EXISTS {table_name}") + create_iceberg_table( + "local", + instance, + table_name, + started_cluster_iceberg_no_spark, + settings={"allow_insert_into_iceberg": 1}, + ) + + # 3. Insert new rows -- this must NOT fail with UNKNOWN_TYPE. + instance.query( + f"INSERT INTO {table_name} VALUES (3, 'charlie', NULL), (4, 'dave', NULL)", + settings={"allow_insert_into_iceberg": 1}, + ) + + # 4. Verify all rows are readable. + result = instance.query( + f"SELECT id, name, placeholder FROM {table_name} ORDER BY id" + ).strip() + expected = ( + "1\talice\t\\N\n" + "2\tbob\t\\N\n" + "3\tcharlie\t\\N\n" + "4\tdave\t\\N" + ) + assert result == expected + + # 5. Verify the unknown column is still Nullable(Nothing) after the write. + describe = instance.query(f"DESCRIBE TABLE {table_name}") + lines = [line.split("\t") for line in describe.strip().split("\n")] + placeholder_row = [row for row in lines if row[0] == "placeholder"] + assert len(placeholder_row) == 1 + assert placeholder_row[0][1] == "Nullable(Nothing)" + + +def _patch_table_with_v3_field(instance, table_name, build_field): + """Bump the table to format-version 3 and append the top-level field returned by + `build_field(last_column_id)`, which returns `(field, new_last_column_id)`.""" + meta, prev_path = _read_latest_metadata(instance, table_name) + meta["format-version"] = 3 + + current_schema_id = meta.get("current-schema-id", 0) + current_schema = None + for schema in meta.get("schemas", []): + if schema.get("schema-id", 0) == current_schema_id: + current_schema = schema + break + assert current_schema is not None + + field, new_last_column_id = build_field(meta.get("last-column-id", 0)) + new_schema_id = max(s.get("schema-id", 0) for s in meta.get("schemas", [])) + 1 + + new_schema = { + "type": "struct", + "schema-id": new_schema_id, + "fields": current_schema["fields"] + [field], + } + meta.setdefault("schemas", []).append(new_schema) + meta["current-schema-id"] = new_schema_id + meta["last-column-id"] = new_last_column_id + + _write_next_metadata(instance, table_name, meta, prev_path) + + +def _patch_table_with_nested_unknown_column(instance, table_name): + """Like `_patch_table_with_unknown_column`, but adds a struct column whose + subfield has the unknown type, so Nothing appears nested inside a Tuple.""" + + def build_field(last_column_id): + field = { + "id": last_column_id + 1, + "name": "nested", + "required": False, + "type": { + "type": "struct", + "fields": [ + { + "id": last_column_id + 2, + "name": "a", + "required": False, + "type": "long", + }, + { + "id": last_column_id + 3, + "name": "u", + "required": False, + "type": "unknown", + }, + ], + }, + } + return field, last_column_id + 3 + + _patch_table_with_v3_field(instance, table_name, build_field) + + +def _patch_table_with_list_unknown_column(instance, table_name): + """Adds a `list` column, whose only content is its list lengths.""" + + def build_field(last_column_id): + field = { + "id": last_column_id + 1, + "name": "tags", + "required": False, + "type": { + "type": "list", + "element-id": last_column_id + 2, + "element": "unknown", + "element-required": False, + }, + } + return field, last_column_id + 2 + + _patch_table_with_v3_field(instance, table_name, build_field) + + +def _recreate_table(instance, table_name, cluster): + instance.query(f"DROP TABLE IF EXISTS {table_name}") + create_iceberg_table( + "local", + instance, + table_name, + cluster, + settings={"allow_insert_into_iceberg": 1}, + ) + + +def test_unknown_type_nested_write(started_cluster_iceberg_no_spark): + """Verify that inserting into a table whose struct column contains an + unknown-typed subfield succeeds: Nothing nested inside a Tuple must not + reach the Parquet writer, which cannot represent it (`UNKNOWN_TYPE`).""" + instance = started_cluster_iceberg_no_spark.instances["node1"] + table_name = "test_unknown_type_nested_write_" + get_uuid_str() + + create_iceberg_table( + "local", + instance, + table_name, + started_cluster_iceberg_no_spark, + "(id Int64, name Nullable(String))", + format_version=2, + ) + instance.query( + f"INSERT INTO {table_name} VALUES (1, 'alice'), (2, 'bob')" + ) + _patch_table_with_nested_unknown_column(instance, table_name) + _recreate_table(instance, table_name, started_cluster_iceberg_no_spark) + + describe = instance.query(f"DESCRIBE TABLE {table_name}") + lines = [line.split("\t") for line in describe.strip().split("\n")] + nested_row = [row for row in lines if row[0] == "nested"] + assert len(nested_row) == 1 + assert "Nothing" in nested_row[0][1], nested_row + + # This must NOT fail with UNKNOWN_TYPE, and must keep the known subfield `a`: + # only the unknown leaf `u` is left out of the data file. + instance.query( + f"INSERT INTO {table_name} (id, name, nested) VALUES (3, 'charlie', (42, NULL)), (4, 'dave', (NULL, NULL))", + settings={"allow_insert_into_iceberg": 1}, + ) + + result = instance.query( + f"SELECT id, name, nested.a, nested.u FROM {table_name} ORDER BY id" + ).strip() + expected = ( + "1\talice\t\\N\t\\N\n" + "2\tbob\t\\N\t\\N\n" + "3\tcharlie\t42\t\\N\n" + "4\tdave\t\\N\t\\N" + ) + assert result == expected + + +def test_unknown_type_list_write(started_cluster_iceberg_no_spark): + """A `list` column has no serialisable leaf, so it is left out of the + data file. That only loses nothing when every list is empty: a non-empty list + such as [NULL] must be rejected instead of silently read back as [].""" + instance = started_cluster_iceberg_no_spark.instances["node1"] + table_name = "test_unknown_type_list_write_" + get_uuid_str() + + create_iceberg_table( + "local", + instance, + table_name, + started_cluster_iceberg_no_spark, + "(id Int64, name Nullable(String))", + format_version=2, + ) + instance.query( + f"INSERT INTO {table_name} VALUES (1, 'alice'), (2, 'bob')" + ) + _patch_table_with_list_unknown_column(instance, table_name) + _recreate_table(instance, table_name, started_cluster_iceberg_no_spark) + + instance.query( + f"INSERT INTO {table_name} (id, name, tags) VALUES (3, 'charlie', [])", + settings={"allow_insert_into_iceberg": 1}, + ) + + error = instance.query_and_get_error( + f"INSERT INTO {table_name} (id, name, tags) VALUES (4, 'dave', [NULL])", + settings={"allow_insert_into_iceberg": 1}, + ) + assert "NOT_IMPLEMENTED" in error, error + assert "Iceberg `unknown` element" in error, error + + result = instance.query( + f"SELECT id, name, tags FROM {table_name} ORDER BY id" + ).strip() + expected = ( + "1\talice\t[]\n" + "2\tbob\t[]\n" + "3\tcharlie\t[]" + ) + assert result == expected + + +def test_unknown_type_rejected_on_v2(started_cluster_iceberg_no_spark): + """`unknown` exists only from Iceberg format version 3, so ClickHouse must not + write it into the metadata of a table with an older format version.""" + instance = started_cluster_iceberg_no_spark.instances["node1"] + + table_name = "test_unknown_type_create_v2_" + get_uuid_str() + error = instance.query_and_get_error( + get_creation_expression( + "local", + table_name, + started_cluster_iceberg_no_spark, + "(id Int64, u Nullable(Nothing))", + 2, + ) + ) + assert "BAD_ARGUMENTS" in error, error + assert "requires format version 3" in error, error + + table_name = "test_unknown_type_add_column_v2_" + get_uuid_str() + create_iceberg_table( + "local", + instance, + table_name, + started_cluster_iceberg_no_spark, + "(id Int64, name Nullable(String))", + format_version=2, + ) + error = instance.query_and_get_error( + f"ALTER TABLE {table_name} ADD COLUMN u Nullable(Nothing)", + settings={"allow_insert_into_iceberg": 1}, + ) + assert "BAD_ARGUMENTS" in error, error + assert "requires format version 3" in error, error + describe = instance.query(f"DESCRIBE TABLE {table_name}") + assert [line.split("\t")[0] for line in describe.strip().split("\n")] == ["id", "name"] + + table_name = "test_unknown_type_create_v3_" + get_uuid_str() + create_iceberg_table( + "local", + instance, + table_name, + started_cluster_iceberg_no_spark, + "(id Int64, u Nullable(Nothing))", + format_version=3, + ) + describe = instance.query(f"DESCRIBE TABLE {table_name}") + u_row = [line.split("\t") for line in describe.strip().split("\n") if line.startswith("u\t")] + assert len(u_row) == 1 and u_row[0][1] == "Nullable(Nothing)", describe + + +def test_unknown_type_write_keeps_bounds(started_cluster_iceberg_no_spark): + """A top-level `unknown` column is left out of the data file, so it must not + contribute statistics: before the fix its null bound made the writer drop + lower/upper bounds for every column, and `column_sizes` described a column + that was never written.""" + instance = started_cluster_iceberg_no_spark.instances["node1"] + table_name = "test_unknown_type_write_keeps_bounds_" + get_uuid_str() + + create_iceberg_table( + "local", + instance, + table_name, + started_cluster_iceberg_no_spark, + "(id Int64, name Nullable(String))", + format_version=2, + ) + instance.query(f"INSERT INTO {table_name} VALUES (1, 'alice'), (2, 'bob')") + _patch_table_with_unknown_column(instance, table_name) + _recreate_table(instance, table_name, started_cluster_iceberg_no_spark) + + meta, _ = _read_latest_metadata(instance, table_name) + current_schema = [ + s for s in meta["schemas"] if s["schema-id"] == meta["current-schema-id"] + ][0] + placeholder_id = [ + f["id"] for f in current_schema["fields"] if f["name"] == "placeholder" + ][0] + + # Two inserts, so each new row range lands in its own data file. + for values in ["(3, 'c', NULL), (4, 'd', NULL)", "(10, 'x', NULL), (11, 'y', NULL)"]: + instance.query( + f"INSERT INTO {table_name} VALUES {values}", + settings={"allow_insert_into_iceberg": 1}, + ) + + files = instance.query( + f"SELECT count(), countIf(mapContains(column_sizes, {placeholder_id})) " + f"FROM system.iceberg_files " + f"WHERE database = currentDatabase() AND table = '{table_name}' AND content = 'DATA'" + ).strip() + assert files == "3\t0", files + + # With bounds on `id` in both new files, min/max pruning skips the files of + # rows 1-2 and 3-4. Without them only the file written before the patch is skipped. + query_id = f"{table_name}_{get_uuid_str()}" + result = instance.query( + f"SELECT id, name FROM {table_name} WHERE id = 10", + query_id=query_id, + settings={ + "use_iceberg_partition_pruning": 1, + "input_format_parquet_bloom_filter_push_down": 0, + "input_format_parquet_filter_push_down": 0, + }, + ).strip() + assert result == "10\tx" + + instance.query("SYSTEM FLUSH LOGS") + pruned = int( + instance.query( + f"SELECT ProfileEvents['IcebergMinMaxIndexPrunedFiles'] FROM system.query_log " + f"WHERE query_id = '{query_id}' AND type = 'QueryFinish'" + ) + ) + assert pruned == 2, pruned + + +def test_unknown_type_only_column_rejected(started_cluster_iceberg_no_spark): + """When every column has only the `unknown` type, nothing can be written: the + insert must fail instead of writing a zero-column Parquet file that later + makes every SELECT fail.""" + instance = started_cluster_iceberg_no_spark.instances["node1"] + table_name = "test_unknown_type_only_column_" + get_uuid_str() + + create_iceberg_table( + "local", + instance, + table_name, + started_cluster_iceberg_no_spark, + "(u Nullable(Nothing))", + format_version=3, + ) + + error = instance.query_and_get_error( + f"INSERT INTO {table_name} VALUES (NULL), (NULL)", + settings={"allow_insert_into_iceberg": 1}, + ) + assert "NOT_IMPLEMENTED" in error, error + assert "every column contains only the Iceberg `unknown` type" in error, error + + assert instance.query(f"SELECT count() FROM {table_name}").strip() == "0"