Skip to content

fix: preserve LightGBM worker attempt identity - #2745

Merged
Rana Singh (ranadeepsingh) merged 8 commits into
microsoft:masterfrom
leninworld:fix-lightgbm-attempt-identity-issue-2699
Oct 1, 2026
Merged

Rana Singh (ranadeepsingh) merged 8 commits into
microsoft:masterfrom
leninworld:fix-lightgbm-attempt-identity-issue-2699

Conversation

@leninworld

@leninworld Lenin Mookiah (leninworld) commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

This PR addresses part of issue #2699.

Related Issues/PRs

Related to #2699

This builds on the first-failure guidance introduced in #2612. It does not close the broader reliability tracker.

What changes are proposed in this pull request?

Distributed LightGBM workers already have Spark stage and task-attempt identity available through TaskContext, but the private topology wire message and first-failure diagnostics retained only part of that identity. This made accepted, stale, duplicate, and unreachable reports difficult to map back to the exact Spark attempt that produced them.

This change:

  • extends the private worker message with optional stageId, taskAttemptId, and task attemptNumber fields alongside the existing stageAttemptNumber;
  • preserves every legacy/current wire layout, including IPv4, bracketed IPv6, and legacy unbracketed IPv6 forms;
  • selects the longest valid metadata layout for current hosts while retaining the shortest valid legacy IPv6 interpretation;
  • snapshots TaskContext once per topology request and reuses that identity in the worker message and unreachable-driver error;
  • adds the complete identity to existing accepted, stale, duplicate/replacement, and unreachable diagnostics without changing their log levels; and
  • explicitly rejects a colon-bearing executor ID before transmission, replacing a previously ambiguous positional parse with a clear error.

The public five-field TaskMessageInfo contract, topology behavior, retry limits, socket timeouts, dependencies, native loading, and estimator APIs are unchanged.

P2 addresses review comment r4129619856: the real-Spark test now obtains Spark through TestBase, resets it to local[1] before the test, and resets the shared provider afterward instead of stopping an unowned session. Production code and the test assertions are unchanged.

Compatibility notes:

  • A defensive message with absent metadata fields is supported by the parser, although that form does not have a dedicated fixture because topology registration runs inside a Spark task.
  • Under the pre-existing contract, a legacy missing stageAttemptNumber and a real value of 0 remain indistinguishable in diagnostics.
  • This is a bounded observability improvement and does not claim to resolve the other failure families tracked in [Tracking] Fabric distributed LightGBM reliability: fixes, regressions, and runtime delivery #2699.

How is this patch tested?

  • I have written tests (not required for typo or doc fix) and confirmed the proposed feature/bug-fix/change works.

Validated in Linux amd64 Docker with Java 11.0.32, Spark 3.5.0, and sbt 1.10.11:

  • lightgbm/compile — passed
  • lightgbm/Test/compile — passed
  • lightgbm/scalastyle — 42 files, 0 errors, 0 warnings
  • lightgbm/Test/scalastyle — 39 files, 0 errors, 0 warnings
  • WorkerWireFormatSuite and DriverSocketRetrySuite — 30 tests passed, 0 failed, 0 aborted
  • git diff --check — passed

The focused tests cover all supported metadata suffix lengths, current and legacy IPv4/IPv6/scoped-IPv6 messages, malformed metadata, exact stale/duplicate/unreachable diagnostic content, cause preservation, the public/private rendering boundary, and a real Spark task that verifies one self-consistent attempt-identity snapshot while preserving the shared TestBase Spark-session lifecycle.

The broader SynapseML suite was not run, and no upstream CI result is claimed. Dependency resolution used ordinary network access.

Validation evidence

SynapseML #2699 P2 validation evidence

