diff --git a/docs/en/antalya/protocol.md b/docs/en/antalya/protocol.md index 2f6266b39aea..177a274dc447 100644 --- a/docs/en/antalya/protocol.md +++ b/docs/en/antalya/protocol.md @@ -24,8 +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 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 4bedfc12502f..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. | +| 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 3a9835e23856..7405ed790693 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 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 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..47e7beaa8f3f 100644 --- a/src/Client/IConnections.h +++ b/src/Client/IConnections.h @@ -30,6 +30,8 @@ class IConnections : boost::noncopyable virtual void sendQueryPlan(const QueryPlan & query_plan) = 0; virtual void sendClusterFunctionReadTaskResponse(const ClusterFunctionReadTaskResponse &) = 0; + + 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..754a83ebc898 100644 --- a/src/Core/AntalyaProtocol.h +++ b/src/Core/AntalyaProtocol.h @@ -7,10 +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; +/// 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_DATA_FILE_META_INFO; +static constexpr auto DBMS_ANTALYA_PROTOCOL_VERSION = DBMS_ANTALYA_PROTOCOL_VERSION_WITH_READ_TASK_EXTENSIONS; 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..c76b60f1d3bb 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; } @@ -194,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) { @@ -205,6 +200,16 @@ void ClusterFunctionReadTaskResponse::serialize(WriteBuffer & out, size_t worker { writeVarUInt(0, out); } + + if (retry_after_us) + { + writeVarUInt(1, out); + writeVarUInt(*retry_after_us, out); + } + else + { + writeVarUInt(0, out); + } } } @@ -260,12 +265,21 @@ 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)); + + 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..7a032e9f9d0e 100644 --- a/src/Interpreters/ClusterFunctionReadTask.h +++ b/src/Interpreters/ClusterFunctionReadTask.h @@ -27,13 +27,15 @@ struct ClusterFunctionReadTaskResponse std::optional iceberg_info; /// File's columns info std::optional file_meta_info; + 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..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,82 @@ 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()); } + +TEST(ClusterFunctionReadTaskResponse, OmitsRetryAfterUsForUpstreamPeer) +{ + ClusterFunctionReadTaskResponse with_retry; + with_retry.path = "data/file.parquet"; + with_retry.retry_after_us = 1000; + + ClusterFunctionReadTaskResponse without_retry; + without_retry.path = with_retry.path; + + String with_retry_bytes; + String without_retry_bytes; + { + WriteBufferFromString out(with_retry_bytes); + with_retry.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, /*antalya_protocol_version*/ 0); + out.finalize(); + } + { + 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) +{ + 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_READ_TASK_EXTENSIONS); + out.finalize(); + + ReadBufferFromString in(serialized); + ClusterFunctionReadTaskResponse deserialized; + 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()); + 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_READ_TASK_EXTENSIONS); + out.finalize(); + + ReadBufferFromString in(serialized); + ClusterFunctionReadTaskResponse deserialized; + 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/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..3e1cff9242bd 100644 --- a/src/QueryPipeline/RemoteQueryExecutor.h +++ b/src/QueryPipeline/RemoteQueryExecutor.h @@ -49,7 +49,7 @@ class TaskIterator { throw Exception(ErrorCodes::NOT_IMPLEMENTED, "Method rescheduleTasksFromReplica is not implemented"); } - virtual ClusterFunctionReadTaskResponsePtr operator()(size_t number_of_current_replica) const = 0; + 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..0bd5e9481924 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_READ_TASK_EXTENSIONS; + 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..ac210b46a7db 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(/*wait_for_retry=*/ false); }, Priority{})); pool.wait(); @@ -2015,10 +1996,7 @@ StorageObjectStorageSource::ReadTaskIterator::ReadTaskIterator( { auto object = object_future.get(); if (object) - { - resolveIcebergObjectStorageIfNeeded(object); buffer.push_back(object); - } } } @@ -2046,13 +2024,12 @@ void StorageObjectStorageSource::ReadTaskIterator::resolveIcebergObjectStorageIf #endif } -ObjectInfoPtr StorageObjectStorageSource::ReadTaskIterator::next(size_t) +ObjectInfoPtr StorageObjectStorageSource::ReadTaskIterator::takeNextObject(bool wait_for_retry) { - size_t current_index = index.fetch_add(1, std::memory_order_relaxed); - - ObjectInfoPtr object_info; - if (current_index >= buffer.size()) + 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"); @@ -2067,13 +2044,38 @@ ObjectInfoPtr StorageObjectStorageSource::ReadTaskIterator::next(size_t) if (!task || task->isEmpty()) return nullptr; - object_info = task->getObjectInfo(); + 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); + 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(/*wait_for_retry=*/ true); + 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..d6d801bd3f74 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageSource.h +++ b/src/Storages/ObjectStorage/StorageObjectStorageSource.h @@ -191,6 +191,11 @@ 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. 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. void resolveIcebergObjectStorageIfNeeded(const ObjectInfoPtr & object); diff --git a/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.cpp b/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.cpp index 44f3d01d1d4b..5d144bfca4e1 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.cpp @@ -41,10 +41,21 @@ 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); + /// 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); { @@ -65,7 +76,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, allow_task_retry); + if (pending.retry_after_us) + return pending; + file = std::move(pending.object); + } if (file) { @@ -82,7 +98,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 +237,9 @@ ObjectInfoPtr StorageObjectStorageStableTaskDistributor::getMatchingFileFromIter return {}; } -ObjectInfoPtr StorageObjectStorageStableTaskDistributor::getAnyUnprocessedFile(size_t number_of_current_replica) +DistributedReadTask StorageObjectStorageStableTaskDistributor::getAnyUnprocessedFile( + size_t number_of_current_replica, + bool allow_task_retry) { /// Limit time of node activity to keep task in queue Poco::Timestamp activity_limit; @@ -240,6 +258,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 + || !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 @@ -257,7 +276,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 +290,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..a300cad5fd75 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h +++ b/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h @@ -14,10 +14,20 @@ #include #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 2. + std::optional retry_after_us; +}; + class StorageObjectStorageStableTaskDistributor { public: @@ -27,7 +37,10 @@ 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 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); /// Insert objects back to unprocessed files void rescheduleTasksFromReplica(size_t number_of_current_replica); @@ -36,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); - ObjectInfoPtr getAnyUnprocessedFile(size_t number_of_current_replica); + DistributedReadTask getAnyUnprocessedFile(size_t number_of_current_replica, bool allow_task_retry); void saveLastNodeActivity(size_t number_of_current_replica); @@ -56,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"); }; 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"