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

Filter by extension

Filter by extension

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

Expand Down
2 changes: 1 addition & 1 deletion docs/en/interfaces/specs/NativeProtocol.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
2 changes: 2 additions & 0 deletions docs/en/sql-reference/distribution-on-cluster.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
2 changes: 2 additions & 0 deletions src/Client/Connection.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
2 changes: 2 additions & 0 deletions src/Client/IConnections.h
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
5 changes: 5 additions & 0 deletions src/Client/MultiplexedConnections.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
7 changes: 4 additions & 3 deletions src/Core/AntalyaProtocol.h
Original file line number Diff line number Diff line change
Expand Up @@ -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
{
Expand Down
67 changes: 0 additions & 67 deletions src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -11,10 +11,6 @@
#include <IO/WriteBufferFromString.h>
#include <Storages/ObjectStorage/DataLakes/IDataLakeMetadata.h>

#include <Poco/JSON/Object.h>
#include <Poco/JSON/Parser.h>
#include <Poco/JSON/JSONException.h>


namespace DB
{
Expand Down Expand Up @@ -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<Poco::JSON::Object::Ptr>();
if (!json)
return;

is_valid = true;

if (json->has("file_path"))
file_path = json->getValue<std::string>("file_path");
if (json->has("retry_after_us"))
retry_after_us = json->getValue<size_t>("retry_after_us");
if (json->has("meta_info"))
file_meta_info = std::make_shared<DataFileMetaInfo>(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();
}

}
49 changes: 2 additions & 47 deletions src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.h
Original file line number Diff line number Diff line change
Expand Up @@ -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<std::string> getFilePath() const { return file_path; }
std::optional<Poco::Timestamp::TimeDiff> getRetryAfterUs() const { return retry_after_us; }
std::optional<DataFileMetaInfoPtr> getFileMetaInfo() const { return file_meta_info; }

private:
bool is_valid = false;
std::optional<std::string> file_path;
std::optional<Poco::Timestamp::TimeDiff> retry_after_us;
std::optional<DataFileMetaInfoPtr> file_meta_info;
};

String relative_path;
/// Object metadata: size, modification time, etc.
std::optional<ObjectMetadata> metadata;
/// Information about columns
std::optional<DataFileMetaInfoPtr> file_meta_info;
/// Retry request after short pause
CommandInTaskResponse command;

RelativePathWithMetadata() = default;

explicit RelativePathWithMetadata(String command_or_path, std::optional<ObjectMetadata> metadata_ = std::nullopt)
: relative_path(std::move(command_or_path))
explicit RelativePathWithMetadata(String relative_path_, std::optional<ObjectMetadata> 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();
}
}
Comment on lines +132 to 136

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P1 Badge Preserve legacy retry commands during rolling upgrades

When a newly upgraded worker is paired with an Antalya v2 initiator, the initiator still encodes a retry as the path {"retry_after_us":...} because it predates the v3 trailer. This constructor now preserves that text as an ordinary object key, so with the default distribution lock and an imbalanced workload the worker tries to fetch that key and the query fails. Keep decoding the legacy command when the peer negotiated an Antalya version below 3 so rolling upgrades remain compatible.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

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

If update initiator first and swarm nodes later, query fall back to behavior with ignoring lock_object_storage_task_distribution_ms setting. I think it is enough.


explicit RelativePathWithMetadata(const DataFileInfo & info, std::optional<ObjectMetadata> metadata_ = std::nullopt);
Expand All @@ -191,7 +147,6 @@ struct RelativePathWithMetadata
void setFileMetaInfo(std::optional<DataFileMetaInfoPtr> file_meta_info_ ) { file_meta_info = file_meta_info_; }
std::optional<DataFileMetaInfoPtr> getFileMetaInfo() const { return file_meta_info; }

const CommandInTaskResponse & getCommand() const { return command; }
std::string getFileNameWithoutExtension() const { return std::filesystem::path(relative_path).stem(); }
};

Expand Down
32 changes: 23 additions & 9 deletions src/Interpreters/ClusterFunctionReadTask.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand Down Expand Up @@ -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)
{
Expand All @@ -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);
}
}
}

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

Expand Down
4 changes: 3 additions & 1 deletion src/Interpreters/ClusterFunctionReadTask.h
Original file line number Diff line number Diff line change
Expand Up @@ -27,13 +27,15 @@ struct ClusterFunctionReadTaskResponse
std::optional<Iceberg::IcebergObjectSerializableInfo> iceberg_info;
/// File's columns info
std::optional<DataFileMetaInfoPtr> file_meta_info;
std::optional<UInt64> 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).
Expand Down
Loading
Loading