diff --git a/docs/en/antalya/part_export.md b/docs/en/antalya/part_export.md index 73d467c5d9b1..35a6368b5ceb 100644 --- a/docs/en/antalya/part_export.md +++ b/docs/en/antalya/part_export.md @@ -47,10 +47,23 @@ SETTINGS allow_experimental_export_merge_tree_part = 1 ## Requirements -Source and destination tables must be 100% compatible: +Source and destination tables must support positional schema conversion: -1. **Identical schemas** - same columns, types, and order -2. **Matching partition keys** - partition expressions must be identical +1. **Positionally compatible schemas** - source columns are matched to destination columns by position, similar to `INSERT INTO dest SELECT * FROM src`. Corresponding types must be safely castable by default. Set `export_merge_tree_part_allow_lossy_cast = 1` to permit lossy casts. +2. **Compatible partitioning** - for destinations other than data lakes, the source and destination `PARTITION BY` expressions must be identical. For Apache Iceberg destinations, the source partition key must be representable as an Iceberg partition spec and must match the destination partition fields and transforms. +3. **Matching partition key column positions and layouts** - it is not enough for the `PARTITION BY` expressions to be textually identical: every top-level column that provides a column or subcolumn used by the source table's partition key must have the same name at the same position in the destination table's schema. If such a column contains a named `Tuple`, its element names must also be declared in the same order. This comparison is recursive through nested tuples and through container types such as `Array` and `Map`. For example, `CREATE TABLE src (a Int32, b Int32) ... PARTITION BY a` and `CREATE TABLE dst (b Int32, a Int32) ... PARTITION BY a` both have the expression `PARTITION BY a`, but `a` is at position 0 in `src` and position 1 in `dst`. The export is rejected with a `BAD_ARGUMENTS` exception whose message includes `Cannot export to : partition key column 'a' is at position 0 in the source table, but the destination's column at that position is named 'b'`. + + This explicit name check only applies to partition key columns. A mismatch in the position of a non-partition-key column is **not** rejected by name - it is only caught if the source and destination types are not castable. If two non-partition-key columns happen to have swapped positions but compatible types, the export succeeds and silently writes values into the wrong destination column, so keep the intended column order rather than relying on type compatibility alone. + + For `PARTITION BY t.a`, this rule applies to the top-level owning column `t`. Exporting from `t Tuple(a Int32, b Int32)` to `t Tuple(b Int32, a Int32)` is rejected, even though `a` is accessed by name. Requiring a stable layout for every partition-key owner also protects positional expressions such as `tupleElement(t, 1)` from changing their meaning after conversion. + + The element-name check only applies when both the source and destination `Tuple` declare explicit names; an unnamed `Tuple` (e.g. `Tuple(Int32, Int32)`) is compared to the destination by element position and type only. For example, exporting from `t Tuple(Int32, Int32)` to `t Tuple(x Int32, y Int32)` is allowed as long as element types match positionally. + + The same rule applies when the named tuple is nested inside a container. For example, `arr Array(Tuple(a Int32, b Int32))` and `arr Array(Tuple(b Int32, a Int32))` are incompatible when `arr` provides an input to the partition key. Likewise, tuple layouts in both the key and value types of `Map` are checked recursively. + + In this case, the export throws a `BAD_ARGUMENTS` exception whose message includes `partition key column 't' has a different Tuple element layout in the source (Tuple(a Int32, b Int32)) and destination (Tuple(b Int32, a Int32)). Tuple element names must be declared in the same order in both tables`. + + For partition expressions containing functions, the check applies to their input columns. For example, `PARTITION BY (toYYYYMM(ts), category)` requires both `ts` and `category` to have the same names at the same top-level positions in both tables. In case a table function is used as the destination, the schema can be omitted and it will be inferred from the source table. diff --git a/docs/en/antalya/partition_export.md b/docs/en/antalya/partition_export.md index 687029b9adc6..7a9773ed293d 100644 --- a/docs/en/antalya/partition_export.md +++ b/docs/en/antalya/partition_export.md @@ -43,6 +43,14 @@ TO TABLE [destination_database.]destination_table - **`partition_id`**: The partition identifier to export (e.g., `'2020'`, `'2021'`) - **`destination_table`**: The target table for the export (typically an S3, Azure, or other object storage table) +## Requirements + +`EXPORT PARTITION` exports each part via the same mechanism as [`EXPORT PART`](/docs/en/antalya/part_export.md#requirements), so the source and destination tables must satisfy the same compatibility requirements, in particular: + +1. **Positionally compatible schemas** - source columns are matched to destination columns by position. Corresponding types must be safely castable unless `export_merge_tree_part_allow_lossy_cast = 1` is set. +2. **Compatible partitioning** - for destinations other than data lakes, the source and destination `PARTITION BY` expressions must be identical. For Apache Iceberg destinations, the source partition key must match the destination partition fields and transforms. +3. **Matching partition key column positions and layouts** - every top-level column that provides a column or subcolumn used by the source table's partition key must have the same name at the same position in the destination table's schema. Named `Tuple` elements within such a column must also be declared in the same order, including tuples nested inside `Array` or `Map`. This applies even if both tables' `PARTITION BY` expressions are textually identical. See [`EXPORT PART` requirements](/docs/en/antalya/part_export.md#requirements) for a worked example and the corresponding exception message. + ## Settings ### Server Settings @@ -251,5 +259,4 @@ WHERE source_table = 'rmt_table' AND destination_table = 's3_table'; ## Related Features -- [ALTER TABLE EXPORT PART](/docs/en/engines/table-engines/mergetree-family/part_export.md) - Export individual parts (non-replicated) - +- [ALTER TABLE EXPORT PART](/docs/en/antalya/part_export.md) - Export individual parts (non-replicated) diff --git a/src/Storages/MergeTree/ExportPartitionUtils.cpp b/src/Storages/MergeTree/ExportPartitionUtils.cpp index 3f85ebc0e1fb..f3c4aadfb31f 100644 --- a/src/Storages/MergeTree/ExportPartitionUtils.cpp +++ b/src/Storages/MergeTree/ExportPartitionUtils.cpp @@ -9,13 +9,21 @@ #include #include #include +#include #include #include #include +#include +#include +#include +#include +#include #include +#include #include #include #include +#include #if USE_AVRO #include @@ -634,6 +642,106 @@ namespace ExportPartitionUtils } #endif + namespace + { + bool haveSameTupleElementLayout(const DataTypePtr & source_type, const DataTypePtr & destination_type) + { + const auto source_type_unwrapped = removeNullable(removeLowCardinality(source_type)); + const auto destination_type_unwrapped = removeNullable(removeLowCardinality(destination_type)); + + const auto * source_tuple = checkAndGetDataType(source_type_unwrapped.get()); + const auto * destination_tuple = checkAndGetDataType(destination_type_unwrapped.get()); + if (source_tuple || destination_tuple) + { + if (!source_tuple || !destination_tuple) + return false; + + if (source_tuple->hasExplicitNames() && destination_tuple->hasExplicitNames()) + { + if (source_tuple->getElementNames() != destination_tuple->getElementNames()) + return false; + } + else if (source_tuple->getElements().size() != destination_tuple->getElements().size()) + return false; + + const auto & source_elements = source_tuple->getElements(); + const auto & destination_elements = destination_tuple->getElements(); + for (size_t i = 0; i < source_elements.size(); ++i) + if (!haveSameTupleElementLayout(source_elements[i], destination_elements[i])) + return false; + + return true; + } + + const auto * source_array = checkAndGetDataType(source_type_unwrapped.get()); + const auto * destination_array = checkAndGetDataType(destination_type_unwrapped.get()); + if (source_array || destination_array) + { + if (!source_array || !destination_array) + return false; + + return haveSameTupleElementLayout(source_array->getNestedType(), destination_array->getNestedType()); + } + + const auto * source_map = checkAndGetDataType(source_type_unwrapped.get()); + const auto * destination_map = checkAndGetDataType(destination_type_unwrapped.get()); + if (source_map || destination_map) + { + if (!source_map || !destination_map) + return false; + + return haveSameTupleElementLayout(source_map->getKeyType(), destination_map->getKeyType()) + && haveSameTupleElementLayout(source_map->getValueType(), destination_map->getValueType()); + } + + return true; + } + + void verifyPartitionKeyColumn( + const ColumnWithTypeAndName & source_column, + const ColumnWithTypeAndName & destination_column, + size_t position, + const StorageID & destination_storage_id) + { + if (source_column.name != destination_column.name) + throw Exception( + ErrorCodes::BAD_ARGUMENTS, + "Cannot export to {}: partition key column '{}' is at position {} in the source " + "table, but the destination's column at that position is named '{}'. EXPORT " + "PART/PARTITION matches columns by position, so partition key columns must be " + "declared at the same position in both tables.", + destination_storage_id.getFullTableName(), + source_column.name, + position, + destination_column.name); + + if (!haveSameTupleElementLayout(source_column.type, destination_column.type)) + throw Exception( + ErrorCodes::BAD_ARGUMENTS, + "Cannot export to {}: partition key column '{}' has a different Tuple element " + "layout in the source ({}) and destination ({}). Tuple element names must be " + "declared in the same order in both tables.", + destination_storage_id.getFullTableName(), + source_column.name, + source_column.type->getName(), + destination_column.type->getName()); + } + } + + void assertPartitionKeyASTAreEqual( + const StorageMetadataPtr & source_metadata, + const StorageMetadataPtr & destination_metadata) + { + constexpr auto query_to_string = [] (const ASTPtr & ast) + { + return ast ? ast->formatWithSecretsOneLine() : ""; + }; + + if (query_to_string(source_metadata->getPartitionKeyAST()) != query_to_string(destination_metadata->getPartitionKeyAST())) + throw Exception(ErrorCodes::BAD_ARGUMENTS, + "Cannot export partition: source and destination tables have different `PARTITION BY` expressions"); + } + void verifyExportSchemaCastable( const StorageMetadataPtr & source_metadata, const StorageMetadataPtr & destination_metadata, @@ -657,15 +765,33 @@ namespace ExportPartitionUtils ActionsDAG::MatchColumnsMode::Position, context); - /// Lossy casts may silently change values, so reject them unless the user opts in. - if (context->getSettingsRef()[Setting::export_merge_tree_part_allow_lossy_cast]) - return; + const auto & source_columns_description = source_metadata->getColumns(); + /// Collect the top-level columns that own columns or subcolumns required by `PARTITION BY`. + /// For example, both `PARTITION BY t.a` and `PARTITION BY (t.a, t.b)` add `t`. + std::unordered_set partition_key_owner_columns; + for (const auto & column_or_subcolumn_name : source_metadata->getColumnsRequiredForPartitionKey()) + { + auto resolved = source_columns_description.tryGetColumnOrSubcolumn( + GetColumnsOptions::All, column_or_subcolumn_name); + const auto & column_name = resolved ? resolved->getNameInStorage() : column_or_subcolumn_name; + partition_key_owner_columns.insert(column_name); + } + + const bool allow_lossy_cast = context->getSettingsRef()[Setting::export_merge_tree_part_allow_lossy_cast]; const size_t num_columns = std::min(source_columns.size(), destination_columns.size()); for (size_t i = 0; i < num_columns; ++i) { const auto & source_column = source_columns[i]; const auto & destination_column = destination_columns[i]; + + if (partition_key_owner_columns.contains(source_column.name)) + verifyPartitionKeyColumn(source_column, destination_column, i, destination_storage_id); + + /// Lossy casts may silently change values, so reject them unless the user opts in. + if (allow_lossy_cast) + continue; + if (!canBeSafelyCast(source_column.type, destination_column.type)) throw Exception(ErrorCodes::INCOMPATIBLE_COLUMNS, "Cannot export to {}: column '{}' requires a lossy cast from {} to {}, " diff --git a/src/Storages/MergeTree/ExportPartitionUtils.h b/src/Storages/MergeTree/ExportPartitionUtils.h index 0bb8acb9bda4..b58ecad9e7bb 100644 --- a/src/Storages/MergeTree/ExportPartitionUtils.h +++ b/src/Storages/MergeTree/ExportPartitionUtils.h @@ -89,6 +89,10 @@ namespace ExportPartitionUtils const std::string & exception_message, const LoggerPtr & log); + void assertPartitionKeyASTAreEqual( + const StorageMetadataPtr & source_metadata, + const StorageMetadataPtr & destination_metadata); + /// Validates that source columns can be exported into the destination with the /// same positional CAST matching as `INSERT INTO dest SELECT * FROM src`. Lossy /// casts are rejected unless `export_merge_tree_part_allow_lossy_cast` is set. diff --git a/src/Storages/MergeTree/MergeTreeData.cpp b/src/Storages/MergeTree/MergeTreeData.cpp index 0ed451b59655..0519730a5ee8 100644 --- a/src/Storages/MergeTree/MergeTreeData.cpp +++ b/src/Storages/MergeTree/MergeTreeData.cpp @@ -6726,11 +6726,6 @@ void MergeTreeData::exportPartToTable( if (!dest_storage->supportsImport(query_context)) throw Exception(ErrorCodes::NOT_IMPLEMENTED, "Destination storage {} does not support MergeTree parts or uses unsupported partitioning", dest_storage->getName()); - auto query_to_string = [] (const ASTPtr & ast) - { - return ast ? ast->formatWithSecretsOneLine() : ""; - }; - auto source_metadata_ptr = getInMemoryMetadataPtr(); auto destination_metadata_ptr = dest_storage->getInMemoryMetadataPtr(); @@ -6791,13 +6786,8 @@ void MergeTreeData::exportPartToTable( ExportPartitionUtils::verifyExportSchemaCastable( source_metadata_ptr, destination_metadata_ptr, dest_storage->getStorageID(), query_context); - /// Iceberg partition compatibility is checked above; here we only need the - /// partition-key ASTs to match (partition-column types follow the lossy-cast gate). if (!dest_storage->isDataLake()) - { - if (query_to_string(source_metadata_ptr->getPartitionKeyAST()) != query_to_string(destination_metadata_ptr->getPartitionKeyAST())) - throw Exception(ErrorCodes::BAD_ARGUMENTS, "Tables have different partition key"); - } + ExportPartitionUtils::assertPartitionKeyASTAreEqual(source_metadata_ptr, destination_metadata_ptr); auto part = getPartIfExists(part_name, {MergeTreeDataPartState::Active, MergeTreeDataPartState::Outdated}); diff --git a/src/Storages/StorageReplicatedMergeTree.cpp b/src/Storages/StorageReplicatedMergeTree.cpp index 5677a5e111a8..7b7533d26d9b 100644 --- a/src/Storages/StorageReplicatedMergeTree.cpp +++ b/src/Storages/StorageReplicatedMergeTree.cpp @@ -8408,11 +8408,6 @@ void StorageReplicatedMergeTree::exportPartitionToTable(const PartitionCommand & if (!dest_storage->supportsImport(query_context)) throw Exception(ErrorCodes::NOT_IMPLEMENTED, "Destination storage {} does not support MergeTree parts or uses unsupported partitioning", dest_storage->getName()); - auto query_to_string = [] (const ASTPtr & ast) - { - return ast ? ast->formatWithSecretsOneLine() : ""; - }; - auto src_snapshot = getInMemoryMetadataPtr(); auto destination_snapshot = dest_storage->getInMemoryMetadataPtr(); @@ -8420,13 +8415,8 @@ void StorageReplicatedMergeTree::exportPartitionToTable(const PartitionCommand & ExportPartitionUtils::verifyExportSchemaCastable( src_snapshot, destination_snapshot, dest_storage->getStorageID(), query_context); - /// Iceberg partition compatibility is checked below; here we only need the - /// partition-key ASTs to match (partition-column types follow the lossy-cast gate). if (!dest_storage->isDataLake()) - { - if (query_to_string(src_snapshot->getPartitionKeyAST()) != query_to_string(destination_snapshot->getPartitionKeyAST())) - throw Exception(ErrorCodes::BAD_ARGUMENTS, "Tables have different partition key"); - } + ExportPartitionUtils::assertPartitionKeyASTAreEqual(src_snapshot, destination_snapshot); zkutil::ZooKeeperPtr zookeeper = getZooKeeperAndAssertNotReadonly(); diff --git a/tests/integration/test_export_merge_tree_part_to_iceberg/test.py b/tests/integration/test_export_merge_tree_part_to_iceberg/test.py index 486adf1f2b17..ceafc0a120fd 100644 --- a/tests/integration/test_export_merge_tree_part_to_iceberg/test.py +++ b/tests/integration/test_export_merge_tree_part_to_iceberg/test.py @@ -13,10 +13,14 @@ test_export_part_with_year_transform_partition – toYearNumSinceEpoch() partition expression test_export_part_with_bucket_partition – icebergBucket(N, col) partition expression test_export_part_partition_key_mismatch_is_rejected – mismatched partition spec rejected synchronously + test_export_part_multi_column_partition_key_success – composite (a, b, c) partition key round-trips + test_export_part_partition_key_mismatch_variants_are_rejected (parametrized) – partition key column reordering, + cardinality mismatches, and transform-expression reordering between src/dst are all rejected synchronously """ import logging import time +from typing import NamedTuple import pytest @@ -462,6 +466,145 @@ def test_export_part_partition_key_mismatch_is_rejected(cluster): node.query(f"DROP TABLE IF EXISTS {iceberg}") +class RejectedPartExportCase(NamedTuple): + src_columns: str + src_partition_by: str + dst_columns: str + dst_partition_by: str + insert_values: str + error_substrings: tuple = () + + +REJECTED_PART_EXPORT_CASES = [ + pytest.param( + RejectedPartExportCase( + src_columns="a Int32, b Int32", + src_partition_by="a", + dst_columns="b Int32, a Int32", + dst_partition_by="a", + insert_values="(1, 1), (1, 2)", + error_substrings=("partition key column",), + ), + id="same_partition_key_different_column_order_single_column", + ), + pytest.param( + RejectedPartExportCase( + src_columns="a Int32, b Int32, c Int32, val String", + src_partition_by="(a, b, c)", + dst_columns="c Int32, b Int32, a Int32, val String", + dst_partition_by="(a, b, c)", + insert_values="(1, 1, 1, 'x'), (1, 1, 1, 'y')", + error_substrings=("partition key column",), + ), + id="same_partition_key_different_column_order_multi_column", + ), + pytest.param( + RejectedPartExportCase( + src_columns="a Int32, b Int32, c Int32, val String", + src_partition_by="(a, b, c)", + dst_columns="a Int32, b Int32, c Int32, val String", + dst_partition_by="(c, b, a)", + insert_values="(1, 2, 3, 'x')", + error_substrings=("partition field 0 mismatch",), + ), + id="multi_column_partition_key_order_mismatch", + ), + pytest.param( + RejectedPartExportCase( + src_columns="a Int32, b Int32, c Int32, val String", + src_partition_by="(a, b, c)", + dst_columns="a Int32, b Int32, c Int32, val String", + dst_partition_by="(a, b)", + insert_values="(1, 2, 3, 'x')", + error_substrings=("partition scheme mismatch",), + ), + id="multi_column_partition_key_fewer_in_destination", + ), + pytest.param( + RejectedPartExportCase( + src_columns="a Int32, b Int32, c Int32, val String", + src_partition_by="(a, b)", + dst_columns="a Int32, b Int32, c Int32, val String", + dst_partition_by="(a, b, c)", + insert_values="(1, 2, 3, 'x')", + error_substrings=("partition scheme mismatch",), + ), + id="multi_column_partition_key_more_in_destination", + ), + pytest.param( + RejectedPartExportCase( + src_columns="other_id Int64, user_id Int64", + src_partition_by="icebergBucket(8, user_id)", + dst_columns="user_id Int64, other_id Int64", + dst_partition_by="icebergBucket(8, user_id)", + insert_values="(1, 42)", + error_substrings=("partition key column",), + ), + id="transform_partition_key_different_column_order", + ), +] + + +@pytest.mark.parametrize("case", REJECTED_PART_EXPORT_CASES) +def test_export_part_partition_key_mismatch_variants_are_rejected(cluster, case): + node = cluster.instances["node1"] + sfx = unique_suffix() + mt = f"mt_rejected_{sfx}" + iceberg = f"iceberg_rejected_{sfx}" + + make_mt(node, mt, case.src_columns, case.src_partition_by) + make_iceberg_s3(node, iceberg, case.dst_columns, case.dst_partition_by) + + node.query(f"INSERT INTO {mt} VALUES {case.insert_values}") + + pid = first_partition_id(node, mt) + part = get_part(node, mt, pid) + + error = node.query_and_get_error( + f"ALTER TABLE {mt} EXPORT PART '{part}' TO TABLE {iceberg} " + f"SETTINGS allow_experimental_export_merge_tree_part = 1, " + f"allow_experimental_insert_into_iceberg = 1" + ) + assert "BAD_ARGUMENTS" in error, f"Expected BAD_ARGUMENTS, got: {error!r}" + for substring in case.error_substrings: + assert substring in error, f"Expected {substring!r} in error, got: {error!r}" + + count = int(node.query(f"SELECT count() FROM {iceberg}").strip()) + assert count == 0, f"Expected 0 rows in Iceberg table after rejected export, got {count}" + + node.query(f"DROP TABLE IF EXISTS {mt} SYNC") + node.query(f"DROP TABLE IF EXISTS {iceberg}") + + +def test_export_part_multi_column_partition_key_success(cluster): + node = cluster.instances["node1"] + sfx = unique_suffix() + mt = f"mt_multi_pkey_ok_{sfx}" + iceberg = f"iceberg_multi_pkey_ok_{sfx}" + + cols = "a Int32, b Int32, c Int32, val String" + make_mt(node, mt, cols, "(a, b, c)") + make_iceberg_s3(node, iceberg, cols, "(a, b, c)") + + node.query(f"INSERT INTO {mt} VALUES (1, 2, 3, 'x'), (1, 2, 3, 'y')") + + pid = first_partition_id(node, mt) + part = get_part(node, mt, pid) + export_part(node, mt, part, iceberg) + wait_for_export_part(node, mt, part) + + count = int(node.query(f"SELECT count() FROM {iceberg}").strip()) + assert count == 2, f"Expected 2 rows in Iceberg table after export, got {count}" + + result = node.query(f"SELECT a, b, c, val FROM {iceberg} ORDER BY val").strip() + assert result == "1\t2\t3\tx\n1\t2\t3\ty", f"Unexpected exported data:\n{result}" + + assert_part_log(node, mt, part) + + node.query(f"DROP TABLE IF EXISTS {mt} SYNC") + node.query(f"DROP TABLE IF EXISTS {iceberg}") + + def test_export_part_with_bucket_partition(cluster): """ Export a part from a MergeTree table partitioned by icebergBucket(8, user_id) @@ -753,3 +896,47 @@ def test_export_part_runtime_cast_failure_propagates_async(cluster): node.query(f"DROP TABLE IF EXISTS {mt} SYNC") node.query(f"DROP TABLE IF EXISTS {iceberg}") + + +def test_export_part_tuple_subcolumn_partition_key_iceberg_rejected(cluster): + node = cluster.instances["node1"] + sfx = unique_suffix() + mt = f"mt_tuple_subcol_{sfx}" + iceberg = f"iceberg_tuple_subcol_{sfx}" + iceberg_partitioned = f"iceberg_tuple_subcol_part_{sfx}" + + create_error = node.query_and_get_error( + f"CREATE TABLE {iceberg_partitioned} (t Tuple(b Int32, a Int32), val String) " + f"ENGINE = IcebergS3('http://minio1:9001/root/data/{iceberg_partitioned}/', 'minio', 'ClickHouse_Minio_P@ssw0rd') " + f"PARTITION BY t.a" + ) + assert "Unknown field to partition" in create_error, ( + f"Expected Iceberg to reject the tuple subcolumn partition key at CREATE time, " + f"got: {create_error!r}" + ) + + make_mt(node, mt, "t Tuple(a Int32, b Int32), val String", "t.a") + make_iceberg_s3(node, iceberg, "t Tuple(b Int32, a Int32), val String", "val") + + node.query(f"INSERT INTO {mt} VALUES ((1, 99), 'x')") + + part = node.query( + f"SELECT name FROM system.parts WHERE database = currentDatabase() " + f"AND table = '{mt}' AND active ORDER BY name LIMIT 1" + ).strip() + + export_error = node.query_and_get_error( + f"ALTER TABLE {mt} EXPORT PART '{part}' TO TABLE {iceberg} " + f"SETTINGS allow_experimental_export_merge_tree_part = 1, " + f"allow_experimental_insert_into_iceberg = 1" + ) + assert "Unknown field to partition" in export_error, ( + f"Expected export validation to reject the tuple subcolumn partition key of {mt}, " + f"got: {export_error!r}" + ) + + count = int(node.query(f"SELECT count() FROM {iceberg}").strip()) + assert count == 0, f"Expected 0 rows in Iceberg table after rejected export, got {count}" + + node.query(f"DROP TABLE IF EXISTS {mt} SYNC") + node.query(f"DROP TABLE IF EXISTS {iceberg}") diff --git a/tests/integration/test_export_merge_tree_part_to_object_storage/test.py b/tests/integration/test_export_merge_tree_part_to_object_storage/test.py index b8c15c26275f..107fc755c269 100644 --- a/tests/integration/test_export_merge_tree_part_to_object_storage/test.py +++ b/tests/integration/test_export_merge_tree_part_to_object_storage/test.py @@ -1,6 +1,7 @@ import logging import time import uuid +from typing import NamedTuple import pytest @@ -312,3 +313,636 @@ def test_pending_patch_parts_skip_before_export(cluster): assert "1\n2\n3" in result, "Export should contain original data before patch" node.query(f"DROP TABLE {mt_table}") + + +class RejectedPartExportCase(NamedTuple): + src_columns: str + src_partition_by: str + dst_columns: str + dst_partition_by: str + insert_values: str + error_substrings: tuple = () + partition_strategy: str = "hive" + + +REJECTED_PART_EXPORT_CASES = [ + pytest.param( + RejectedPartExportCase( + src_columns="a Int32, b Int32", + src_partition_by="a", + dst_columns="b Int32, a Int32", + dst_partition_by="a", + insert_values="(1, 1), (1, 2)", + error_substrings=( + "partition key column 'a' is at position 0 in the source table", + ), + ), + id="same_partition_key_different_column_order_single_column", + ), + pytest.param( + RejectedPartExportCase( + src_columns="a Int32, b Int32, c Int32, val String", + src_partition_by="(a, b, c)", + dst_columns="c Int32, b Int32, a Int32, val String", + dst_partition_by="(a, b, c)", + insert_values="(1, 1, 1, 'x'), (1, 1, 1, 'y')", + error_substrings=( + "partition key column 'a' is at position 0 in the source table", + ), + ), + id="same_partition_key_different_column_order_multi_column", + ), + pytest.param( + RejectedPartExportCase( + src_columns="a Int32, b Int32, c Int32, val String", + src_partition_by="(a, b, c)", + dst_columns="a Int32, b Int32, c Int32, val String", + dst_partition_by="(c, b, a)", + insert_values="(1, 2, 3, 'x')", + error_substrings=( + "source and destination tables have different `PARTITION BY` expressions", + ), + ), + id="multi_column_partition_key_order_mismatch", + ), + pytest.param( + RejectedPartExportCase( + src_columns="a Int32, b Int32, c Int32, val String", + src_partition_by="(a, b, c)", + dst_columns="a Int32, b Int32, c Int32, val String", + dst_partition_by="(a, b)", + insert_values="(1, 2, 3, 'x')", + error_substrings=( + "source and destination tables have different `PARTITION BY` expressions", + ), + ), + id="multi_column_partition_key_fewer_in_destination", + ), + pytest.param( + RejectedPartExportCase( + src_columns="a Int32, b Int32, c Int32, val String", + src_partition_by="(a, b)", + dst_columns="a Int32, b Int32, c Int32, val String", + dst_partition_by="(a, b, c)", + insert_values="(1, 2, 3, 'x')", + error_substrings=( + "source and destination tables have different `PARTITION BY` expressions", + ), + ), + id="multi_column_partition_key_more_in_destination", + ), + pytest.param( + RejectedPartExportCase( + src_columns="ts DateTime, category String, decoy DateTime, val String", + src_partition_by="(toYYYYMM(ts), category)", + dst_columns="decoy DateTime, category String, ts DateTime, val String", + dst_partition_by="(toYYYYMM(ts), category)", + insert_values=( + "('2024-03-05 15:00:00', 'category', " + "'2024-03-06 15:00:00', 'x')" + ), + error_substrings=( + "partition key column 'ts' is at position 0 in the source table", + ), + partition_strategy="wildcard", + ), + id="function_and_column_partition_key_owner_reordered", + ), + pytest.param( + RejectedPartExportCase( + src_columns=( + "t Tuple(ts DateTime, value Int32), category String, " + "decoy Tuple(ts DateTime, value Int32), val String" + ), + src_partition_by="(toYYYYMM(t.ts), category)", + dst_columns=( + "decoy Tuple(ts DateTime, value Int32), category String, " + "t Tuple(ts DateTime, value Int32), val String" + ), + dst_partition_by="(toYYYYMM(t.ts), category)", + insert_values=( + "(('2024-03-05 15:00:00', 1), 'category', " + "('2024-03-06 15:00:00', 2), 'x')" + ), + error_substrings=( + "partition key column 't' is at position 0 in the source table", + ), + partition_strategy="wildcard", + ), + id="function_over_subcolumn_partition_key_owner_reordered", + ), +] + + +@pytest.mark.parametrize("case", REJECTED_PART_EXPORT_CASES) +def test_export_part_partition_key_mismatch_variants_are_rejected(cluster, case): + skip_if_remote_database_disk_enabled(cluster) + node = cluster.instances["node1"] + + postfix = str(uuid.uuid4()).replace("-", "_") + mt_table = f"rejected_mt_table_{postfix}" + s3_table = f"rejected_s3_table_{postfix}" + + node.query(f""" + CREATE TABLE {mt_table} ({case.src_columns}) + ENGINE = MergeTree() + PARTITION BY {case.src_partition_by} + ORDER BY tuple() + SETTINGS enable_block_number_column = 1, enable_block_offset_column = 1 + """) + + filename = ( + f"{s3_table}/{{_partition_id}}/{{_file}}" + if case.partition_strategy == "wildcard" + else s3_table + ) + node.query(f""" + CREATE TABLE {s3_table} ({case.dst_columns}) + ENGINE = S3(s3_conn, filename='{filename}', format=Parquet, partition_strategy='{case.partition_strategy}') + PARTITION BY {case.dst_partition_by} + """) + + node.query(f"INSERT INTO {mt_table} VALUES {case.insert_values}") + + part_name = node.query( + f"SELECT name FROM system.parts WHERE database = currentDatabase() " + f"AND table = '{mt_table}' AND active ORDER BY name LIMIT 1" + ).strip() + + error = node.query_and_get_error(f"ALTER TABLE {mt_table} EXPORT PART '{part_name}' TO TABLE {s3_table}") + assert "BAD_ARGUMENTS" in error, f"Expected BAD_ARGUMENTS, got: {error}" + for substring in case.error_substrings: + assert substring in error, f"Expected {substring!r} in error, got: {error}" + + if case.partition_strategy == "hive": + count = int(node.query(f"SELECT count() FROM {s3_table}").strip()) + assert count == 0, ( + f"Expected 0 rows in destination after rejected export, got {count}" + ) + + +def test_export_part_multi_column_partition_key_success(cluster): + skip_if_remote_database_disk_enabled(cluster) + node = cluster.instances["node1"] + + postfix = str(uuid.uuid4()).replace("-", "_") + mt_table = f"multi_pkey_ok_mt_table_{postfix}" + s3_table = f"multi_pkey_ok_s3_table_{postfix}" + + node.query(f""" + CREATE TABLE {mt_table} (a Int32, b Int32, c Int32, val String) + ENGINE = MergeTree() + PARTITION BY (a, b, c) + ORDER BY tuple() + SETTINGS enable_block_number_column = 1, enable_block_offset_column = 1 + """) + + node.query(f""" + CREATE TABLE {s3_table} (a Int32, b Int32, c Int32, val String) + ENGINE = S3(s3_conn, filename='{s3_table}', format=Parquet, partition_strategy='hive') + PARTITION BY (a, b, c) + """) + + node.query(f"INSERT INTO {mt_table} VALUES (1, 2, 3, 'x'), (1, 2, 3, 'y')") + + part_name = node.query( + f"SELECT name FROM system.parts WHERE database = currentDatabase() " + f"AND table = '{mt_table}' AND active ORDER BY name LIMIT 1" + ).strip() + + node.query(f"ALTER TABLE {mt_table} EXPORT PART '{part_name}' TO TABLE {s3_table}") + + time.sleep(5) + + count = int(node.query(f"SELECT count() FROM {s3_table}").strip()) + assert count == 2, f"Expected 2 rows in destination after export, got {count}" + + result = node.query(f"SELECT a, b, c, val FROM {s3_table} ORDER BY val").strip() + assert result == "1\t2\t3\tx\n1\t2\t3\ty", f"Unexpected exported data:\n{result}" + + +@pytest.mark.parametrize( + "owner_name, source_type, destination_type, partition_by, insert_value", + [ + pytest.param( + "t", + "Tuple(a Int32, b Int32)", + "Tuple(b Int32, a Int32)", + "t.a", + "(1, 99)", + id="named_subcolumn", + ), + pytest.param( + "t", + "Tuple(a Int32, b Int32)", + "Tuple(b Int32, a Int32)", + "tupleElement(t, 1)", + "(1, 99)", + id="positional_tuple_element", + ), + pytest.param( + "arr", + "Array(Tuple(a Int32, b Int32))", + "Array(Tuple(b Int32, a Int32))", + "tupleElement(arr[1], 'a')", + "[(1, 99)]", + id="tuple_nested_in_array", + ), + pytest.param( + "m", + "Map(String, Tuple(a Int32, b Int32))", + "Map(String, Tuple(b Int32, a Int32))", + "tupleElement(m['key'], 'a')", + "map('key', (1, 99))", + id="tuple_nested_in_map_value", + ), + ], +) +def test_export_part_tuple_fields_reordered_for_partition_key_is_rejected( + cluster, + owner_name, + source_type, + destination_type, + partition_by, + insert_value, +): + skip_if_remote_database_disk_enabled(cluster) + node = cluster.instances["node1"] + + postfix = str(uuid.uuid4()).replace("-", "_") + mt_table = f"reordered_tuple_mt_table_{postfix}" + s3_table = f"reordered_tuple_s3_table_{postfix}" + + node.query(f""" + CREATE TABLE {mt_table} ({owner_name} {source_type}, val String) + ENGINE = MergeTree() + PARTITION BY {partition_by} + ORDER BY tuple() + SETTINGS enable_block_number_column = 1, enable_block_offset_column = 1 + """) + + node.query(f""" + CREATE TABLE {s3_table} ({owner_name} {destination_type}, val String) + ENGINE = S3(s3_conn, filename='{s3_table}/{{_partition_id}}/{{_file}}', format=Parquet, partition_strategy='wildcard') + PARTITION BY {partition_by} + """) + + node.query(f"INSERT INTO {mt_table} VALUES ({insert_value}, 'x')") + + part_name = node.query( + f"SELECT name FROM system.parts WHERE database = currentDatabase() " + f"AND table = '{mt_table}' AND active ORDER BY name LIMIT 1" + ).strip() + + error = node.query_and_get_error( + f"ALTER TABLE {mt_table} EXPORT PART '{part_name}' TO TABLE {s3_table}" + ) + assert "BAD_ARGUMENTS" in error and "different Tuple element layout" in error, ( + f"Expected export to reject reordered named `Tuple` fields used by " + f"`PARTITION BY {partition_by}`, got: {error!r}" + ) + + +def test_export_part_unnamed_tuple_partition_key_owner_matching_named_destination_is_allowed(cluster): + skip_if_remote_database_disk_enabled(cluster) + node = cluster.instances["node1"] + + postfix = str(uuid.uuid4()).replace("-", "_") + mt_table = f"unnamed_tuple_ok_mt_table_{postfix}" + s3_table = f"unnamed_tuple_ok_s3_table_{postfix}" + + node.query(f""" + CREATE TABLE {mt_table} (t Tuple(Int32, Int32), val String) + ENGINE = MergeTree() + PARTITION BY tupleElement(t, 1) + ORDER BY tuple() + SETTINGS enable_block_number_column = 1, enable_block_offset_column = 1 + """) + + node.query(f""" + CREATE TABLE {s3_table} (t Tuple(x Int32, y Int32), val String) + ENGINE = S3(s3_conn, filename='{s3_table}/{{_partition_id}}/{{_file}}', format=Parquet, partition_strategy='wildcard') + PARTITION BY tupleElement(t, 1) + """) + + node.query(f"INSERT INTO {mt_table} VALUES ((1, 99), 'x')") + + part_name = node.query( + f"SELECT name FROM system.parts WHERE database = currentDatabase() " + f"AND table = '{mt_table}' AND active ORDER BY name LIMIT 1" + ).strip() + + node.query(f"ALTER TABLE {mt_table} EXPORT PART '{part_name}' TO TABLE {s3_table}") + + +def test_export_part_subcolumn_partition_key_different_subcolumn_is_rejected(cluster): + skip_if_remote_database_disk_enabled(cluster) + node = cluster.instances["node1"] + + postfix = str(uuid.uuid4()).replace("-", "_") + mt_table = f"subcol_diff_subcol_mt_table_{postfix}" + s3_table = f"subcol_diff_subcol_s3_table_{postfix}" + + node.query(f""" + CREATE TABLE {mt_table} (a Tuple(b Int32, c Int32), val String) + ENGINE = MergeTree() + PARTITION BY a.b + ORDER BY tuple() + SETTINGS enable_block_number_column = 1, enable_block_offset_column = 1 + """) + + node.query(f""" + CREATE TABLE {s3_table} (a Tuple(b Int32, c Int32), val String) + ENGINE = S3(s3_conn, filename='{s3_table}/{{_partition_id}}/{{_file}}', format=Parquet, partition_strategy='wildcard') + PARTITION BY a.c + """) + + node.query(f"INSERT INTO {mt_table} VALUES ((1, 2), 'x')") + + part_name = node.query( + f"SELECT name FROM system.parts WHERE database = currentDatabase() " + f"AND table = '{mt_table}' AND active ORDER BY name LIMIT 1" + ).strip() + + error = node.query_and_get_error( + f"ALTER TABLE {mt_table} EXPORT PART '{part_name}' TO TABLE {s3_table}" + ) + assert ( + "BAD_ARGUMENTS" in error + and "source and destination tables have different `PARTITION BY` expressions" + in error + ), ( + f"Both tables declare `a` as the same Tuple(b Int32, c Int32) (so the column-cast " + f"check passes and the owner-name-only `partition_key_owner_columns` contains " + f"only `a`, so `verifyExportSchemaCastable` cannot distinguish `a.b` from " + f"`a.c`), but the source " + f"partitions by `a.b` and the destination by `a.c` — a genuinely different " + f"partition key that must be caught by the `PARTITION BY` AST comparison; " + f"got: {error!r}" + ) + + +def test_export_part_tuple_subcolumn_partition_key_owner_column_reordered_is_rejected(cluster): + skip_if_remote_database_disk_enabled(cluster) + node = cluster.instances["node1"] + + postfix = str(uuid.uuid4()).replace("-", "_") + mt_table = f"tuple_subcol_owner_mt_table_{postfix}" + s3_table = f"tuple_subcol_owner_s3_table_{postfix}" + + node.query(f""" + CREATE TABLE {mt_table} (t Tuple(a Int32, b Int32), decoy Tuple(a Int32, b Int32), val String) + ENGINE = MergeTree() + PARTITION BY t.a + ORDER BY tuple() + SETTINGS enable_block_number_column = 1, enable_block_offset_column = 1 + """) + + node.query(f""" + CREATE TABLE {s3_table} (decoy Tuple(a Int32, b Int32), t Tuple(a Int32, b Int32), val String) + ENGINE = S3(s3_conn, filename='{s3_table}/{{_partition_id}}/{{_file}}', format=Parquet, partition_strategy='wildcard') + PARTITION BY t.a + """) + + node.query(f"INSERT INTO {mt_table} VALUES ((1, 100), (2, 200), 'x')") + + part_name = node.query( + f"SELECT name FROM system.parts WHERE database = currentDatabase() " + f"AND table = '{mt_table}' AND active ORDER BY name LIMIT 1" + ).strip() + + error = node.query_and_get_error( + f"ALTER TABLE {mt_table} EXPORT PART '{part_name}' TO TABLE {s3_table}" + ) + assert "BAD_ARGUMENTS" in error and "partition key column" in error, ( + f"Expected export to reject `t` and `decoy` swapping positions around the " + f"partition key column `t.a`, the same way a plain (non-tuple) partition key " + f"column position swap is rejected; got: {error!r}" + ) + + +def test_export_part_multiple_partition_key_subcolumns_with_same_owner_reordered_is_rejected(cluster): + skip_if_remote_database_disk_enabled(cluster) + node = cluster.instances["node1"] + + postfix = str(uuid.uuid4()).replace("-", "_") + mt_table = f"same_owner_subcolumns_mt_table_{postfix}" + s3_table = f"same_owner_subcolumns_s3_table_{postfix}" + + node.query(f""" + CREATE TABLE {mt_table} ( + t Tuple(a Int32, b Int32), + decoy Tuple(a Int32, b Int32), + val String + ) + ENGINE = MergeTree() + PARTITION BY (t.a, t.b) + ORDER BY tuple() + SETTINGS enable_block_number_column = 1, enable_block_offset_column = 1 + """) + + node.query(f""" + CREATE TABLE {s3_table} ( + decoy Tuple(a Int32, b Int32), + t Tuple(a Int32, b Int32), + val String + ) + ENGINE = S3(s3_conn, filename='{s3_table}/{{_partition_id}}/{{_file}}', format=Parquet, partition_strategy='wildcard') + PARTITION BY (t.a, t.b) + """) + + node.query(f"INSERT INTO {mt_table} VALUES ((1, 10), (2, 20), 'x')") + + part_name = node.query( + f"SELECT name FROM system.parts WHERE database = currentDatabase() " + f"AND table = '{mt_table}' AND active ORDER BY name LIMIT 1" + ).strip() + + error = node.query_and_get_error( + f"ALTER TABLE {mt_table} EXPORT PART '{part_name}' TO TABLE {s3_table}" + ) + assert "BAD_ARGUMENTS" in error and "partition key column 't'" in error, ( + f"Expected both `t.a` and `t.b` to resolve to the same top-level owner `t` " + f"and reject swapping `t` with `decoy`; got: {error!r}" + ) + + +def test_export_part_multi_level_subcolumn_partition_key_owner_reordered_is_rejected(cluster): + skip_if_remote_database_disk_enabled(cluster) + node = cluster.instances["node1"] + + postfix = str(uuid.uuid4()).replace("-", "_") + mt_table = f"nested_subcol_owner_mt_table_{postfix}" + s3_table = f"nested_subcol_owner_s3_table_{postfix}" + + node.query(f""" + CREATE TABLE {mt_table} ( + t Tuple(x Tuple(a Int32, b Int32), c Int32), + decoy Tuple(x Tuple(a Int32, b Int32), c Int32), + val String + ) + ENGINE = MergeTree() + PARTITION BY t.x.a + ORDER BY tuple() + SETTINGS enable_block_number_column = 1, enable_block_offset_column = 1 + """) + + node.query(f""" + CREATE TABLE {s3_table} ( + decoy Tuple(x Tuple(a Int32, b Int32), c Int32), + t Tuple(x Tuple(a Int32, b Int32), c Int32), + val String + ) + ENGINE = S3(s3_conn, filename='{s3_table}/{{_partition_id}}/{{_file}}', format=Parquet, partition_strategy='wildcard') + PARTITION BY t.x.a + """) + + node.query(f"INSERT INTO {mt_table} VALUES ((((1, 100), 1000)), (((2, 200), 2000)), 'x')") + + part_name = node.query( + f"SELECT name FROM system.parts WHERE database = currentDatabase() " + f"AND table = '{mt_table}' AND active ORDER BY name LIMIT 1" + ).strip() + + error = node.query_and_get_error( + f"ALTER TABLE {mt_table} EXPORT PART '{part_name}' TO TABLE {s3_table}" + ) + assert "BAD_ARGUMENTS" in error and "partition key column" in error, ( + f"Expected export to reject `t` and `decoy` swapping positions around the " + f"two-level-deep partition key column `t.x.a`. This only works if " + f"`getNameInStorage` resolves all the way to the top-level column `t`, not to " + f"the intermediate level `t.x`; got: {error!r}" + ) + + +def test_export_part_multiple_subcolumn_partition_keys_owner_reordered_is_rejected(cluster): + skip_if_remote_database_disk_enabled(cluster) + node = cluster.instances["node1"] + + postfix = str(uuid.uuid4()).replace("-", "_") + mt_table = f"multi_subcol_key_mt_table_{postfix}" + s3_table = f"multi_subcol_key_s3_table_{postfix}" + + node.query(f""" + CREATE TABLE {mt_table} ( + t Tuple(a Int32, x Int32), + u Tuple(b Int32, y Int32), + decoy Tuple(b Int32, y Int32), + val String + ) + ENGINE = MergeTree() + PARTITION BY (t.a, u.b) + ORDER BY tuple() + SETTINGS enable_block_number_column = 1, enable_block_offset_column = 1 + """) + + node.query(f""" + CREATE TABLE {s3_table} ( + t Tuple(a Int32, x Int32), + decoy Tuple(b Int32, y Int32), + u Tuple(b Int32, y Int32), + val String + ) + ENGINE = S3(s3_conn, filename='{s3_table}/{{_partition_id}}/{{_file}}', format=Parquet, partition_strategy='wildcard') + PARTITION BY (t.a, u.b) + """) + + node.query(f"INSERT INTO {mt_table} VALUES ((1, 10), (2, 20), (3, 30), 'x')") + + part_name = node.query( + f"SELECT name FROM system.parts WHERE database = currentDatabase() " + f"AND table = '{mt_table}' AND active ORDER BY name LIMIT 1" + ).strip() + + error = node.query_and_get_error( + f"ALTER TABLE {mt_table} EXPORT PART '{part_name}' TO TABLE {s3_table}" + ) + assert "BAD_ARGUMENTS" in error and "partition key column 'u'" in error, ( + f"`t` (owner of key part `t.a`) stays at position 0 on both sides, so the guard " + f"must independently catch `u` (owner of key part `u.b`) swapping positions " + f"with `decoy` — a partition key with two subcolumn-owning columns must have " + f"both validated, not just the first one encountered; got: {error!r}" + ) + + +def test_export_part_mixed_flat_and_subcolumn_partition_key_flat_part_reordered_is_rejected(cluster): + skip_if_remote_database_disk_enabled(cluster) + node = cluster.instances["node1"] + + postfix = str(uuid.uuid4()).replace("-", "_") + mt_table = f"mixed_key_mt_table_{postfix}" + s3_table = f"mixed_key_s3_table_{postfix}" + + node.query(f""" + CREATE TABLE {mt_table} (a Int32, t Tuple(b Int32, c Int32), decoy Int32, val String) + ENGINE = MergeTree() + PARTITION BY (a, t.b) + ORDER BY tuple() + SETTINGS enable_block_number_column = 1, enable_block_offset_column = 1 + """) + + node.query(f""" + CREATE TABLE {s3_table} (decoy Int32, t Tuple(b Int32, c Int32), a Int32, val String) + ENGINE = S3(s3_conn, filename='{s3_table}/{{_partition_id}}/{{_file}}', format=Parquet, partition_strategy='wildcard') + PARTITION BY (a, t.b) + """) + + node.query(f"INSERT INTO {mt_table} VALUES (1, (2, 3), 4, 'x')") + + part_name = node.query( + f"SELECT name FROM system.parts WHERE database = currentDatabase() " + f"AND table = '{mt_table}' AND active ORDER BY name LIMIT 1" + ).strip() + + error = node.query_and_get_error( + f"ALTER TABLE {mt_table} EXPORT PART '{part_name}' TO TABLE {s3_table}" + ) + assert "BAD_ARGUMENTS" in error and "partition key column 'a'" in error, ( + f"`t` (owner of key part `t.b`) stays at position 1 on both sides, so the guard " + f"must independently catch the plain, non-tuple key part `a` swapping positions " + f"with `decoy` — the pre-existing flat-column check and the new subcolumn-owner " + f"resolution must both keep working when combined in one `PARTITION BY` " + f"expression; got: {error!r}" + ) + + +def test_export_part_subcolumn_partition_key_owner_reordered_rejected_even_with_allow_lossy_cast(cluster): + skip_if_remote_database_disk_enabled(cluster) + node = cluster.instances["node1"] + + postfix = str(uuid.uuid4()).replace("-", "_") + mt_table = f"lossy_owner_mt_table_{postfix}" + s3_table = f"lossy_owner_s3_table_{postfix}" + + node.query(f""" + CREATE TABLE {mt_table} (t Tuple(a Int32, b Int32), decoy Tuple(a Int32, b Int32), val String) + ENGINE = MergeTree() + PARTITION BY t.a + ORDER BY tuple() + SETTINGS enable_block_number_column = 1, enable_block_offset_column = 1 + """) + + node.query(f""" + CREATE TABLE {s3_table} (decoy Tuple(a Int32, b Int32), t Tuple(a Int32, b Int32), val String) + ENGINE = S3(s3_conn, filename='{s3_table}/{{_partition_id}}/{{_file}}', format=Parquet, partition_strategy='wildcard') + PARTITION BY t.a + """) + + node.query(f"INSERT INTO {mt_table} VALUES ((1, 100), (2, 200), 'x')") + + part_name = node.query( + f"SELECT name FROM system.parts WHERE database = currentDatabase() " + f"AND table = '{mt_table}' AND active ORDER BY name LIMIT 1" + ).strip() + + error = node.query_and_get_error( + f"ALTER TABLE {mt_table} EXPORT PART '{part_name}' TO TABLE {s3_table} " + f"SETTINGS export_merge_tree_part_allow_lossy_cast = 1" + ) + assert "BAD_ARGUMENTS" in error and "partition key column" in error, ( + f"The partition-key position/name guard is checked before the " + f"`allow_lossy_cast` early-continue in verifyExportSchemaCastable, so setting " + f"`export_merge_tree_part_allow_lossy_cast = 1` must not suppress the rejection " + f"of `t`/`decoy` swapping positions around the partition key column `t.a`; " + f"got: {error!r}" + ) diff --git a/tests/integration/test_export_replicated_mt_partition_to_iceberg/test.py b/tests/integration/test_export_replicated_mt_partition_to_iceberg/test.py index ad2deba8de19..53e18c3df66b 100644 --- a/tests/integration/test_export_replicated_mt_partition_to_iceberg/test.py +++ b/tests/integration/test_export_replicated_mt_partition_to_iceberg/test.py @@ -3,6 +3,7 @@ import logging import re import time +from typing import NamedTuple import pytest from avro.datafile import DataFileReader @@ -754,6 +755,7 @@ def check_accepted(mt, iceberg, description): f"ALTER TABLE {mt} EXPORT PARTITION ID '{pid}' TO TABLE {iceberg}", settings={"allow_insert_into_iceberg": 1}, ) + return pid # 1. Compound identity: (year, region) cols = "id Int64, year Int32, region String" @@ -761,7 +763,12 @@ def check_accepted(mt, iceberg, description): make_rmt(node, t, cols, "(year, region)") node.query(f"INSERT INTO {t} VALUES (1, 2023, 'EU')") make_iceberg_s3(node, i, cols, "(year, region)") - check_accepted(t, i, "compound identity (year, region)") + pid = check_accepted(t, i, "compound identity (year, region)") + wait_for_export_status(node, t, i, pid, "COMPLETED") + count = int(node.query(f"SELECT count() FROM {i}").strip()) + assert count == 1, f"[compound identity (year, region)] Expected 1 row in Iceberg table, got {count}" + result = node.query(f"SELECT id, year, region FROM {i}").strip() + assert result == "1\t2023\tEU", f"[compound identity (year, region)] Unexpected exported data:\n{result}" # 2. Year transform cols = "id Int64, event_date Date" @@ -837,6 +844,8 @@ def assert_rejected(mt, iceberg, description): node.query(f"INSERT INTO {t} VALUES (1, 2020, 'EU')") make_iceberg_s3(node, i, cols, "(region, year)") assert_rejected(t, i, "compound field order reversed") + count = int(node.query(f"SELECT count() FROM {i}").strip()) + assert count == 0, f"[compound field order reversed] Expected 0 rows in destination, got {count}" # 2. Transform mismatch: MergeTree year-transform, Iceberg identity on same Date col cols = "id Int64, event_date Date" @@ -869,6 +878,8 @@ def assert_rejected(mt, iceberg, description): node.query(f"INSERT INTO {t} VALUES (1, 2020, 'EU')") make_iceberg_s3(node, i, cols, "year") assert_rejected(t, i, "2-field MergeTree vs 1-field Iceberg") + count = int(node.query(f"SELECT count() FROM {i}").strip()) + assert count == 0, f"[2-field MergeTree vs 1-field Iceberg] Expected 0 rows in destination, got {count}" # 6. Unsupported MergeTree expression: intDiv(year, 100) is not an Iceberg transform cols = "id Int64, year Int32" @@ -1324,6 +1335,128 @@ def test_export_partition_with_renamed_destination_column(cluster): ) +class RejectedPartitionExportCase(NamedTuple): + src_columns: str + src_partition_by: str + dst_columns: str + dst_partition_by: str + insert_values: str + error_substrings: tuple = () + + +REJECTED_PARTITION_EXPORT_CASES = [ + pytest.param( + RejectedPartitionExportCase( + src_columns="a Int32, b Int32", + src_partition_by="a", + dst_columns="b Int32, a Int32", + dst_partition_by="a", + insert_values="(1, 1), (1, 2)", + error_substrings=("partition key column",), + ), + id="same_partition_key_different_column_order_single_column", + ), + pytest.param( + RejectedPartitionExportCase( + src_columns="a Int32, b Int32, c Int32, val String", + src_partition_by="(a, b, c)", + dst_columns="c Int32, b Int32, a Int32, val String", + dst_partition_by="(a, b, c)", + insert_values="(1, 1, 1, 'x'), (1, 1, 1, 'y')", + error_substrings=("partition key column",), + ), + id="same_partition_key_different_column_order_multi_column", + ), + pytest.param( + RejectedPartitionExportCase( + src_columns="a Int32, b Int32, c Int32, val String", + src_partition_by="(a, b)", + dst_columns="a Int32, b Int32, c Int32, val String", + dst_partition_by="(a, b, c)", + insert_values="(1, 2, 3, 'x')", + error_substrings=("partition scheme mismatch",), + ), + id="multi_column_partition_key_more_in_destination", + ), + pytest.param( + RejectedPartitionExportCase( + src_columns="other_id Int64, user_id Int64", + src_partition_by="icebergBucket(8, user_id)", + dst_columns="user_id Int64, other_id Int64", + dst_partition_by="icebergBucket(8, user_id)", + insert_values="(1, 42)", + error_substrings=("partition key column",), + ), + id="transform_partition_key_different_column_order", + ), +] + + +@pytest.mark.parametrize("case", REJECTED_PARTITION_EXPORT_CASES) +def test_export_partition_partition_key_mismatch_variants_are_rejected(cluster, case): + node = cluster.instances["replica1"] + + uid = unique_suffix() + mt_table = f"mt_rejected_{uid}" + iceberg_table = f"iceberg_rejected_{uid}" + + make_rmt(node, mt_table, case.src_columns, case.src_partition_by, replica_name="replica1") + make_iceberg_s3(node, iceberg_table, case.dst_columns, partition_by=case.dst_partition_by) + + node.query(f"INSERT INTO {mt_table} VALUES {case.insert_values}") + + pid = first_partition_id(node, mt_table) + error = node.query_and_get_error( + f"ALTER TABLE {mt_table} EXPORT PARTITION ID '{pid}' TO TABLE {iceberg_table}", + settings={"allow_insert_into_iceberg": 1}, + ) + assert "BAD_ARGUMENTS" in error, f"Expected BAD_ARGUMENTS, got: {error}" + for substring in case.error_substrings: + assert substring in error, f"Expected {substring!r} in error, got: {error}" + + error_all = node.query_and_get_error( + f"ALTER TABLE {mt_table} EXPORT PARTITION ALL TO TABLE {iceberg_table}", + settings={"allow_insert_into_iceberg": 1}, + ) + assert "BAD_ARGUMENTS" in error_all, f"Expected BAD_ARGUMENTS, got: {error_all}" + + count = int(node.query(f"SELECT count() FROM {iceberg_table}").strip()) + assert count == 0, f"Expected 0 rows in destination after rejected export, got {count}" + + +def test_export_partition_multi_column_partition_key_success_all(cluster): + node = cluster.instances["replica1"] + + uid = unique_suffix() + mt_table = f"mt_multi_pkey_ok_all_{uid}" + iceberg_table = f"iceberg_multi_pkey_ok_all_{uid}" + + cols = "a Int32, b Int32, c Int32, val String" + make_rmt(node, mt_table, cols, "(a, b, c)", replica_name="replica1") + make_iceberg_s3(node, iceberg_table, cols, partition_by="(a, b, c)") + + node.query(f"INSERT INTO {mt_table} VALUES (1, 2, 3, 'x'), (4, 5, 6, 'y')") + + partition_ids = node.query( + f"SELECT DISTINCT partition_id FROM system.parts WHERE database = currentDatabase() " + f"AND table = '{mt_table}' AND active ORDER BY partition_id" + ).strip().split("\n") + + node.query( + f"ALTER TABLE {mt_table} EXPORT PARTITION ALL TO TABLE {iceberg_table}", + settings={"allow_insert_into_iceberg": 1}, + ) + + for pid in partition_ids: + wait_for_export_status(node, mt_table, iceberg_table, pid, "COMPLETED") + + count = int(node.query(f"SELECT count() FROM {iceberg_table}").strip()) + assert count == 2, f"Expected 2 rows in destination after export, got {count}" + + result = node.query(f"SELECT a, b, c, val FROM {iceberg_table} ORDER BY val").strip() + assert result == "1\t2\t3\tx\n4\t5\t6\ty", f"Unexpected exported data:\n{result}" + + def test_export_partition_with_castable_widening(cluster): """A lossless widening of both a data column (id Int32 -> Int64) and the partition column (year Int32 -> Int64) round-trips.""" diff --git a/tests/integration/test_export_replicated_mt_partition_to_object_storage/test.py b/tests/integration/test_export_replicated_mt_partition_to_object_storage/test.py index 8d4589292e3c..aeb738496b86 100644 --- a/tests/integration/test_export_replicated_mt_partition_to_object_storage/test.py +++ b/tests/integration/test_export_replicated_mt_partition_to_object_storage/test.py @@ -1,6 +1,7 @@ import logging import time import uuid +from typing import NamedTuple import pytest @@ -1747,3 +1748,196 @@ def test_export_partition_all_failure_modes(cluster): f"ALTER TABLE {mt_table} EXPORT PARTITION ALL TO TABLE {s3_table}" f" SETTINGS export_merge_tree_partition_all_on_error = 'skip_conflicts'" ) + + +class RejectedPartitionExportCase(NamedTuple): + src_columns: str + src_partition_by: str + dst_columns: str + dst_partition_by: str + insert_values: str + error_substrings: tuple = () + + +REJECTED_PARTITION_EXPORT_CASES = [ + pytest.param( + RejectedPartitionExportCase( + src_columns="a Int32, b Int32", + src_partition_by="a", + dst_columns="b Int32, a Int32", + dst_partition_by="a", + insert_values="(1, 1), (1, 2)", + error_substrings=("partition key column",), + ), + id="same_partition_key_different_column_order_single_column", + ), + pytest.param( + RejectedPartitionExportCase( + src_columns="a Int32, b Int32, c Int32, val String", + src_partition_by="(a, b, c)", + dst_columns="c Int32, b Int32, a Int32, val String", + dst_partition_by="(a, b, c)", + insert_values="(1, 1, 1, 'x'), (1, 1, 1, 'y')", + error_substrings=("partition key column",), + ), + id="same_partition_key_different_column_order_multi_column", + ), + pytest.param( + RejectedPartitionExportCase( + src_columns="a Int32, b Int32, c Int32, val String", + src_partition_by="(a, b, c)", + dst_columns="a Int32, b Int32, c Int32, val String", + dst_partition_by="(c, b, a)", + insert_values="(1, 2, 3, 'x')", + error_substrings=( + "source and destination tables have different `PARTITION BY` expressions", + ), + ), + id="multi_column_partition_key_order_mismatch", + ), + pytest.param( + RejectedPartitionExportCase( + src_columns="a Int32, b Int32, c Int32, val String", + src_partition_by="(a, b, c)", + dst_columns="a Int32, b Int32, c Int32, val String", + dst_partition_by="(a, b)", + insert_values="(1, 2, 3, 'x')", + error_substrings=( + "source and destination tables have different `PARTITION BY` expressions", + ), + ), + id="multi_column_partition_key_fewer_in_destination", + ), + pytest.param( + RejectedPartitionExportCase( + src_columns="a Int32, b Int32, c Int32, val String", + src_partition_by="(a, b)", + dst_columns="a Int32, b Int32, c Int32, val String", + dst_partition_by="(a, b, c)", + insert_values="(1, 2, 3, 'x')", + error_substrings=( + "source and destination tables have different `PARTITION BY` expressions", + ), + ), + id="multi_column_partition_key_more_in_destination", + ), +] + + +@pytest.mark.parametrize("case", REJECTED_PARTITION_EXPORT_CASES) +def test_export_partition_partition_key_mismatch_variants_are_rejected(cluster, case): + skip_if_remote_database_disk_enabled(cluster) + node = cluster.instances["replica1"] + + postfix = str(uuid.uuid4()).replace("-", "_") + mt_table = f"rejected_mt_table_{postfix}" + s3_table = f"rejected_s3_table_{postfix}" + + node.query(f""" + CREATE TABLE {mt_table} ({case.src_columns}) + ENGINE = ReplicatedMergeTree('/clickhouse/tables/{mt_table}', 'replica1') + PARTITION BY {case.src_partition_by} + ORDER BY tuple() + """) + + node.query(f""" + CREATE TABLE {s3_table} ({case.dst_columns}) + ENGINE = S3(s3_conn, filename='{s3_table}', format=Parquet, partition_strategy='hive') + PARTITION BY {case.dst_partition_by} + """) + + node.query(f"INSERT INTO {mt_table} VALUES {case.insert_values}") + + partition_id = node.query( + f"SELECT partition_id FROM system.parts WHERE database = currentDatabase() " + f"AND table = '{mt_table}' AND active ORDER BY name LIMIT 1" + ).strip() + + error = node.query_and_get_error(f"ALTER TABLE {mt_table} EXPORT PARTITION ID '{partition_id}' TO TABLE {s3_table}") + assert "BAD_ARGUMENTS" in error, f"Expected BAD_ARGUMENTS, got: {error}" + for substring in case.error_substrings: + assert substring in error, f"Expected {substring!r} in error, got: {error}" + + error_all = node.query_and_get_error(f"ALTER TABLE {mt_table} EXPORT PARTITION ALL TO TABLE {s3_table}") + assert "BAD_ARGUMENTS" in error_all, f"Expected BAD_ARGUMENTS, got: {error_all}" + + count = int(node.query(f"SELECT count() FROM {s3_table}").strip()) + assert count == 0, f"Expected 0 rows in destination after rejected export, got {count}" + + +def test_export_partition_multi_column_partition_key_success(cluster): + skip_if_remote_database_disk_enabled(cluster) + node = cluster.instances["replica1"] + + postfix = str(uuid.uuid4()).replace("-", "_") + mt_table = f"multi_pkey_ok_mt_table_{postfix}" + s3_table = f"multi_pkey_ok_s3_table_{postfix}" + + node.query(f""" + CREATE TABLE {mt_table} (a Int32, b Int32, c Int32, val String) + ENGINE = ReplicatedMergeTree('/clickhouse/tables/{mt_table}', 'replica1') + PARTITION BY (a, b, c) + ORDER BY tuple() + """) + + node.query(f""" + CREATE TABLE {s3_table} (a Int32, b Int32, c Int32, val String) + ENGINE = S3(s3_conn, filename='{s3_table}', format=Parquet, partition_strategy='hive') + PARTITION BY (a, b, c) + """) + + node.query(f"INSERT INTO {mt_table} VALUES (1, 2, 3, 'x'), (1, 2, 3, 'y')") + + partition_id = node.query( + f"SELECT partition_id FROM system.parts WHERE database = currentDatabase() " + f"AND table = '{mt_table}' AND active ORDER BY name LIMIT 1" + ).strip() + + node.query(f"ALTER TABLE {mt_table} EXPORT PARTITION ID '{partition_id}' TO TABLE {s3_table}") + wait_for_export_status(node, mt_table, s3_table, partition_id, "COMPLETED") + + count = int(node.query(f"SELECT count() FROM {s3_table}").strip()) + assert count == 2, f"Expected 2 rows in destination after export, got {count}" + + result = node.query(f"SELECT a, b, c, val FROM {s3_table} ORDER BY val").strip() + assert result == "1\t2\t3\tx\n1\t2\t3\ty", f"Unexpected exported data:\n{result}" + + +def test_export_partition_multi_column_partition_key_success_all(cluster): + skip_if_remote_database_disk_enabled(cluster) + node = cluster.instances["replica1"] + + postfix = str(uuid.uuid4()).replace("-", "_") + mt_table = f"multi_pkey_ok_all_mt_table_{postfix}" + s3_table = f"multi_pkey_ok_all_s3_table_{postfix}" + + node.query(f""" + CREATE TABLE {mt_table} (a Int32, b Int32, c Int32, val String) + ENGINE = ReplicatedMergeTree('/clickhouse/tables/{mt_table}', 'replica1') + PARTITION BY (a, b, c) + ORDER BY tuple() + """) + + node.query(f""" + CREATE TABLE {s3_table} (a Int32, b Int32, c Int32, val String) + ENGINE = S3(s3_conn, filename='{s3_table}', format=Parquet, partition_strategy='hive') + PARTITION BY (a, b, c) + """) + + node.query(f"INSERT INTO {mt_table} VALUES (1, 2, 3, 'x'), (4, 5, 6, 'y')") + + partition_ids = node.query( + f"SELECT DISTINCT partition_id FROM system.parts WHERE database = currentDatabase() " + f"AND table = '{mt_table}' AND active ORDER BY partition_id" + ).strip().split("\n") + + node.query(f"ALTER TABLE {mt_table} EXPORT PARTITION ALL TO TABLE {s3_table}") + + for pid in partition_ids: + wait_for_export_status(node, mt_table, s3_table, pid, "COMPLETED") + + count = int(node.query(f"SELECT count() FROM {s3_table}").strip()) + assert count == 2, f"Expected 2 rows in destination after export, got {count}" + + result = node.query(f"SELECT a, b, c, val FROM {s3_table} ORDER BY val").strip() + assert result == "1\t2\t3\tx\n4\t5\t6\ty", f"Unexpected exported data:\n{result}"