diff --git a/docs/en/antalya/protocol.md b/docs/en/antalya/protocol.md index be10b9014a31..2f6266b39aea 100644 --- a/docs/en/antalya/protocol.md +++ b/docs/en/antalya/protocol.md @@ -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} diff --git a/docs/en/interfaces/specs/NativeProtocol.md b/docs/en/interfaces/specs/NativeProtocol.md index 8a2c8f6ae40b..4bedfc12502f 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 | +| 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 | diff --git a/src/Client/Connection.cpp b/src/Client/Connection.cpp index 71b11aded41c..3fa7e133741c 100644 --- a/src/Client/Connection.cpp +++ b/src/Client/Connection.cpp @@ -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(); } diff --git a/src/Core/AntalyaProtocol.h b/src/Core/AntalyaProtocol.h index 53cf8d1ceb68..c79ca44e306a 100644 --- a/src/Core/AntalyaProtocol.h +++ b/src/Core/AntalyaProtocol.h @@ -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 { diff --git a/src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.cpp b/src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.cpp index c763863d1a72..06fb85e75440 100644 --- a/src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.cpp +++ b/src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.cpp @@ -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] != '{') diff --git a/src/Interpreters/ClusterFunctionReadTask.cpp b/src/Interpreters/ClusterFunctionReadTask.cpp index 3adc8cfabcec..0dc62e3351ab 100644 --- a/src/Interpreters/ClusterFunctionReadTask.cpp +++ b/src/Interpreters/ClusterFunctionReadTask.cpp @@ -2,6 +2,7 @@ #include #include #include +#include #include #include #include @@ -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; } @@ -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(); @@ -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(worker_protocol_version), static_cast(DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION)); @@ -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); @@ -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::deserialize(in)); + } } } diff --git a/src/Interpreters/ClusterFunctionReadTask.h b/src/Interpreters/ClusterFunctionReadTask.h index c41fa8293e75..e6fbb06984f0 100644 --- a/src/Interpreters/ClusterFunctionReadTask.h +++ b/src/Interpreters/ClusterFunctionReadTask.h @@ -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; diff --git a/src/Interpreters/tests/gtest_cluster_function_read_task.cpp b/src/Interpreters/tests/gtest_cluster_function_read_task.cpp index b471b53ed4b6..cef49f02151b 100644 --- a/src/Interpreters/tests/gtest_cluster_function_read_task.cpp +++ b/src/Interpreters/tests/gtest_cluster_function_read_task.cpp @@ -2,6 +2,7 @@ #include #include +#include #include #include #include @@ -9,6 +10,7 @@ #include #include #include +#include #include #if USE_PARQUET @@ -312,3 +314,87 @@ TEST(ClusterFunctionReadTaskResponse, RoundTripsFileBucketInfoOnSupportedProtoco } #endif + +static ClusterFunctionReadTaskResponse makeResponseWithFileMeta() +{ + ClusterFunctionReadTaskResponse response; + response.path = "data/file.parquet"; + auto meta = std::make_shared(); + 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(), 1); + EXPECT_EQ(column.hyperrectangle->right.safeGet(), 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()); +} diff --git a/src/Server/TCPHandler.cpp b/src/Server/TCPHandler.cpp index e890daed3a7a..574301bda35a 100644 --- a/src/Server/TCPHandler.cpp +++ b/src/Server/TCPHandler.cpp @@ -2250,7 +2250,7 @@ ClusterFunctionReadTaskResponsePtr TCPHandler::receiveClusterFunctionReadTaskRes case Protocol::Client::ReadTaskResponse: { auto task = std::make_shared(); - task->deserialize(*in); + task->deserialize(*in, client_antalya_protocol_version); return task; } diff --git a/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp b/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp index 201742db474f..9e3ad4677022 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageCluster.cpp @@ -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 @@ -614,14 +613,12 @@ class TaskDistributor : public TaskIterator std::vector && 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; } @@ -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) }; } diff --git a/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.cpp b/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.cpp index 26e6713b9146..44f3d01d1d4b 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.cpp +++ b/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.cpp @@ -24,15 +24,13 @@ StorageObjectStorageStableTaskDistributor::StorageObjectStorageStableTaskDistrib std::shared_ptr iterator_, std::vector && 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(); @@ -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; { diff --git a/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h b/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h index 3a5a16998be3..652bbb6a03e0 100644 --- a/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h +++ b/src/Storages/ObjectStorage/StorageObjectStorageStableTaskDistributor.h @@ -25,8 +25,7 @@ class StorageObjectStorageStableTaskDistributor std::shared_ptr iterator_, std::vector && 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); @@ -57,7 +56,6 @@ class StorageObjectStorageStableTaskDistributor std::mutex mutex; bool iterator_exhausted = false; - bool iceberg_read_optimization_enabled = 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 1df88cb6adf8..244a6a1f210f 100644 --- a/src/Storages/ObjectStorage/tests/gtest_rendezvous_hashing.cpp +++ b/src/Storages/ObjectStorage/tests/gtest_rendezvous_hashing.cpp @@ -102,7 +102,7 @@ TEST(RendezvousHashing, SingleNode) { auto iterator = makeIterator(); std::vector replicas = {"replica0", "replica1", "replica2", "replica3"}; - StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0, false); + StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0); std::vector paths; ASSERT_TRUE(extractNForReplica(distributor, paths, 0, 10)); ASSERT_TRUE(checkHead(paths, {6})); @@ -111,7 +111,7 @@ TEST(RendezvousHashing, SingleNode) { auto iterator = makeIterator(); std::vector replicas = {"replica0", "replica1", "replica2", "replica3"}; - StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0, false); + StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0); std::vector paths; ASSERT_TRUE(extractNForReplica(distributor, paths, 1, 10)); ASSERT_TRUE(checkHead(paths, {0, 2, 4})); @@ -120,7 +120,7 @@ TEST(RendezvousHashing, SingleNode) { auto iterator = makeIterator(); std::vector replicas = {"replica0", "replica1", "replica2", "replica3"}; - StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0, false); + StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0); std::vector paths; ASSERT_TRUE(extractNForReplica(distributor, paths, 2, 10)); ASSERT_TRUE(checkHead(paths, {1, 5, 7, 8})); @@ -129,7 +129,7 @@ TEST(RendezvousHashing, SingleNode) { auto iterator = makeIterator(); std::vector replicas = {"replica0", "replica1", "replica2", "replica3"}; - StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0, false); + StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0); std::vector paths; ASSERT_TRUE(extractNForReplica(distributor, paths, 3, 10)); ASSERT_TRUE(checkHead(paths, {3, 9})); @@ -140,7 +140,7 @@ TEST(RendezvousHashing, MultipleNodes) { auto iterator = makeIterator(); std::vector replicas = {"replica0", "replica1", "replica2", "replica3"}; - StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0, false); + StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0); { std::vector paths; @@ -172,7 +172,7 @@ TEST(RendezvousHashing, SingleNodeReducedCluster) { auto iterator = makeIterator(); std::vector replicas = {"replica2", "replica1"}; - StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0, false); + StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0); std::vector paths; ASSERT_TRUE(extractNForReplica(distributor, paths, 0, 10)); ASSERT_TRUE(checkHead(paths, {1, 5, 6, 7, 8, 9})); @@ -181,7 +181,7 @@ TEST(RendezvousHashing, SingleNodeReducedCluster) { auto iterator = makeIterator(); std::vector replicas = {"replica2", "replica1"}; - StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0, false); + StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0); std::vector paths; ASSERT_TRUE(extractNForReplica(distributor, paths, 1, 10)); ASSERT_TRUE(checkHead(paths, {0, 2, 3, 4})); @@ -192,7 +192,7 @@ TEST(RendezvousHashing, MultipleNodesReducedCluster) { auto iterator = makeIterator(); std::vector replicas = {"replica2", "replica1"}; - StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0, false); + StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0); { std::vector paths; @@ -211,7 +211,7 @@ TEST(RendezvousHashing, MultipleNodesReducedClusterOneByOne) { auto iterator = makeIterator(); std::vector replicas = {"replica2", "replica1"}; - StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0, false); + StorageObjectStorageStableTaskDistributor distributor(iterator, std::move(replicas), false, 0); std::vector paths0; std::vector paths1; 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 new file mode 100644 index 000000000000..273ec80e1c58 --- /dev/null +++ b/tests/integration/test_iceberg_mixed_antalya_upstream/configs/config.d/cluster.xml @@ -0,0 +1,20 @@ + + + + + + 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 new file mode 100644 index 000000000000..779c3e37a7cb --- /dev/null +++ b/tests/integration/test_iceberg_mixed_antalya_upstream/test.py @@ -0,0 +1,80 @@ +import pytest + +from helpers.cluster import ClickHouseCluster +from helpers.config_cluster import minio_access_key, minio_secret_key + +# Four-part tag. `latest` and `26.6` move when a new build is published. +UPSTREAM_IMAGE = "clickhouse/clickhouse-server" +UPSTREAM_TAG = "26.6.8.7" + +cluster = ClickHouseCluster(__file__) +antalya = cluster.add_instance( + "antalya", + main_configs=["configs/config.d/cluster.xml"], + with_minio=True, +) +upstream = cluster.add_instance( + "upstream", + main_configs=["configs/config.d/cluster.xml"], + image=UPSTREAM_IMAGE, + tag=UPSTREAM_TAG, + with_installed_binary=True, + with_remote_database_disk=False, +) + +TABLE_URL = "http://minio1:9001/root/mixed_antalya_upstream/" + + +@pytest.fixture(scope="module") +def started_cluster(): + try: + cluster.start() + antalya.query( + f""" + CREATE TABLE mixed_proto (id Int32, tag String) + ENGINE = IcebergS3('{TABLE_URL}', '{minio_access_key}', '{minio_secret_key}') + """ + ) + # Two files, each with a constant `tag`, so the Antalya read optimization has + # something to send. The other side must still return these values. + antalya.query( + "INSERT INTO mixed_proto VALUES (1, 'alpha')", + settings={"allow_insert_into_iceberg": 1}, + ) + antalya.query( + "INSERT INTO mixed_proto VALUES (2, 'beta')", + settings={"allow_insert_into_iceberg": 1}, + ) + yield cluster + finally: + cluster.shutdown() + + +@pytest.mark.parametrize( + ("initiator_name", "cluster_name", "settings"), + [ + 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", + ), + ], +) +def test_mixed_iceberg_cluster_read(started_cluster, initiator_name, cluster_name, settings): + 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"