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
70 changes: 42 additions & 28 deletions docs/en/antalya/partition_export.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,25 +9,23 @@ The `ALTER TABLE EXPORT PARTITION` command exports entire partitions from `Merge

The set of parts that are exported is based on the list of parts the replica that received the export command sees. On `Replicated*MergeTree`, the other replicas will assist in the export process if they have those parts locally. Otherwise they will ignore it.

The partition export tasks of both engines can be observed through `system.partition_exports`.
Exporting a partition that was exported before is allowed and exports its parts again; Preventing duplicates is up to the user.

`system.replicated_partition_exports` is kept as an alias of `system.partition_exports` for backwards compatibility. It returns exactly the same rows, including exports of plain `MergeTree` tables.
Export tasks of both engines, including those of a [`TTL ... EXPORT TO TABLE`](/docs/en/antalya/ttl_export.md) expression, can be observed through `system.distributed_exports`, one row per task. Manual exports are independent of the `EXPORT` TTL: they neither read nor update what the TTL has exported.

The same partition can not be exported to the same destination more than once. This behavior can be overriden with `export_merge_tree_partition_force_export`.

The export task can be killed by issuing the kill command: `KILL EXPORT PARTITION <where predicate for system.partition_exports>`.
The export task can be killed by issuing the kill command: `KILL EXPORT <where predicate for system.distributed_exports>`.

The task is persistent - it should be resumed after crashes, failures and etc.

A part with no surviving rows writes no file. This happens when every row of the part was removed by a lightweight delete: the part is still exported, and counts as done, but it contributes nothing to the destination. If that is true of every part of the partition, the export produces no files at all and there is nothing to commit, so the task reaches `COMPLETED` without touching the destination. Such a part therefore has no entry in the `destination_file_paths` column of `system.partition_exports`.
A part with no surviving rows writes no file. This happens when every row of the part was removed by a lightweight delete: the part is still exported, and counts as done, but it contributes nothing to the destination. If that is true of every part of the partition, the export produces no files at all and there is nothing to commit, so the task reaches `COMPLETED` without touching the destination. Such a part therefore has no entry in the `destination_file_paths` column of `system.distributed_exports`.

### On Apache Iceberg storage exports:

Each MergeTree part that has surviving rows will become a separate file (or more depending on `max_bytes` and `max_rows` settings) following the engine naming convention. Once all parts have been exported, new snapshots / manifest files are generated and the data is comitted using the Apache Iceberg commit mechanism.

The manifest file produced by the commit contains a summary field `clickhouse.export-partition-transaction-id` that stores the transaction id. This field is used to implement idempotency and avoid data duplication. Some Apache Iceberg storage managers employ old manifests cleanup, ClickHouse does not.

**IMPORTANT**: In case the storage is managed by a 3rd party application that cleans up old manifest files, it is important that the TTL of such files are greater than the timeout of export partition tasks. If it is not configured in such a way, it is possible to accidentally duplicate data in the extremely rare case a ClickHouse node is the only node working on a given export task, commits the data to Iceberg, crashes before marking the task as done and only boots up after the manifest cleanup has deleted the commit manifest. In such scenario, ClickHouse would attempt to commit those files again producing duplicates. The task timeout on ClickHouse side is controlled by the setting `export_merge_tree_partition_task_timeout_seconds`.
**IMPORTANT**: In case the storage is managed by a 3rd party application that cleans up old manifest files, it is important that the TTL of such files are greater than the timeout of export partition tasks. If it is not configured in such a way, it is possible to accidentally duplicate data in the extremely rare case a ClickHouse node is the only node working on a given export task, commits the data to Iceberg, crashes before marking the task as done and only boots up after the manifest cleanup has deleted the commit manifest. In such scenario, ClickHouse would attempt to commit those files again producing duplicates. The task timeout on ClickHouse side is controlled by the setting `export_merge_tree_task_timeout_seconds`.

The Iceberg manifest files contain statistics about the data. Exporting a merge tree partition is a non ephemeral long running task, in which nodes can be turned off and turned on. This means the stats of individual files need to be persisted somewhere in order to produce the final manifest. This is implemented through sidecars. Each data file exported will contain a "sibling" sidecar file named `<data_file_name>_clickhouse_export_part_sidecar.avro`. ClickHouse does not clean up these files, and they can be safely deleted once the data is comitted.

Expand All @@ -40,7 +38,7 @@ The source partition must not be split in the destination. This is validated at

### On plain object storage exports:

Each MergeTree part will become a separate file with the following name convention: `<table_directory>/<partitioning>/<data_part_name>_<merge_tree_part_checksum>.<format>`. To ensure atomicity, a commit file containing the relative paths of all exported parts is also shipped. A data file should only be considered part of the dataset if a commit file references it. The commit file will be named using the following convention: `<table_directory>/commit_<partition_id>_<transaction_id>`.
Each MergeTree part will become a separate file with the following name convention: `<table_directory>/<partitioning>/<data_part_name>_<merge_tree_part_checksum>.<format>`. To ensure atomicity, a commit file containing the relative paths of all exported parts is also shipped. A data file should only be considered part of the dataset if a commit file references it. The commit file will be named using the following convention: `<table_directory>/commit_<transaction_id>`.

## Plain (non-replicated) MergeTree {#plain-non-replicated-mergetree}

Expand Down Expand Up @@ -85,23 +83,26 @@ TO TABLE [destination_database.]destination_table

### Query Settings

#### `export_merge_tree_partition_force_export` (Optional)

- **Type**: `Bool`
- **Default**: `false`
- **Description**: Ignore existing partition export and overwrite the ZooKeeper entry. Allows re-exporting a partition that was already exported to the same destination. **IMPORTANT:** this is dangerous because it can lead to duplicated data, use it with caution.

#### `export_merge_tree_partition_retry_initial_backoff_seconds` (Optional)
#### `export_merge_tree_retry_initial_backoff_seconds` (Optional)

- **Type**: `UInt64`
- **Default**: `5`
- **Description**: Initial delay (in seconds) before retrying a failed part export. The delay grows exponentially with the per-replica retry count (`delay = min(initial << (attempts - 1), max)`). The back-off is per-replica in-memory state: it only spaces this replica's retries out in time and never prevents another replica from attempting the same part. Retryable failures (transient memory/network/object-storage/Keeper errors) are retried until the task succeeds or `export_merge_tree_partition_task_timeout_seconds` elapses, while non-retryable failures (e.g. schema/type incompatibilities) fail the task immediately.
- **Description**: Initial delay (in seconds) before retrying a failed part export. The delay grows exponentially with the per-replica retry count (`delay = min(initial << (attempts - 1), max)`). The back-off is per-replica in-memory state: it only spaces this replica's retries out in time and never prevents another replica from attempting the same part. Retryable failures (transient memory/network/object-storage/Keeper errors) are retried until the task succeeds or `export_merge_tree_task_timeout_seconds` elapses, while non-retryable failures (e.g. schema/type incompatibilities) fail the task immediately.

#### `export_merge_tree_partition_retry_max_backoff_seconds` (Optional)
#### `export_merge_tree_retry_max_backoff_seconds` (Optional)

- **Type**: `UInt64`
- **Default**: `300`
- **Description**: Maximum delay (in seconds) between retries of a failed part export. Caps the exponential growth controlled by `export_merge_tree_partition_retry_initial_backoff_seconds`.
- **Description**: Maximum delay (in seconds) between retries of a failed part export. Caps the exponential growth controlled by `export_merge_tree_retry_initial_backoff_seconds`.

#### `export_merge_tree_partition_all_on_error` (Optional) {#export-merge-tree-partition-all-on-error}

- **Type**: `ExportPartitionAllOnError`
- **Default**: `throw_first`
- **Description**: How `EXPORT PARTITION ALL` handles a partition that cannot be exported. Possible values:
- `throw_first` - stop at the first failing partition and throw
- `collect` - try every partition, then throw one error listing the failing ones
- `skip_conflicts` - behaves like `throw_first`: re-exporting a partition is no longer refused, so there are no conflicts to skip

#### `export_merge_tree_part_file_already_exists_policy` (Optional)

Expand Down Expand Up @@ -130,12 +131,12 @@ TO TABLE [destination_database.]destination_table
- **Default**: `{part_name}_{checksum}`
- **Description**: Pattern for the filename of the exported merge tree part. The `part_name` and `checksum` are calculated and replaced on the fly. Additional macros are supported.

### `export_merge_tree_partition_task_timeout_seconds` (Optional)
### `export_merge_tree_task_timeout_seconds` (Optional)

- **Type**: `UInt64`
- **Default**: `3600`
- **Default**: `86400`
- **Description**: The timeout is measured from the manifest's create_time. Set to 0 to disable the timeout.
When the timeout is exceeded the task transitions to KILLED (same terminal state as `KILL QUERY ... EXPORT PARTITION`), and a `last_exception_per_replica` entry on the replica that fires the timeout is populated with a timeout reason.
When the timeout is exceeded the task transitions to KILLED (same terminal state as `KILL EXPORT`), and a `last_exception_per_replica` entry on the replica that fires the timeout is populated with a timeout reason.

Notes:
- Enforcement is best-effort: actual kill latency is bounded by one manifest-updater poll cycle (~30s) plus ZooKeeper watch propagation.
Expand Down Expand Up @@ -186,13 +187,14 @@ PARTITION BY year;
INSERT INTO rmt_table VALUES (1, 2020), (2, 2020), (3, 2020), (4, 2021);

ALTER TABLE rmt_table EXPORT PARTITION ID '2020' TO TABLE s3_table;
```

