Repository navigation
Conversation
`retry_after_us` for `lock_object_storage_task_distribution_ms` is appended to `ReadTaskResponse` only when both peers negotiated version 3, so an upstream worker still receives a plain object path. Co-authored-by: Cursor <cursoragent@cursor.com>
If any replica did not negotiate Antalya protocol version 3, `lock_object_storage_task_distribution_ms` is skipped for the query and files are assigned immediately. Co-authored-by: Cursor <cursoragent@cursor.com>
|
@codex review |
Codex Review SummaryThis comment shows the latest Codex review activity on this pull request.
ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings. |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: f475189edf
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| 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(); | ||
| } | ||
| } |
There was a problem hiding this comment.
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 👍 / 👎.
There was a problem hiding this comment.
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.
| objects.push_back(pool_scheduler([this]() -> ObjectInfoPtr | ||
| { | ||
| auto task = callback(); | ||
| if (!task || task->isEmpty()) | ||
| return nullptr; | ||
| return task->getObjectInfo(); | ||
| return takeNextObject(); | ||
| }, Priority{})); |
There was a problem hiding this comment.
Avoid waiting for retries during eager prefetch
If a worker receives some real tasks among its max_threads prefetches but another prefetch receives retry_after_us because files remain assigned to another replica, takeNextObject loops until that retry yields a task or end marker. Since the constructor waits for every future before exposing the already fetched objects, none of this worker's valid tasks can begin processing during the lock interval; a large lock_object_storage_task_distribution_ms can therefore stall an entire worker rather than only its idle reader. The initial prefetch needs to defer retry handling without blocking construction of all readers.
Useful? React with 👍 / 👎.
| auto wait_time = std::min<UInt64>(100000, *task->retry_after_us); | ||
| ProfileEvents::increment(ProfileEvents::ObjectStorageClusterWaitingMicroseconds, wait_time); | ||
| sleepForMicroseconds(wait_time); | ||
| continue; |
There was a problem hiding this comment.
Recheck swarm mode after each retry
When SYSTEM STOP SWARM MODE is issued while this iterator is processing retry responses, the mode is checked only before entering the loop. Each retry then sleeps and continues directly to another callback, so the worker can keep requesting and accept a new task after swarm mode has been disabled, defeating the command's graceful-shutdown guarantee. Recheck isSwarmModeEnabled before every subsequent request.
Useful? React with 👍 / 👎.
A `retry_after_us` answer in `ReadTaskIterator` construction is skipped so tasks already fetched by other threads can start. The wait happens later, in the reader that asks for another file. Co-authored-by: Cursor <cursoragent@cursor.com>
`SYSTEM STOP SWARM MODE` during a `retry_after_us` wait in `takeNextObject` stops further task requests. Co-authored-by: Cursor <cursoragent@cursor.com>
|
@codex review |
|
Codex Review: Didn't find any major issues. More of your lovely PRs please. Reviewed commit: ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
If Codex has suggestions, it will comment; otherwise it will react with 👍. Codex can also answer questions or update the PR. Try commenting "@codex address that feedback". |
zvonand
left a comment
There was a problem hiding this comment.
fine in general, just needs a cleanup + protocol version bump revert. it is OK to rename previous anchor instead.
|
|
||
| ```text | ||
| client -> server "ClickHouse client (antalya:2)" | ||
| server -> client "ClickHouse (antalya:2)" |
There was a problem hiding this comment.
actually no need to bump the protocol. previous change was not yet released. keep version 2 pls
| | 7 | Scalar | not specified | Scalar data block | | ||
| | 8 | IgnoredPartUUIDs | not specified | Parts to exclude from query | | ||
| | 9 | ReadTaskResponse | not specified | S3 cluster read response. When both peers negotiated [Antalya protocol](/antalya/protocol) version 2 or newer, the initiator appends optional Iceberg column statistics after the upstream payload. | | ||
| | 9 | ReadTaskResponse | not specified | S3 cluster read response. When both peers negotiated [Antalya protocol](/antalya/protocol) version 2 or newer, the initiator appends optional Iceberg column statistics after the upstream payload. Version 3 or newer also appends optional `retry_after_us` for `lock_object_storage_task_distribution_ms`. If any worker negotiated a lower Antalya version, that field is omitted and the setting is ignored for the query. | |
|
|
||
| Each node begins processing files for which it is the primary node. After completing its assigned files, a node may take tasks from other nodes, either immediately or after waiting for `lock_object_storage_task_distribution_ms` milliseconds if the primary node does not request new files during that interval. The default value of `lock_object_storage_task_distribution_ms` is 500 milliseconds. This setting balances between caching efficiency and workload redistribution when nodes are imbalanced. | ||
|
|
||
| The wait is a field of `ReadTaskResponse` and is sent only when every worker negotiated Antalya protocol version 3. If any worker negotiated a lower version, including an upstream node, the setting is ignored for that query and files are assigned immediately. |
| virtual void sendClusterFunctionReadTaskResponse(const ClusterFunctionReadTaskResponse &) = 0; | ||
|
|
||
| /// Negotiated Antalya protocol version of the connection that produced the last packet. | ||
| /// `0` when the peer is not an Antalya build, or when this connection type has no current replica. |
| /// 2 — `ReadTaskResponse` carries optional `DataFileMetaInfo` after the upstream payload. | ||
| static constexpr auto DBMS_ANTALYA_PROTOCOL_VERSION_WITH_DATA_FILE_META_INFO = 2; | ||
| /// 3 — `ReadTaskResponse` carries optional `retry_after_us` after the version 2 trailer. | ||
| static constexpr auto DBMS_ANTALYA_PROTOCOL_VERSION_WITH_TASK_RETRY = 3; |
| std::optional<Iceberg::IcebergObjectSerializableInfo> iceberg_info; | ||
| /// File's columns info | ||
| std::optional<DataFileMetaInfoPtr> file_meta_info; | ||
| /// Microseconds the worker should wait before asking for another task. |
There was a problem hiding this comment.
this comment explains obvious behavior
| throw Exception(ErrorCodes::NOT_IMPLEMENTED, "Method rescheduleTasksFromReplica is not implemented"); | ||
| } | ||
| virtual ClusterFunctionReadTaskResponsePtr operator()(size_t number_of_current_replica) const = 0; | ||
| /// `antalya_protocol_version` is the negotiated version of the worker that asked for a task (`0` for an upstream peer). |
There was a problem hiding this comment.
one more statement of the obvious :)
Both trailers ship in the same release, so they share `DBMS_ANTALYA_PROTOCOL_VERSION_WITH_READ_TASK_EXTENSIONS`. Co-authored-by: Cursor <cursoragent@cursor.com>
Changelog category (leave one):
Changelog entry (a user-readable short description of the changes that goes to CHANGELOG.md):
Use Antalya protocol for
lock_object_storage_task_distribution_msinstead of JSONDocumentation entry for user-facing changes
Followed #2482, completely removed JSON in object path string, use Antalya protocol instead.
CI/CD Options
Exclude tests:
Regression jobs to run: