diff --git a/docs/en/antalya/part_export.md b/docs/en/antalya/part_export.md index 73d467c5d9b1..2ff4a7301674 100644 --- a/docs/en/antalya/part_export.md +++ b/docs/en/antalya/part_export.md @@ -51,6 +51,9 @@ Source and destination tables must be 100% compatible: 1. **Identical schemas** - same columns, types, and order 2. **Matching partition keys** - partition expressions must be identical +3. **Partition key columns at the same position** - columns are matched by position, similar to `INSERT INTO dest SELECT * FROM src`. It is not enough for the `PARTITION BY` expressions to be textually identical: every column that is part of the source table's partition key must also sit at the same position in the destination table's schema. 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`, so the export is rejected with `BAD_ARGUMENTS: 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 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 full column order identical (per point 1) rather than relying on this check alone. 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..ccd18914845b 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/engines/table-engines/mergetree-family/part_export.md#requirements), so the source and destination tables must satisfy the same compatibility requirements, in particular: + +1. **Identical schemas** - same columns, types, and order +2. **Matching partition keys** - partition expressions must be identical +3. **Partition key columns at the same position** - columns are matched by position, so every column that is part of the source table's partition key must also sit at the same position in the destination table's schema, even if both tables' `PARTITION BY` expressions are textually identical. See [`EXPORT PART` requirements](/docs/en/engines/table-engines/mergetree-family/part_export.md#requirements) for a worked example and the exact error message. + ## Settings ### Server Settings diff --git a/src/Storages/MergeTree/ExportPartitionUtils.cpp b/src/Storages/MergeTree/ExportPartitionUtils.cpp index 3f85ebc0e1fb..433c663a445b 100644 --- a/src/Storages/MergeTree/ExportPartitionUtils.cpp +++ b/src/Storages/MergeTree/ExportPartitionUtils.cpp @@ -9,10 +9,14 @@ #include #include #include +#include #include #include #include +#include +#include #include +#include #include #include #include @@ -634,6 +638,32 @@ namespace ExportPartitionUtils } #endif + namespace + { + std::optional getDateTimeTimeZoneName(const DataTypePtr & type) + { + if (const auto * datetime_type = checkAndGetDataType(type.get())) + return datetime_type->getTimeZone().getTimeZone(); + if (const auto * datetime64_type = checkAndGetDataType(type.get())) + return datetime64_type->getTimeZone().getTimeZone(); + return {}; + } + } + + void verifyMergeTreePartitionCompatibility( + 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 +687,51 @@ 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; + auto partition_key_columns = source_metadata->getColumnsRequiredForPartitionKey(); + const std::unordered_set partition_key_column_set( + std::make_move_iterator(partition_key_columns.begin()), + std::make_move_iterator(partition_key_columns.end())); + + 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_column_set.contains(source_column.name) && 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, + i, + destination_column.name); + + if (partition_key_column_set.contains(source_column.name)) + { + const auto source_time_zone = getDateTimeTimeZoneName(source_column.type); + const auto destination_time_zone = getDateTimeTimeZoneName(destination_column.type); + if (source_time_zone && destination_time_zone && *source_time_zone != *destination_time_zone) + throw Exception(ErrorCodes::BAD_ARGUMENTS, + "Cannot export to {}: partition key column '{}' is {} in the source table " + "but {} in the destination. The destination's hive-style partition path is " + "rendered from the source value without converting the timezone, so this " + "would silently shift the exported value by the timezone offset. Use the " + "same timezone in both tables' partition key column.", + destination_storage_id.getFullTableName(), + destination_column.name, + source_column.type->getName(), + destination_column.type->getName()); + } + + /// 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..dd1ae4c18094 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 verifyMergeTreePartitionCompatibility( + 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..ac8ad03c3b63 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::verifyMergeTreePartitionCompatibility(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..759736dc06be 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::verifyMergeTreePartitionCompatibility(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..e1dd26ea3820 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) 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..b4a1160313aa 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,216 @@ 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 = () + + +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=("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=("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=("different `PARTITION BY` expressions",), + ), + id="multi_column_partition_key_more_in_destination", + ), + pytest.param( + RejectedPartExportCase( + src_columns="id Int64, ts DateTime('UTC')", + src_partition_by="ts", + dst_columns="id Int64, ts DateTime('Asia/Tokyo')", + dst_partition_by="ts", + insert_values="(1, '2024-03-05 15:00:00')", + error_substrings=("timezone",), + ), + id="partition_key_timezone_mismatch", + ), +] + + +@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 + """) + + 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}") + + 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}" + + 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}" + + +def test_export_part_non_partition_key_timezone_mismatch_is_allowed(cluster): + skip_if_remote_database_disk_enabled(cluster) + node = cluster.instances["node1"] + + postfix = str(uuid.uuid4()).replace("-", "_") + mt_table = f"tz_ok_mt_table_{postfix}" + s3_table = f"tz_ok_s3_table_{postfix}" + + node.query(f""" + CREATE TABLE {mt_table} (id Int64, ts DateTime('UTC')) + ENGINE = MergeTree() + PARTITION BY id + ORDER BY tuple() + SETTINGS enable_block_number_column = 1, enable_block_offset_column = 1 + """) + + node.query(f""" + CREATE TABLE {s3_table} (id Int64, ts DateTime('Asia/Tokyo')) + ENGINE = S3(s3_conn, filename='{s3_table}', format=Parquet, partition_strategy='hive') + PARTITION BY id + """) + + node.query(f"INSERT INTO {mt_table} VALUES (1, '2024-03-05 15:00:00')") + + 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 == 1, f"Expected 1 row in destination after export, got {count}" + + source_ts = node.query(f"SELECT ts FROM {mt_table}").strip() + assert source_ts == "2024-03-05 15:00:00", f"Unexpected source value: {source_ts}" + + dest_ts = node.query(f"SELECT ts FROM {s3_table}").strip() + assert dest_ts == "2024-03-06 00:00:00", ( + f"Expected the exported value to be the same instant displayed in the " + f"destination's Asia/Tokyo timezone ('2024-03-06 00:00:00'), got: {dest_ts}" + ) + + source_unix_ts = int(node.query(f"SELECT toUnixTimestamp(ts) FROM {mt_table}").strip()) + dest_unix_ts = int(node.query(f"SELECT toUnixTimestamp(ts) FROM {s3_table}").strip()) + assert source_unix_ts == dest_unix_ts, ( + f"Expected exported DateTime value to be preserved regardless of the " + f"destination column's timezone, got source={source_unix_ts}, dest={dest_unix_ts}" + ) + + 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..0b38f97df4a5 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,201 @@ 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=("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=("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=("different `PARTITION BY` expressions",), + ), + id="multi_column_partition_key_more_in_destination", + ), + pytest.param( + RejectedPartitionExportCase( + src_columns="id Int64, ts DateTime('UTC')", + src_partition_by="ts", + dst_columns="id Int64, ts DateTime('Asia/Tokyo')", + dst_partition_by="ts", + insert_values="(1, '2024-03-05 15:00:00')", + error_substrings=("timezone",), + ), + id="partition_key_timezone_mismatch", + ), +] + + +@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}"