## Killing Exports

You can cancel in-progress partition exports using the `KILL EXPORT PARTITION` command:
You can cancel in-progress exports using the `KILL EXPORT` command:

```sql
KILL EXPORT PARTITION
KILL EXPORT
WHERE partition_id = '2020'
AND source_table = 'rmt_table'
AND destination_table = 's3_table'
Expand All @@ -202,13 +204,13 @@ WHERE partition_id = '2020'

### Active and Completed Exports

Monitor partition exports using the `system.partition_exports` table:
Monitor exports using the `system.distributed_exports` table:

```sql
arthur :) select * from system.partition_exports Format Vertical;
arthur :) select * from system.distributed_exports Format Vertical;

SELECT *
FROM system.partition_exports
FROM system.distributed_exports
FORMAT Vertical

Query id: 9efc271a-a501-44d1-834f-bc4d20156164
Expand All @@ -234,7 +236,10 @@ destination_file_paths: {'2022_0_0_0':['data/year=2022/2022_0_0_0_<hash>.par
committed_metadata_file:
committed_manifest_list:
committed_manifest_file:
committed_marker_file: data/commit_2022_9b2c1e5a-3f47-4c8e-8a1d-6f0b2d4e7c31
committed_marker_file: data/commit_9b2c1e5a-3f47-4c8e-8a1d-6f0b2d4e7c31
local_backoff_per_part: []
source: query
retry_of: []

Row 2:
──────
Expand All @@ -258,6 +263,9 @@ committed_metadata_file: data/metadata/v3.metadata.json
committed_manifest_list: data/metadata/snap-4029103741930112856-1-<uuid>.avro
committed_manifest_file: data/metadata/<uuid>-m0.avro
committed_marker_file:
local_backoff_per_part: []
source: query
retry_of: []

2 rows in set. Elapsed: 0.019 sec.

Expand Down Expand Up @@ -286,15 +294,21 @@ Status values include:
- `committed_manifest_file` — for Iceberg destinations: path of the manifest file referenced by `committed_manifest_list`. Empty under the same conditions as `committed_metadata_file`.
- `committed_marker_file` — for plain object storage destinations: path of the per-transaction commit marker file written by the destination. Empty for Iceberg destinations and for tasks that have not committed yet.

### Source columns {#source-columns}

- `source` — `query` for a task of `EXPORT PARTITION`, `ttl` for a task of the table's `TTL ... EXPORT TO TABLE` expression.
- `retry_of` — for a task of the `EXPORT` TTL: transaction ids of earlier tasks that failed to export some of its parts after exporting all of theirs, so their commit may still land. Its commit checks whether any of them landed at the destination after all.

To pick the latest exception across replicas:

```sql
SELECT
arraySort(x -> -x.time, last_exception_per_replica)[1] AS latest_exception
FROM system.partition_exports
FROM system.distributed_exports
WHERE source_table = 'rmt_table' AND destination_table = 's3_table';
```

## Related Features

- [ALTER TABLE EXPORT PART](/docs/en/antalya/part_export.md) - Export individual parts (non-replicated)
- [TTL ... EXPORT TO TABLE](/docs/en/antalya/ttl_export.md) - Export parts in the background once their TTL is due
Loading
Loading