Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions docs/en/antalya/part_export.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
8 changes: 8 additions & 0 deletions docs/en/antalya/partition_export.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
72 changes: 69 additions & 3 deletions src/Storages/MergeTree/ExportPartitionUtils.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -9,10 +9,14 @@
#include <Storages/MergeTree/MergeTreeData.h>
#include <filesystem>
#include <thread>
#include <unordered_map>
#include <unordered_set>
#include <Core/Block.h>
#include <Core/Settings.h>
#include <DataTypes/DataTypeDateTime.h>
#include <DataTypes/DataTypeDateTime64.h>
#include <DataTypes/Utils.h>
#include <Functions/FunctionHelpers.h>
#include <Interpreters/ActionsDAG.h>
#include <Interpreters/Context.h>
#include <Interpreters/ExpressionActions.h>
Expand Down Expand Up @@ -634,6 +638,32 @@ namespace ExportPartitionUtils
}
#endif

namespace
{
std::optional<String> getDateTimeTimeZoneName(const DataTypePtr & type)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I found method getExplicitTimeZoneOfDateTimeArgument like this.
And fragment from method extractTimeZoneNameFromFunctionArguments.
May be possible to reuse something?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Seems they are a slightly defferent. But I could use checkAndGetDataType from them.

{
if (const auto * datetime_type = checkAndGetDataType<DataTypeDateTime>(type.get()))
return datetime_type->getTimeZone().getTimeZone();
if (const auto * datetime64_type = checkAndGetDataType<DataTypeDateTime64>(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,
Expand All @@ -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<String> partition_key_column_set(
std::make_move_iterator(partition_key_columns.begin()),
std::make_move_iterator(partition_key_columns.end()));
Comment on lines +690 to +693

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Normalize subcolumns before checking partition positions

When PARTITION BY references a subcolumn such as t.a, getColumnsRequiredForPartitionKey returns the subcolumn name, while source_columns contains only the top-level readable column t; consequently, this set never matches and the new positional check is skipped. For example, source columns (t Tuple(a UInt32), u Tuple(a UInt32)) and destination columns (u Tuple(a UInt32), t Tuple(a UInt32)), both partitioned by t.a, have positionally compatible types and pass validation, but export source u into destination t while constructing the partition from source t.a, silently mispartitioning the data. Map required subcolumns back to their top-level storage columns before constructing this set.

Useful? React with 👍 / 👎.


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)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

What if partitioning is for several columns, and destination has same columns but in different order?
Both have columns 'a' and 'b', in same order, but source with partition by a,b, and destination partition by b,a.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Or when destination has more columns in partition list than source.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I added the tests for simular cases:

What if partitioning is for several columns, and destination has same columns but in different order?
https://github.com/Altinity/ClickHouse/pull/2134/changes#diff-28d3e30160a3442e2b9ef01e3bf0b10b4c27de6a39efea3175ccec0a3f2d8949R1786

Or when destination has more columns in partition list than source.
https://github.com/Altinity/ClickHouse/pull/2134/changes#diff-28d3e30160a3442e2b9ef01e3bf0b10b4c27de6a39efea3175ccec0a3f2d8949R1799

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 {}, "
Expand Down
4 changes: 4 additions & 0 deletions src/Storages/MergeTree/ExportPartitionUtils.h
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
12 changes: 1 addition & 11 deletions src/Storages/MergeTree/MergeTreeData.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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();

Expand Down Expand Up @@ -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});

Expand Down
12 changes: 1 addition & 11 deletions src/Storages/StorageReplicatedMergeTree.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -8408,25 +8408,15 @@ 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();

/// Positional CAST matching, like `INSERT INTO dest SELECT * FROM src`.
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();

Expand Down
143 changes: 143 additions & 0 deletions tests/integration/test_export_merge_tree_part_to_iceberg/test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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)
Expand Down
Loading
Loading