Skip to content

[ML] Surface CCS skipped-cluster stats on datafeed extraction failure - #157567

Merged
edsavage merged 8 commits into
elastic:mainfrom
edsavage:ml-datafeed-ccs-stats-on-extraction-failure
Aug 27, 2026
Merged

edsavage merged 8 commits into
elastic:mainfrom
edsavage:ml-datafeed-ccs-stats-on-extraction-failure

Conversation

@edsavage

@edsavage edsavage commented Aug 24, 2026 •

Copy link
Copy Markdown
Contributor

Summary

When a datafeed search against a remote cluster throws ResourceNotFoundException (one or more
remote clusters were skipped), the extraction-failure code path in DatafeedJob.run() re-throws
before reaching crossClusterSearchStats.update(). As a result skipped_clusters is always 0
in the running datafeed stats API response even when every search round skips remote clusters.

This PR wires cluster-state metadata from the point of failure back to the stats update:

  • DataExtractor.getLinkedClusterStates() — new default interface method returning List.of(). Each extractor overrides it to expose the cluster states it has observed so far.
  • ScrollDataExtractor — updates lastLinkedClusterStates inside the checkForSkippedClusters catch block before re-throwing, so state is not lost when the scroll response is released.
  • AbstractAggregationDataExtractor / CompositeAggregationDataExtractor — the static executeSearchRequest() wraps ResourceNotFoundException in a new package-private SkippedClustersException carrying the extracted cluster states. The instance search() method catches this, updates lastLinkedClusterStates, then re-throws.
  • ChunkedDataExtractor — catches ResourceNotFoundException from the temporary summary extractor in setUpChunkedSearch(), merges its skipped cluster states into lastLinkedClusterStates before re-throwing, so states are not lost when the summary extractor is discarded. Also added a test for this path in ChunkedDataExtractorTests.
  • DatafeedJob.run() — calls dataExtractor.getLinkedClusterStates() in the extraction catch block and passes non-empty results to crossClusterSearchStats.update() before re-throwing.

Test plan

  • Existing CrossClusterSearchStatsTests and DatafeedJobTests continue to pass
  • Run full ML unit test suite: ./gradlew :x-pack:plugin:ml:test
  • Manual verification: configure a datafeed against an unavailable remote cluster and confirm skipped_clusters is non-zero in GET _ml/datafeeds/<id>/_stats

When a datafeed search throws ResourceNotFoundException (skipped remote
clusters), DatafeedJob's extraction catch block never reached the
crossClusterSearchStats.update() call, leaving skipped_clusters at zero
in the running datafeed stats API response.

Fix the pipeline so that cluster-state metadata observed before the
exception is propagated out:

* Each extractor stores cluster states in lastLinkedClusterStates and
  exposes them via a new DataExtractor.getLinkedClusterStates() default
  method.
* ScrollDataExtractor updates lastLinkedClusterStates in the
  checkForSkippedClusters catch block before re-throwing.
* AbstractAggregationDataExtractor and CompositeAggregationDataExtractor
  wrap the static executeSearchRequest() ResourceNotFoundException in a
  new package-private SkippedClustersException that carries the states
  through the static-method boundary, then update lastLinkedClusterStates
  when catching it in search().
* ChunkedDataExtractor.getLinkedClusterStates() delegates to the inner
  extractor so failures inside a chunk window propagate outward.
* DatafeedJob.run() calls dataExtractor.getLinkedClusterStates() in the
  extraction catch block and passes non-empty results to
  crossClusterSearchStats.update() before re-throwing.
@edsavage edsavage added >bug :ml Machine learning Team:ML Meta label for the ML team auto-backport Automatically create backport pull requests when merged v9.6.0 v8.19.21 v9.4.6 v9.5.3 labels Aug 24, 2026
@elasticsearchmachine

Copy link
Copy Markdown
Collaborator

Hi @edsavage, I've created a changelog YAML for you.

@github-actions

github-actions Bot commented Aug 24, 2026 •

Copy link
Copy Markdown
Contributor

🔍 Preview links for changed docs

⏳ Building and deploying preview... View progress

This comment will be updated with preview links when the build is complete.

@github-actions

Copy link
Copy Markdown
Contributor

ℹ️ Important: Docs version tagging

👋 Thanks for updating the docs! Just a friendly reminder that our docs are now cumulative. This means all 9.x versions are documented on the same page and published off of the main branch, instead of creating separate pages for each minor version.

We use applies_to tags to mark version-specific features and changes.

Expand for a quick overview

When to use applies_to tags:

✅ At the page level to indicate which products/deployments the content applies to (mandatory)
✅ When features change state (e.g. preview, ga) in a specific version
✅ When availability differs across deployments and environments

What NOT to do:

❌ Don't remove or replace information that applies to an older version
❌ Don't add new information that applies to a specific version without an applies_to tag
❌ Don't forget that applies_to tags can be used at the page, section, and inline level

🤔 Need help?

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.

Pull request overview

This PR fixes ML datafeed running stats reporting skipped_clusters: 0 when CCS searches fail due to skipped remote clusters, by propagating per-cluster metadata from extractor failure paths back to CrossClusterSearchStats.

Changes:

  • Add DataExtractor#getLinkedClusterStates() (default empty) and implement it in scroll/aggregation/chunked extractors to expose the most recently observed CCS cluster states (including on failures).
  • Capture linked cluster states before releasing SearchResponse when ResourceNotFoundException is thrown due to skipped clusters (scroll + aggregation paths).
  • Update DatafeedJob to update crossClusterSearchStats in the extraction failure catch block using extractor-provided partial states; add changelog entry.

Reviewed changes

Copilot reviewed 8 out of 8 changed files in this pull request and generated 4 comments.

Show a summary per file
File Description
x-pack/plugin/ml/src/main/java/org/elasticsearch/xpack/ml/datafeed/extractor/scroll/ScrollDataExtractor.java Captures linked cluster states on skipped-cluster failure before releasing the scroll response; exposes cached states via getLinkedClusterStates().
x-pack/plugin/ml/src/main/java/org/elasticsearch/xpack/ml/datafeed/extractor/DataExtractor.java Adds default getLinkedClusterStates() to allow failure-path CCS state propagation.
x-pack/plugin/ml/src/main/java/org/elasticsearch/xpack/ml/datafeed/extractor/chunked/ChunkedDataExtractor.java Delegates getLinkedClusterStates() to the active inner extractor so chunked extraction can surface failure-path states.
x-pack/plugin/ml/src/main/java/org/elasticsearch/xpack/ml/datafeed/extractor/aggregation/SkippedClustersException.java New internal exception used to carry linked cluster states extracted from a response that is about to be released.
x-pack/plugin/ml/src/main/java/org/elasticsearch/xpack/ml/datafeed/extractor/aggregation/CompositeAggregationDataExtractor.java Captures linked cluster states when the new skipped-cluster wrapper exception is thrown; exposes cached states.
x-pack/plugin/ml/src/main/java/org/elasticsearch/xpack/ml/datafeed/extractor/aggregation/AbstractAggregationDataExtractor.java Wraps skipped-cluster ResourceNotFoundException to retain linked cluster states prior to response release; exposes cached states.
x-pack/plugin/ml/src/main/java/org/elasticsearch/xpack/ml/datafeed/DatafeedJob.java Updates CCS stats even when extraction fails by pulling partial linked cluster states from the extractor in the exception path.
docs/changelog/157567.yaml Adds ML changelog entry for the bug fix.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

- Re-throw original ResourceNotFoundException from SkippedClustersException
  catch in AbstractAggregationDataExtractor and
  CompositeAggregationDataExtractor so the exception type and HTTP status
  code seen by callers and unwrapCause() are unchanged.
- Persist the richer linked-cluster state back to lastLinkedClusterStates
  in ChunkedDataExtractor.getLinkedClusterStates() so subsequent calls
  return the merged snapshot rather than recomputing from a potentially
  cleared inner extractor.
- Add DatafeedJobTests.testSkippedClustersStatsUpdatedOnExtractionFailure
  covering the failure path where the extractor throws while returning
  non-empty cluster states.
@edsavage
edsavage marked this pull request as ready for review August 25, 2026 01:17
@elasticsearchmachine

Copy link
Copy Markdown
Collaborator

Pinging @elastic/ml-core (Team:ML)

@edsavage
edsavage requested review from prwhelan and valeriy42 August 25, 2026 22:52
@edsavage
edsavage merged commit 3ce0b34 into elastic:main Aug 27, 2026
37 checks passed
@elasticsearchmachine

Copy link
Copy Markdown
Collaborator

💔 Backport failed

You can use sqren/backport to manually backport by running backport --upstream elastic/elasticsearch --pr 157567

@edsavage edsavage removed >bug backport pending :ml Machine learning Team:ML Meta label for the ML team auto-backport Automatically create backport pull requests when merged v9.6.0 v8.19.21 v9.4.6 v9.5.3 labels Aug 27, 2026
edsavage added a commit that referenced this pull request Aug 28, 2026
…ilure (#157933)

Backport of #157567 to 9.5.

Conflict resolution: DatafeedJobTests.java on 9.5 was missing the two new imports (ElasticsearchSecurityException, ResourceNotFoundException) and the three new test methods introduced on main alongside the fix. Both were cleanly added.
edsavage added a commit that referenced this pull request Aug 28, 2026
michalborek pushed a commit to michalborek/elasticsearch that referenced this pull request Sep 1, 2026
…elastic#157567)

When a datafeed search against a remote cluster throws ResourceNotFoundException (one or more
remote clusters were skipped), the extraction-failure code path in DatafeedJob.run() re-throws
before reaching crossClusterSearchStats.update(). As a result skipped_clusters is always 0
in the running datafeed stats API response even when every search round skips remote clusters.

This PR wires cluster-state metadata from the point of failure back to the stats update
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.

4 participants