Maintainer additions (folded in from #2746)

Rana Singh (@ranadeepsingh) added the numTasks fix from #2746 on top of the two commits above, so both parts of #2699 land together with credit to the original work. The contributor's commits are unchanged. The added commits are:

  • 70a853c7ef fix: expand LightGBM partitions to numTasks and explain missing-task timeouts. Without barrier mode, the driver waits for exactly numTasks tasks. coalesce can't add partitions, so with 1 input partition and numTasks=2 only one task ran and the driver waited 20 minutes before a misleading "Connection refused" error. The input is now repartitioned up to numTasks (by group column for LightGBMRanker). When numTasks is larger than the tasks Spark can run at once, the driver now fails with LightGBMMissingTasksException, which names the missing partitions, and the executor-side "could not reach the driver" message mentions that cause.
  • 09f82ca3f5 logs the extra shuffle at warn level and documents its cost in the LightGBM overview.
  • fbbbb22890 makes the new missing-partition unit test start its short accept timeout only after the first report is recorded, so slow agents can't flake it.
  • 898e99712c moves the unreachable-driver message builder into DriverUnreachableFailure.scala with no wording change. Combining both PRs put NetworkManager.scala at 821 lines, over the 800-line style limit.
  • ecd71ff2dc routes every input partition-count read (automatic numTasks, barrier mode, the ranker, and the new check) through one documented inputPartitionCount helper, after review flagged the direct df.rdd call. No behavior change.
  • fcb89ce1dd extends this PR's executor-id guard after review: besides :, it now rejects an empty id, = (the executor partition list delimiter), and control characters such as CR and LF before transmission. These already broke the topology messages, but only failed later with a confusing parse error.

The unreachable-driver message keeps this PR's full attempt identity and adds the missing-tasks hint for first attempts.

What doesn't change: automatic numTasks, explicit numTasks at or below the partition count (still coalesce), and barrier mode. There are no public signature or parameter changes.

Cost: an explicit numTasks on uncached AQE input reads the partition count, which runs a pending upstream shuffle one extra time. On Fabric, a 20M-row aggregate fit took 10.2s median versus 7.9s on master. Automatic numTasks was about 9.5s on both. Caching the input avoids it.

Validation of the combined branch (JDK 11, 2 CPUs):

  • lightgbm/scalastyle and lightgbm/Test/scalastyle are clean.
  • DriverSocketRetrySuite, LightGBMNumTasksSuite, and WorkerWireFormatSuite passed 39/39 in 4 consecutive runs.
  • After ecd71ff2dc: LightGBMNumTasksSuite, LightGBMRankerPartitionSuite, and DriverSocketRetrySuite passed 30/30.
  • After fcb89ce1dd: WorkerWireFormatSuite, DriverSocketRetrySuite, and NetworkManagerSuite passed 48/48.
  • Azure build 238451992 on fcb89ce1dd (current head): 408 of 409 LightGBM tests passed, with no failures, including WorkerWireFormatSuite 16/16, DriverSocketRetrySuite 17/17, and LightGBMNumTasksSuite 7/7. No test failed anywhere in the build. UnitTests io1 failed on an Azure login certificate error, and the spark4.1 release compatibility check failed as it does on other current PRs.
  • Azure build 238439461 on ecd71ff2dc: 407 of 408 LightGBM tests passed, with no failures. The other is the ignored performance test. The only failed job, UnitTests speech1, hit an Azure login certificate error.
  • Azure build 238431530 on 898e99712c: 374 of 375 LightGBM tests passed and none failed (the other is "Performance testing", which is marked ignore in the source), including LightGBMNumTasksSuite 7/7, DriverSocketRetrySuite 17/17, and WorkerWireFormatSuite 15/15. The failed jobs were infrastructure: an Azure login certificate error stopped UnitTests lightgbm5 before tests ran, and ImageFeaturizerSuite failed on a corrupted cached ONNX model download.

Fabric end-to-end on this branch (head 898e99712c, jars ...-232-898e9971-SNAPSHOT, Spark 3.5.5, 1 executor with 8 slots, AQE on, bulk mode): the classifier, ranker, and empty-partition cases passed, and numTasks=9 on 8 slots failed with LightGBMMissingTasksException naming partition 8. The driver log shows the new attempt-identity worker messages being sent and parsed, for example enabledTask:10.3.160.5:12411:2:1:0:5:5:0. The earlier master comparison, run on the #2746 jars:

Case master #2746
Classifier, 1 partition, numTasks=4 "could not reach the driver" passes, accuracy 1.0
Ranker, 1 partition, numTasks=4, grouping off same error passes
Ranker with empty partitions same error passes
numTasks=9 on 8 slots generic retry error "received network reports from 8 of 9 ... Missing partitions: 8"

Streaming mode without slotNames fails on this Fabric runtime for both master and this branch. That's the existing #2242.

Does this PR change any dependencies?

  • No. You can skip this section.
  • Yes. Make sure the dependencies are resolved correctly, and list changes here.

Does this PR add a new feature? If so, have you added samples on website?

  • No. You can skip this section.
  • Yes. Make sure you have added samples following below steps.
  1. Find the corresponding markdown file for your new feature in website/docs/documentation folder.
    Make sure you choose the correct class estimators/transformers and namespace.
  2. Follow the pattern in markdown file and add another section for your new API, including pyspark, scala (and .NET potentially) samples.
  3. Make sure the DocTable points to correct API link.
  4. Navigate to website folder, and run npm start to make sure the website renders correctly.
  5. Don't forget to add <!--pytest-codeblocks:cont--> before each python code blocks to enable auto-tests for python samples.
  6. Make sure the WebsiteSamplesTests job pass in the pipeline.

Copilot AI balanced review requested due to automatic review settings September 29, 2026 04:23
@azure-pipelines

Copy link
Copy Markdown
Azure Pipelines:
There may be pipelines that require an authorized user to comment /azp run to run.

@github-actions

Copy link
Copy Markdown

Hey Lenin Mookiah (@leninworld) 👋!
Thank you so much for contributing to our repository 🙌.
Someone from SynapseML Team will be reviewing this pull request soon.

We use semantic commit messages to streamline the release process.
Before your pull request can be merged, you should make sure your first commit and PR title start with a semantic prefix.
This helps us to create release messages and credit you for your hard work!

Examples of commit messages with semantic prefixes:

  • fix: Fix LightGBM crashes with empty partitions
  • feat: Make HTTP on Spark back-offs configurable
  • docs: Update Spark Serving usage
  • build: Add codecov support
  • perf: improve LightGBM memory usage
  • refactor: make python code generation rely on classes
  • style: Remove nulls from CNTKModel
  • test: Add test coverage for CNTKModel

To test your commit locally, please follow our guild on building from source.
Check out the developer guide for additional guidance on testing your change.

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Copilot review overview

🟡 Changes recommended

The new test can stop the shared Spark context and cause order-dependent failures in later suites.

Review effort: Balanced
Findings: 1 High severity

Open (1)
What changed in this PR

Preserves complete Spark attempt identity in LightGBM worker messages and diagnostics while maintaining wire compatibility.

Changes:

  • Extends worker wire metadata with stage and task attempt identifiers.
  • Reuses a single task-identity snapshot across registration and failure diagnostics.
  • Adds parsing, compatibility, and diagnostic tests.
File Description
WorkerMessage.scala Extends wire parsing, formatting, and identity summaries.
NetworkManager.scala Captures and logs complete worker identity.
WorkerWireFormatSuite.scala Tests wire compatibility and malformed metadata.
DriverSocketRetrySuite.scala Tests diagnostics and real Spark task identity.

💡 Configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Copilot AI review requested due to automatic review settings September 29, 2026 05:21

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Copilot review overview

🔵 Needs a closer look

The compatibility-sensitive socket protocol and distributed retry behavior warrant final human review despite strong focused coverage.

Review effort: Balanced
Findings: None

Resolved since last review (1)

Ranadeep Singh and others added 4 commits September 30, 2026 23:37
…timeouts

Without barrier execution mode, the LightGBM driver waits for exactly numTasks
training tasks to report. Two cases left it waiting until the timeout, after
which tasks failed with a misleading "Connection refused" error:

- An explicit numTasks larger than the input partition count. coalesce can
  only merge partitions, so fewer tasks ran than the driver expected. The
  input is now repartitioned up to numTasks. The ranker repartitions by its
  grouping column so query groups stay whole. The partition count is only
  read for an explicit numTasks, so automatic sizing adds no extra work.
- numTasks larger than the tasks Spark can run at once. This can't be fixed
  automatically, so the driver now fails with LightGBMMissingTasksException,
  which lists how many tasks reported and which partitions are missing. The
  job failure is attached as a suppressed exception.

Adds regression tests that fail on master, and a troubleshooting note in the
LightGBM overview. Addresses part of microsoft#2699.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
…extra pass

Log the numTasks expansion at warn level, because it adds a shuffle the
user did not ask for, and say how to avoid it. Document that an explicit
numTasks now reads the input partition count, like automatic numTasks and
barrier mode already do, which runs a pending adaptive shuffle once more
on uncached input.

Fabric E2E (runtime Spark 3.5.5, 1 executor x 8 slots, AQE on) measured
explicit numTasks fits at 10.2s median versus 7.9s on master for an
uncached 20M-row aggregate, matching one extra upstream pass.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
The missing-partition test set a 500ms accept timeout before the fake
worker connected, so a slow agent could time out with zero reports and
fail the assertion on partition 1. The driver accepts one connection at a
time, so the helper now applies the short timeout on the second accept,
after the first report is recorded, and the test waits for that point.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Combining the attempt-identity changes with the numTasks expansion fix
pushed NetworkManager.scala past the 800-line style limit. The message
builder for a task that cannot reach the driver is self-contained, so it
now lives in DriverUnreachableFailure.scala with no change in behavior or
wording.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Copilot AI balanced review requested due to automatic review settings September 30, 2026 23:46
@ranadeepsingh

Copy link
Copy Markdown
Collaborator

/azp run

@azure-pipelines

Copy link
Copy Markdown
Azure Pipelines:
Successfully started running 1 pipeline(s).

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Copilot review overview

🟡 Changes recommended

The new direct RDD API usage violates the repository’s Spark Connect and managed-mode compatibility requirement.

Review effort: Balanced
Findings: 1 Medium severity

Open (1)

Review flagged the new numTasks check for calling df.rdd directly. Spark
has no Dataset-level partition count, and the same read was already used
in three other places. All four now go through inputPartitionCount,
which documents that it only inspects the planned input on the driver
and can run pending adaptive shuffle stages early. No behavior change.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Copilot AI balanced review requested due to automatic review settings October 1, 2026 01:03
@ranadeepsingh

Copy link
Copy Markdown
Collaborator

/azp run

@azure-pipelines

Copy link
Copy Markdown
Azure Pipelines:
Successfully started running 1 pipeline(s).

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Copilot review overview

🟡 Changes recommended

Executor IDs containing control characters or = can still corrupt the topology wire protocols.

Review effort: Balanced
Findings: 1 High severity

Open (1)
Resolved since last review (1)
Previously missed (1)

In code that hasn't changed since last review

Low severity Correct cost statement for the non-barrier ranker path

docs/​Explore Algorithms/​LightGBM/​Overview.md:425

This cost statement is inaccurate for the default non-barrier ranker path: LightGBMRanker.scala:95-100 performs the grouping repartition directly and deliberately skips inputPartitionCount. Qualify the statement so ranker users are not told to expect an extra input-partition read that does not occur.

…sages

The executor id is sent verbatim in the worker line protocol and in the
driver's executor=partitions list. The existing guard only rejected ':'.
An empty id, '=' or a control character such as CR or LF also breaks
those messages, and previously failed later with a confusing parse error.
They are now rejected before transmission with a clear message. The ':'
message is unchanged.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Copilot AI balanced review requested due to automatic review settings October 1, 2026 02:29
@ranadeepsingh

Copy link
Copy Markdown
Collaborator

/azp run

@azure-pipelines

Copy link
Copy Markdown
Azure Pipelines:
Successfully started running 1 pipeline(s).

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Copilot review overview

🔵 Needs a closer look

Distributed protocol, scheduling, and repartitioning changes warrant final human review despite extensive targeted validation.

Review effort: Balanced
Findings: None

Resolved since last review (1)

@ranadeepsingh

Copy link
Copy Markdown
Collaborator

Thanks Lenin Mookiah (@leninworld) for this. Having the full Spark attempt identity on every worker report makes LightGBM topology failures much easier to trace back to the exact task attempt, and your guard against a : in the executor id fixed a real ambiguity in the wire parser.

To keep the #2699 work in one place and with your name on it, I folded my follow-up (#2746) into this PR instead of merging it separately. Your two commits are unchanged. I added these on top:

  • 70a853c, 09f82ca, fbbbb22: without barrier mode, an explicit numTasks above the input partition count now expands the input instead of hanging for 20 minutes, and when tasks can't all run at once the driver fails with an error that names the missing partitions. I merged the unreachable-driver message so it keeps your attempt identity and adds that hint.
  • 898e997: moved that message builder into DriverUnreachableFailure.scala, because the two changes together put NetworkManager.scala over the 800-line style limit.
  • ecd71ff: one helper for reading the input partition count, after review flagged a direct df.rdd call.
  • fcb89ce: extended your executor id guard to also reject an empty id, =, and control characters, which review pointed out also break the protocol.

On Azure build 238451992 for fcb89ce, 408 of 409 LightGBM tests passed and none failed (the other is marked ignore), including WorkerWireFormatSuite 16/16. No test failed anywhere in the build. Two jobs did fail: UnitTests io1 on an Azure login certificate error, and the spark4.1 release compatibility check, which also fails on other current PRs. I also ran this branch on Microsoft Fabric. The numTasks cases passed, and the driver log shows your extended worker messages being sent and parsed.

Please take a look and sign off if this fits what you had in mind. If you'd rather keep this PR to the identity change only, say so and I'll move my commits back out to #2746.

A maintainer approval is still needed before merge.

@leninworld

Copy link
Copy Markdown
Contributor Author

Thanks Lenin Mookiah (Lenin Mookiah (@leninworld)) for this. Having the full Spark attempt identity on every worker report makes LightGBM topology failures much easier to trace back to the exact task attempt, and your guard against a : in the executor id fixed a real ambiguity in the wire parser.

To keep the #2699 work in one place and with your name on it, I folded my follow-up (#2746) into this PR instead of merging it separately. Your two commits are unchanged. I added these on top:

  • 70a853c, 09f82ca, fbbbb22: without barrier mode, an explicit numTasks above the input partition count now expands the input instead of hanging for 20 minutes, and when tasks can't all run at once the driver fails with an error that names the missing partitions. I merged the unreachable-driver message so it keeps your attempt identity and adds that hint.
  • 898e997: moved that message builder into DriverUnreachableFailure.scala, because the two changes together put NetworkManager.scala over the 800-line style limit.
  • ecd71ff: one helper for reading the input partition count, after review flagged a direct df.rdd call.
  • fcb89ce: extended your executor id guard to also reject an empty id, =, and control characters, which review pointed out also break the protocol.

On Azure build 238451992 for fcb89ce, 408 of 409 LightGBM tests passed and none failed (the other is marked ignore), including WorkerWireFormatSuite 16/16. No test failed anywhere in the build. Two jobs did fail: UnitTests io1 on an Azure login certificate error, and the spark4.1 release compatibility check, which also fails on other current PRs. I also ran this branch on Microsoft Fabric. The numTasks cases passed, and the driver log shows your extended worker messages being sent and parsed.

Please take a look and sign off if this fits what you had in mind. If you'd rather keep this PR to the identity change only, say so and I'll move my commits back out to #2746.

A maintainer approval is still needed before merge.

Thanks Rana Singh (@ranadeepsingh). The combined scope fits what I had in mind, and I’m happy to keep your follow-up commits in this PR rather than separate them into #2746.

The numTasks expansion, missing-partition diagnostics, compatibility-safe partition-count helper, and stricter executor-id validation complement the original worker-attempt identity change well. The LightGBM test results and Fabric validation also address my concerns.

Please proceed with the remaining maintainer review and approval. Thanks for preserving my original commits and consolidating the #2699 work here.

@ranadeepsingh

Rana Singh (ranadeepsingh) commented Oct 1, 2026 •

Copy link
Copy Markdown
Collaborator

I reviewed every file at head fcb89ce1dd against the current base.

File by file

File Change Checked
docs/.../LightGBM/Overview.md New troubleshooting section for missing tasks, plus a note on the extra cost of an explicit numTasks Accurate, matches the code
DriverUnreachableFailure.scala Moved out of NetworkManager to keep it under the line limit. Same message, now with the attempt identity and the missing-tasks hint Wording preserved, the old helper was package-private
LightGBMBase.scala An explicit numTasks above the input partition count now repartitions so every task gets data. Partition count reads go through one inputPartitionCount helper. A missing-tasks failure closes the driver sockets and fails fast Automatic numTasks and barrier mode unchanged. Sockets are still closed in the outer finally
LightGBMMissingTasksException.scala New exception with a clear message about which tasks never reported Constructor is package-private
LightGBMRanker.scala Grouped ranking repartitions by the group column, so groups stay together when partitions expand Same result as the #2744 path. Barrier mode unchanged
NetworkManager.scala Workers report the full Spark attempt identity. A non-barrier socket timeout raises the missing-tasks error TaskContext is read once per worker
WorkerMessage.scala Parser accepts the old and new message layouts, including IPv6 hosts. Executor IDs with :, =, control characters, or an empty value are rejected with a clear error Old worker messages still parse. The rejected IDs already broke parsing before this PR
DriverSocketRetrySuite.scala Covers the missing-tasks failure path Shared Spark session restored in finally
LightGBMNumTasksSuite.scala Fits with numTasks above the partition count for the classifier, regressor and grouped ranker Exercises the public fit path
WorkerWireFormatSuite.scala Round trips for every message layout and the new executor ID checks Includes legacy and IPv6 cases

Safety

  • Only Scala sources, tests and one doc page. No build, dependency, workflow or pipeline changes.
  • No binary files, non-ASCII or hidden characters, process execution, reflection, deserialization, new network endpoints, environment reads or secrets.
  • Commit authorship is clean: 2 commits from Lenin Mookiah (@leninworld), 6 maintainer commits on top, no history rewrites.

Regressions

  • No public method or parameter was removed or changed, so no Python wrapper changes.
  • The only cost is one extra pass over uncached input when a user sets numTasks explicitly. On Fabric that fit took 10.2s vs 7.9s on master. Automatic numTasks was about 9.5s on both. The docs recommend caching for this case.
  • master has only CodeQL workflow bumps since the base, and the PR still merges cleanly.

Validation

  • Local, exact head: scalastyle clean for main and test, 70 of 70 tests passed across the six LightGBM network, numTasks, ranker and wire format suites.
  • Azure build 238451992: LightGBM 408 of 409 passed (the other is the ignored performance test). The only red jobs were an infrastructure failure in io1 and the Spark 4.1 compatibility check. That check fails because spark4.1 has not yet received fix: preserve LightGBM ranker grouping under AQE #2744, so the LightGBMRanker.scala patch conflicts. It needs the usual master sync to spark4.1, not a code change here.
  • Fabric end to end: the job fails on master and passes with this PR. The driver log shows the new worker messages.
  • Lenin Mookiah (@leninworld) signed off on the combined scope.

Recommended for the next SynapseML release.

@ranadeepsingh
Rana Singh (ranadeepsingh) merged commit 6ce51f7 into microsoft:master Oct 1, 2026
80 of 81 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants