From e8c8b9fb705460a612f7f3d121e63022f0a07c8c Mon Sep 17 00:00:00 2001 From: Anton Ivashkin Date: Tue, 6 Oct 2026 11:17:37 +0200 Subject: [PATCH 1/6] Send task retry delay in the Antalya protocol trailer. `retry_after_us` for `lock_object_storage_task_distribution_ms` is appended to `ReadTaskResponse` only when both peers negotiated version 3, so an upstream worker still receives a plain object path. Co-authored-by: Cursor --- docs/en/antalya/protocol.md | 8 ++- docs/en/interfaces/specs/NativeProtocol.md | 2 +- .../sql-reference/distribution-on-cluster.md | 2 + src/Client/Connection.h | 2 + src/Client/IConnections.h | 4 ++ src/Client/MultiplexedConnections.h | 5 ++ src/Core/AntalyaProtocol.h | 4 +- .../ObjectStorages/IObjectStorage.cpp | 67 ------------------ .../ObjectStorages/IObjectStorage.h | 49 +------------ src/Interpreters/ClusterFunctionReadTask.cpp | 45 ++++++++++-- src/Interpreters/ClusterFunctionReadTask.h | 6 +- .../gtest_cluster_function_read_task.cpp | 58 ++++++++++++++++ src/QueryPipeline/RemoteQueryExecutor.cpp | 4 +- src/QueryPipeline/RemoteQueryExecutor.h | 3 +- .../StorageObjectStorageCluster.cpp | 16 +++-- .../StorageObjectStorageSource.cpp | 69 +++++++++---------- .../StorageObjectStorageSource.h | 3 + ...rageObjectStorageStableTaskDistributor.cpp | 29 +++++--- ...torageObjectStorageStableTaskDistributor.h | 15 +++- .../tests/gtest_rendezvous_hashing.cpp | 2 +- src/Storages/StorageFileCluster.cpp | 2 +- src/Storages/StorageURLCluster.cpp | 2 +- .../configs/config.d/cluster.xml | 12 ++++ .../test.py | 24 ++----- 24 files changed, 233 insertions(+), 200 deletions(-) diff --git a/docs/en/antalya/protocol.md b/docs/en/antalya/protocol.md index 2f6266b39aea..06b4885ae012 100644 --- a/docs/en/antalya/protocol.md +++ b/docs/en/antalya/protocol.md @@ -14,8 +14,8 @@ upstream ClickHouse cannot reach, defined in `src/Core/AntalyaProtocol.h`. Both the name string of their `Hello`, on every connection: ```text -client -> server "ClickHouse client (antalya:2)" -server -> client "ClickHouse (antalya:2)" +client -> server "ClickHouse client (antalya:3)" +server -> client "ClickHouse (antalya:3)" ``` Each side parses the peer's suffix, caps the value with `min(own, peer)` and keeps the result. `0` @@ -27,6 +27,10 @@ Version 1 is the advertisement itself. Version 2 appends optional Iceberg column statistics (`DataFileMetaInfo`) to `ReadTaskResponse`, after the upstream cluster-processing payload. +Version 3 appends optional `retry_after_us` after that trailer. The initiator sends it when +`lock_object_storage_task_distribution_ms` asks a worker to wait before taking another node's file. +A peer that negotiated a lower version is given the file immediately, so the task path stays an object key. + ## Adding an Antalya-only wire change {#adding-a-wire-change} - Bump `DBMS_ANTALYA_PROTOCOL_VERSION` by one and gate the change on the negotiated value. diff --git a/docs/en/interfaces/specs/NativeProtocol.md b/docs/en/interfaces/specs/NativeProtocol.md index 4bedfc12502f..d00f88a66c02 100644 --- a/docs/en/interfaces/specs/NativeProtocol.md +++ b/docs/en/interfaces/specs/NativeProtocol.md @@ -833,7 +833,7 @@ External clients that don't use SSH auth never see packets 11, 12, or 18 — the | 6 | KeepAlive | not specified | Connection keepalive | | 7 | Scalar | not specified | Scalar data block | | 8 | IgnoredPartUUIDs | not specified | Parts to exclude from query | -| 9 | ReadTaskResponse | not specified | S3 cluster read response. When both peers negotiated [Antalya protocol](/antalya/protocol) version 2 or newer, the initiator appends optional Iceberg column statistics after the upstream payload. | +| 9 | ReadTaskResponse | not specified | S3 cluster read response. When both peers negotiated [Antalya protocol](/antalya/protocol) version 2 or newer, the initiator appends optional Iceberg column statistics after the upstream payload. Version 3 or newer also appends optional `retry_after_us` for `lock_object_storage_task_distribution_ms`. | | 10 | MergeTreeReadTaskResponse | not specified | Parallel read task response | | 11 | SSHChallengeRequest | [SSH auth](#ssh-authentication) | SSH auth challenge request | | 12 | SSHChallengeResponse | [SSH auth](#ssh-authentication) | SSH auth challenge response | diff --git a/docs/en/sql-reference/distribution-on-cluster.md b/docs/en/sql-reference/distribution-on-cluster.md index 3a9835e23856..c68de6772749 100644 --- a/docs/en/sql-reference/distribution-on-cluster.md +++ b/docs/en/sql-reference/distribution-on-cluster.md @@ -16,6 +16,8 @@ This improves cache efficiency by minimizing data movement among nodes. Each node begins processing files for which it is the primary node. After completing its assigned files, a node may take tasks from other nodes, either immediately or after waiting for `lock_object_storage_task_distribution_ms` milliseconds if the primary node does not request new files during that interval. The default value of `lock_object_storage_task_distribution_ms` is 500 milliseconds. This setting balances between caching efficiency and workload redistribution when nodes are imbalanced. +The wait is a field of `ReadTaskResponse` and is sent only when the worker negotiated Antalya protocol version 3. A worker that negotiated a lower version, including an upstream node, receives the file immediately. + ## `SYSTEM STOP SWARM MODE` command If a node needs to shut down gracefully, the command `SYSTEM STOP SWARM MODE` prevents the node from receiving new tasks for *Cluster-family queries. The node finishes processing already assigned files before it can safely shut down without errors. diff --git a/src/Client/Connection.h b/src/Client/Connection.h index d2a4a5504a35..8d66462d15eb 100644 --- a/src/Client/Connection.h +++ b/src/Client/Connection.h @@ -187,6 +187,8 @@ class Connection : public IServerConnection UInt64 getParallelReplicasProtocolVersion() const { return server_parallel_replicas_protocol_version; } + UInt64 getServerAntalyaProtocolVersion() const { return server_antalya_protocol_version; } + private: String host; UInt16 port; diff --git a/src/Client/IConnections.h b/src/Client/IConnections.h index 1769c4a37a0e..4ca0fcd1a306 100644 --- a/src/Client/IConnections.h +++ b/src/Client/IConnections.h @@ -30,6 +30,10 @@ class IConnections : boost::noncopyable virtual void sendQueryPlan(const QueryPlan & query_plan) = 0; virtual void sendClusterFunctionReadTaskResponse(const ClusterFunctionReadTaskResponse &) = 0; + + /// Negotiated Antalya protocol version of the connection that produced the last packet. + /// `0` when the peer is not an Antalya build, or when this connection type has no current replica. + virtual UInt64 getServerAntalyaProtocolVersion() const { return 0; } virtual void sendMergeTreeReadTaskResponse(const ParallelReadResponse & response) = 0; /// Get packet from any replica. diff --git a/src/Client/MultiplexedConnections.h b/src/Client/MultiplexedConnections.h index 8b83eeb97e7a..bbe3f552fed0 100644 --- a/src/Client/MultiplexedConnections.h +++ b/src/Client/MultiplexedConnections.h @@ -42,6 +42,11 @@ class MultiplexedConnections final : public IConnections void sendQueryPlan(const QueryPlan & query_plan) override; void sendClusterFunctionReadTaskResponse(const ClusterFunctionReadTaskResponse & response) override; + + UInt64 getServerAntalyaProtocolVersion() const override + { + return current_connection ? current_connection->getServerAntalyaProtocolVersion() : 0; + } void sendMergeTreeReadTaskResponse(const ParallelReadResponse & response) override; Packet receivePacket() override; diff --git a/src/Core/AntalyaProtocol.h b/src/Core/AntalyaProtocol.h index c79ca44e306a..e05a8bb593be 100644 --- a/src/Core/AntalyaProtocol.h +++ b/src/Core/AntalyaProtocol.h @@ -9,8 +9,10 @@ namespace DB /// 2 — `ReadTaskResponse` carries optional `DataFileMetaInfo` after the upstream payload. static constexpr auto DBMS_ANTALYA_PROTOCOL_VERSION_WITH_DATA_FILE_META_INFO = 2; +/// 3 — `ReadTaskResponse` carries optional `retry_after_us` after the version 2 trailer. +static constexpr auto DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY = 3; /// Bump for every Antalya-only wire protocol change. See `docs/en/antalya/protocol.md`. -static constexpr auto DBMS_ANTALYA_PROTOCOL_VERSION = DBMS_ANTALYA_PROTOCOL_VERSION_WITH_DATA_FILE_META_INFO; +static constexpr auto DBMS_ANTALYA_PROTOCOL_VERSION = DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY; namespace AntalyaProtocol { diff --git a/src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.cpp b/src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.cpp index 06fb85e75440..f51f274c0ba4 100644 --- a/src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.cpp +++ b/src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.cpp @@ -11,10 +11,6 @@ #include #include -#include -#include -#include - namespace DB { @@ -160,67 +156,4 @@ RelativePathWithMetadata::RelativePathWithMetadata(const DataFileInfo & info, st file_meta_info = info.file_meta_info; } -RelativePathWithMetadata::CommandInTaskResponse::CommandInTaskResponse(const std::string & task) -{ - /// TEMPORARY WORKAROUND. This constructor runs for every string passed to `RelativePathWithMetadata`, - /// which includes every key returned by a `listObjects` / `iterate` call of every object storage, not only - /// the task-distributor answers it exists for (`{"retry_after_us": N}` from - /// `StorageObjectStorageStableTaskDistributor`). Parsing an ordinary object key as JSON throws and catches - /// one `JSONException` per listed key. Besides the cost, under ASan the fake stack frame of `parseImpl` - /// that exits by exception is never released (it is re-entered at the same stack depth, and `FakeStack::GC` - /// frees only frames strictly below the next allocation), so every listed key leaks one frame per thread; - /// once the size class is full every `__asan_stack_malloc_1` scans all 8192 slots and every small function - /// on that thread becomes ~100x slower. See https://github.com/Altinity/ClickHouse/issues/2362. - /// - /// Only try to parse strings that can be a JSON object. Object keys never start with `{`; the distributor - /// answer always does. Column statistics no longer travel in this JSON: they are an Antalya protocol - /// trailer on `ReadTaskResponse`. `retry_after_us` is still encoded here until that command moves too. - { - const auto first = task.find_first_not_of(" \t\r\n"); - if (first == std::string::npos || task[first] != '{') - return; - } - - Poco::JSON::Parser parser; - try - { - auto json = parser.parse(task).extract(); - if (!json) - return; - - is_valid = true; - - if (json->has("file_path")) - file_path = json->getValue("file_path"); - if (json->has("retry_after_us")) - retry_after_us = json->getValue("retry_after_us"); - if (json->has("meta_info")) - file_meta_info = std::make_shared(json->getObject("meta_info")); - } - catch (const Poco::JSON::JSONException &) - { /// Not a JSON - return; - } - catch (const Poco::SyntaxException &) - { /// Not a JSON - return; - } -} - -std::string RelativePathWithMetadata::CommandInTaskResponse::toString() const -{ - Poco::JSON::Object json; - if (file_path.has_value()) - json.set("file_path", file_path.value()); - if (retry_after_us.has_value()) - json.set("retry_after_us", retry_after_us.value()); - if (file_meta_info.has_value()) - json.set("meta_info", file_meta_info.value()->toJson()); - - std::ostringstream oss; - oss.exceptions(std::ios::failbit); - Poco::JSON::Stringifier::stringify(json, oss); - return oss.str(); -} - } diff --git a/src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.h b/src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.h index 436f8c05bd51..4408464dc947 100644 --- a/src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.h +++ b/src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.h @@ -121,62 +121,18 @@ struct DataLakeObjectMetadata; struct RelativePathWithMetadata { - class CommandInTaskResponse - { - public: - CommandInTaskResponse() = default; - explicit CommandInTaskResponse(const std::string & task); - - bool isValid() const { return is_valid; } - void setFilePath(const std::string & file_path_ ) - { - file_path = file_path_; - is_valid = true; - } - void setRetryAfterUs(Poco::Timestamp::TimeDiff time_us) - { - retry_after_us = time_us; - is_valid = true; - } - void setFileMetaInfo(DataFileMetaInfoPtr file_meta_info_ ) - { - file_meta_info = file_meta_info_; - is_valid = true; - } - - std::string toString() const; - - std::optional getFilePath() const { return file_path; } - std::optional getRetryAfterUs() const { return retry_after_us; } - std::optional getFileMetaInfo() const { return file_meta_info; } - - private: - bool is_valid = false; - std::optional file_path; - std::optional retry_after_us; - std::optional file_meta_info; - }; - String relative_path; /// Object metadata: size, modification time, etc. std::optional metadata; /// Information about columns std::optional file_meta_info; - /// Retry request after short pause - CommandInTaskResponse command; RelativePathWithMetadata() = default; - explicit RelativePathWithMetadata(String command_or_path, std::optional metadata_ = std::nullopt) - : relative_path(std::move(command_or_path)) + explicit RelativePathWithMetadata(String relative_path_, std::optional metadata_ = std::nullopt) + : relative_path(std::move(relative_path_)) , metadata(std::move(metadata_)) - , command(relative_path) { - if (command.isValid()) - { - relative_path = command.getFilePath().value_or(""); - file_meta_info = command.getFileMetaInfo(); - } } explicit RelativePathWithMetadata(const DataFileInfo & info, std::optional metadata_ = std::nullopt); @@ -191,7 +147,6 @@ struct RelativePathWithMetadata void setFileMetaInfo(std::optional file_meta_info_ ) { file_meta_info = file_meta_info_; } std::optional getFileMetaInfo() const { return file_meta_info; } - const CommandInTaskResponse & getCommand() const { return command; } std::string getFileNameWithoutExtension() const { return std::filesystem::path(relative_path).stem(); } }; diff --git a/src/Interpreters/ClusterFunctionReadTask.cpp b/src/Interpreters/ClusterFunctionReadTask.cpp index 0dc62e3351ab..cab05f07ff6a 100644 --- a/src/Interpreters/ClusterFunctionReadTask.cpp +++ b/src/Interpreters/ClusterFunctionReadTask.cpp @@ -44,13 +44,8 @@ ClusterFunctionReadTaskResponse::ClusterFunctionReadTaskResponse(ObjectInfoPtr o if (context->getSettingsRef()[Setting::allow_experimental_iceberg_read_optimization]) file_meta_info = object->relative_path_with_metadata.file_meta_info; - if (object->relative_path_with_metadata.getCommand().isValid()) - path = object->relative_path_with_metadata.getCommand().toString(); - else - { - const bool send_over_whole_archive = !context->getSettingsRef()[Setting::cluster_function_process_archive_on_multiple_nodes]; - path = send_over_whole_archive ? object->getPathOrPathToArchiveIfArchive() : object->getPath(); - } + const bool send_over_whole_archive = !context->getSettingsRef()[Setting::cluster_function_process_archive_on_multiple_nodes]; + path = send_over_whole_archive ? object->getPathOrPathToArchiveIfArchive() : object->getPath(); file_bucket_info = object->file_bucket_info; } @@ -133,6 +128,17 @@ void ClusterFunctionReadTaskResponse::serialize(WriteBuffer & out, size_t worker DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION_WITH_ICEBERG_METADATA); } + /// Fail closed: dropping `retry_after_us` would send an empty path, and the worker would stop reading. + if (retry_after_us && antalya_protocol_version < DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY) + { + throw Exception( + ErrorCodes::UNKNOWN_PROTOCOL, + "Worker Antalya protocol version {} cannot carry `retry_after_us`, which is required for " + "`lock_object_storage_task_distribution_ms` (minimum Antalya protocol version: {})", + antalya_protocol_version, + DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY); + } + /// Fail closed: protocol < 4 omits `file_bucket_info`, so each bucket task becomes a full-file /// read and bucket-split cluster queries return duplicated rows. if (protocol_version < DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION_WITH_FILE_BUCKETS_INFO @@ -206,6 +212,19 @@ void ClusterFunctionReadTaskResponse::serialize(WriteBuffer & out, size_t worker writeVarUInt(0, out); } } + + if (antalya_protocol_version >= DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY) + { + if (retry_after_us) + { + writeVarUInt(1, out); + writeVarUInt(*retry_after_us, out); + } + else + { + writeVarUInt(0, out); + } + } } void ClusterFunctionReadTaskResponse::deserialize(ReadBuffer & in, size_t antalya_protocol_version) @@ -267,6 +286,18 @@ void ClusterFunctionReadTaskResponse::deserialize(ReadBuffer & in, size_t antaly if (has_file_meta_info) file_meta_info = std::make_shared(DataFileMetaInfo::deserialize(in)); } + + if (antalya_protocol_version >= DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY) + { + UInt64 has_retry_after_us = 0; + readVarUInt(has_retry_after_us, in); + if (has_retry_after_us) + { + UInt64 retry = 0; + readVarUInt(retry, in); + retry_after_us = retry; + } + } } } diff --git a/src/Interpreters/ClusterFunctionReadTask.h b/src/Interpreters/ClusterFunctionReadTask.h index e6fbb06984f0..1e233775fdd0 100644 --- a/src/Interpreters/ClusterFunctionReadTask.h +++ b/src/Interpreters/ClusterFunctionReadTask.h @@ -27,13 +27,17 @@ struct ClusterFunctionReadTaskResponse std::optional iceberg_info; /// File's columns info std::optional file_meta_info; + /// Microseconds the worker should wait before asking for another task. + /// Sent only when both peers negotiated Antalya protocol version 3. + std::optional retry_after_us; /// Convert received response into ObjectInfo. ObjectInfoPtr getObjectInfo() const; /// Whether response is empty. /// It is used to identify an end of processing. - bool isEmpty() const { return path.empty(); } + /// A retry request has no path and is not the end of processing. + bool isEmpty() const { return path.empty() && !retry_after_us; } /// Serialize according to the cluster-processing protocol version. /// `antalya_protocol_version` is the negotiated Antalya version of this hop (`0` for an upstream peer). diff --git a/src/Interpreters/tests/gtest_cluster_function_read_task.cpp b/src/Interpreters/tests/gtest_cluster_function_read_task.cpp index cef49f02151b..1b4e1023a651 100644 --- a/src/Interpreters/tests/gtest_cluster_function_read_task.cpp +++ b/src/Interpreters/tests/gtest_cluster_function_read_task.cpp @@ -398,3 +398,61 @@ TEST(ClusterFunctionReadTaskResponse, AntalyaProtocolWithoutFileMetaInfoConsumes EXPECT_FALSE(deserialized.file_meta_info.has_value()); EXPECT_TRUE(in.eof()); } + +TEST(ClusterFunctionReadTaskResponse, RejectsRetryAfterUsForUpstreamPeer) +{ + ClusterFunctionReadTaskResponse response; + response.retry_after_us = 1000; + + String serialized; + WriteBufferFromString out(serialized); + try + { + response.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, /*antalya_protocol_version*/ 0); + FAIL() << "Expected exception"; + } + catch (const Exception & e) + { + EXPECT_EQ(e.code(), ErrorCodes::UNKNOWN_PROTOCOL); + EXPECT_NE(e.message().find("retry_after_us"), std::string::npos); + } +} + +TEST(ClusterFunctionReadTaskResponse, RoundTripsRetryAfterUsOnAntalyaProtocol) +{ + ClusterFunctionReadTaskResponse response; + response.retry_after_us = 2500; + EXPECT_FALSE(response.isEmpty()); + + String serialized; + WriteBufferFromString out(serialized); + response.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY); + out.finalize(); + + ReadBufferFromString in(serialized); + ClusterFunctionReadTaskResponse deserialized; + deserialized.deserialize(in, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY); + ASSERT_TRUE(deserialized.retry_after_us.has_value()); + EXPECT_EQ(*deserialized.retry_after_us, 2500); + EXPECT_TRUE(deserialized.path.empty()); + EXPECT_FALSE(deserialized.isEmpty()); + EXPECT_TRUE(in.eof()); +} + +TEST(ClusterFunctionReadTaskResponse, AntalyaProtocolWithoutRetryConsumesPresenceFlag) +{ + ClusterFunctionReadTaskResponse response; + response.path = "data/file.parquet"; + + String serialized; + WriteBufferFromString out(serialized); + response.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY); + out.finalize(); + + ReadBufferFromString in(serialized); + ClusterFunctionReadTaskResponse deserialized; + deserialized.deserialize(in, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY); + EXPECT_FALSE(deserialized.retry_after_us.has_value()); + EXPECT_EQ(deserialized.path, "data/file.parquet"); + EXPECT_TRUE(in.eof()); +} diff --git a/src/QueryPipeline/RemoteQueryExecutor.cpp b/src/QueryPipeline/RemoteQueryExecutor.cpp index 429ec6998d51..c15334dc3d29 100644 --- a/src/QueryPipeline/RemoteQueryExecutor.cpp +++ b/src/QueryPipeline/RemoteQueryExecutor.cpp @@ -773,7 +773,9 @@ void RemoteQueryExecutor::processReadTaskRequest() ProfileEvents::increment(ProfileEvents::ReadTaskRequestsReceived); - auto response = (*extension->task_iterator)(extension->replica_info->number_of_current_replica); + auto response = (*extension->task_iterator)( + extension->replica_info->number_of_current_replica, + connections->getServerAntalyaProtocolVersion()); connections->sendClusterFunctionReadTaskResponse(*response); } diff --git a/src/QueryPipeline/RemoteQueryExecutor.h b/src/QueryPipeline/RemoteQueryExecutor.h index b83afc2c51c8..e0243d02a8fd 100644 --- a/src/QueryPipeline/RemoteQueryExecutor.h +++ b/src/QueryPipeline/RemoteQueryExecutor.h @@ -49,7 +49,8 @@ class TaskIterator { throw Exception(ErrorCodes::NOT_IMPLEMENTED, "Method rescheduleTasksFromReplica is not implemented"); } - virtual ClusterFunctionReadTaskResponsePtr operator()(size_t number_of_current_replica) const = 0; + /// `antalya_protocol_version` is the negotiated version of the worker that asked for a task (`0` for an upstream peer). + virtual ClusterFunctionReadTaskResponsePtr operator()(size_t number_of_current_replica, UInt64 antalya_protocol_version = 0) const = 0; }; /// This class allows one to launch queries on remote replicas of one shard and get results diff --git a/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp b/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp index 9e3ad4677022..0729118d3647 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp @@ -9,6 +9,7 @@ #include #include +#include #include #include #include @@ -627,16 +628,23 @@ class TaskDistributor : public TaskIterator task_distributor.rescheduleTasksFromReplica(number_of_current_replica); } - ClusterFunctionReadTaskResponsePtr operator()(size_t number_of_current_replica) const override + ClusterFunctionReadTaskResponsePtr operator()(size_t number_of_current_replica, UInt64 antalya_protocol_version) const override { fiu_do_on(FailPoints::storage_cluster_read_sleep, { sleepForSeconds(10); }); - auto task = task_distributor.getNextTask(number_of_current_replica); - if (task) - return std::make_shared(std::move(task), context); + const bool worker_supports_task_retry = antalya_protocol_version >= DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY; + auto task = task_distributor.getNextTask(number_of_current_replica, worker_supports_task_retry); + if (task.retry_after_us) + { + auto response = std::make_shared(); + response->retry_after_us = task.retry_after_us; + return response; + } + if (task.object) + return std::make_shared(std::move(task.object), context); return std::make_shared(); } diff --git a/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp b/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp index 79b63f6dc1b7..f54401cfa13e 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp @@ -850,22 +850,6 @@ StorageObjectStorageSource::ReaderHolder StorageObjectStorageSource::createReade if (!object_info) return {}; - if (object_info->relative_path_with_metadata.getCommand().isValid()) - { - auto retry_after_us = object_info->relative_path_with_metadata.getCommand().getRetryAfterUs(); - if (retry_after_us.has_value()) - { - /// TODO: Make asyncronous waiting without sleep in thread - /// Now this sleep is on executor node in worker thread - /// Does not block query initiator - auto wait_time = std::min(Poco::Timestamp::TimeDiff(100000ul), retry_after_us.value()); - ProfileEvents::increment(ProfileEvents::ObjectStorageClusterWaitingMicroseconds, wait_time); - sleepForMicroseconds(wait_time); - continue; - } - object_info->relative_path_with_metadata.setFileMetaInfo(object_info->relative_path_with_metadata.getCommand().getFileMetaInfo()); - } - if (object_info->getPath().empty()) return {}; if (!object_info->getObjectMetadata()) @@ -2003,10 +1987,7 @@ StorageObjectStorageSource::ReadTaskIterator::ReadTaskIterator( for (size_t i = 0; i < max_threads_count; ++i) objects.push_back(pool_scheduler([this]() -> ObjectInfoPtr { - auto task = callback(); - if (!task || task->isEmpty()) - return nullptr; - return task->getObjectInfo(); + return takeNextObject(); }, Priority{})); pool.wait(); @@ -2015,10 +1996,7 @@ StorageObjectStorageSource::ReadTaskIterator::ReadTaskIterator( { auto object = object_future.get(); if (object) - { - resolveIcebergObjectStorageIfNeeded(object); buffer.push_back(object); - } } } @@ -2046,19 +2024,16 @@ void StorageObjectStorageSource::ReadTaskIterator::resolveIcebergObjectStorageIf #endif } -ObjectInfoPtr StorageObjectStorageSource::ReadTaskIterator::next(size_t) +ObjectInfoPtr StorageObjectStorageSource::ReadTaskIterator::takeNextObject() { - size_t current_index = index.fetch_add(1, std::memory_order_relaxed); - - ObjectInfoPtr object_info; - if (current_index >= buffer.size()) + if (!getContext()->isSwarmModeEnabled()) { - if (!getContext()->isSwarmModeEnabled()) - { - LOG_DEBUG(getLogger("StorageObjectStorageSource"), "STOP SWARM MODE called, stop getting new tasks"); - return nullptr; - } + LOG_DEBUG(getLogger("StorageObjectStorageSource"), "STOP SWARM MODE called, stop getting new tasks"); + return nullptr; + } + while (true) + { auto task = callback(); if (auto query_status = getContext()->getProcessListElement()) @@ -2067,13 +2042,35 @@ ObjectInfoPtr StorageObjectStorageSource::ReadTaskIterator::next(size_t) if (!task || task->isEmpty()) return nullptr; - object_info = task->getObjectInfo(); + if (task->retry_after_us) + { + /// The wait stays on the worker so the initiator thread is not blocked. + /// TODO: wait asynchronously instead of sleeping in this thread. + auto wait_time = std::min(100000, *task->retry_after_us); + ProfileEvents::increment(ProfileEvents::ObjectStorageClusterWaitingMicroseconds, wait_time); + sleepForMicroseconds(wait_time); + continue; + } + + auto object_info = task->getObjectInfo(); resolveIcebergObjectStorageIfNeeded(object_info); + return object_info; } - else +} + +ObjectInfoPtr StorageObjectStorageSource::ReadTaskIterator::next(size_t) +{ + size_t current_index = index.fetch_add(1, std::memory_order_relaxed); + + ObjectInfoPtr object_info; + if (current_index >= buffer.size()) { - object_info = buffer[current_index]; + object_info = takeNextObject(); + if (!object_info) + return nullptr; } + else + object_info = buffer[current_index]; if (!is_archive) return object_info; diff --git a/src/Storages/ObjectStorage/StorageObjectStorageSource.h b/src/Storages/ObjectStorage/StorageObjectStorageSource.h index 29de8bb7a99b..04bac2473f1a 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageSource.h +++ b/src/Storages/ObjectStorage/StorageObjectStorageSource.h @@ -191,6 +191,9 @@ class StorageObjectStorageSource::ReadTaskIterator : public IObjectIterator, pri private: ObjectInfoPtr createObjectInfoInArchive(const std::string & path_to_archive, const std::string & path_in_archive); + /// Fetch the next object, waiting and asking again when the initiator sends `retry_after_us`. + ObjectInfoPtr takeNextObject(); + /// For Iceberg objects: resolve which storage the file lives in (possibly a secondary storage) /// from the raw metadata path and record it on the object. No-op for non-Iceberg objects. void resolveIcebergObjectStorageIfNeeded(const ObjectInfoPtr & object); diff --git a/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.cpp b/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.cpp index 44f3d01d1d4b..f20818f5ed53 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.cpp @@ -41,7 +41,7 @@ StorageObjectStorageStableTaskDistributor::StorageObjectStorageStableTaskDistrib } } -ObjectInfoPtr StorageObjectStorageStableTaskDistributor::getNextTask(size_t number_of_current_replica) +DistributedReadTask StorageObjectStorageStableTaskDistributor::getNextTask(size_t number_of_current_replica, bool worker_supports_task_retry) { LOG_TRACE(log, "Received request from replica {} looking for a file", number_of_current_replica); @@ -65,7 +65,12 @@ ObjectInfoPtr StorageObjectStorageStableTaskDistributor::getNextTask(size_t numb file = getMatchingFileFromIterator(number_of_current_replica); // 3. Process unprocessed files if iterator is exhausted if (!file) - file = getAnyUnprocessedFile(number_of_current_replica); + { + auto pending = getAnyUnprocessedFile(number_of_current_replica, worker_supports_task_retry); + if (pending.retry_after_us) + return pending; + file = std::move(pending.object); + } if (file) { @@ -82,7 +87,7 @@ ObjectInfoPtr StorageObjectStorageStableTaskDistributor::getNextTask(size_t numb processed_file_list_ptr->second.push_back(file); } - return file; + return DistributedReadTask{.object = std::move(file), .retry_after_us = std::nullopt}; } size_t StorageObjectStorageStableTaskDistributor::getReplicaForFile(const String & file_path) @@ -221,7 +226,9 @@ ObjectInfoPtr StorageObjectStorageStableTaskDistributor::getMatchingFileFromIter return {}; } -ObjectInfoPtr StorageObjectStorageStableTaskDistributor::getAnyUnprocessedFile(size_t number_of_current_replica) +DistributedReadTask StorageObjectStorageStableTaskDistributor::getAnyUnprocessedFile( + size_t number_of_current_replica, + bool worker_supports_task_retry) { /// Limit time of node activity to keep task in queue Poco::Timestamp activity_limit; @@ -240,6 +247,7 @@ ObjectInfoPtr StorageObjectStorageStableTaskDistributor::getAnyUnprocessedFile(s auto number_of_matched_replica = it->second.second; auto last_activity = last_node_activity.find(number_of_matched_replica); if (lock_object_storage_task_distribution_us <= 0 // file deferring is turned off + || !worker_supports_task_retry // peer cannot be asked to wait || it->second.second == number_of_current_replica // file is matching with current replica || last_activity == last_node_activity.end() // msut never be happen, last_activity is filled for each replica on start || activity_limit > last_activity->second) // matched replica did not ask for a new files for a while @@ -257,7 +265,7 @@ ObjectInfoPtr StorageObjectStorageStableTaskDistributor::getAnyUnprocessedFile(s ); ProfileEvents::increment(ProfileEvents::ObjectStorageClusterSentToNonMatchedReplica); - return next_file; + return DistributedReadTask{.object = std::move(next_file), .retry_after_us = std::nullopt}; } oldest_activity = std::min(oldest_activity, last_activity->second); @@ -271,11 +279,12 @@ ObjectInfoPtr StorageObjectStorageStableTaskDistributor::getAnyUnprocessedFile(s oldest_activity - activity_limit ); - /// All unprocessed files owned by alive replicas with recenlty activity - /// Need to retry after (oldest_activity - activity_limit) microseconds - RelativePathWithMetadata::CommandInTaskResponse response; - response.setRetryAfterUs(oldest_activity - activity_limit); - return std::make_shared(response.toString()); + /// All unprocessed files owned by alive replicas with recently activity. + /// The worker waits `(oldest_activity - activity_limit)` microseconds, then asks again. + auto retry_after_us = oldest_activity - activity_limit; + if (retry_after_us < 0) + retry_after_us = 0; + return DistributedReadTask{.object = nullptr, .retry_after_us = static_cast(retry_after_us)}; } return {}; diff --git a/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h b/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h index 652bbb6a03e0..2b9acd794539 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h +++ b/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h @@ -14,10 +14,19 @@ #include #include #include +#include namespace DB { +struct DistributedReadTask +{ + ObjectInfoPtr object; + /// Microseconds the worker should wait before asking again. + /// Set only when `object` is empty and the worker understands Antalya protocol version 3. + std::optional retry_after_us; +}; + class StorageObjectStorageStableTaskDistributor { public: @@ -27,7 +36,9 @@ class StorageObjectStorageStableTaskDistributor bool send_over_whole_archive_, uint64_t lock_object_storage_task_distribution_ms_); - ObjectInfoPtr getNextTask(size_t number_of_current_replica); + /// `worker_supports_task_retry` is false for an upstream peer. That peer is given a file immediately, + /// because it cannot be asked to wait without putting a command into the object path. + DistributedReadTask getNextTask(size_t number_of_current_replica, bool worker_supports_task_retry = true); /// Insert objects back to unprocessed files void rescheduleTasksFromReplica(size_t number_of_current_replica); @@ -36,7 +47,7 @@ class StorageObjectStorageStableTaskDistributor size_t getReplicaForFile(const String & file_path); ObjectInfoPtr getPreQueuedFile(size_t number_of_current_replica); ObjectInfoPtr getMatchingFileFromIterator(size_t number_of_current_replica); - ObjectInfoPtr getAnyUnprocessedFile(size_t number_of_current_replica); + DistributedReadTask getAnyUnprocessedFile(size_t number_of_current_replica, bool worker_supports_task_retry); void saveLastNodeActivity(size_t number_of_current_replica); diff --git a/src/Storages/ObjectStorage/tests/gtest_rendezvous_hashing.cpp b/src/Storages/ObjectStorage/tests/gtest_rendezvous_hashing.cpp index 244a6a1f210f..3b6b54d39e2d 100644 --- a/src/Storages/ObjectStorage/tests/gtest_rendezvous_hashing.cpp +++ b/src/Storages/ObjectStorage/tests/gtest_rendezvous_hashing.cpp @@ -53,7 +53,7 @@ namespace { for (size_t i = 0; i < n; ++i) { - ObjectInfoPtr object = distributor.getNextTask(replica_id); + ObjectInfoPtr object = distributor.getNextTask(replica_id).object; if (!object) return false; auto path = object->getPath(); diff --git a/src/Storages/StorageFileCluster.cpp b/src/Storages/StorageFileCluster.cpp index 36494b16991e..797f86acbfc7 100644 --- a/src/Storages/StorageFileCluster.cpp +++ b/src/Storages/StorageFileCluster.cpp @@ -126,7 +126,7 @@ class FileTaskIterator : public TaskIterator ~FileTaskIterator() override = default; - ClusterFunctionReadTaskResponsePtr operator()(size_t /* number_of_current_replica */) const override + ClusterFunctionReadTaskResponsePtr operator()(size_t /* number_of_current_replica */, UInt64 /* antalya_protocol_version */) const override { auto file = iterator.next(); if (file.empty()) diff --git a/src/Storages/StorageURLCluster.cpp b/src/Storages/StorageURLCluster.cpp index 2a2437778614..1ba03ce76248 100644 --- a/src/Storages/StorageURLCluster.cpp +++ b/src/Storages/StorageURLCluster.cpp @@ -155,7 +155,7 @@ class UrlTaskIterator : public TaskIterator ~UrlTaskIterator() override = default; - ClusterFunctionReadTaskResponsePtr operator()(size_t /* number_of_current_replica */) const override + ClusterFunctionReadTaskResponsePtr operator()(size_t /* number_of_current_replica */, UInt64 /* antalya_protocol_version */) const override { auto url = iterator.next(); if (url.empty()) diff --git a/tests/integration/test_iceberg_mixed_antalya_upstream/configs/config.d/cluster.xml b/tests/integration/test_iceberg_mixed_antalya_upstream/configs/config.d/cluster.xml index 273ec80e1c58..f01efc87a639 100644 --- a/tests/integration/test_iceberg_mixed_antalya_upstream/configs/config.d/cluster.xml +++ b/tests/integration/test_iceberg_mixed_antalya_upstream/configs/config.d/cluster.xml @@ -16,5 +16,17 @@ + + + + antalya + 9000 + + + upstream + 9000 + + + diff --git a/tests/integration/test_iceberg_mixed_antalya_upstream/test.py b/tests/integration/test_iceberg_mixed_antalya_upstream/test.py index 779c3e37a7cb..bf0531478c4a 100644 --- a/tests/integration/test_iceberg_mixed_antalya_upstream/test.py +++ b/tests/integration/test_iceberg_mixed_antalya_upstream/test.py @@ -51,30 +51,20 @@ def started_cluster(): @pytest.mark.parametrize( - ("initiator_name", "cluster_name", "settings"), + ("initiator_name", "cluster_name"), [ - pytest.param( - "antalya", - "upstream_only", - # Default 500 makes the Antalya initiator send a JSON retry command in the task - # path. An upstream worker would open that text as an object key. - {"lock_object_storage_task_distribution_ms": 0}, - id="antalya_initiator_upstream_worker", - ), - pytest.param( - "upstream", - "antalya_only", - None, - id="upstream_initiator_antalya_worker", - ), + pytest.param("antalya", "upstream_only", id="antalya_initiator_upstream_worker"), + pytest.param("upstream", "antalya_only", id="upstream_initiator_antalya_worker"), + # Both workers, default `lock_object_storage_task_distribution_ms`. The Antalya initiator + # must not put a retry command into the object path: the upstream worker would open that text as a key. + pytest.param("antalya", "mixed", id="antalya_initiator_mixed_workers"), ], ) -def test_mixed_iceberg_cluster_read(started_cluster, initiator_name, cluster_name, settings): +def test_mixed_iceberg_cluster_read(started_cluster, initiator_name, cluster_name): expected = antalya.query(f"SELECT id, tag FROM icebergS3('{TABLE_URL}', '{minio_access_key}', '{minio_secret_key}') ORDER BY id") initiator = started_cluster.instances[initiator_name] got = initiator.query( f"SELECT id, tag FROM icebergS3Cluster('{cluster_name}', '{TABLE_URL}', '{minio_access_key}', '{minio_secret_key}') ORDER BY id", - settings=settings, ) assert got == expected assert got == "1\talpha\n2\tbeta\n" From f475189edf513d0272d3b82ec191e3a81389ee7b Mon Sep 17 00:00:00 2001 From: Anton Ivashkin Date: Tue, 6 Oct 2026 11:34:03 +0200 Subject: [PATCH 2/6] Ignore task-distribution lock when a worker cannot wait. If any replica did not negotiate Antalya protocol version 3, `lock_object_storage_task_distribution_ms` is skipped for the query and files are assigned immediately. Co-authored-by: Cursor --- docs/en/antalya/protocol.md | 3 +- docs/en/interfaces/specs/NativeProtocol.md | 2 +- .../sql-reference/distribution-on-cluster.md | 2 +- src/Interpreters/ClusterFunctionReadTask.cpp | 11 ------ .../gtest_cluster_function_read_task.cpp | 34 +++++++++++++------ ...rageObjectStorageStableTaskDistributor.cpp | 17 ++++++++-- ...torageObjectStorageStableTaskDistributor.h | 10 ++++-- 7 files changed, 48 insertions(+), 31 deletions(-) diff --git a/docs/en/antalya/protocol.md b/docs/en/antalya/protocol.md index 06b4885ae012..852037b3cd3f 100644 --- a/docs/en/antalya/protocol.md +++ b/docs/en/antalya/protocol.md @@ -29,7 +29,8 @@ after the upstream cluster-processing payload. Version 3 appends optional `retry_after_us` after that trailer. The initiator sends it when `lock_object_storage_task_distribution_ms` asks a worker to wait before taking another node's file. -A peer that negotiated a lower version is given the file immediately, so the task path stays an object key. +If any worker negotiated a lower version, the initiator ignores that setting for the query and +assigns files immediately, so every task path stays an object key. ## Adding an Antalya-only wire change {#adding-a-wire-change} diff --git a/docs/en/interfaces/specs/NativeProtocol.md b/docs/en/interfaces/specs/NativeProtocol.md index d00f88a66c02..26fa37db1c88 100644 --- a/docs/en/interfaces/specs/NativeProtocol.md +++ b/docs/en/interfaces/specs/NativeProtocol.md @@ -833,7 +833,7 @@ External clients that don't use SSH auth never see packets 11, 12, or 18 — the | 6 | KeepAlive | not specified | Connection keepalive | | 7 | Scalar | not specified | Scalar data block | | 8 | IgnoredPartUUIDs | not specified | Parts to exclude from query | -| 9 | ReadTaskResponse | not specified | S3 cluster read response. When both peers negotiated [Antalya protocol](/antalya/protocol) version 2 or newer, the initiator appends optional Iceberg column statistics after the upstream payload. Version 3 or newer also appends optional `retry_after_us` for `lock_object_storage_task_distribution_ms`. | +| 9 | ReadTaskResponse | not specified | S3 cluster read response. When both peers negotiated [Antalya protocol](/antalya/protocol) version 2 or newer, the initiator appends optional Iceberg column statistics after the upstream payload. Version 3 or newer also appends optional `retry_after_us` for `lock_object_storage_task_distribution_ms`. If any worker negotiated a lower Antalya version, that field is omitted and the setting is ignored for the query. | | 10 | MergeTreeReadTaskResponse | not specified | Parallel read task response | | 11 | SSHChallengeRequest | [SSH auth](#ssh-authentication) | SSH auth challenge request | | 12 | SSHChallengeResponse | [SSH auth](#ssh-authentication) | SSH auth challenge response | diff --git a/docs/en/sql-reference/distribution-on-cluster.md b/docs/en/sql-reference/distribution-on-cluster.md index c68de6772749..0304bd7c7f8e 100644 --- a/docs/en/sql-reference/distribution-on-cluster.md +++ b/docs/en/sql-reference/distribution-on-cluster.md @@ -16,7 +16,7 @@ This improves cache efficiency by minimizing data movement among nodes. Each node begins processing files for which it is the primary node. After completing its assigned files, a node may take tasks from other nodes, either immediately or after waiting for `lock_object_storage_task_distribution_ms` milliseconds if the primary node does not request new files during that interval. The default value of `lock_object_storage_task_distribution_ms` is 500 milliseconds. This setting balances between caching efficiency and workload redistribution when nodes are imbalanced. -The wait is a field of `ReadTaskResponse` and is sent only when the worker negotiated Antalya protocol version 3. A worker that negotiated a lower version, including an upstream node, receives the file immediately. +The wait is a field of `ReadTaskResponse` and is sent only when every worker negotiated Antalya protocol version 3. If any worker negotiated a lower version, including an upstream node, the setting is ignored for that query and files are assigned immediately. ## `SYSTEM STOP SWARM MODE` command diff --git a/src/Interpreters/ClusterFunctionReadTask.cpp b/src/Interpreters/ClusterFunctionReadTask.cpp index cab05f07ff6a..44e18b50894c 100644 --- a/src/Interpreters/ClusterFunctionReadTask.cpp +++ b/src/Interpreters/ClusterFunctionReadTask.cpp @@ -128,17 +128,6 @@ void ClusterFunctionReadTaskResponse::serialize(WriteBuffer & out, size_t worker DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION_WITH_ICEBERG_METADATA); } - /// Fail closed: dropping `retry_after_us` would send an empty path, and the worker would stop reading. - if (retry_after_us && antalya_protocol_version < DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY) - { - throw Exception( - ErrorCodes::UNKNOWN_PROTOCOL, - "Worker Antalya protocol version {} cannot carry `retry_after_us`, which is required for " - "`lock_object_storage_task_distribution_ms` (minimum Antalya protocol version: {})", - antalya_protocol_version, - DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY); - } - /// Fail closed: protocol < 4 omits `file_bucket_info`, so each bucket task becomes a full-file /// read and bucket-split cluster queries return duplicated rows. if (protocol_version < DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION_WITH_FILE_BUCKETS_INFO diff --git a/src/Interpreters/tests/gtest_cluster_function_read_task.cpp b/src/Interpreters/tests/gtest_cluster_function_read_task.cpp index 1b4e1023a651..c0b9007a7762 100644 --- a/src/Interpreters/tests/gtest_cluster_function_read_task.cpp +++ b/src/Interpreters/tests/gtest_cluster_function_read_task.cpp @@ -399,23 +399,35 @@ TEST(ClusterFunctionReadTaskResponse, AntalyaProtocolWithoutFileMetaInfoConsumes EXPECT_TRUE(in.eof()); } -TEST(ClusterFunctionReadTaskResponse, RejectsRetryAfterUsForUpstreamPeer) +TEST(ClusterFunctionReadTaskResponse, OmitsRetryAfterUsForUpstreamPeer) { - ClusterFunctionReadTaskResponse response; - response.retry_after_us = 1000; + ClusterFunctionReadTaskResponse with_retry; + with_retry.path = "data/file.parquet"; + with_retry.retry_after_us = 1000; - String serialized; - WriteBufferFromString out(serialized); - try + ClusterFunctionReadTaskResponse without_retry; + without_retry.path = with_retry.path; + + String with_retry_bytes; + String without_retry_bytes; { - response.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, /*antalya_protocol_version*/ 0); - FAIL() << "Expected exception"; + WriteBufferFromString out(with_retry_bytes); + with_retry.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, /*antalya_protocol_version*/ 0); + out.finalize(); } - catch (const Exception & e) { - EXPECT_EQ(e.code(), ErrorCodes::UNKNOWN_PROTOCOL); - EXPECT_NE(e.message().find("retry_after_us"), std::string::npos); + WriteBufferFromString out(without_retry_bytes); + without_retry.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, /*antalya_protocol_version*/ 0); + out.finalize(); } + EXPECT_EQ(with_retry_bytes, without_retry_bytes); + + ReadBufferFromString in(with_retry_bytes); + ClusterFunctionReadTaskResponse deserialized; + deserialized.deserialize(in, /*antalya_protocol_version*/ 0); + EXPECT_FALSE(deserialized.retry_after_us.has_value()); + EXPECT_EQ(deserialized.path, "data/file.parquet"); + EXPECT_TRUE(in.eof()); } TEST(ClusterFunctionReadTaskResponse, RoundTripsRetryAfterUsOnAntalyaProtocol) diff --git a/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.cpp b/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.cpp index f20818f5ed53..5d144bfca4e1 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.cpp @@ -45,6 +45,17 @@ DistributedReadTask StorageObjectStorageStableTaskDistributor::getNextTask(size_ { LOG_TRACE(log, "Received request from replica {} looking for a file", number_of_current_replica); + /// One peer that cannot be asked to wait makes the lock useless: that peer would take another + /// node's file immediately. Run the rest of the query as if the setting were 0. + if (!worker_supports_task_retry && !task_retry_disabled.exchange(true, std::memory_order_relaxed)) + { + LOG_INFO( + log, + "Replica {} does not support task retry, ignoring `lock_object_storage_task_distribution_ms` for this query", + number_of_current_replica); + } + const bool allow_task_retry = worker_supports_task_retry && !task_retry_disabled.load(std::memory_order_relaxed); + saveLastNodeActivity(number_of_current_replica); { @@ -66,7 +77,7 @@ DistributedReadTask StorageObjectStorageStableTaskDistributor::getNextTask(size_ // 3. Process unprocessed files if iterator is exhausted if (!file) { - auto pending = getAnyUnprocessedFile(number_of_current_replica, worker_supports_task_retry); + auto pending = getAnyUnprocessedFile(number_of_current_replica, allow_task_retry); if (pending.retry_after_us) return pending; file = std::move(pending.object); @@ -228,7 +239,7 @@ ObjectInfoPtr StorageObjectStorageStableTaskDistributor::getMatchingFileFromIter DistributedReadTask StorageObjectStorageStableTaskDistributor::getAnyUnprocessedFile( size_t number_of_current_replica, - bool worker_supports_task_retry) + bool allow_task_retry) { /// Limit time of node activity to keep task in queue Poco::Timestamp activity_limit; @@ -247,7 +258,7 @@ DistributedReadTask StorageObjectStorageStableTaskDistributor::getAnyUnprocessed auto number_of_matched_replica = it->second.second; auto last_activity = last_node_activity.find(number_of_matched_replica); if (lock_object_storage_task_distribution_us <= 0 // file deferring is turned off - || !worker_supports_task_retry // peer cannot be asked to wait + || !allow_task_retry // this request must not wait: peer is old, or the query already saw one || it->second.second == number_of_current_replica // file is matching with current replica || last_activity == last_node_activity.end() // msut never be happen, last_activity is filled for each replica on start || activity_limit > last_activity->second) // matched replica did not ask for a new files for a while diff --git a/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h b/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h index 2b9acd794539..52859ba8d27a 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h +++ b/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h @@ -15,6 +15,7 @@ #include #include #include +#include namespace DB { @@ -36,8 +37,9 @@ class StorageObjectStorageStableTaskDistributor bool send_over_whole_archive_, uint64_t lock_object_storage_task_distribution_ms_); - /// `worker_supports_task_retry` is false for an upstream peer. That peer is given a file immediately, - /// because it cannot be asked to wait without putting a command into the object path. + /// `worker_supports_task_retry` is false for a peer that did not negotiate Antalya protocol version 3. + /// The first such peer disables `lock_object_storage_task_distribution_ms` for the rest of the query: + /// every replica is then given a file immediately, with no `retry_after_us`. DistributedReadTask getNextTask(size_t number_of_current_replica, bool worker_supports_task_retry = true); /// Insert objects back to unprocessed files @@ -47,7 +49,7 @@ class StorageObjectStorageStableTaskDistributor size_t getReplicaForFile(const String & file_path); ObjectInfoPtr getPreQueuedFile(size_t number_of_current_replica); ObjectInfoPtr getMatchingFileFromIterator(size_t number_of_current_replica); - DistributedReadTask getAnyUnprocessedFile(size_t number_of_current_replica, bool worker_supports_task_retry); + DistributedReadTask getAnyUnprocessedFile(size_t number_of_current_replica, bool allow_task_retry); void saveLastNodeActivity(size_t number_of_current_replica); @@ -67,6 +69,8 @@ class StorageObjectStorageStableTaskDistributor std::mutex mutex; bool iterator_exhausted = false; + /// Set when any worker in this query cannot carry `retry_after_us`. Further tasks skip the lock. + std::atomic task_retry_disabled{false}; LoggerPtr log = getLogger("StorageClusterTaskDistributor"); }; From 11d8b1dcc2942a1ab11dfce2833038f9fe3e091f Mon Sep 17 00:00:00 2001 From: Anton Ivashkin Date: Tue, 6 Oct 2026 14:20:11 +0200 Subject: [PATCH 3/6] Do not wait for a task retry during reader prefetch. A `retry_after_us` answer in `ReadTaskIterator` construction is skipped so tasks already fetched by other threads can start. The wait happens later, in the reader that asks for another file. Co-authored-by: Cursor --- .../ObjectStorage/StorageObjectStorageSource.cpp | 9 ++++++--- src/Storages/ObjectStorage/StorageObjectStorageSource.h | 6 ++++-- 2 files changed, 10 insertions(+), 5 deletions(-) diff --git a/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp b/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp index f54401cfa13e..f801068dee48 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp @@ -1987,7 +1987,7 @@ StorageObjectStorageSource::ReadTaskIterator::ReadTaskIterator( for (size_t i = 0; i < max_threads_count; ++i) objects.push_back(pool_scheduler([this]() -> ObjectInfoPtr { - return takeNextObject(); + return takeNextObject(/*wait_for_retry=*/ false); }, Priority{})); pool.wait(); @@ -2024,7 +2024,7 @@ void StorageObjectStorageSource::ReadTaskIterator::resolveIcebergObjectStorageIf #endif } -ObjectInfoPtr StorageObjectStorageSource::ReadTaskIterator::takeNextObject() +ObjectInfoPtr StorageObjectStorageSource::ReadTaskIterator::takeNextObject(bool wait_for_retry) { if (!getContext()->isSwarmModeEnabled()) { @@ -2044,6 +2044,9 @@ ObjectInfoPtr StorageObjectStorageSource::ReadTaskIterator::takeNextObject() if (task->retry_after_us) { + if (!wait_for_retry) + return nullptr; + /// The wait stays on the worker so the initiator thread is not blocked. /// TODO: wait asynchronously instead of sleeping in this thread. auto wait_time = std::min(100000, *task->retry_after_us); @@ -2065,7 +2068,7 @@ ObjectInfoPtr StorageObjectStorageSource::ReadTaskIterator::next(size_t) ObjectInfoPtr object_info; if (current_index >= buffer.size()) { - object_info = takeNextObject(); + object_info = takeNextObject(/*wait_for_retry=*/ true); if (!object_info) return nullptr; } diff --git a/src/Storages/ObjectStorage/StorageObjectStorageSource.h b/src/Storages/ObjectStorage/StorageObjectStorageSource.h index 04bac2473f1a..d6d801bd3f74 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageSource.h +++ b/src/Storages/ObjectStorage/StorageObjectStorageSource.h @@ -191,8 +191,10 @@ class StorageObjectStorageSource::ReadTaskIterator : public IObjectIterator, pri private: ObjectInfoPtr createObjectInfoInArchive(const std::string & path_to_archive, const std::string & path_in_archive); - /// Fetch the next object, waiting and asking again when the initiator sends `retry_after_us`. - ObjectInfoPtr takeNextObject(); + /// Fetch the next object. When `wait_for_retry` is set, a `retry_after_us` response is waited out + /// on this thread and another task is requested. The constructor prefetch passes false so one + /// waiting slot cannot hold back tasks the other slots already received. + ObjectInfoPtr takeNextObject(bool wait_for_retry); /// For Iceberg objects: resolve which storage the file lives in (possibly a secondary storage) /// from the raw metadata path and record it on the object. No-op for non-Iceberg objects. From af45bd6e8d3f8bb389c42d494b7763736fc3a574 Mon Sep 17 00:00:00 2001 From: Anton Ivashkin Date: Tue, 6 Oct 2026 14:22:19 +0200 Subject: [PATCH 4/6] Recheck swarm mode before each task retry request. `SYSTEM STOP SWARM MODE` during a `retry_after_us` wait in `takeNextObject` stops further task requests. Co-authored-by: Cursor --- .../ObjectStorage/StorageObjectStorageSource.cpp | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp b/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp index f801068dee48..ac210b46a7db 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageSource.cpp @@ -2026,14 +2026,16 @@ void StorageObjectStorageSource::ReadTaskIterator::resolveIcebergObjectStorageIf ObjectInfoPtr StorageObjectStorageSource::ReadTaskIterator::takeNextObject(bool wait_for_retry) { - if (!getContext()->isSwarmModeEnabled()) - { - LOG_DEBUG(getLogger("StorageObjectStorageSource"), "STOP SWARM MODE called, stop getting new tasks"); - return nullptr; - } - while (true) { + /// Recheck after each `retry_after_us` wait. `SYSTEM STOP SWARM MODE` during that wait + /// must not be followed by another task request. + if (!getContext()->isSwarmModeEnabled()) + { + LOG_DEBUG(getLogger("StorageObjectStorageSource"), "STOP SWARM MODE called, stop getting new tasks"); + return nullptr; + } + auto task = callback(); if (auto query_status = getContext()->getProcessListElement()) From 6d0d320c08a5e7eb6b357ea6f561efa67eb6b29c Mon Sep 17 00:00:00 2001 From: Anton Ivashkin Date: Tue, 6 Oct 2026 17:34:48 +0200 Subject: [PATCH 5/6] Keep `retry_after_us` on Antalya protocol version 2. Both trailers ship in the same release, so they share `DBMS_ANTALYA_PROTOCOL_VERSION_WITH_READ_TASK_EXTENSIONS`. Co-authored-by: Cursor --- docs/en/antalya/protocol.md | 16 +++++++--------- docs/en/interfaces/specs/NativeProtocol.md | 2 +- docs/en/sql-reference/distribution-on-cluster.md | 2 +- src/Core/AntalyaProtocol.h | 9 ++++----- src/Interpreters/ClusterFunctionReadTask.cpp | 10 ++-------- src/Interpreters/ClusterFunctionReadTask.h | 2 +- .../tests/gtest_cluster_function_read_task.cpp | 16 ++++++++-------- .../StorageObjectStorageCluster.cpp | 2 +- .../StorageObjectStorageStableTaskDistributor.h | 4 ++-- 9 files changed, 27 insertions(+), 36 deletions(-) diff --git a/docs/en/antalya/protocol.md b/docs/en/antalya/protocol.md index 852037b3cd3f..177a274dc447 100644 --- a/docs/en/antalya/protocol.md +++ b/docs/en/antalya/protocol.md @@ -14,8 +14,8 @@ upstream ClickHouse cannot reach, defined in `src/Core/AntalyaProtocol.h`. Both the name string of their `Hello`, on every connection: ```text -client -> server "ClickHouse client (antalya:3)" -server -> client "ClickHouse (antalya:3)" +client -> server "ClickHouse client (antalya:2)" +server -> client "ClickHouse (antalya:2)" ``` Each side parses the peer's suffix, caps the value with `min(own, peer)` and keeps the result. `0` @@ -24,13 +24,11 @@ worker and worker to worker negotiate independently. Version 1 is the advertisement itself. -Version 2 appends optional Iceberg column statistics (`DataFileMetaInfo`) to `ReadTaskResponse`, -after the upstream cluster-processing payload. - -Version 3 appends optional `retry_after_us` after that trailer. The initiator sends it when -`lock_object_storage_task_distribution_ms` asks a worker to wait before taking another node's file. -If any worker negotiated a lower version, the initiator ignores that setting for the query and -assigns files immediately, so every task path stays an object key. +Version 2 appends two optional fields to `ReadTaskResponse`, after the upstream cluster-processing +payload. The first is Iceberg column statistics (`DataFileMetaInfo`). The second is `retry_after_us`, +sent when `lock_object_storage_task_distribution_ms` asks a worker to wait before taking another +node's file. If any worker negotiated a lower version, the initiator ignores that setting for the +query and assigns files immediately, so every task path stays an object key. ## Adding an Antalya-only wire change {#adding-a-wire-change} diff --git a/docs/en/interfaces/specs/NativeProtocol.md b/docs/en/interfaces/specs/NativeProtocol.md index 26fa37db1c88..c1eb12703fb7 100644 --- a/docs/en/interfaces/specs/NativeProtocol.md +++ b/docs/en/interfaces/specs/NativeProtocol.md @@ -833,7 +833,7 @@ External clients that don't use SSH auth never see packets 11, 12, or 18 — the | 6 | KeepAlive | not specified | Connection keepalive | | 7 | Scalar | not specified | Scalar data block | | 8 | IgnoredPartUUIDs | not specified | Parts to exclude from query | -| 9 | ReadTaskResponse | not specified | S3 cluster read response. When both peers negotiated [Antalya protocol](/antalya/protocol) version 2 or newer, the initiator appends optional Iceberg column statistics after the upstream payload. Version 3 or newer also appends optional `retry_after_us` for `lock_object_storage_task_distribution_ms`. If any worker negotiated a lower Antalya version, that field is omitted and the setting is ignored for the query. | +| 9 | ReadTaskResponse | not specified | S3 cluster read response. When both peers negotiated [Antalya protocol](/antalya/protocol) version 2 or newer, the initiator appends optional Iceberg column statistics and then optional `retry_after_us` for `lock_object_storage_task_distribution_ms` after the upstream payload. If any worker negotiated a lower Antalya version, `retry_after_us` is omitted and that setting is ignored for the query. | | 10 | MergeTreeReadTaskResponse | not specified | Parallel read task response | | 11 | SSHChallengeRequest | [SSH auth](#ssh-authentication) | SSH auth challenge request | | 12 | SSHChallengeResponse | [SSH auth](#ssh-authentication) | SSH auth challenge response | diff --git a/docs/en/sql-reference/distribution-on-cluster.md b/docs/en/sql-reference/distribution-on-cluster.md index 0304bd7c7f8e..7405ed790693 100644 --- a/docs/en/sql-reference/distribution-on-cluster.md +++ b/docs/en/sql-reference/distribution-on-cluster.md @@ -16,7 +16,7 @@ This improves cache efficiency by minimizing data movement among nodes. Each node begins processing files for which it is the primary node. After completing its assigned files, a node may take tasks from other nodes, either immediately or after waiting for `lock_object_storage_task_distribution_ms` milliseconds if the primary node does not request new files during that interval. The default value of `lock_object_storage_task_distribution_ms` is 500 milliseconds. This setting balances between caching efficiency and workload redistribution when nodes are imbalanced. -The wait is a field of `ReadTaskResponse` and is sent only when every worker negotiated Antalya protocol version 3. If any worker negotiated a lower version, including an upstream node, the setting is ignored for that query and files are assigned immediately. +The wait is a field of `ReadTaskResponse` and is sent only when every worker negotiated Antalya protocol version 2. If any worker negotiated a lower version, including an upstream node, the setting is ignored for that query and files are assigned immediately. ## `SYSTEM STOP SWARM MODE` command diff --git a/src/Core/AntalyaProtocol.h b/src/Core/AntalyaProtocol.h index e05a8bb593be..754a83ebc898 100644 --- a/src/Core/AntalyaProtocol.h +++ b/src/Core/AntalyaProtocol.h @@ -7,12 +7,11 @@ namespace DB { -/// 2 — `ReadTaskResponse` carries optional `DataFileMetaInfo` after the upstream payload. -static constexpr auto DBMS_ANTALYA_PROTOCOL_VERSION_WITH_DATA_FILE_META_INFO = 2; -/// 3 — `ReadTaskResponse` carries optional `retry_after_us` after the version 2 trailer. -static constexpr auto DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY = 3; +/// 2 — `ReadTaskResponse` carries optional `DataFileMetaInfo`, then optional `retry_after_us`, +/// after the upstream payload. Both fields ship in the same release. +static constexpr auto DBMS_ANTALYA_PROTOCOL_VERSION_WITH_READ_TASK_EXTENSIONS = 2; /// Bump for every Antalya-only wire protocol change. See `docs/en/antalya/protocol.md`. -static constexpr auto DBMS_ANTALYA_PROTOCOL_VERSION = DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY; +static constexpr auto DBMS_ANTALYA_PROTOCOL_VERSION = DBMS_ANTALYA_PROTOCOL_VERSION_WITH_READ_TASK_EXTENSIONS; namespace AntalyaProtocol { diff --git a/src/Interpreters/ClusterFunctionReadTask.cpp b/src/Interpreters/ClusterFunctionReadTask.cpp index 44e18b50894c..c76b60f1d3bb 100644 --- a/src/Interpreters/ClusterFunctionReadTask.cpp +++ b/src/Interpreters/ClusterFunctionReadTask.cpp @@ -189,7 +189,7 @@ void ClusterFunctionReadTaskResponse::serialize(WriteBuffer & out, size_t worker } } - if (antalya_protocol_version >= DBMS_ANTALYA_PROTOCOL_VERSION_WITH_DATA_FILE_META_INFO) + if (antalya_protocol_version >= DBMS_ANTALYA_PROTOCOL_VERSION_WITH_READ_TASK_EXTENSIONS) { if (file_meta_info && *file_meta_info) { @@ -200,10 +200,7 @@ void ClusterFunctionReadTaskResponse::serialize(WriteBuffer & out, size_t worker { writeVarUInt(0, out); } - } - if (antalya_protocol_version >= DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY) - { if (retry_after_us) { writeVarUInt(1, out); @@ -268,16 +265,13 @@ void ClusterFunctionReadTaskResponse::deserialize(ReadBuffer & in, size_t antaly } } - if (antalya_protocol_version >= DBMS_ANTALYA_PROTOCOL_VERSION_WITH_DATA_FILE_META_INFO) + if (antalya_protocol_version >= DBMS_ANTALYA_PROTOCOL_VERSION_WITH_READ_TASK_EXTENSIONS) { UInt64 has_file_meta_info = 0; readVarUInt(has_file_meta_info, in); if (has_file_meta_info) file_meta_info = std::make_shared(DataFileMetaInfo::deserialize(in)); - } - if (antalya_protocol_version >= DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY) - { UInt64 has_retry_after_us = 0; readVarUInt(has_retry_after_us, in); if (has_retry_after_us) diff --git a/src/Interpreters/ClusterFunctionReadTask.h b/src/Interpreters/ClusterFunctionReadTask.h index 1e233775fdd0..54f1f7644fcd 100644 --- a/src/Interpreters/ClusterFunctionReadTask.h +++ b/src/Interpreters/ClusterFunctionReadTask.h @@ -28,7 +28,7 @@ struct ClusterFunctionReadTaskResponse /// File's columns info std::optional file_meta_info; /// Microseconds the worker should wait before asking for another task. - /// Sent only when both peers negotiated Antalya protocol version 3. + /// Sent only when both peers negotiated Antalya protocol version 2. std::optional retry_after_us; /// Convert received response into ObjectInfo. diff --git a/src/Interpreters/tests/gtest_cluster_function_read_task.cpp b/src/Interpreters/tests/gtest_cluster_function_read_task.cpp index c0b9007a7762..6c369dca0df9 100644 --- a/src/Interpreters/tests/gtest_cluster_function_read_task.cpp +++ b/src/Interpreters/tests/gtest_cluster_function_read_task.cpp @@ -363,12 +363,12 @@ TEST(ClusterFunctionReadTaskResponse, RoundTripsFileMetaInfoOnAntalyaProtocol) String serialized; WriteBufferFromString out(serialized); - response.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_DATA_FILE_META_INFO); + response.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_READ_TASK_EXTENSIONS); out.finalize(); ReadBufferFromString in(serialized); ClusterFunctionReadTaskResponse deserialized; - deserialized.deserialize(in, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_DATA_FILE_META_INFO); + deserialized.deserialize(in, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_READ_TASK_EXTENSIONS); ASSERT_TRUE(deserialized.file_meta_info && *deserialized.file_meta_info); const auto & column = (*deserialized.file_meta_info)->columns_info.at("id"); @@ -389,12 +389,12 @@ TEST(ClusterFunctionReadTaskResponse, AntalyaProtocolWithoutFileMetaInfoConsumes String serialized; WriteBufferFromString out(serialized); - response.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_DATA_FILE_META_INFO); + response.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_READ_TASK_EXTENSIONS); out.finalize(); ReadBufferFromString in(serialized); ClusterFunctionReadTaskResponse deserialized; - deserialized.deserialize(in, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_DATA_FILE_META_INFO); + deserialized.deserialize(in, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_READ_TASK_EXTENSIONS); EXPECT_FALSE(deserialized.file_meta_info.has_value()); EXPECT_TRUE(in.eof()); } @@ -438,12 +438,12 @@ TEST(ClusterFunctionReadTaskResponse, RoundTripsRetryAfterUsOnAntalyaProtocol) String serialized; WriteBufferFromString out(serialized); - response.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY); + response.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_READ_TASK_EXTENSIONS); out.finalize(); ReadBufferFromString in(serialized); ClusterFunctionReadTaskResponse deserialized; - deserialized.deserialize(in, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY); + deserialized.deserialize(in, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_READ_TASK_EXTENSIONS); ASSERT_TRUE(deserialized.retry_after_us.has_value()); EXPECT_EQ(*deserialized.retry_after_us, 2500); EXPECT_TRUE(deserialized.path.empty()); @@ -458,12 +458,12 @@ TEST(ClusterFunctionReadTaskResponse, AntalyaProtocolWithoutRetryConsumesPresenc String serialized; WriteBufferFromString out(serialized); - response.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY); + response.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_READ_TASK_EXTENSIONS); out.finalize(); ReadBufferFromString in(serialized); ClusterFunctionReadTaskResponse deserialized; - deserialized.deserialize(in, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY); + deserialized.deserialize(in, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_READ_TASK_EXTENSIONS); EXPECT_FALSE(deserialized.retry_after_us.has_value()); EXPECT_EQ(deserialized.path, "data/file.parquet"); EXPECT_TRUE(in.eof()); diff --git a/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp b/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp index 0729118d3647..0bd5e9481924 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp @@ -635,7 +635,7 @@ class TaskDistributor : public TaskIterator sleepForSeconds(10); }); - const bool worker_supports_task_retry = antalya_protocol_version >= DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY; + const bool worker_supports_task_retry = antalya_protocol_version >= DBMS_ANTALYA_PROTOCOL_VERSION_WITH_READ_TASK_EXTENSIONS; auto task = task_distributor.getNextTask(number_of_current_replica, worker_supports_task_retry); if (task.retry_after_us) { diff --git a/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h b/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h index 52859ba8d27a..a300cad5fd75 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h +++ b/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h @@ -24,7 +24,7 @@ struct DistributedReadTask { ObjectInfoPtr object; /// Microseconds the worker should wait before asking again. - /// Set only when `object` is empty and the worker understands Antalya protocol version 3. + /// Set only when `object` is empty and the worker understands Antalya protocol version 2. std::optional retry_after_us; }; @@ -37,7 +37,7 @@ class StorageObjectStorageStableTaskDistributor bool send_over_whole_archive_, uint64_t lock_object_storage_task_distribution_ms_); - /// `worker_supports_task_retry` is false for a peer that did not negotiate Antalya protocol version 3. + /// `worker_supports_task_retry` is false for a peer that did not negotiate Antalya protocol version 2. /// The first such peer disables `lock_object_storage_task_distribution_ms` for the rest of the query: /// every replica is then given a file immediately, with no `retry_after_us`. DistributedReadTask getNextTask(size_t number_of_current_replica, bool worker_supports_task_retry = true); From 0147b7edc5611c721359849ed3a2d3b6bb46591e Mon Sep 17 00:00:00 2001 From: Anton Ivashkin Date: Tue, 6 Oct 2026 17:39:46 +0200 Subject: [PATCH 6/6] Remove useless comments --- src/Client/IConnections.h | 2 -- src/Interpreters/ClusterFunctionReadTask.h | 2 -- src/QueryPipeline/RemoteQueryExecutor.h | 1 - 3 files changed, 5 deletions(-) diff --git a/src/Client/IConnections.h b/src/Client/IConnections.h index 4ca0fcd1a306..47e7beaa8f3f 100644 --- a/src/Client/IConnections.h +++ b/src/Client/IConnections.h @@ -31,8 +31,6 @@ class IConnections : boost::noncopyable virtual void sendClusterFunctionReadTaskResponse(const ClusterFunctionReadTaskResponse &) = 0; - /// Negotiated Antalya protocol version of the connection that produced the last packet. - /// `0` when the peer is not an Antalya build, or when this connection type has no current replica. virtual UInt64 getServerAntalyaProtocolVersion() const { return 0; } virtual void sendMergeTreeReadTaskResponse(const ParallelReadResponse & response) = 0; diff --git a/src/Interpreters/ClusterFunctionReadTask.h b/src/Interpreters/ClusterFunctionReadTask.h index 54f1f7644fcd..7a032e9f9d0e 100644 --- a/src/Interpreters/ClusterFunctionReadTask.h +++ b/src/Interpreters/ClusterFunctionReadTask.h @@ -27,8 +27,6 @@ struct ClusterFunctionReadTaskResponse std::optional iceberg_info; /// File's columns info std::optional file_meta_info; - /// Microseconds the worker should wait before asking for another task. - /// Sent only when both peers negotiated Antalya protocol version 2. std::optional retry_after_us; /// Convert received response into ObjectInfo. diff --git a/src/QueryPipeline/RemoteQueryExecutor.h b/src/QueryPipeline/RemoteQueryExecutor.h index e0243d02a8fd..3e1cff9242bd 100644 --- a/src/QueryPipeline/RemoteQueryExecutor.h +++ b/src/QueryPipeline/RemoteQueryExecutor.h @@ -49,7 +49,6 @@ class TaskIterator { throw Exception(ErrorCodes::NOT_IMPLEMENTED, "Method rescheduleTasksFromReplica is not implemented"); } - /// `antalya_protocol_version` is the negotiated version of the worker that asked for a task (`0` for an upstream peer). virtual ClusterFunctionReadTaskResponsePtr operator()(size_t number_of_current_replica, UInt64 antalya_protocol_version = 0) const = 0; };