Repository navigation
[ML] Surface CCS skipped-cluster stats on datafeed extraction failure - #157567
Conversation
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.
|
Hi @edsavage, I've created a changelog YAML for you. |
🔍 Preview links for changed docs⏳ Building and deploying preview... View progress This comment will be updated with preview links when the build is complete. |
ℹ️ 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 overviewWhen to use applies_to tags:✅ At the page level to indicate which products/deployments the content applies to (mandatory) What NOT to do:❌ Don't remove or replace information that applies to an older version 🤔 Need help?
|
There was a problem hiding this comment.
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
SearchResponsewhenResourceNotFoundExceptionis thrown due to skipped clusters (scroll + aggregation paths). - Update
DatafeedJobto updatecrossClusterSearchStatsin 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.
|
Pinging @elastic/ml-core (Team:ML) |
💔 Backport failedYou can use sqren/backport to manually backport by running |
…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
Summary
When a datafeed search against a remote cluster throws
ResourceNotFoundException(one or moreremote clusters were skipped), the extraction-failure code path in
DatafeedJob.run()re-throwsbefore reaching
crossClusterSearchStats.update(). As a resultskipped_clustersis always 0in 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 returningList.of(). Each extractor overrides it to expose the cluster states it has observed so far.ScrollDataExtractor— updateslastLinkedClusterStatesinside thecheckForSkippedClusterscatch block before re-throwing, so state is not lost when the scroll response is released.AbstractAggregationDataExtractor/CompositeAggregationDataExtractor— the staticexecuteSearchRequest()wrapsResourceNotFoundExceptionin a new package-privateSkippedClustersExceptioncarrying the extracted cluster states. The instancesearch()method catches this, updateslastLinkedClusterStates, then re-throws.ChunkedDataExtractor— catchesResourceNotFoundExceptionfrom the temporary summary extractor insetUpChunkedSearch(), merges its skipped cluster states intolastLinkedClusterStatesbefore re-throwing, so states are not lost when the summary extractor is discarded. Also added a test for this path inChunkedDataExtractorTests.DatafeedJob.run()— callsdataExtractor.getLinkedClusterStates()in the extraction catch block and passes non-empty results tocrossClusterSearchStats.update()before re-throwing.Test plan
CrossClusterSearchStatsTestsandDatafeedJobTestscontinue to pass./gradlew :x-pack:plugin:ml:testskipped_clustersis non-zero inGET _ml/datafeeds/<id>/_stats