Skip to content

Commit bb12dec

Browse files
Merge branch 'antalya-26.6' into feature/antalya-26.6/iceberg-puffin-deletion-vectors-read-2
Resolve conflicts by keeping deletion-vector read paths and count fail-closed behavior, while adopting multistorage path resolution from antalya-26.6. Co-authored-by: Cursor <cursoragent@cursor.com>
2 parents 79d773f + d42f80a commit bb12dec

104 files changed

Lines changed: 4381 additions & 518 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎docs/en/antalya/partition_export.md‎

Lines changed: 52 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -167,37 +167,49 @@ Query id: 9efc271a-a501-44d1-834f-bc4d20156164
167167
168168
Row 1:
169169
──────
170-
source_database: default
171-
source_table: replicated_source
172-
destination_database: default
173-
destination_table: replicated_destination
174-
create_time: 2025-11-21 18:21:51
175-
partition_id: 2022
176-
transaction_id: 7397746091717128192
177-
source_replica: r1
178-
parts: ['2022_0_0_0','2022_1_1_0','2022_2_2_0']
179-
parts_count: 3
180-
parts_to_do: 0
181-
status: COMPLETED
170+
source_database: default
171+
source_table: replicated_source
172+
destination_database: default
173+
destination_table: s3_destination
174+
create_time: 2025-11-21 18:21:51
175+
partition_id: 2022
176+
transaction_id: 9b2c1e5a-3f47-4c8e-8a1d-6f0b2d4e7c31
177+
query_id: 3fa3c8d3-7d6b-4f8b-9aa2-2c1f1ad0a111
178+
source_replica: r1
179+
parts: ['2022_0_0_0','2022_1_1_0','2022_2_2_0']
180+
parts_count: 3
181+
parts_to_do: 0
182+
status: COMPLETED
182183
last_exception_per_replica: []
183-
exception_count: 0
184+
exception_count: 0
185+
destination_file_paths: {'2022_0_0_0':['data/year=2022/2022_0_0_0_<hash>.parquet'],'2022_1_1_0':['data/year=2022/2022_1_1_0_<hash>.parquet'],'2022_2_2_0':['data/year=2022/2022_2_2_0_<hash>.parquet']}
186+
committed_metadata_file:
187+
committed_manifest_list:
188+
committed_manifest_file:
189+
committed_marker_file: data/commit_2022_9b2c1e5a-3f47-4c8e-8a1d-6f0b2d4e7c31
184190
185191
Row 2:
186192
──────
187-
source_database: default
188-
source_table: replicated_source
189-
destination_database: default
190-
destination_table: replicated_destination
191-
create_time: 2025-11-21 18:20:35
192-
partition_id: 2021
193-
transaction_id: 7397745772618674176
194-
source_replica: r1
195-
parts: ['2021_0_0_0']
196-
parts_count: 1
197-
parts_to_do: 0
198-
status: COMPLETED
199-
last_exception_per_replica: []
200-
exception_count: 0
193+
source_database: default
194+
source_table: replicated_source
195+
destination_database: default
196+
destination_table: iceberg_destination
197+
create_time: 2025-11-21 18:20:35
198+
partition_id: 2021
199+
transaction_id: d0e4f7a2-8c19-4b6d-9e3a-1f5c7b2e9d40
200+
query_id: 1c8e0fd0-6a3a-4d6e-9bd6-bdf64adfe118
201+
source_replica: r2
202+
parts: ['2021_0_0_0']
203+
parts_count: 1
204+
parts_to_do: 0
205+
status: COMPLETED
206+
last_exception_per_replica: [('r1','Code: 999. Coordination::Exception: Session expired','2021_0_0_0','2025-11-21 18:20:42',1)]
207+
exception_count: 1
208+
destination_file_paths: {'2021_0_0_0':['data/year=2021/2021_0_0_0_<hash>.parquet']}
209+
committed_metadata_file: data/metadata/v3.metadata.json
210+
committed_manifest_list: data/metadata/snap-4029103741930112856-1-<uuid>.avro
211+
committed_manifest_file: data/metadata/<uuid>-m0.avro
212+
committed_marker_file:
201213
202214
2 rows in set. Elapsed: 0.019 sec.
203215
@@ -215,6 +227,19 @@ Status values include:
215227
- `last_exception_per_replica` is an `Array(Tuple(replica String, message String, part String, time DateTime, count UInt64))`. Each tuple is the most recent exception observed by a single replica plus a best-effort within-replica `count`. Replicas that have never reported an exception are omitted.
216228
- `exception_count` is the sum of every `count` in `last_exception_per_replica`. Each replica owns its own counter, so cross-replica updates do not race; the sum is exact w.r.t. the snapshot returned. Within a single replica concurrent failing writers may under-count by one.
217229

230+
### Per-part destination file paths
231+
232+
- `destination_file_paths` is a `Map(String, Array(String))` keyed by source part name. Each value is the list of file paths written to the destination object storage when that part was exported (a single part can produce multiple files depending on `max_bytes` / `max_rows`). If a refresh cannot read a processed entry from ZooKeeper, the affected key holds the sentinel `<failed to read from zk>` instead of silently under-counting.
233+
234+
### Commit info columns
235+
236+
These columns surface paths produced by the destination storage during commit, so it is possible to inspect what was written without consulting the destination directly:
237+
238+
- `committed_metadata_file` — for Iceberg destinations: path of the new `vN.metadata.json` written by the commit. Empty for non-Iceberg destinations and before the commit lands. If the commit was already finished by a previous run (detected via the transaction id stored in the snapshot summary), this column carries a human-readable sentinel string instead of a path because the original committer's paths are not recoverable from inside the impl.
239+
- `committed_manifest_list` — for Iceberg destinations: path of the manifest list file (`snap-*.avro`) referenced by the new snapshot. Empty under the same conditions as `committed_metadata_file`.
240+
- `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`.
241+
- `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.
242+
218243
To pick the latest exception across replicas:
219244
220245
```sql

‎docs/en/engines/database-engines/datalake.md‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,7 @@ The following settings are supported:
5454
| `storage_endpoint` | Endpoint URL for the underlying storage |
5555
| `oauth_server_uri` | URI of the OAuth2 authorization server for authentication |
5656
| `vended_credentials` | Boolean indicating whether to use vended credentials from the catalog (supports AWS S3 and Azure ADLS Gen2) |
57+
| `vended_credentials_cache_ttl` | Maximum cache entry lifetime (in seconds) for vended credentials (REST catalogs only). Default `300`; `0` disables caching. |
5758
| `aws_access_key_id` | AWS access key ID for S3/Glue access (if not using vended credentials) |
5859
| `aws_secret_access_key` | AWS secret access key for S3/Glue access (if not using vended credentials) |
5960
| `region` | AWS region for the service (e.g., `us-east-1`) |

‎programs/client/Client.cpp‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,8 @@
1717
#include <Common/Config/ConfigProcessor.h>
1818
#include <Common/Config/getClientConfigPath.h>
1919
#include <Common/CurrentThread.h>
20+
#include <Common/DateLUT.h>
21+
#include <Common/DateLUTImpl.h>
2022
#include <Common/QueryScope.h>
2123
#include <Common/Exception.h>
2224
#include <Common/TerminalSize.h>
@@ -546,6 +548,13 @@ void Client::connect()
546548
UInt64 server_version_minor = 0;
547549
UInt64 server_version_patch = 0;
548550

551+
/// Capture the client local time zone before the branch below may switch the process default
552+
/// to the server time zone. `serverTimezoneInstance()` reads the process default directly and
553+
/// ignores `session_timezone`; `instance()` would fold in an explicit `--session_timezone` and
554+
/// cache the wrong zone. `connect()` can run again on reconnect, so only capture once.
555+
if (client_local_timezone.empty())
556+
client_local_timezone = DateLUT::serverTimezoneInstance().getTimeZone();
557+
549558
if (hosts_and_ports.empty())
550559
{
551560
String host = config().getString("host", "localhost");

‎src/Client/ClientBase.cpp‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -150,6 +150,8 @@ namespace Setting
150150
extern const SettingsFloatAuto promql_evaluation_time;
151151
extern const SettingsBool into_outfile_create_parent_directories;
152152
extern const SettingsBool ignore_format_null_for_explain;
153+
extern const SettingsBool use_client_time_zone;
154+
extern const SettingsTimezone session_timezone;
153155
}
154156

155157
namespace ErrorCodes
@@ -2513,6 +2515,16 @@ void ClientBase::processParsedSingleQuery(
25132515

25142516
applySettingsFromServerIfNeeded(); // after connect() and applySettingsFromQuery()
25152517

2518+
/// With `use_client_time_zone`, DateTime string literals must be interpreted in the client time
2519+
/// zone. The client parses synchronous INSERT literals itself, but literals interpreted server-side
2520+
/// (asynchronous INSERT, SELECT) rely on `session_timezone`. Seed it with the client time zone unless
2521+
/// the user set `session_timezone` explicitly. This is transient (reverted with the other query
2522+
/// settings below), so it tracks per-query `use_client_time_zone` changes in both directions.
2523+
if (!client_local_timezone.empty()
2524+
&& client_context->getSettingsRef()[Setting::use_client_time_zone]
2525+
&& !client_context->getSettingsRef().isChanged("session_timezone"))
2526+
client_context->setSetting("session_timezone", client_local_timezone);
2527+
25162528
ASTPtr input_function;
25172529
const auto * insert = parsed_query->as<ASTInsertQuery>();
25182530
if (insert && insert->select)

‎src/Client/ClientBase.h‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -326,6 +326,11 @@ class ClientBase
326326
ContextMutablePtr global_context;
327327
ContextMutablePtr client_context;
328328

329+
/// The client local time zone, captured on the first connect() before it may switch the
330+
/// process default to the server time zone. Used to seed `session_timezone` per query when
331+
/// `use_client_time_zone` is set, so server-side literal parsing matches the client side.
332+
String client_local_timezone;
333+
329334
String default_database;
330335
String query_id;
331336
Int32 suggestion_limit{};

‎src/Common/FailPoint.cpp‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -170,6 +170,7 @@ static struct InitFiu
170170
ONCE(iceberg_export_after_commit_before_zk_completed) \
171171
REGULAR(export_partition_commit_always_throw) \
172172
ONCE(export_partition_status_change_throw) \
173+
REGULAR(export_partition_processed_paths_sync_fail) \
173174
REGULAR(export_part_non_retryable_throw) \
174175
REGULAR(export_part_retryable_throw) \
175176
ONCE(backup_add_empty_memory_table) \

‎src/Common/ProfileEvents.cpp‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1541,6 +1541,8 @@ The server successfully detected this situation and will download merged part fr
15411541
M(ObjectStorageListObjectsCacheMisses, "Number of times object storage list objects operation miss the cache.", ValueType::Number) \
15421542
M(ObjectStorageListObjectsCacheExactMatchHits, "Number of times object storage list objects operation hit the cache with an exact match.", ValueType::Number) \
15431543
M(ObjectStorageListObjectsCachePrefixMatchHits, "Number of times object storage list objects operation miss the cache using prefix matching.", ValueType::Number) \
1544+
M(DataLakeRestCatalogCredentialsVended, "Number of table metadata requests to REST catalog asking to vend storage credentials.", ValueType::Number) \
1545+
M(DataLakeRestCatalogCredentialsCacheHits, "Number of table metadata requests to REST catalog reusing cached storage credentials.", ValueType::Number) \
15441546
\
15451547
M(DataLakeRestCatalogLoadConfig, "Number of 'load config' requests to Iceberg REST catalog.", ValueType::Number) \
15461548
M(DataLakeRestCatalogLoadConfigMicroseconds, "Total time of 'load config' requests to Iceberg REST catalog.", ValueType::Microseconds) \

‎src/Core/ProtocolDefines.h‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -39,7 +39,8 @@ static constexpr auto DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION_WITH_ICEBERG_META
3939
static constexpr auto DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION_WITH_FILE_BUCKETS_INFO = 4;
4040
static constexpr auto DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION_WITH_EXCLUDED_ROWS = 5;
4141
static constexpr auto DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION_WITH_ICEBERG_FILE_STATS = 6;
42-
static constexpr auto DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION = DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION_WITH_ICEBERG_FILE_STATS;
42+
static constexpr auto DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION_WITH_ICEBERG_ABSOLUTE_PATH = 9;
43+
static constexpr auto DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION = DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION_WITH_ICEBERG_ABSOLUTE_PATH;
4344

4445
static constexpr auto DATA_LAKE_TABLE_STATE_SNAPSHOT_PROTOCOL_VERSION = 1;
4546

‎src/Core/Settings.cpp‎

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -641,6 +641,9 @@ Use multiple threads for azure multipart upload.
641641
)", 0) \
642642
DECLARE(Bool, s3_throw_on_zero_files_match, false, R"(
643643
Throw an error, when ListObjects request cannot match any files
644+
)", 0) \
645+
DECLARE(Bool, object_storage_propagate_credentials_to_other_storages, false, R"(
646+
Reuse base-storage credentials for a secondary object storage. For `S3`, credentials are reused when the endpoint matches; when this setting is enabled, they are also reused across different endpoints, including less secure connections (for example, from `https` to plain `http`). For `Azure`, reads stay within the base account.
644647
)", 0) \
645648
DECLARE(Bool, hdfs_throw_on_zero_files_match, false, R"(
646649
Throw an error if matched zero files according to glob expansion rules.
@@ -8070,11 +8073,15 @@ To survive a long transient outage (e.g. object storage downtime), raise `export
80708073
DECLARE(UInt64, export_merge_tree_partition_retry_max_backoff_seconds, 300, R"(
80718074
Maximum delay (in seconds) between retries of a failed part export in an export partition task. Caps the exponential growth controlled by `export_merge_tree_partition_retry_initial_backoff_seconds`.
80728075
)", 0) \
8073-
DECLARE(UInt64, export_merge_tree_partition_task_timeout_seconds, 3600, R"(
8076+
DECLARE(UInt64, export_merge_tree_partition_task_timeout_seconds, 86400, R"(
80748077
Maximum wall-clock duration (in seconds) an export partition task is allowed to remain in the PENDING state before it is auto-killed by the background cleanup loop.
80758078
The timeout is measured from the manifest's create_time. Set to 0 to disable the timeout.
80768079
When the timeout is exceeded the task transitions to KILLED (same terminal state as `KILL QUERY ... EXPORT PARTITION`), and `last_exception` is populated with a timeout reason.
80778080
8081+
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.
8082+
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.
8083+
In such scenario, ClickHouse would attempt to commit those files again producing duplicates.
8084+
80788085
Notes:
80798086
- Enforcement is best-effort: actual kill latency is bounded by one manifest-updater poll cycle (~30s) plus ZooKeeper watch propagation.
80808087
)", 0) \

‎src/Core/SettingsChangesHistory.cpp‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -283,7 +283,7 @@ const VersionToSettingsChangesMap & getSettingsChangesHistory()
283283
addSettingsChanges(settings_changes_history, "26.1.3.20001.altinityantalya",
284284
{
285285
{"iceberg_partition_timezone", "", "", "New setting."},
286-
// {"s3_propagate_credentials_to_other_storages", false, false, "New setting"},
286+
{"object_storage_propagate_credentials_to_other_storages", false, false, "New setting"},
287287
{"export_merge_tree_part_filename_pattern", "", "{part_name}_{checksum}", "New setting"},
288288
// {"use_parquet_metadata_cache", false, true, "Enables cache of parquet file metadata."},
289289
// {"input_format_parquet_use_metadata_cache", true, false, "Obsolete. No-op"}, // https://github.com/Altinity/ClickHouse/pull/586

0 commit comments

Comments
 (0)