Skip to content

Java SDK: Change coordinator comm to use multiplexing - #69254

Merged
jason810496 merged 1 commit into
apache:mainfrom
FrankYang0529:airflow-java-sdk-multiplexing
Jul 30, 2026
Merged

jason810496 merged 1 commit into
apache:mainfrom
FrankYang0529:airflow-java-sdk-multiplexing

Conversation

@FrankYang0529

Copy link
Copy Markdown
Member

Why

#69080 made the coordinator comm thread-safe by holding a single lock across each client's entire send → receive → deserialize round trip. It's correct, but only one request can be in flight at a time, so a task that issues concurrent calls over the shared comm socket pays a full, serialized round trip per call.

How

Every frame already carries a request id, so a single background dispatcher (readLoop) becomes the sole socket reader and routes each response to the waiter registered under that id — the lock now guards only the write, so concurrent calls can overlap.


Was generative AI tooling used to co-author this PR?
  • Yes - Claude Code with Opus 4.8

  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Server.kt Outdated
Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Comm.kt Outdated
Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Comm.kt Outdated
@FrankYang0529
FrankYang0529 force-pushed the airflow-java-sdk-multiplexing branch from e159284 to d847241 Compare July 3, 2026 07:46

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

Sorry for the late review as it took me a longer time than I expected to understand all the kotlin specific concurrency concept and syntax in this PR. Here're some comments from CC that still make sense to me.

Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Comm.kt Outdated
Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Server.kt Outdated
@FrankYang0529
FrankYang0529 force-pushed the airflow-java-sdk-multiplexing branch 2 times, most recently from a267532 to 0c6bf38 Compare July 26, 2026 04:48
Serializing each client's whole send → receive → deserialize round trip
behind one lock (apache#69080) allows only a single request in flight at a
time, so a task issuing concurrent calls over the shared comm socket
pays a full, serialized round trip per call.
Every frame already carries a request id, so a single background
dispatcher can read the socket and route each response by id, with the
lock guarding only the write.

Signed-off-by: PoAn Yang <payang@apache.org>
@FrankYang0529
FrankYang0529 force-pushed the airflow-java-sdk-multiplexing branch from 0c6bf38 to 89fe9f5 Compare July 26, 2026 06:07

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

Thanks.

@jason810496
jason810496 merged commit abd837a into apache:main Jul 30, 2026
90 checks passed
@FrankYang0529
FrankYang0529 deleted the airflow-java-sdk-multiplexing branch July 30, 2026 06:18
dabla pushed a commit to dabla/airflow that referenced this pull request Aug 14, 2026
Serializing each client's whole send → receive → deserialize round trip
behind one lock (apache#69080) allows only a single request in flight at a
time, so a task issuing concurrent calls over the shared comm socket
pays a full, serialized round trip per call.
Every frame already carries a request id, so a single background
dispatcher can read the socket and route each response by id, with the
lock guarding only the write.

Signed-off-by: PoAn Yang <payang@apache.org>
imrichardwu pushed a commit to imrichardwu/airflow that referenced this pull request Sep 11, 2026
Serializing each client's whole send → receive → deserialize round trip
behind one lock (apache#69080) allows only a single request in flight at a
time, so a task issuing concurrent calls over the shared comm socket
pays a full, serialized round trip per call.
Every frame already carries a request id, so a single background
dispatcher can read the socket and route each response by id, with the
lock guarding only the write.

Signed-off-by: PoAn Yang <payang@apache.org>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants