Skip to content

Use Antalya protocol for lock_object_storage_task_distribution_ms instead of JSON - #2489

Open
ianton-ru wants to merge 6 commits into
antalya-26.6from
feature/antalya-26.6/task-retry-antalya-protocol
Open

ianton-ru wants to merge 6 commits into
antalya-26.6from
feature/antalya-26.6/task-retry-antalya-protocol

Conversation

@ianton-ru

Copy link
Copy Markdown

Changelog category (leave one):

  • Improvement

Changelog entry (a user-readable short description of the changes that goes to CHANGELOG.md):

Use Antalya protocol for lock_object_storage_task_distribution_ms instead of JSON

Documentation entry for user-facing changes

Followed #2482, completely removed JSON in object path string, use Antalya protocol instead.

CI/CD Options

Exclude tests:

  • Fast test
  • Integration Tests
  • Stateless tests
  • Stateful tests
  • Unit tests
  • Performance tests
  • Aarch64 tests
  • All with ASAN
  • All with TSAN
  • All with MSAN
  • All with UBSAN
  • All with Coverage
  • All Regression
  • Disable CI Cache

Regression jobs to run:

  • Fast suites (mostly <1h)
  • Aggregate Functions (2h)
  • Alter (1.5h)
  • Benchmark (30m)
  • CAS (content-addressed storage; Antalya only)
  • ClickHouse Keeper (1h)
  • Iceberg (2h)
  • LDAP (1h)
  • OAuth (5m)
  • Parquet (1.5h)
  • RBAC (1.5h)
  • SSL Server (1h)
  • S3 (2h)
  • S3 Export (2h)
  • Swarms (30m)
  • Tiered Storage (2h)

ianton-ru and others added 2 commits October 6, 2026 11:17
`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>
@ianton-ru
ianton-ru marked this pull request as ready for review October 6, 2026 09:45
@github-actions

github-actions Bot commented Oct 6, 2026 •

Copy link
Copy Markdown

Workflow [PR], commit [0147b7e]

@ianton-ru

Copy link
Copy Markdown
Author

@codex review

@chatgpt-codex-connector

chatgpt-codex-connector Bot commented Oct 6, 2026 •

Copy link
Copy Markdown

Codex Review Summary

This comment shows the latest Codex review activity on this pull request.

Review Status Commit Review trigger
📝 Code Review ✅ Completed 2026-10-06T12:33:18.360498Z af45bd6 Manual request
ℹ️ 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" or "@codex security review".

Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings.

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

💡 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".

Comment on lines +132 to 136
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();
}
}

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.

Comment on lines 1988 to 1991
objects.push_back(pool_scheduler([this]() -> ObjectInfoPtr
{
auto task = callback();
if (!task || task->isEmpty())
return nullptr;
return task->getObjectInfo();
return takeNextObject();
}, Priority{}));

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge 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 👍 / 👎.

Comment on lines +2049 to +2052
auto wait_time = std::min<UInt64>(100000, *task->retry_after_us);
ProfileEvents::increment(ProfileEvents::ObjectStorageClusterWaitingMicroseconds, wait_time);
sleepForMicroseconds(wait_time);
continue;

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge 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 👍 / 👎.

ianton-ru and others added 2 commits October 6, 2026 14:20
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>
@ianton-ru

Copy link
Copy Markdown
Author

@codex review

@chatgpt-codex-connector

Copy link
Copy Markdown

Codex Review: Didn't find any major issues. More of your lovely PRs please.

Reviewed commit: af45bd6e8d

ℹ️ 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".

@zvonand zvonand added port-antalya PRs to be ported to all new Antalya releases antalya-26.6 antalya labels Oct 6, 2026

@zvonand zvonand left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

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)"

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

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. |

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

see above


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.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

same note about version

Comment thread src/Client/IConnections.h Outdated
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.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

does not deserve a comment :)

Comment thread src/Core/AntalyaProtocol.h Outdated
/// 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;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

same, no need to bump

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.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

this comment explains obvious behavior

Comment thread src/QueryPipeline/RemoteQueryExecutor.h Outdated
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).

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

one more statement of the obvious :)

ianton-ru and others added 2 commits October 6, 2026 17:34
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>

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

antalya antalya-26.6 port-antalya PRs to be ported to all new Antalya releases

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants