Skip to content
Merged
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
9 changes: 6 additions & 3 deletions docs/en/antalya/protocol.md
Original file line number Diff line number Diff line change
Expand Up @@ -14,15 +14,18 @@ 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:1)"
server -> client "ClickHouse (antalya:1)"
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`
means the peer is not an Antalya build. Negotiation is per hop and not transitive: initiator to
worker and worker to worker negotiate independently.

Version 1 is the advertisement itself. Nothing is gated on it yet.
Version 1 is the advertisement itself.

Version 2 appends optional Iceberg column statistics (`DataFileMetaInfo`) to `ReadTaskResponse`,
after the upstream cluster-processing payload.

## 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 |
| 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. |
| 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: 1 addition & 1 deletion src/Client/Connection.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -1151,7 +1151,7 @@ void Connection::sendData(const Block & block, const String & name, bool scalar)
void Connection::sendClusterFunctionReadTaskResponse(const ClusterFunctionReadTaskResponse & response)
{
writeVarUInt(Protocol::Client::ReadTaskResponse, *out);
response.serialize(*out, worker_cluster_function_protocol_version);
response.serialize(*out, worker_cluster_function_protocol_version, server_antalya_protocol_version);
out->finishChunk();
out->next();
}
Expand Down
4 changes: 3 additions & 1 deletion src/Core/AntalyaProtocol.h
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,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;
/// Bump for every Antalya-only wire protocol change. See `docs/en/antalya/protocol.md`.
static constexpr auto DBMS_ANTALYA_PROTOCOL_VERSION = 1;
static constexpr auto DBMS_ANTALYA_PROTOCOL_VERSION = DBMS_ANTALYA_PROTOCOL_VERSION_WITH_DATA_FILE_META_INFO;

namespace AntalyaProtocol
{
Expand Down
5 changes: 2 additions & 3 deletions src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -173,9 +173,8 @@ RelativePathWithMetadata::CommandInTaskResponse::CommandInTaskResponse(const std
/// 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. The proper fix is to stop multiplexing the command into the path field (a separate
/// `ObjectInfo` kind or the versioned cluster-function protocol, see
/// https://github.com/Altinity/ClickHouse/pull/1360), after which this probe goes away entirely.
/// 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] != '{')
Expand Down
30 changes: 27 additions & 3 deletions src/Interpreters/ClusterFunctionReadTask.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
#include <Interpreters/SetSerialization.h>
#include <Interpreters/Context.h>
#include <AggregateFunctions/AggregateFunctionGroupBitmapData.h>
#include <Core/AntalyaProtocol.h>
#include <Core/Settings.h>
#include <Core/ProtocolDefines.h>
#include <Common/Exception.h>
Expand All @@ -23,6 +24,7 @@ namespace ErrorCodes
}
namespace Setting
{
extern const SettingsBool allow_experimental_iceberg_read_optimization;
extern const SettingsBool cluster_function_process_archive_on_multiple_nodes;
}

Expand All @@ -39,7 +41,8 @@ ClusterFunctionReadTaskResponse::ClusterFunctionReadTaskResponse(ObjectInfoPtr o
iceberg_info = iceberg_object->info;
#endif

file_meta_info = object->relative_path_with_metadata.file_meta_info;
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();
Expand Down Expand Up @@ -86,7 +89,7 @@ ObjectInfoPtr ClusterFunctionReadTaskResponse::getObjectInfo() const
return object;
}

void ClusterFunctionReadTaskResponse::serialize(WriteBuffer & out, size_t worker_protocol_version) const
void ClusterFunctionReadTaskResponse::serialize(WriteBuffer & out, size_t worker_protocol_version, size_t antalya_protocol_version) const
{
auto protocol_version
= std::min(static_cast<UInt64>(worker_protocol_version), static_cast<UInt64>(DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION));
Expand Down Expand Up @@ -190,9 +193,22 @@ void ClusterFunctionReadTaskResponse::serialize(WriteBuffer & out, size_t worker
writeVarUInt(0, out);
}
}

if (antalya_protocol_version >= DBMS_ANTALYA_PROTOCOL_VERSION_WITH_DATA_FILE_META_INFO)
{
if (file_meta_info && *file_meta_info)
{
writeVarUInt(1, out);
(*file_meta_info)->serialize(out);
}
else
{
writeVarUInt(0, out);
}
}
}

void ClusterFunctionReadTaskResponse::deserialize(ReadBuffer & in)
void ClusterFunctionReadTaskResponse::deserialize(ReadBuffer & in, size_t antalya_protocol_version)
{
size_t protocol_version = 0;
readVarUInt(protocol_version, in);
Expand Down Expand Up @@ -243,6 +259,14 @@ void ClusterFunctionReadTaskResponse::deserialize(ReadBuffer & in)
iceberg_info->deserializeForClusterFunctionProtocol(in, protocol_version);
}
}

if (antalya_protocol_version >= DBMS_ANTALYA_PROTOCOL_VERSION_WITH_DATA_FILE_META_INFO)
{
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));
}
}

}
11 changes: 6 additions & 5 deletions src/Interpreters/ClusterFunctionReadTask.h
Original file line number Diff line number Diff line change
Expand Up @@ -35,11 +35,12 @@ struct ClusterFunctionReadTaskResponse
/// It is used to identify an end of processing.
bool isEmpty() const { return path.empty(); }

/// Serialize according to the protocol version.
void serialize(WriteBuffer & out, size_t worker_protocol_version) const;
/// Deserialize. Protocol version will be received from `in`
/// and the result will be deserialized accordingly.
void deserialize(ReadBuffer & in);
/// Serialize according to the cluster-processing protocol version.
/// `antalya_protocol_version` is the negotiated Antalya version of this hop (`0` for an upstream peer).
void serialize(WriteBuffer & out, size_t worker_protocol_version, size_t antalya_protocol_version = 0) const;
/// Deserialize. The cluster-processing protocol version is read from `in`.
/// `antalya_protocol_version` must be the same negotiated value `serialize` was given.
void deserialize(ReadBuffer & in, size_t antalya_protocol_version = 0);
};

using ClusterFunctionReadTaskResponsePtr = std::shared_ptr<ClusterFunctionReadTaskResponse>;
Expand Down
86 changes: 86 additions & 0 deletions src/Interpreters/tests/gtest_cluster_function_read_task.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -2,13 +2,15 @@

#include <AggregateFunctions/AggregateFunctionGroupBitmapData.h>
#include <Common/Exception.h>
#include <Core/AntalyaProtocol.h>
#include <Core/NamesAndTypes.h>
#include <Core/ProtocolDefines.h>
#include <DataTypes/DataTypesNumber.h>
#include <IO/ReadBufferFromString.h>
#include <IO/WriteBufferFromString.h>
#include <Interpreters/ActionsDAG.h>
#include <Interpreters/ClusterFunctionReadTask.h>
#include <Storages/ObjectStorage/DataLakes/IDataLakeMetadata.h>
#include <config.h>

#if USE_PARQUET
Expand Down Expand Up @@ -312,3 +314,87 @@ TEST(ClusterFunctionReadTaskResponse, RoundTripsFileBucketInfoOnSupportedProtoco
}

#endif

static ClusterFunctionReadTaskResponse makeResponseWithFileMeta()
{
ClusterFunctionReadTaskResponse response;
response.path = "data/file.parquet";
auto meta = std::make_shared<DataFileMetaInfo>();
DataFileMetaInfo::ColumnInfo column;
column.rows_count = 10;
column.nulls_count = 0;
column.hyperrectangle = Range(Field(Int64(1)), true, Field(Int64(5)), true);
meta->columns_info.emplace("id", std::move(column));
response.file_meta_info = std::move(meta);
return response;
}

TEST(ClusterFunctionReadTaskResponse, OmitsFileMetaInfoForUpstreamPeer)
{
auto with_meta = makeResponseWithFileMeta();
ClusterFunctionReadTaskResponse without_meta;
without_meta.path = with_meta.path;

String with_meta_bytes;
String without_meta_bytes;
{
WriteBufferFromString out(with_meta_bytes);
with_meta.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, /*antalya_protocol_version*/ 0);
out.finalize();
}
{
WriteBufferFromString out(without_meta_bytes);
without_meta.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, /*antalya_protocol_version*/ 0);
out.finalize();
}
EXPECT_EQ(with_meta_bytes, without_meta_bytes);

ReadBufferFromString in(with_meta_bytes);
ClusterFunctionReadTaskResponse deserialized;
deserialized.deserialize(in, /*antalya_protocol_version*/ 0);
EXPECT_FALSE(deserialized.file_meta_info.has_value());
EXPECT_EQ(deserialized.path, "data/file.parquet");
EXPECT_TRUE(in.eof());
}

TEST(ClusterFunctionReadTaskResponse, RoundTripsFileMetaInfoOnAntalyaProtocol)
{
auto response = makeResponseWithFileMeta();

String serialized;
WriteBufferFromString out(serialized);
response.serialize(out, DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_DATA_FILE_META_INFO);
out.finalize();

ReadBufferFromString in(serialized);
ClusterFunctionReadTaskResponse deserialized;
deserialized.deserialize(in, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_DATA_FILE_META_INFO);

ASSERT_TRUE(deserialized.file_meta_info && *deserialized.file_meta_info);
const auto & column = (*deserialized.file_meta_info)->columns_info.at("id");
ASSERT_TRUE(column.rows_count.has_value());
ASSERT_TRUE(column.nulls_count.has_value());
EXPECT_EQ(*column.rows_count, 10);
EXPECT_EQ(*column.nulls_count, 0);
ASSERT_TRUE(column.hyperrectangle.has_value());
EXPECT_EQ(column.hyperrectangle->left.safeGet<Int64>(), 1);
EXPECT_EQ(column.hyperrectangle->right.safeGet<Int64>(), 5);
EXPECT_TRUE(in.eof());
}

TEST(ClusterFunctionReadTaskResponse, AntalyaProtocolWithoutFileMetaInfoConsumesPresenceFlag)
{
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_DATA_FILE_META_INFO);
out.finalize();

ReadBufferFromString in(serialized);
ClusterFunctionReadTaskResponse deserialized;
deserialized.deserialize(in, DBMS_ANTALYA_PROTOCOL_VERSION_WITH_DATA_FILE_META_INFO);
EXPECT_FALSE(deserialized.file_meta_info.has_value());
EXPECT_TRUE(in.eof());
}
2 changes: 1 addition & 1 deletion src/Server/TCPHandler.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -2250,7 +2250,7 @@ ClusterFunctionReadTaskResponsePtr TCPHandler::receiveClusterFunctionReadTaskRes
case Protocol::Client::ReadTaskResponse:
{
auto task = std::make_shared<ClusterFunctionReadTaskResponse>();
task->deserialize(*in);
task->deserialize(*in, client_antalya_protocol_version);
return task;
}

Expand Down
10 changes: 3 additions & 7 deletions src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,6 @@ namespace Setting
extern const SettingsInt64 delta_lake_snapshot_start_version;
extern const SettingsInt64 delta_lake_snapshot_end_version;
extern const SettingsUInt64 lock_object_storage_task_distribution_ms;
extern const SettingsBool allow_experimental_iceberg_read_optimization;
}

namespace ErrorCodes
Expand Down Expand Up @@ -614,14 +613,12 @@ class TaskDistributor : public TaskIterator
std::vector<std::string> && ids_of_hosts,
bool send_over_whole_archive,
uint64_t lock_object_storage_task_distribution_ms,
ContextPtr context_,
bool iceberg_read_optimization_enabled)
ContextPtr context_)
: task_distributor(
iterator,
std::move(ids_of_hosts),
send_over_whole_archive,
lock_object_storage_task_distribution_ms,
iceberg_read_optimization_enabled)
lock_object_storage_task_distribution_ms)
, context(context_) {}
~TaskDistributor() override = default;
bool supportRerunTask() const override { return true; }
Expand Down Expand Up @@ -710,8 +707,7 @@ RemoteQueryExecutor::Extension StorageObjectStorageCluster::getTaskIteratorExten
std::move(ids_of_hosts),
/* send_over_whole_archive */!local_context->getSettingsRef()[Setting::cluster_function_process_archive_on_multiple_nodes],
lock_object_storage_task_distribution_ms,
local_context,
/* iceberg_read_optimization_enabled */local_context->getSettingsRef()[Setting::allow_experimental_iceberg_read_optimization]);
local_context);

return RemoteQueryExecutor::Extension{ .task_iterator = std::move(callback) };
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,15 +24,13 @@ StorageObjectStorageStableTaskDistributor::StorageObjectStorageStableTaskDistrib
std::shared_ptr<IObjectIterator> iterator_,
std::vector<std::string> && ids_of_nodes_,
bool send_over_whole_archive_,
uint64_t lock_object_storage_task_distribution_ms_,
bool iceberg_read_optimization_enabled_)
uint64_t lock_object_storage_task_distribution_ms_)
: iterator(std::move(iterator_))
, send_over_whole_archive(send_over_whole_archive_)
, connection_to_files(ids_of_nodes_.size())
, ids_of_nodes(std::move(ids_of_nodes_))
, lock_object_storage_task_distribution_us(lock_object_storage_task_distribution_ms_ * 1000)
, iterator_exhausted(false)
, iceberg_read_optimization_enabled(iceberg_read_optimization_enabled_)
{
Poco::Timestamp now;
size_t nodes = ids_of_nodes.size();
Expand Down Expand Up @@ -187,17 +185,6 @@ ObjectInfoPtr StorageObjectStorageStableTaskDistributor::getMatchingFileFromIter

String file_identifier = getFileIdentifier(object_info, true);

if (iceberg_read_optimization_enabled)
{
auto file_meta_info = object_info->relative_path_with_metadata.getFileMetaInfo();
if (file_meta_info.has_value())
{
auto file_path = send_over_whole_archive ? object_info->getPathOrPathToArchiveIfArchive() : object_info->getPath();
object_info->relative_path_with_metadata.command.setFilePath(file_path);
object_info->relative_path_with_metadata.command.setFileMetaInfo(file_meta_info.value());
}
}

size_t file_replica_idx;

{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,8 +25,7 @@ class StorageObjectStorageStableTaskDistributor
std::shared_ptr<IObjectIterator> iterator_,
std::vector<std::string> && ids_of_nodes_,
bool send_over_whole_archive_,
uint64_t lock_object_storage_task_distribution_ms_,
bool iceberg_read_optimization_enabled_);
uint64_t lock_object_storage_task_distribution_ms_);

ObjectInfoPtr getNextTask(size_t number_of_current_replica);

Expand Down Expand Up @@ -57,7 +56,6 @@ class StorageObjectStorageStableTaskDistributor

std::mutex mutex;
bool iterator_exhausted = false;
bool iceberg_read_optimization_enabled = false;

LoggerPtr log = getLogger("StorageClusterTaskDistributor");
};
Expand Down
Loading
Loading