From e3b94bcfd93e5f0b3dd5a58161b77947290da94e Mon Sep 17 00:00:00 2001 From: Kanthi Subramanian Date: Tue, 15 Sep 2026 05:10:47 +0200 Subject: [PATCH 1/8] Added support for iceberg v3 unknown data type --- .../DataLakes/Iceberg/Constant.h | 1 + .../DataLakes/Iceberg/SchemaProcessor.cpp | 3 + .../ObjectStorage/DataLakes/Iceberg/Utils.cpp | 2 + .../gtest_iceberg_metadata_generator.cpp | 35 ++++++ .../tests/gtest_iceberg_schema_processor.cpp | 13 ++ .../test_unknown_type.py | 116 ++++++++++++++++++ 6 files changed, 170 insertions(+) create mode 100644 tests/integration/test_storage_iceberg_no_spark/test_unknown_type.py diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Constant.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/Constant.h index c0c3c4e78996..d1e33b2fcae8 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/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..30016d735287 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp @@ -574,6 +574,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: 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/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..a074ddc64619 --- /dev/null +++ b/tests/integration/test_storage_iceberg_no_spark/test_unknown_type.py @@ -0,0 +1,116 @@ +#!/usr/bin/env python3 + +""" +Integration test for Iceberg v3 `unknown` primitive type. + +Creates an Iceberg v3 table with an `unknown`-typed column via PyIceberg, +writes some rows, then reads the table from ClickHouse and asserts: +- The unknown column is mapped to Nullable(Nothing) +- All values in the unknown column are NULL +- Non-unknown columns read correctly +""" + +from pyiceberg.catalog import load_catalog +from pyiceberg.schema import Schema +from pyiceberg.types import NestedField, LongType, StringType, UnknownType +from pyiceberg.partitioning import PartitionSpec +from pyiceberg.table.sorting import SortOrder +from helpers.config_cluster import minio_secret_key, minio_access_key +import pyarrow as pa +import uuid + +BASE_URL = "http://rest:8181/v1" +CATALOG_NAME = "demo" + + +def load_catalog_impl(started_cluster): + base_url_local_raw = f"http://localhost:{started_cluster.iceberg_rest_catalog_port}" + return load_catalog( + CATALOG_NAME, + **{ + "uri": base_url_local_raw, + "type": "rest", + "s3.endpoint": f"http://{started_cluster.minio_ip}:{started_cluster.minio_port}", + "s3.access-key-id": minio_access_key, + "s3.secret-access-key": minio_secret_key, + }, + ) + + +def create_clickhouse_iceberg_database(started_cluster, node, name): + settings = { + "catalog_type": "rest", + "warehouse": "demo", + "storage_endpoint": "http://minio1:9001/warehouse-rest", + } + + node.query( + f""" +DROP DATABASE IF EXISTS {name}; +SET allow_database_iceberg=true; +SET write_full_path_in_iceberg_metadata=1; +CREATE DATABASE {name} ENGINE = DataLakeCatalog('{BASE_URL}', 'minio', '{minio_secret_key}') +SETTINGS {",".join((k+"="+repr(v) for k, v in settings.items()))} + """ + ) + + +def test_unknown_type_read(started_cluster_iceberg_no_spark): + """Create an Iceberg v3 table with an unknown-typed column and verify ClickHouse reads it.""" + instance = started_cluster_iceberg_no_spark.instances["node1"] + root_namespace = f"clickhouse_{uuid.uuid4()}" + catalog = load_catalog_impl(started_cluster_iceberg_no_spark) + + schema = Schema( + NestedField(field_id=1, name="id", field_type=LongType(), required=True), + NestedField(field_id=2, name="name", field_type=StringType(), required=False), + NestedField( + field_id=3, name="placeholder", field_type=UnknownType(), required=False + ), + ) + + table = catalog.create_table( + identifier=f"{root_namespace}.test_unknown", + schema=schema, + location="s3://warehouse-rest/data", + partition_spec=PartitionSpec(), + sort_order=SortOrder(), + properties={"format-version": "3"}, + ) + + # Write data; the unknown column is all NULLs in the Arrow table + data = pa.table( + { + "id": pa.array([1, 2, 3], type=pa.int64()), + "name": pa.array(["alice", "bob", "charlie"], type=pa.large_string()), + "placeholder": pa.nulls(3), + } + ) + table.append(data) + + create_clickhouse_iceberg_database( + started_cluster_iceberg_no_spark, instance, CATALOG_NAME + ) + + fqtn = f"{CATALOG_NAME}.`{root_namespace}.test_unknown`" + + # Verify row count + assert instance.query(f"SELECT count() FROM {fqtn}").strip() == "3" + + # Verify the unknown column type is Nullable(Nothing) + describe = instance.query(f"DESCRIBE TABLE {fqtn}") + 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)" + + # Verify all values in the unknown column are NULL + result = instance.query(f"SELECT placeholder FROM {fqtn}").strip() + assert result == "\\N\n\\N\n\\N" + + # Verify non-unknown columns read correctly alongside the unknown column + result = instance.query( + f"SELECT id, name, placeholder FROM {fqtn} ORDER BY id" + ).strip() + expected = "1\talice\t\\N\n2\tbob\t\\N\n3\tcharlie\t\\N" + assert result == expected From bb1b0962f6d7caaffa000242082ccd273fb6873f Mon Sep 17 00:00:00 2001 From: Kanthi Subramanian Date: Wed, 16 Sep 2026 18:25:43 +0200 Subject: [PATCH 2/8] Changed integration test to modify manifest --- .../test_unknown_type.py | 185 ++++++++++-------- 1 file changed, 103 insertions(+), 82 deletions(-) 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 index a074ddc64619..8501a2a5fb99 100644 --- a/tests/integration/test_storage_iceberg_no_spark/test_unknown_type.py +++ b/tests/integration/test_storage_iceberg_no_spark/test_unknown_type.py @@ -3,114 +3,135 @@ """ Integration test for Iceberg v3 `unknown` primitive type. -Creates an Iceberg v3 table with an `unknown`-typed column via PyIceberg, -writes some rows, then reads the table from ClickHouse and asserts: +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 """ -from pyiceberg.catalog import load_catalog -from pyiceberg.schema import Schema -from pyiceberg.types import NestedField, LongType, StringType, UnknownType -from pyiceberg.partitioning import PartitionSpec -from pyiceberg.table.sorting import SortOrder -from helpers.config_cluster import minio_secret_key, minio_access_key -import pyarrow as pa -import uuid - -BASE_URL = "http://rest:8181/v1" -CATALOG_NAME = "demo" - - -def load_catalog_impl(started_cluster): - base_url_local_raw = f"http://localhost:{started_cluster.iceberg_rest_catalog_port}" - return load_catalog( - CATALOG_NAME, - **{ - "uri": base_url_local_raw, - "type": "rest", - "s3.endpoint": f"http://{started_cluster.minio_ip}:{started_cluster.minio_port}", - "s3.access-key-id": minio_access_key, - "s3.secret-access-key": minio_secret_key, - }, - ) +import json +import re +from helpers.iceberg_utils import create_iceberg_table, get_uuid_str -def create_clickhouse_iceberg_database(started_cluster, node, name): - settings = { - "catalog_type": "rest", - "warehouse": "demo", - "storage_endpoint": "http://minio1:9001/warehouse-rest", - } - node.query( - f""" -DROP DATABASE IF EXISTS {name}; -SET allow_database_iceberg=true; -SET write_full_path_in_iceberg_metadata=1; -CREATE DATABASE {name} ENGINE = DataLakeCatalog('{BASE_URL}', 'minio', '{minio_secret_key}') -SETTINGS {",".join((k+"="+repr(v) for k, v in settings.items()))} - """ +def _metadata_dir(table_name): + return ( + f"/var/lib/clickhouse/user_files/iceberg_data/default/{table_name}/metadata" ) -def test_unknown_type_read(started_cluster_iceberg_no_spark): - """Create an Iceberg v3 table with an unknown-typed column and verify ClickHouse reads it.""" - instance = started_cluster_iceberg_no_spark.instances["node1"] - root_namespace = f"clickhouse_{uuid.uuid4()}" - catalog = load_catalog_impl(started_cluster_iceberg_no_spark) - - schema = Schema( - NestedField(field_id=1, name="id", field_type=LongType(), required=True), - NestedField(field_id=2, name="name", field_type=StringType(), required=False), - NestedField( - field_id=3, name="placeholder", field_type=UnknownType(), required=False - ), +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"] ) - table = catalog.create_table( - identifier=f"{root_namespace}.test_unknown", - schema=schema, - location="s3://warehouse-rest/data", - partition_spec=PartitionSpec(), - sort_order=SortOrder(), - properties={"format-version": "3"}, - ) - # Write data; the unknown column is all NULLs in the Arrow table - data = pa.table( - { - "id": pa.array([1, 2, 3], type=pa.int64()), - "name": pa.array(["alice", "bob", "charlie"], type=pa.large_string()), - "placeholder": pa.nulls(3), - } +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, ) - table.append(data) - - create_clickhouse_iceberg_database( - started_cluster_iceberg_no_spark, instance, CATALOG_NAME + instance.query( + f"INSERT INTO {table_name} VALUES (1, 'alice'), (2, 'bob'), (3, 'charlie')" ) - fqtn = f"{CATALOG_NAME}.`{root_namespace}.test_unknown`" + # 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, + ) - # Verify row count - assert instance.query(f"SELECT count() FROM {fqtn}").strip() == "3" + # 6. Verify row count. + assert instance.query(f"SELECT count() FROM {table_name}").strip() == "3" - # Verify the unknown column type is Nullable(Nothing) - describe = instance.query(f"DESCRIBE TABLE {fqtn}") + # 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 len(placeholder_row) == 1, ( + f"Expected one 'placeholder' column, got: {lines}" + ) assert placeholder_row[0][1] == "Nullable(Nothing)" - # Verify all values in the unknown column are NULL - result = instance.query(f"SELECT placeholder FROM {fqtn}").strip() + # 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" - # Verify non-unknown columns read correctly alongside the unknown column + # 9. Verify non-unknown columns read correctly alongside the unknown column. result = instance.query( - f"SELECT id, name, placeholder FROM {fqtn} ORDER BY id" + 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 From 7e4d5753a8bd82ebf99fbebb124cdfcde7772e0c Mon Sep 17 00:00:00 2001 From: Kanthi Subramanian Date: Sun, 20 Sep 2026 00:39:02 +0200 Subject: [PATCH 3/8] Fix write path for unknown iceberg data type with parquet --- .../DataLakes/Iceberg/MultipleFileWriter.cpp | 51 +++++++++- .../DataLakes/Iceberg/MultipleFileWriter.h | 7 ++ .../test_unknown_type.py | 96 +++++++++++++++++++ 3 files changed, 152 insertions(+), 2 deletions(-) diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp index 3414bf04e4c5..ab8475356c5a 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp @@ -1,5 +1,7 @@ #include +#include +#include #include #include #include @@ -39,6 +41,33 @@ MultipleFileWriter::MultipleFileWriter( , new_file_path_callback(std::move(new_file_path_callback_)) { column_mapper->setStorageColumnEncoding(Iceberg::IcebergSchemaProcessor::traverseSchema(schema_)); + + /// Iceberg `unknown` type maps to Nullable(Nothing) in ClickHouse. No + /// serialisation format (Parquet, ORC, Avro) can represent the Nothing type, + /// and the column is guaranteed to contain only NULLs, so we strip it from the + /// block that is passed to the format writer. The column still appears in the + /// Iceberg schema metadata and is read back as NULLs on the read path. + for (size_t i = 0; i < sample_block->columns(); ++i) + { + auto inner_type = removeNullable(sample_block->getByPosition(i).type); + if (isNothing(inner_type)) + nothing_column_indices.push_back(i); + } + + if (nothing_column_indices.empty()) + { + filtered_sample_block = sample_block; + } + else + { + Block filtered; + for (size_t i = 0; i < sample_block->columns(); ++i) + { + if (std::find(nothing_column_indices.begin(), nothing_column_indices.end(), i) == nothing_column_indices.end()) + filtered.insert(sample_block->getByPosition(i)); + } + filtered_sample_block = std::make_shared(std::move(filtered)); + } } void MultipleFileWriter::startNewFile() @@ -69,7 +98,7 @@ 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); } void MultipleFileWriter::consume(const Chunk & chunk) @@ -78,7 +107,25 @@ void MultipleFileWriter::consume(const Chunk & chunk) { startNewFile(); } - output_format->write(sample_block->cloneWithColumns(chunk.getColumns())); + + if (nothing_column_indices.empty()) + { + output_format->write(sample_block->cloneWithColumns(chunk.getColumns())); + } + else + { + /// Strip Nullable(Nothing) columns before passing to the format writer. + auto columns = chunk.getColumns(); + Columns filtered_columns; + filtered_columns.reserve(columns.size() - nothing_column_indices.size()); + for (size_t i = 0; i < columns.size(); ++i) + { + if (std::find(nothing_column_indices.begin(), nothing_column_indices.end(), i) == nothing_column_indices.end()) + filtered_columns.push_back(columns[i]); + } + output_format->write(filtered_sample_block->cloneWithColumns(std::move(filtered_columns))); + } + 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..51c0a5c56198 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.h @@ -92,6 +92,13 @@ class MultipleFileWriter std::optional format_settings; const String& write_format; SharedHeader sample_block; + /// Indices of Nullable(Nothing) columns (Iceberg `unknown` type) that must be + /// stripped before the data reaches the Parquet/format writer, because no + /// serialisation format supports the Nothing type. The columns exist only in + /// the Iceberg schema metadata and are always read back as NULLs. + std::vector nothing_column_indices; + /// `sample_block` with the 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/tests/integration/test_storage_iceberg_no_spark/test_unknown_type.py b/tests/integration/test_storage_iceberg_no_spark/test_unknown_type.py index 8501a2a5fb99..c9694581abe6 100644 --- a/tests/integration/test_storage_iceberg_no_spark/test_unknown_type.py +++ b/tests/integration/test_storage_iceberg_no_spark/test_unknown_type.py @@ -135,3 +135,99 @@ def test_unknown_type_read(started_cluster_iceberg_no_spark): ).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)" From e77daf606009b8205b73a95fc8dac75221db9f5f Mon Sep 17 00:00:00 2001 From: Kanthi Subramanian Date: Mon, 28 Sep 2026 00:21:56 +0200 Subject: [PATCH 4/8] Added logic to create a helper function to check for unknown in nested types --- .../DataLakes/Iceberg/MultipleFileWriter.cpp | 62 +++++++--- .../DataLakes/Iceberg/MultipleFileWriter.h | 11 +- .../test_unknown_type.py | 109 ++++++++++++++++++ 3 files changed, 159 insertions(+), 23 deletions(-) diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp index ab8475356c5a..0eb698b0a014 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp @@ -1,7 +1,10 @@ #include +#include +#include #include #include +#include #include #include #include @@ -14,6 +17,36 @@ namespace DB #if USE_AVRO +namespace +{ + +/// Iceberg `unknown` is a primitive type, so it can also appear nested inside +/// structs, lists, and maps. No serialisation format (Parquet, ORC, Avro) can +/// represent Nothing at any nesting level (the Parquet writer throws +/// `UNKNOWN_TYPE`), so such columns must be stripped from the writer's block. +bool containsNothing(const DataTypePtr & type) +{ + auto inner_type = removeNullable(type); + if (isNothing(inner_type)) + return true; + if (const auto * tuple_type = typeid_cast(inner_type.get())) + { + for (const auto & elem : tuple_type->getElements()) + { + if (containsNothing(elem)) + return true; + } + return false; + } + if (const auto * array_type = typeid_cast(inner_type.get())) + return containsNothing(array_type->getNestedType()); + if (const auto * map_type = typeid_cast(inner_type.get())) + return containsNothing(map_type->getKeyType()) || containsNothing(map_type->getValueType()); + return false; +} + +} + MultipleFileWriter::MultipleFileWriter( UInt64 max_data_file_num_rows_, UInt64 max_data_file_num_bytes_, @@ -49,23 +82,19 @@ MultipleFileWriter::MultipleFileWriter( /// Iceberg schema metadata and is read back as NULLs on the read path. for (size_t i = 0; i < sample_block->columns(); ++i) { - auto inner_type = removeNullable(sample_block->getByPosition(i).type); - if (isNothing(inner_type)) - nothing_column_indices.push_back(i); + if (!containsNothing(sample_block->getByPosition(i).type)) + kept_column_indices.push_back(i); } - if (nothing_column_indices.empty()) + if (kept_column_indices.size() == sample_block->columns()) { filtered_sample_block = sample_block; } else { Block filtered; - for (size_t i = 0; i < sample_block->columns(); ++i) - { - if (std::find(nothing_column_indices.begin(), nothing_column_indices.end(), i) == nothing_column_indices.end()) - filtered.insert(sample_block->getByPosition(i)); - } + for (size_t i : kept_column_indices) + filtered.insert(sample_block->getByPosition(i)); filtered_sample_block = std::make_shared(std::move(filtered)); } } @@ -108,21 +137,18 @@ void MultipleFileWriter::consume(const Chunk & chunk) startNewFile(); } - if (nothing_column_indices.empty()) + if (kept_column_indices.size() == sample_block->columns()) { output_format->write(sample_block->cloneWithColumns(chunk.getColumns())); } else { - /// Strip Nullable(Nothing) columns before passing to the format writer. - auto columns = chunk.getColumns(); + /// Strip columns containing Nothing before passing to the format writer. + const auto & columns = chunk.getColumns(); Columns filtered_columns; - filtered_columns.reserve(columns.size() - nothing_column_indices.size()); - for (size_t i = 0; i < columns.size(); ++i) - { - if (std::find(nothing_column_indices.begin(), nothing_column_indices.end(), i) == nothing_column_indices.end()) - filtered_columns.push_back(columns[i]); - } + filtered_columns.reserve(kept_column_indices.size()); + for (size_t i : kept_column_indices) + filtered_columns.push_back(columns[i]); output_format->write(filtered_sample_block->cloneWithColumns(std::move(filtered_columns))); } diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.h index 51c0a5c56198..3c7a5a240bba 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.h @@ -92,11 +92,12 @@ class MultipleFileWriter std::optional format_settings; const String& write_format; SharedHeader sample_block; - /// Indices of Nullable(Nothing) columns (Iceberg `unknown` type) that must be - /// stripped before the data reaches the Parquet/format writer, because no - /// serialisation format supports the Nothing type. The columns exist only in - /// the Iceberg schema metadata and are always read back as NULLs. - std::vector nothing_column_indices; + /// Indices of `sample_block` columns that are passed to the format writer. + /// Columns containing Nothing at any nesting level (Iceberg `unknown` type) + /// are excluded, because no serialisation format supports the Nothing type. + /// Those columns exist only in the Iceberg schema metadata and are always + /// read back as NULLs. + std::vector kept_column_indices; /// `sample_block` with the Nothing columns removed; used by the format writer. SharedHeader filtered_sample_block; UInt64 total_bytes = 0; 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 index c9694581abe6..cf1e849bd082 100644 --- a/tests/integration/test_storage_iceberg_no_spark/test_unknown_type.py +++ b/tests/integration/test_storage_iceberg_no_spark/test_unknown_type.py @@ -231,3 +231,112 @@ def test_unknown_type_write(started_cluster_iceberg_no_spark): 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_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.""" + 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) + struct_field_id = last_column_id + 1 + known_subfield_id = last_column_id + 2 + unknown_subfield_id = last_column_id + 3 + 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": struct_field_id, + "name": "nested", + "required": False, + "type": { + "type": "struct", + "fields": [ + { + "id": known_subfield_id, + "name": "a", + "required": False, + "type": "long", + }, + { + "id": unknown_subfield_id, + "name": "u", + "required": False, + "type": "unknown", + }, + ], + }, + } + ], + } + meta.setdefault("schemas", []).append(new_schema) + meta["current-schema-id"] = new_schema_id + meta["last-column-id"] = unknown_subfield_id + + _write_next_metadata(instance, table_name, meta, prev_path) + + +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) + + 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}, + ) + + 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. + instance.query( + f"INSERT INTO {table_name} (id, name) VALUES (3, 'charlie'), (4, 'dave')", + settings={"allow_insert_into_iceberg": 1}, + ) + + result = instance.query( + f"SELECT id, name FROM {table_name} ORDER BY id" + ).strip() + expected = ( + "1\talice\n" + "2\tbob\n" + "3\tcharlie\n" + "4\tdave" + ) + assert result == expected From f67f1e1ee7aff8fef84f2dc0dc1132ad76af3d21 Mon Sep 17 00:00:00 2001 From: Kanthi Subramanian Date: Thu, 1 Oct 2026 19:36:04 +0200 Subject: [PATCH 5/8] Added logic to string unknown data type in fileWriter --- .../DataLakes/Iceberg/MultipleFileWriter.cpp | 232 ++++++++++++++---- .../DataLakes/Iceberg/MultipleFileWriter.h | 34 ++- .../tests/gtest_iceberg_strip_nothing.cpp | 115 +++++++++ 3 files changed, 331 insertions(+), 50 deletions(-) create mode 100644 src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_strip_nothing.cpp diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp index 0eb698b0a014..0231c35a5666 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp @@ -1,5 +1,10 @@ #include +#include +#include +#include +#include +#include #include #include #include @@ -15,38 +20,149 @@ namespace DB { -#if USE_AVRO +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 { -/// Iceberg `unknown` is a primitive type, so it can also appear nested inside -/// structs, lists, and maps. No serialisation format (Parquet, ORC, Avro) can -/// represent Nothing at any nesting level (the Parquet writer throws -/// `UNKNOWN_TYPE`), so such columns must be stripped from the writer's block. -bool containsNothing(const DataTypePtr & type) +ColumnPtr stripNothingColumnImpl(const ColumnPtr & column, const DataTypePtr & type) { - auto inner_type = removeNullable(type); - if (isNothing(inner_type)) - return true; - if (const auto * tuple_type = typeid_cast(inner_type.get())) + 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())) { - for (const auto & elem : tuple_type->getElements()) + 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 (containsNothing(elem)) - return true; + if (stripNothing(elements[i])) + kept_columns.push_back(stripNothingColumnImpl(tuple_column.getColumnPtr(i), elements[i])); } - return false; + 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()); } - if (const auto * array_type = typeid_cast(inner_type.get())) - return containsNothing(array_type->getNestedType()); - if (const auto * map_type = typeid_cast(inner_type.get())) - return containsNothing(map_type->getKeyType()) || containsNothing(map_type->getValueType()); - return false; + + 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( UInt64 max_data_file_num_rows_, UInt64 max_data_file_num_bytes_, @@ -75,26 +191,27 @@ MultipleFileWriter::MultipleFileWriter( { column_mapper->setStorageColumnEncoding(Iceberg::IcebergSchemaProcessor::traverseSchema(schema_)); - /// Iceberg `unknown` type maps to Nullable(Nothing) in ClickHouse. No - /// serialisation format (Parquet, ORC, Avro) can represent the Nothing type, - /// and the column is guaranteed to contain only NULLs, so we strip it from the - /// block that is passed to the format writer. The column still appears in the - /// Iceberg schema metadata and is read back as NULLs on the read path. - for (size_t i = 0; i < sample_block->columns(); ++i) + written_column_types.reserve(sample_block->columns()); + for (const auto & column : *sample_block) { - if (!containsNothing(sample_block->getByPosition(i).type)) - kept_column_indices.push_back(i); + written_column_types.push_back(Iceberg::stripNothing(column.type)); + has_nothing_leaves |= written_column_types.back() != column.type; } - if (kept_column_indices.size() == sample_block->columns()) + if (!has_nothing_leaves) { filtered_sample_block = sample_block; } else { Block filtered; - for (size_t i : kept_column_indices) - filtered.insert(sample_block->getByPosition(i)); + 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)); } } @@ -130,27 +247,56 @@ void MultipleFileWriter::startNewFile() write_format, *buffer, *filtered_sample_block, context, format_settings, format_filter_info); } +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. + 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(); } - if (kept_column_indices.size() == sample_block->columns()) - { - output_format->write(sample_block->cloneWithColumns(chunk.getColumns())); - } + if (filtered_columns) + output_format->write(filtered_sample_block->cloneWithColumns(std::move(*filtered_columns))); else - { - /// Strip columns containing Nothing before passing to the format writer. - const auto & columns = chunk.getColumns(); - Columns filtered_columns; - filtered_columns.reserve(kept_column_indices.size()); - for (size_t i : kept_column_indices) - filtered_columns.push_back(columns[i]); - output_format->write(filtered_sample_block->cloneWithColumns(std::move(filtered_columns))); - } + output_format->write(sample_block->cloneWithColumns(chunk.getColumns())); output_format->flush(); *current_file_num_rows += chunk.getNumRows(); diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.h index 3c7a5a240bba..1c082e1e26f3 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,9 @@ 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; + UInt64 max_data_file_num_rows; UInt64 max_data_file_num_bytes; Poco::JSON::Array::Ptr schema; @@ -92,13 +112,13 @@ class MultipleFileWriter std::optional format_settings; const String& write_format; SharedHeader sample_block; - /// Indices of `sample_block` columns that are passed to the format writer. - /// Columns containing Nothing at any nesting level (Iceberg `unknown` type) - /// are excluded, because no serialisation format supports the Nothing type. - /// Those columns exist only in the Iceberg schema metadata and are always - /// read back as NULLs. - std::vector kept_column_indices; - /// `sample_block` with the Nothing columns removed; used by the format writer. + /// 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/tests/gtest_iceberg_strip_nothing.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_strip_nothing.cpp new file mode 100644 index 000000000000..b6f7a3600eb3 --- /dev/null +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_strip_nothing.cpp @@ -0,0 +1,115 @@ +#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))})})})); + } +} From 67cbdb61ff7969fafc48f3b04808ed9b7a4608a2 Mon Sep 17 00:00:00 2001 From: Kanthi Subramanian Date: Thu, 1 Oct 2026 22:39:09 +0200 Subject: [PATCH 6/8] Address findings, reject unknown on iceberg v2 tables, dropped unknown column will be ignored in file stats, inserts with nothing will be rejected --- .../DataLakes/Iceberg/DataFileStatistics.cpp | 23 +- .../DataLakes/Iceberg/DataFileStatistics.h | 6 + .../DataLakes/Iceberg/MetadataGenerator.cpp | 2 + .../DataLakes/Iceberg/MultipleFileWriter.cpp | 18 + .../DataLakes/Iceberg/MultipleFileWriter.h | 2 + .../ObjectStorage/DataLakes/Iceberg/Utils.cpp | 14 + .../ObjectStorage/DataLakes/Iceberg/Utils.h | 5 + .../tests/gtest_iceberg_strip_nothing.cpp | 85 +++++ .../test_unknown_type.py | 317 +++++++++++++++--- 9 files changed, 419 insertions(+), 53 deletions(-) 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 0231c35a5666..71674f8d6eef 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.cpp @@ -204,6 +204,8 @@ MultipleFileWriter::MultipleFileWriter( } else { + aggregate_stats.excludeColumns(getStatisticsExcludedColumns()); + Block filtered; for (size_t i = 0; i < sample_block->columns(); ++i) { @@ -224,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(); @@ -247,6 +251,14 @@ void MultipleFileWriter::startNewFile() 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; @@ -284,6 +296,12 @@ Columns MultipleFileWriter::filterColumns(const Columns & columns) const 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()); diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.h b/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.h index 1c082e1e26f3..00712850182a 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.h +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/MultipleFileWriter.h @@ -87,6 +87,8 @@ class MultipleFileWriter 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; diff --git a/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/Utils.cpp index 30016d735287..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) { @@ -1060,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_strip_nothing.cpp b/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_strip_nothing.cpp index b6f7a3600eb3..42a1617b8d4f 100644 --- a/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_strip_nothing.cpp +++ b/src/Storages/ObjectStorage/DataLakes/Iceberg/tests/gtest_iceberg_strip_nothing.cpp @@ -3,6 +3,9 @@ #include #include #include +#include +#include +#include #include using namespace DB; @@ -113,3 +116,85 @@ TEST(IcebergStripNothing, ArrayAndMapColumnsKeepOffsetsAndKeys) 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 index cf1e849bd082..0e076f526c2b 100644 --- a/tests/integration/test_storage_iceberg_no_spark/test_unknown_type.py +++ b/tests/integration/test_storage_iceberg_no_spark/test_unknown_type.py @@ -14,7 +14,11 @@ import json import re -from helpers.iceberg_utils import create_iceberg_table, get_uuid_str +from helpers.iceberg_utils import ( + create_iceberg_table, + get_creation_expression, + get_uuid_str, +) def _metadata_dir(table_name): @@ -233,9 +237,9 @@ def test_unknown_type_write(started_cluster_iceberg_no_spark): assert placeholder_row[0][1] == "Nullable(Nothing)" -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 _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 @@ -247,48 +251,84 @@ def _patch_table_with_nested_unknown_column(instance, table_name): break assert current_schema is not None - last_column_id = meta.get("last-column-id", 0) - struct_field_id = last_column_id + 1 - known_subfield_id = last_column_id + 2 - unknown_subfield_id = last_column_id + 3 + 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"] - + [ - { - "id": struct_field_id, - "name": "nested", - "required": False, - "type": { - "type": "struct", - "fields": [ - { - "id": known_subfield_id, - "name": "a", - "required": False, - "type": "long", - }, - { - "id": unknown_subfield_id, - "name": "u", - "required": False, - "type": "unknown", - }, - ], - }, - } - ], + "fields": current_schema["fields"] + [field], } meta.setdefault("schemas", []).append(new_schema) meta["current-schema-id"] = new_schema_id - meta["last-column-id"] = unknown_subfield_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 @@ -308,35 +348,214 @@ def test_unknown_type_nested_write(started_cluster_iceberg_no_spark): 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() - 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}, + "(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) - 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. instance.query( - f"INSERT INTO {table_name} (id, name) VALUES (3, 'charlie'), (4, 'dave')", + 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 FROM {table_name} ORDER BY id" + f"SELECT id, name, tags FROM {table_name} ORDER BY id" ).strip() expected = ( - "1\talice\n" - "2\tbob\n" - "3\tcharlie\n" - "4\tdave" + "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" From da969e7256ac0898bd774c7c3685dfe2ed274a02 Mon Sep 17 00:00:00 2001 From: Kanthi Subramanian Date: Fri, 2 Oct 2026 17:42:07 +0200 Subject: [PATCH 7/8] Fix unknown type integration test --- .../StorageObjectStorageSource.cpp | 41 +++++ .../test_tuple_element_read.py | 39 +++++ .../test_unknown_type.py | 43 ++++-- .../configs/config.d/cache_62964.xml | 9 ++ .../test_writes_field_ids_spark_read.py | 140 ++++++++++++++++++ 5 files changed, 256 insertions(+), 16 deletions(-) create mode 100644 tests/integration/test_storage_iceberg_no_spark/test_tuple_element_read.py create mode 100644 tests/integration/test_storage_iceberg_with_spark/configs/config.d/cache_62964.xml create mode 100644 tests/integration/test_storage_iceberg_with_spark/test_writes_field_ids_spark_read.py diff --git a/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp b/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp index 79b63f6dc1b7..191d8db7db8e 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp @@ -34,6 +34,7 @@ #include #include #include +#include #include #include #include @@ -1065,6 +1066,12 @@ StorageObjectStorageSource::ReaderHolder StorageObjectStorageSource::createReade if (column_name.starts_with("materialize(") && column_name.ends_with(")")) continue; + /// Tuple element reads request nested fields like `t.a`. Writers don't have to store + /// statistics for nested fields (ClickHouse stores them only for top-level columns), + /// so missing statistics don't prove that the field is absent from the file. + if (column_name.contains('.')) + continue; + /// Skip columns produced by prewhere or row-level filter expressions — /// they are computed at read time, not stored in the file. if (format_filter_info @@ -1470,6 +1477,40 @@ StorageObjectStorageSource::ReaderHolder StorageObjectStorageSource::createReade for (const auto & column_name : row_lineage_columns) columns_to_extract.emplace_back(column_name, rowLineageColumnType()); + /// Tuple element reads request nested fields like `t.a` as separate columns. If the data lake + /// schema transform has rebuilt the whole `t` instead (e.g. the file was written before `t` + /// was added), extract `t.a` as a subcolumn of `t`. + const Block & extract_header = builder.getHeader(); + bool need_materialize = false; + for (auto & column : columns_to_extract) + { + const String name_in_storage = column.getNameInStorage(); + if (extract_header.has(name_in_storage)) + continue; + + for (size_t pos = name_in_storage.rfind('.'); pos != String::npos && pos != 0; pos = name_in_storage.rfind('.', pos - 1)) + { + const auto * storage_column = extract_header.findByName(name_in_storage.substr(0, pos)); + const String subcolumn_name = column.name.substr(pos + 1); + if (storage_column && storage_column->type->tryGetSubcolumnType(subcolumn_name)) + { + column = NameAndTypePair(storage_column->name, subcolumn_name, storage_column->type, column.type); + /// Subcolumns can't be extracted from a constant column, e.g. the default the schema transform + /// supplies for a column missing in the file. + need_materialize |= storage_column->column && isColumnConst(*storage_column->column); + break; + } + } + } + + if (need_materialize) + { + builder.addSimpleTransform([](const SharedHeader & header) + { + return std::make_shared(header, /*remove_special_representations=*/ false); + }); + } + builder.addSimpleTransform([&](const SharedHeader & header) { return std::make_shared(header, columns_to_extract); diff --git a/tests/integration/test_storage_iceberg_no_spark/test_tuple_element_read.py b/tests/integration/test_storage_iceberg_no_spark/test_tuple_element_read.py new file mode 100644 index 000000000000..f23055f13f02 --- /dev/null +++ b/tests/integration/test_storage_iceberg_no_spark/test_tuple_element_read.py @@ -0,0 +1,39 @@ +import pytest + +from helpers.iceberg_utils import ( + create_iceberg_table, + get_uuid_str +) + + +@pytest.mark.parametrize("read_optimization", [0, 1]) +def test_tuple_element_read(started_cluster_iceberg_no_spark, read_optimization): + """Manifests written by ClickHouse carry statistics only for top-level columns, + so a tuple element such as `t.a` has none. That must not make the read + optimization treat the element as absent from the file and return NULL.""" + instance = started_cluster_iceberg_no_spark.instances["node1"] + TABLE_NAME = "test_tuple_element_read_" + get_uuid_str() + + create_iceberg_table( + "local", + instance, + TABLE_NAME, + started_cluster_iceberg_no_spark, + "(id Int64, t Tuple(a Nullable(Int64), b Nullable(String)))", + format_version=2, + ) + instance.query( + f"INSERT INTO {TABLE_NAME} VALUES (1, (42, 'x')), (2, (NULL, 'y'))", + settings={"allow_insert_into_iceberg": 1}, + ) + + settings = {"allow_experimental_iceberg_read_optimization": read_optimization} + assert ( + instance.query(f"SELECT id, t.a, t.b FROM {TABLE_NAME} ORDER BY id", settings=settings) + == "1\t42\tx\n2\t\\N\ty\n" + ) + assert instance.query(f"SELECT t.a FROM {TABLE_NAME} ORDER BY id", settings=settings) == "42\n\\N\n" + assert ( + instance.query(f"SELECT id, t, t.b FROM {TABLE_NAME} ORDER BY id", settings=settings) + == "1\t(42,'x')\tx\n2\t(NULL,'y')\ty\n" + ) 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 index 0e076f526c2b..e235e9170d27 100644 --- a/tests/integration/test_storage_iceberg_no_spark/test_unknown_type.py +++ b/tests/integration/test_storage_iceberg_no_spark/test_unknown_type.py @@ -363,16 +363,23 @@ def test_unknown_type_nested_write(started_cluster_iceberg_no_spark): 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 + for read_optimization in (0, 1): + result = instance.query( + f"SELECT id, name, nested.a, nested.u FROM {table_name} ORDER BY id", + settings={"allow_experimental_iceberg_read_optimization": read_optimization}, + ).strip() + assert result == expected, read_optimization + + result = instance.query( + f"SELECT id, nested FROM {table_name} WHERE id >= 3 ORDER BY id" + ).strip() + assert result == "3\t(42,NULL)\n4\t(NULL,NULL)" def test_unknown_type_list_write(started_cluster_iceberg_no_spark): @@ -421,8 +428,12 @@ def test_unknown_type_list_write(started_cluster_iceberg_no_spark): 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.""" + write it into the metadata of a table with an older format version. + + ClickHouse itself refuses a top-level `Nullable(Nothing)` column in any table, + but `Nothing` nested in a `Tuple` is accepted, so that is what reaches Iceberg.""" instance = started_cluster_iceberg_no_spark.instances["node1"] + nested_unknown = "Tuple(a Nullable(Int64), u Nullable(Nothing))" table_name = "test_unknown_type_create_v2_" + get_uuid_str() error = instance.query_and_get_error( @@ -430,30 +441,30 @@ def test_unknown_type_rejected_on_v2(started_cluster_iceberg_no_spark): "local", table_name, started_cluster_iceberg_no_spark, - "(id Int64, u Nullable(Nothing))", + f"(id Int64, nested {nested_unknown})", 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() + table_name = "test_unknown_type_modify_column_v2_" + get_uuid_str() create_iceberg_table( "local", instance, table_name, started_cluster_iceberg_no_spark, - "(id Int64, name Nullable(String))", + "(id Int64, nested Tuple(a Nullable(Int64)))", format_version=2, ) + describe_before = instance.query(f"DESCRIBE TABLE {table_name}") error = instance.query_and_get_error( - f"ALTER TABLE {table_name} ADD COLUMN u Nullable(Nothing)", + f"ALTER TABLE {table_name} MODIFY COLUMN nested {nested_unknown}", 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"] + assert instance.query(f"DESCRIBE TABLE {table_name}") == describe_before table_name = "test_unknown_type_create_v3_" + get_uuid_str() create_iceberg_table( @@ -461,12 +472,12 @@ def test_unknown_type_rejected_on_v2(started_cluster_iceberg_no_spark): instance, table_name, started_cluster_iceberg_no_spark, - "(id Int64, u Nullable(Nothing))", + f"(id Int64, nested {nested_unknown})", 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 + nested_row = [line.split("\t") for line in describe.strip().split("\n") if line.startswith("nested\t")] + assert len(nested_row) == 1 and "Nothing" in nested_row[0][1], describe def test_unknown_type_write_keeps_bounds(started_cluster_iceberg_no_spark): @@ -547,12 +558,12 @@ def test_unknown_type_only_column_rejected(started_cluster_iceberg_no_spark): instance, table_name, started_cluster_iceberg_no_spark, - "(u Nullable(Nothing))", + "(t Tuple(u Nullable(Nothing)))", format_version=3, ) error = instance.query_and_get_error( - f"INSERT INTO {table_name} VALUES (NULL), (NULL)", + f"INSERT INTO {table_name} SELECT tuple(NULL) FROM numbers(2)", settings={"allow_insert_into_iceberg": 1}, ) assert "NOT_IMPLEMENTED" in error, error diff --git a/tests/integration/test_storage_iceberg_with_spark/configs/config.d/cache_62964.xml b/tests/integration/test_storage_iceberg_with_spark/configs/config.d/cache_62964.xml new file mode 100644 index 000000000000..55996d877400 --- /dev/null +++ b/tests/integration/test_storage_iceberg_with_spark/configs/config.d/cache_62964.xml @@ -0,0 +1,9 @@ + + + + + 1Gi + cache_62964 + + + diff --git a/tests/integration/test_storage_iceberg_with_spark/test_writes_field_ids_spark_read.py b/tests/integration/test_storage_iceberg_with_spark/test_writes_field_ids_spark_read.py new file mode 100644 index 000000000000..d1b50ec9a1fc --- /dev/null +++ b/tests/integration/test_storage_iceberg_with_spark/test_writes_field_ids_spark_read.py @@ -0,0 +1,140 @@ +import glob + +from helpers.iceberg_utils import ( + create_iceberg_table, + get_uuid_str, + default_download_directory, +) + + +def _avro_field_ids(node_schema): + """Collect every {path: field-id} pair from an avro record schema tree, + descending through array/map wrappers exactly like the writer does.""" + import avro.schema + + result = {} + + def walk(schema, path): + if isinstance(schema, avro.schema.RecordSchema): + for field in schema.fields: + child_path = f"{path}.{field.name}" if path else field.name + fid = field.other_props.get("field-id") + if fid is not None: + result[child_path] = fid + walk(field.type, child_path) + elif isinstance(schema, avro.schema.UnionSchema): + for member in schema.schemas: + walk(member, path) + elif isinstance(schema, avro.schema.ArraySchema): + walk(schema.items, f"{path}.element") + elif isinstance(schema, avro.schema.MapSchema): + walk(schema.values, f"{path}.value") + + walk(node_schema, "") + return result + + +def _avro_metadata_schemas(path): + """Return (manifest_list_schema, [manifest_schemas]) writer schemas from the + Avro files under path/metadata/. Manifest lists are the snap-*.avro files; + the remaining *.avro files are manifests.""" + import avro.datafile + import avro.io + + def _writer_schema(avro_path): + with open(avro_path, "rb") as f: + reader = avro.datafile.DataFileReader(f, avro.io.DatumReader()) + return reader.datum_reader.writers_schema + + manifest_lists = glob.glob(f"{path}/metadata/snap-*.avro") + manifests = [ + p + for p in glob.glob(f"{path}/metadata/*.avro") + if "/snap-" not in p and not p.rsplit("/", 1)[-1].startswith("snap-") + ] + assert manifest_lists, "no manifest-list (snap-*.avro) was written" + assert manifests, "no manifest (*.avro) was written" + return _writer_schema(manifest_lists[0]), [_writer_schema(p) for p in manifests] + + +# Regression test for https://github.com/ClickHouse/ClickHouse/issues/111763. +# The Avro schemas embedded in ClickHouse-written manifest-list and manifest +# files must carry the Iceberg spec field-ids (manifest_path=500, status=0, +# data_file=2, ...). The bundled avro-cpp JSON compiler drops the field-id/ +# element-id attributes, so before the fix the header schema serialized by +# DataFileWriter omitted them and PyIceberg/Spark could not plan a scan of a +# ClickHouse-written table (ValueError: Cannot convert field, missing field-id). +# ClickHouse itself reads via the Iceberg `schema` metadata key, so a ClickHouse +# round-trip cannot catch this - an external reader (Spark here) is required. +def test_writes_manifest_field_ids_spark_read(started_cluster_iceberg_with_spark): + instance = started_cluster_iceberg_with_spark.instances["node1"] + spark = started_cluster_iceberg_with_spark.spark_session + storage_type = "local" + TABLE_NAME = "test_manifest_field_ids_" + get_uuid_str() + local_path = f"/var/lib/clickhouse/user_files/iceberg_data/default/{TABLE_NAME}" + + # Partitioned so the manifest `data_file.partition` struct is populated and + # its field-id must match the persisted partition spec. + create_iceberg_table( + storage_type, + instance, + TABLE_NAME, + started_cluster_iceberg_with_spark, + "(id Int32, label String, score Float64)", + 2, + partition_by="id", + format="Avro", + ) + + instance.query( + f"INSERT INTO {TABLE_NAME} VALUES (1, 'alice', 1.5), (2, 'bob', 2.5), (3, 'charlie', 3.5)", + settings={"allow_insert_into_iceberg": 1}, + ) + + default_download_directory( + started_cluster_iceberg_with_spark, + storage_type, + f"{local_path}/", + f"{local_path}/", + ) + + manifest_list_schema, manifest_schemas = _avro_metadata_schemas(local_path) + + # Manifest-list (`manifest_file`) spec field-ids. + ml_ids = _avro_field_ids(manifest_list_schema) + assert ml_ids.get("manifest_path") == 500, ml_ids + assert ml_ids.get("manifest_length") == 501, ml_ids + assert ml_ids.get("partition_spec_id") == 502, ml_ids + assert ml_ids.get("added_snapshot_id") == 503, ml_ids + # `partitions` is a list; its field-summary subfields carry ids too. + assert ml_ids.get("partitions") == 507, ml_ids + assert ml_ids.get("partitions.element.contains_null") == 509, ml_ids + + # The partition spec field-id the manifest `partition` struct must reuse. The + # single partition column `id` is field-id 1 in the schema, so the spec assigns + # it partition field-id 1001 (ClickHouse numbers partition fields from 1001). + expected_partition_field_id = 1001 + + # Manifest (`manifest_entry`) spec field-ids, on every manifest written. + for manifest_schema in manifest_schemas: + m_ids = _avro_field_ids(manifest_schema) + assert m_ids.get("status") == 0, m_ids + assert m_ids.get("snapshot_id") == 1, m_ids + assert m_ids.get("data_file") == 2, m_ids + assert m_ids.get("data_file.file_path") == 100, m_ids + assert m_ids.get("data_file.record_count") == 103, m_ids + # The map subfields (column_sizes etc.) are Avro logicalType=map arrays; + # their key/value carry Iceberg key-id/value-id, reached via `.element`. + assert m_ids.get("data_file.column_sizes.element.key") == 117, m_ids + assert m_ids.get("data_file.column_sizes.element.value") == 118, m_ids + # The manifest `partition` struct field-id must equal the partition spec's. + assert m_ids.get("data_file.partition.id") == expected_partition_field_id, m_ids + + # End-to-end: Spark plans the scan via the manifest field-ids (exactly what + # was broken) and reads the ClickHouse-written rows back. + rows = spark.read.format("iceberg").load(local_path).orderBy("id").collect() + assert [(r.id, r.label, r.score) for r in rows] == [ + (1, "alice", 1.5), + (2, "bob", 2.5), + (3, "charlie", 3.5), + ] From fac315cd112e6f81c5a24a2146911d8204090aa6 Mon Sep 17 00:00:00 2001 From: Kanthi Subramanian Date: Fri, 2 Oct 2026 18:41:53 +0200 Subject: [PATCH 8/8] Deleted cache file --- .../configs/config.d/cache_62964.xml | 9 --------- 1 file changed, 9 deletions(-) delete mode 100644 tests/integration/test_storage_iceberg_with_spark/configs/config.d/cache_62964.xml diff --git a/tests/integration/test_storage_iceberg_with_spark/configs/config.d/cache_62964.xml b/tests/integration/test_storage_iceberg_with_spark/configs/config.d/cache_62964.xml deleted file mode 100644 index 55996d877400..000000000000 --- a/tests/integration/test_storage_iceberg_with_spark/configs/config.d/cache_62964.xml +++ /dev/null @@ -1,9 +0,0 @@ - - - - - 1Gi - cache_62964 - - -