Skip to content
Merged
Show file tree
Hide file tree
Changes from 31 commits
Commits
Show all changes
447 commits
Select commit Hold shift + click to select a range
6fcc1ca
use a queue for async consumption of lucene generated bytes
drempapis Jan 20, 2026
36d0dd3
create tests form the Producer-Consumer pattern
drempapis Jan 21, 2026
61191cc
Merge branch 'chunked_fetch_phase' of github.com:drempapis/elasticsea…
drempapis Jan 21, 2026
15b2775
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 21, 2026
1e35a61
Merge branch 'chunked_fetch_phase' of github.com:drempapis/elasticsea…
drempapis Jan 21, 2026
768b4f3
apply spot and update transport version
drempapis Jan 21, 2026
e4e4aeb
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 21, 2026
c9be35b
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 21, 2026
821ab00
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 21, 2026
663d412
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 22, 2026
fc2b941
fix test
drempapis Jan 22, 2026
c72e681
Use a ThrottledTaskRunner rather than a custom producer/consumer impl…
drempapis Jan 23, 2026
26ed0eb
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 23, 2026
8527069
[CI] Auto commit changes from spotless
Jan 23, 2026
cd5b924
Add tests and code improvements
drempapis Jan 23, 2026
32fb8e9
Merge branch 'chunked_fetch_phase' of github.com:drempapis/elasticsea…
drempapis Jan 23, 2026
29fc625
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 23, 2026
fea70b4
Merge branch 'chunked_fetch_phase' of github.com:drempapis/elasticsea…
drempapis Jan 23, 2026
34e1e6b
update transport version + spotless
drempapis Jan 23, 2026
9dfb471
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 23, 2026
de64257
update test code
drempapis Jan 26, 2026
c41913a
Merge branch 'chunked_fetch_phase' of github.com:drempapis/elasticsea…
drempapis Jan 26, 2026
7f6f607
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 26, 2026
5464e71
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 26, 2026
b67c080
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 26, 2026
71d9267
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 26, 2026
8fbbc7a
update code for counting on cb-bytes
drempapis Jan 26, 2026
896f6c7
add test
drempapis Jan 26, 2026
35a62f7
merge master
drempapis Jan 26, 2026
704dcaa
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 26, 2026
154d081
[CI] Auto commit changes from spotless
Jan 26, 2026
807eca5
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 26, 2026
b61635f
update tests
drempapis Jan 26, 2026
d252561
Merge branch 'chunked_fetch_phase' of github.com:drempapis/elasticsea…
drempapis Jan 26, 2026
1e6b79b
fix tests
drempapis Jan 27, 2026
6836419
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 27, 2026
014fadb
update comment
drempapis Jan 27, 2026
d5b6bdd
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 27, 2026
9912c64
Merge branch 'chunked_fetch_phase' of github.com:drempapis/elasticsea…
drempapis Jan 27, 2026
9aa11f0
update transport version
drempapis Jan 27, 2026
e43612f
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 28, 2026
f1d8849
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 28, 2026
10dca8b
revert changes
drempapis Jan 28, 2026
f801944
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 28, 2026
e0bec49
spotless apply
drempapis Jan 28, 2026
07f3ac6
Merge branch 'chunked_fetch_phase' of github.com:drempapis/elasticsea…
drempapis Jan 28, 2026
7c81453
add test
drempapis Jan 28, 2026
5104d2d
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 28, 2026
1bcafe0
Fix Leak test
drempapis Jan 28, 2026
ed3d896
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 28, 2026
942db6b
[CI] Auto commit changes from spotless
Jan 28, 2026
aeef791
make configurable a parameter
drempapis Jan 29, 2026
c2dc153
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 29, 2026
e050efc
update transport version|
drempapis Jan 29, 2026
e0c2233
[CI] Auto commit changes from spotless
Jan 29, 2026
b85301c
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 29, 2026
75c24d3
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 29, 2026
5fe6c00
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 30, 2026
6f6a820
add transport version
drempapis Jan 30, 2026
b7cec57
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 30, 2026
c9ed3c9
Revert code
drempapis Jan 30, 2026
4504583
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 30, 2026
e0b5c28
Merge branch 'main' into chunked_fetch_phase
drempapis Jan 30, 2026
57c7ce4
Remove unnecessary newline in warnings method
drempapis Jan 30, 2026
7526348
Remove unnecessary newline in warnings method
drempapis Jan 30, 2026
c041721
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 2, 2026
4a1f246
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 2, 2026
70715f8
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 3, 2026
d52b2aa
update transport version
drempapis Feb 3, 2026
d325219
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 3, 2026
0f29d58
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 5, 2026
9eea215
update transport version
drempapis Feb 5, 2026
20c0214
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 5, 2026
bd3b38e
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 5, 2026
6240ebf
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 6, 2026
a0f9e56
update after review
drempapis Feb 6, 2026
8d02d05
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 6, 2026
bc75863
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 6, 2026
ad0c461
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 10, 2026
dbbb1e5
update transport version
drempapis Feb 10, 2026
427e9b3
update test
drempapis Feb 10, 2026
64da35d
update test
drempapis Feb 10, 2026
de183d0
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 18, 2026
f6b5700
update transport version
drempapis Feb 18, 2026
0a2f51a
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 21, 2026
eb88283
update transport version
drempapis Feb 21, 2026
736f68e
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 22, 2026
dfe50a2
update transport version
drempapis Feb 22, 2026
ef2d4f7
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 23, 2026
eb22906
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 24, 2026
b2d236a
update transport version
drempapis Feb 24, 2026
dede2e8
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 26, 2026
9d5b44d
update transport version
drempapis Feb 26, 2026
a792941
Convert ResponseStreamKey to a record into ActiveFetchPhaseTasks
drempapis Feb 26, 2026
5eb105e
Remove unused standard mode from TransportFetchPhaseResponseChunkAction
drempapis Feb 26, 2026
531d9b3
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 26, 2026
4ec63d3
update javadoc
drempapis Feb 26, 2026
c21c027
update javadoc
drempapis Feb 26, 2026
71ba8ee
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 26, 2026
33c45c8
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 26, 2026
193757f
update javadoc
drempapis Feb 26, 2026
f8d4da5
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 26, 2026
8d0923f
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 26, 2026
a288a34
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 26, 2026
c621cff
Merge branch 'main' into chunked_fetch_phase
drempapis Feb 26, 2026
3175d5f
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 3, 2026
582a741
update transport version
drempapis Mar 3, 2026
96c459e
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 3, 2026
75165b8
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 4, 2026
9323988
update transport version
drempapis Mar 4, 2026
467fe07
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 4, 2026
3312778
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 4, 2026
6e3c29d
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 4, 2026
c499cdc
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 4, 2026
b7656f6
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 5, 2026
e4e3d39
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 5, 2026
c44611a
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 5, 2026
7641b39
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 6, 2026
04d2593
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 6, 2026
84a1cd1
update after review
drempapis Mar 6, 2026
0c855aa
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 6, 2026
5e41578
update after review
drempapis Mar 9, 2026
19cd32e
Remove Type.HITS enum
drempapis Mar 9, 2026
9c7145e
[CI] Auto commit changes from spotless
Mar 9, 2026
6cab3fa
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 10, 2026
955c108
update transport version
drempapis Mar 10, 2026
8856115
update after review
drempapis Mar 10, 2026
cc9b26d
update after review
drempapis Mar 10, 2026
dbb6dde
make method more readable
drempapis Mar 10, 2026
5702c5d
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 10, 2026
27aa07a
[CI] Auto commit changes from spotless
Mar 10, 2026
dd81129
update test
drempapis Mar 10, 2026
6af47ca
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 10, 2026
4fbfaa8
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 10, 2026
79abf17
update after review
drempapis Mar 10, 2026
987bcb5
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 10, 2026
90d865d
Merge branch 'chunked_fetch_phase' of github.com:drempapis/elasticsea…
drempapis Mar 10, 2026
54ed6d8
update after review
drempapis Mar 10, 2026
c22510d
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 10, 2026
839ada3
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 11, 2026
24eb525
update transport versionb
drempapis Mar 11, 2026
011d81c
update after review
drempapis Mar 11, 2026
075a797
Track chunked fetch stream allocations on request breaker
drempapis Mar 11, 2026
b8085dd
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 11, 2026
32a85b8
[CI] Auto commit changes from spotless
Mar 11, 2026
57f8f4b
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 11, 2026
955d254
remove redundant close
drempapis Mar 11, 2026
11a5cf4
Use ActionListener helpers for the FetchPhase
drempapis Mar 12, 2026
9f5eda2
update after review
drempapis Mar 12, 2026
dd05ba0
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 12, 2026
c9dfe26
update transport version
drempapis Mar 12, 2026
2b01ecd
update after review
drempapis Mar 12, 2026
a7ab9c7
update after review
drempapis Mar 12, 2026
d4706c4
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 12, 2026
dd54b4a
[CI] Auto commit changes from spotless
Mar 12, 2026
328dc1b
update after review
drempapis Mar 12, 2026
99f1fe7
Merge branch 'chunked_fetch_phase' of github.com:drempapis/elasticsea…
drempapis Mar 12, 2026
23a2c69
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 12, 2026
8e37435
Merge branch 'chunked_fetch_phase' of github.com:drempapis/elasticsea…
drempapis Mar 12, 2026
ea6c433
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 12, 2026
d06125a
update after review
drempapis Mar 12, 2026
a1b26be
Merge branch 'chunked_fetch_phase' of github.com:drempapis/elasticsea…
drempapis Mar 12, 2026
a82402a
[CI] Auto commit changes from spotless
Mar 12, 2026
1d740a2
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 12, 2026
abe0131
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 13, 2026
aef0bdd
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 13, 2026
45360b4
[CI] Auto commit changes from spotless
Mar 13, 2026
c6f2c25
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 16, 2026
638d634
Add tests and update transport version
drempapis Mar 16, 2026
b4b030a
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 16, 2026
0572cc3
[CI] Auto commit changes from spotless
Mar 16, 2026
b38157c
Add Unit tests to excercise various scenarios
drempapis Mar 16, 2026
40f51f2
[CI] Auto commit changes from spotless
Mar 16, 2026
4a92a3d
Add more ITests
drempapis Mar 16, 2026
eddc5b7
Merge branch 'chunked_fetch_phase' of github.com:drempapis/elasticsea…
drempapis Mar 16, 2026
f004e1b
[CI] Auto commit changes from spotless
Mar 16, 2026
a10b57c
Remove unused timestamp and redundant from metadata from fetch chunks
drempapis Mar 16, 2026
034584d
Remove unused timestamp and redundant from metadata from fetch chunks
drempapis Mar 16, 2026
2ff5b22
[CI] Auto commit changes from spotless
Mar 16, 2026
983c11c
update test limits
drempapis Mar 16, 2026
5e0c5ce
update test limits
drempapis Mar 16, 2026
8b96036
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 16, 2026
2c1777e
update transport version
drempapis Mar 16, 2026
8cc3966
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 16, 2026
16b3cee
update after review
drempapis Mar 16, 2026
575f192
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 17, 2026
efd3856
update transport version
drempapis Mar 17, 2026
b0e6413
Reject negative shard ids in fetch chunks
drempapis Mar 17, 2026
bdace3b
[CI] Auto commit changes from spotless
Mar 17, 2026
faa3842
Remove shard id validation from fetch chunks
drempapis Mar 17, 2026
95de3c2
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 17, 2026
512cac4
Merge branch 'chunked_fetch_phase' of github.com:drempapis/elasticsea…
drempapis Mar 17, 2026
e0cb817
[CI] Auto commit changes from spotless
Mar 17, 2026
b20dadc
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 17, 2026
4193d6f
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 17, 2026
dcaa330
revert code in Streaming fetch for when reading Lucene related data
drempapis Mar 18, 2026
2f89ea5
[CI] Auto commit changes from spotless
Mar 18, 2026
b67b5ce
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 18, 2026
cc85fc0
update transport version
drempapis Mar 18, 2026
83d7329
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 18, 2026
367a920
update after review
drempapis Mar 18, 2026
70842ef
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 18, 2026
caeba1f
update after review
drempapis Mar 18, 2026
bbf27b0
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 18, 2026
c729198
update after release
drempapis Mar 18, 2026
056f27c
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 18, 2026
4fec8eb
[CI] Auto commit changes from spotless
Mar 18, 2026
a9e237c
update after review
drempapis Mar 18, 2026
95dea71
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 18, 2026
532a44a
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 19, 2026
a01efc5
Merge branch 'chunked_fetch_phase' of github.com:drempapis/elasticsea…
drempapis Mar 19, 2026
e3b0296
update transport version
drempapis Mar 19, 2026
faa056d
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 19, 2026
6e4eb95
update transport version
drempapis Mar 19, 2026
51aa5f1
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 19, 2026
a987b96
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 19, 2026
f2ad449
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 19, 2026
1f3ca15
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 19, 2026
c591ca8
update after merge
drempapis Mar 19, 2026
6509139
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 19, 2026
f064ae1
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 19, 2026
7ce7e71
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 19, 2026
fd471f4
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 19, 2026
2dfd675
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 20, 2026
10e2a92
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 20, 2026
c79fb54
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 20, 2026
27ca22b
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 20, 2026
0ca27b9
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 20, 2026
ba57235
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 20, 2026
c644e32
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 21, 2026
0705998
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 21, 2026
97d3ed8
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 22, 2026
f1fb579
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 22, 2026
57b89c4
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 24, 2026
cf53b1c
update transport version
drempapis Mar 24, 2026
935f673
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 24, 2026
ea0bf5e
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 24, 2026
93fe7f0
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 25, 2026
14ca541
update transport version
drempapis Mar 25, 2026
b09529e
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 25, 2026
4d5d9f1
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 25, 2026
79d57a8
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 25, 2026
4ce2b1b
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 25, 2026
1b77ed6
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 25, 2026
e861874
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 27, 2026
d497b5b
update code
drempapis Mar 27, 2026
1528d63
update after review
drempapis Mar 27, 2026
7f90408
[CI] Auto commit changes from spotless
Mar 27, 2026
4976248
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 27, 2026
3916b53
Merge branch 'main' into chunked_fetch_phase
drempapis Mar 27, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -423,6 +423,8 @@
import org.elasticsearch.rest.action.synonyms.RestGetSynonymsSetsAction;
import org.elasticsearch.rest.action.synonyms.RestPutSynonymRuleAction;
import org.elasticsearch.rest.action.synonyms.RestPutSynonymsAction;
import org.elasticsearch.search.fetch.chunk.ActiveFetchPhaseTasks;
import org.elasticsearch.search.fetch.chunk.TransportFetchPhaseResponseChunkAction;
import org.elasticsearch.snapshots.TransportUpdateSnapshotStatusAction;
import org.elasticsearch.tasks.Task;
import org.elasticsearch.telemetry.TelemetryProvider;
Expand Down Expand Up @@ -761,6 +763,7 @@ public <Request extends ActionRequest, Response extends ActionResponse> void reg
actions.register(TransportMultiSearchAction.TYPE, TransportMultiSearchAction.class);
actions.register(TransportExplainAction.TYPE, TransportExplainAction.class);
actions.register(TransportClearScrollAction.TYPE, TransportClearScrollAction.class);
actions.register(TransportFetchPhaseResponseChunkAction.TYPE, TransportFetchPhaseResponseChunkAction.class);
actions.register(RecoveryAction.INSTANCE, TransportRecoveryAction.class);
actions.register(TransportNodesReloadSecureSettingsAction.TYPE, TransportNodesReloadSecureSettingsAction.class);
actions.register(AutoCreateAction.INSTANCE, AutoCreateAction.TransportAction.class);
Expand Down Expand Up @@ -1090,6 +1093,7 @@ protected void configure() {
bind(new TypeLiteral<RequestValidators<IndicesAliasesRequest>>() {
}).toInstance(indicesAliasesRequestRequestValidators);
bind(AutoCreateIndex.class).toInstance(autoCreateIndex);
bind(ActiveFetchPhaseTasks.class).asEagerSingleton();

// register ActionType -> transportAction Map used by NodeClient
@SuppressWarnings("rawtypes")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import org.elasticsearch.search.dfs.AggregatedDfs;
import org.elasticsearch.search.dfs.DfsKnnResults;
import org.elasticsearch.search.dfs.DfsSearchResult;
import org.elasticsearch.search.fetch.chunk.TransportFetchPhaseCoordinationAction;
import org.elasticsearch.search.internal.ShardSearchRequest;
import org.elasticsearch.search.query.QuerySearchRequest;
import org.elasticsearch.search.query.QuerySearchResult;
Expand Down Expand Up @@ -56,18 +57,25 @@ class DfsQueryPhase extends SearchPhase {
private final AbstractSearchAsyncAction<?> context;
private final SearchProgressListener progressListener;
private long phaseStartTimeInNanos;
private final TransportFetchPhaseCoordinationAction fetchCoordinationAction;

DfsQueryPhase(SearchPhaseResults<SearchPhaseResult> queryResult, Client client, AbstractSearchAsyncAction<?> context) {
DfsQueryPhase(
SearchPhaseResults<SearchPhaseResult> queryResult,
Client client,
AbstractSearchAsyncAction<?> context,
TransportFetchPhaseCoordinationAction fetchCoordinationAction
) {
super(NAME);
this.progressListener = context.getTask().getProgressListener();
this.queryResult = queryResult;
this.client = client;
this.context = context;
this.fetchCoordinationAction = fetchCoordinationAction;
}

// protected for testing
protected SearchPhase nextPhase(AggregatedDfs dfs) {
return SearchQueryThenFetchAsyncAction.nextPhase(client, context, queryResult, dfs);
return SearchQueryThenFetchAsyncAction.nextPhase(client, context, queryResult, dfs, fetchCoordinationAction);
}

@SuppressWarnings("unchecked")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@

import org.apache.logging.log4j.Logger;
import org.apache.lucene.search.ScoreDoc;
import org.elasticsearch.action.ActionListener;
import org.elasticsearch.cluster.node.DiscoveryNode;
import org.elasticsearch.common.util.concurrent.AbstractRunnable;
import org.elasticsearch.common.util.concurrent.AtomicArray;
import org.elasticsearch.core.Nullable;
Expand All @@ -18,6 +20,7 @@
import org.elasticsearch.search.dfs.AggregatedDfs;
import org.elasticsearch.search.fetch.FetchSearchResult;
import org.elasticsearch.search.fetch.ShardFetchSearchRequest;
import org.elasticsearch.search.fetch.chunk.TransportFetchPhaseCoordinationAction;
import org.elasticsearch.search.internal.ShardSearchContextId;
import org.elasticsearch.search.rank.RankDoc;
import org.elasticsearch.search.rank.RankDocShardInfo;
Expand All @@ -28,6 +31,8 @@
import java.util.List;
import java.util.Map;

import static org.elasticsearch.search.fetch.chunk.TransportFetchPhaseCoordinationAction.CHUNKED_FETCH_PHASE;

/**
* This search phase merges the query results from the previous phase together and calculates the topN hits for this search.
* Then it reaches out to all relevant shards to fetch the topN hits.
Expand All @@ -44,12 +49,14 @@ class FetchSearchPhase extends SearchPhase {
@Nullable
private final SearchPhaseResults<SearchPhaseResult> resultConsumer;
private final SearchPhaseController.ReducedQueryPhase reducedQueryPhase;
private final TransportFetchPhaseCoordinationAction fetchCoordinationAction;

FetchSearchPhase(
SearchPhaseResults<SearchPhaseResult> resultConsumer,
AggregatedDfs aggregatedDfs,
AbstractSearchAsyncAction<?> context,
@Nullable SearchPhaseController.ReducedQueryPhase reducedQueryPhase
@Nullable SearchPhaseController.ReducedQueryPhase reducedQueryPhase,
TransportFetchPhaseCoordinationAction fetchCoordinationAction
) {
super(NAME);
if (context.getNumShards() != resultConsumer.getNumShards()) {
Expand All @@ -67,6 +74,7 @@ class FetchSearchPhase extends SearchPhase {
this.progressListener = context.getTask().getProgressListener();
this.reducedQueryPhase = reducedQueryPhase;
this.resultConsumer = reducedQueryPhase == null ? resultConsumer : null;
this.fetchCoordinationAction = fetchCoordinationAction;
}

// protected for tests
Expand Down Expand Up @@ -98,7 +106,7 @@ private void innerRun() throws Exception {
final int numShards = context.getNumShards();
// Usually when there is a single shard, we force the search type QUERY_THEN_FETCH. But when there's kNN, we might
// still use DFS_QUERY_THEN_FETCH, which does not perform the "query and fetch" optimization during the query phase.
final boolean queryAndFetchOptimization = numShards == 1
boolean queryAndFetchOptimization = numShards == 1
&& context.getRequest().hasKnnSearch() == false
&& reducedQueryPhase.queryPhaseRankCoordinatorContext() == null
&& (context.getRequest().source() == null || context.getRequest().source().rankBuilder() == null);
Expand Down Expand Up @@ -212,6 +220,8 @@ private void executeFetch(
final ShardSearchContextId contextId = shardPhaseResult.queryResult() != null
? shardPhaseResult.queryResult().getContextId()
: shardPhaseResult.rankFeatureResult().getContextId();

// Create the listener that handles the fetch result
var listener = new SearchActionListener<FetchSearchResult>(shardTarget, shardIndex) {
@Override
public void innerOnResponse(FetchSearchResult result) {
Expand All @@ -237,29 +247,43 @@ public void onFailure(Exception e) {
}
}
};

// Get connection to the target node
final Transport.Connection connection;
try {
connection = context.getConnection(shardTarget.getClusterAlias(), shardTarget.getNodeId());
} catch (Exception e) {
listener.onFailure(e);
return;
}
context.getSearchTransport()
.sendExecuteFetch(
connection,
new ShardFetchSearchRequest(
context.getOriginalIndices(shardPhaseResult.getShardIndex()),
contextId,
shardPhaseResult.getShardSearchRequest(),
entry,
rankDocs,
lastEmittedDocForShard,
shardPhaseResult.getRescoreDocIds(),
aggregatedDfs
),

// Create the fetch request
final ShardFetchSearchRequest shardFetchRequest = new ShardFetchSearchRequest(
context.getOriginalIndices(shardPhaseResult.getShardIndex()),
contextId,
shardPhaseResult.getShardSearchRequest(),
entry,
rankDocs,
lastEmittedDocForShard,
shardPhaseResult.getRescoreDocIds(),
aggregatedDfs
);

if (connection.getTransportVersion().supports(CHUNKED_FETCH_PHASE)) {
shardFetchRequest.setCoordinatingNode(context.getSearchTransport().transportService().getLocalNode());
shardFetchRequest.setCoordinatingTaskId(context.getTask().getId());

DiscoveryNode targetNode = connection.getNode();

// Execute via coordination action
fetchCoordinationAction.execute(
context.getTask(),
listener
new TransportFetchPhaseCoordinationAction.Request(shardFetchRequest, targetNode),
ActionListener.wrap(response -> listener.onResponse(response.getResult()), listener::onFailure)
);
} else {
context.getSearchTransport().sendExecuteFetch(connection, shardFetchRequest, context.getTask(), listener);
}
}

private void moveToNextPhase(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
import org.elasticsearch.search.SearchShardTarget;
import org.elasticsearch.search.builder.SearchSourceBuilder;
import org.elasticsearch.search.dfs.AggregatedDfs;
import org.elasticsearch.search.fetch.chunk.TransportFetchPhaseCoordinationAction;
import org.elasticsearch.search.internal.ShardSearchContextId;
import org.elasticsearch.search.rank.RankDoc;
import org.elasticsearch.search.rank.context.RankFeaturePhaseRankCoordinatorContext;
Expand Down Expand Up @@ -48,12 +49,14 @@ public class RankFeaturePhase extends SearchPhase {
private final AggregatedDfs aggregatedDfs;
private final SearchProgressListener progressListener;
private final RankFeaturePhaseRankCoordinatorContext rankFeaturePhaseRankCoordinatorContext;
private final TransportFetchPhaseCoordinationAction fetchCoordinationAction;

RankFeaturePhase(
SearchPhaseResults<SearchPhaseResult> queryPhaseResults,
AggregatedDfs aggregatedDfs,
AbstractSearchAsyncAction<?> context,
RankFeaturePhaseRankCoordinatorContext rankFeaturePhaseRankCoordinatorContext
RankFeaturePhaseRankCoordinatorContext rankFeaturePhaseRankCoordinatorContext,
TransportFetchPhaseCoordinationAction fetchCoordinationAction
) {
super(NAME);
assert rankFeaturePhaseRankCoordinatorContext != null;
Expand All @@ -72,6 +75,7 @@ public class RankFeaturePhase extends SearchPhase {
this.rankPhaseResults = new ArraySearchPhaseResults<>(context.getNumShards());
context.addReleasable(rankPhaseResults);
this.progressListener = context.getTask().getProgressListener();
this.fetchCoordinationAction = fetchCoordinationAction;
}

@Override
Expand Down Expand Up @@ -267,6 +271,9 @@ private float maxScore(ScoreDoc[] scoreDocs) {
}

void moveToNextPhase(SearchPhaseResults<SearchPhaseResult> phaseResults, SearchPhaseController.ReducedQueryPhase reducedQueryPhase) {
context.executeNextPhase(NAME, () -> new FetchSearchPhase(phaseResults, aggregatedDfs, context, reducedQueryPhase));
context.executeNextPhase(
NAME,
() -> new FetchSearchPhase(phaseResults, aggregatedDfs, context, reducedQueryPhase, fetchCoordinationAction)
);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
import org.elasticsearch.search.SearchPhaseResult;
import org.elasticsearch.search.SearchShardTarget;
import org.elasticsearch.search.dfs.DfsSearchResult;
import org.elasticsearch.search.fetch.chunk.TransportFetchPhaseCoordinationAction;
import org.elasticsearch.search.internal.AliasFilter;
import org.elasticsearch.transport.Transport;

Expand All @@ -31,6 +32,7 @@ final class SearchDfsQueryThenFetchAsyncAction extends AbstractSearchAsyncAction
private final SearchPhaseResults<SearchPhaseResult> queryPhaseResultConsumer;
private final SearchProgressListener progressListener;
private final Client client;
private final TransportFetchPhaseCoordinationAction fetchCoordinationAction;

SearchDfsQueryThenFetchAsyncAction(
Logger logger,
Expand All @@ -50,7 +52,8 @@ final class SearchDfsQueryThenFetchAsyncAction extends AbstractSearchAsyncAction
SearchResponse.Clusters clusters,
Client client,
SearchResponseMetrics searchResponseMetrics,
Map<String, Object> searchRequestAttributes
Map<String, Object> searchRequestAttributes,
TransportFetchPhaseCoordinationAction fetchCoordinationAction
) {
super(
"dfs",
Expand Down Expand Up @@ -81,6 +84,7 @@ final class SearchDfsQueryThenFetchAsyncAction extends AbstractSearchAsyncAction
notifyListShards(progressListener, clusters, request, shardsIts);
}
this.client = client;
this.fetchCoordinationAction = fetchCoordinationAction;
}

@Override
Expand All @@ -94,7 +98,7 @@ protected void executePhaseOnShard(

@Override
protected SearchPhase getNextPhase() {
return new DfsQueryPhase(queryPhaseResultConsumer, client, this);
return new DfsQueryPhase(queryPhaseResultConsumer, client, this, fetchCoordinationAction);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@
import org.elasticsearch.search.SearchShardTarget;
import org.elasticsearch.search.builder.PointInTimeBuilder;
import org.elasticsearch.search.dfs.AggregatedDfs;
import org.elasticsearch.search.fetch.chunk.TransportFetchPhaseCoordinationAction;
import org.elasticsearch.search.internal.AliasFilter;
import org.elasticsearch.search.internal.SearchContext;
import org.elasticsearch.search.internal.ShardSearchContextId;
Expand Down Expand Up @@ -98,6 +99,7 @@ public class SearchQueryThenFetchAsyncAction extends AbstractSearchAsyncAction<S
private final Client client;
private final boolean batchQueryPhase;
private long phaseStartTimeNanos;
private final TransportFetchPhaseCoordinationAction fetchPhaseCoordinationAction;

SearchQueryThenFetchAsyncAction(
Logger logger,
Expand All @@ -118,7 +120,8 @@ public class SearchQueryThenFetchAsyncAction extends AbstractSearchAsyncAction<S
Client client,
boolean batchQueryPhase,
SearchResponseMetrics searchResponseMetrics,
Map<String, Object> searchRequestAttributes
Map<String, Object> searchRequestAttributes,
TransportFetchPhaseCoordinationAction fetchPhaseCoordinationAction
) {
super(
"query",
Expand Down Expand Up @@ -150,6 +153,7 @@ public class SearchQueryThenFetchAsyncAction extends AbstractSearchAsyncAction<S
if (progressListener != SearchProgressListener.NOOP) {
notifyListShards(progressListener, clusters, request, shardsIts);
}
this.fetchPhaseCoordinationAction = fetchPhaseCoordinationAction;
}

@Override
Expand Down Expand Up @@ -208,18 +212,19 @@ static SearchPhase nextPhase(
Client client,
AbstractSearchAsyncAction<?> context,
SearchPhaseResults<SearchPhaseResult> queryResults,
AggregatedDfs aggregatedDfs
AggregatedDfs aggregatedDfs,
TransportFetchPhaseCoordinationAction fetchCoordinationAction
) {
var rankFeaturePhaseCoordCtx = RankFeaturePhase.coordinatorContext(context.getRequest().source(), client);
if (rankFeaturePhaseCoordCtx == null) {
return new FetchSearchPhase(queryResults, aggregatedDfs, context, null);
return new FetchSearchPhase(queryResults, aggregatedDfs, context, null, fetchCoordinationAction);
}
return new RankFeaturePhase(queryResults, aggregatedDfs, context, rankFeaturePhaseCoordCtx);
return new RankFeaturePhase(queryResults, aggregatedDfs, context, rankFeaturePhaseCoordCtx, fetchCoordinationAction);
}

@Override
protected SearchPhase getNextPhase() {
return nextPhase(client, this, results, null);
return nextPhase(client, this, results, null, fetchPhaseCoordinationAction);
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,8 @@
import org.elasticsearch.search.fetch.ScrollQueryFetchSearchResult;
import org.elasticsearch.search.fetch.ShardFetchRequest;
import org.elasticsearch.search.fetch.ShardFetchSearchRequest;
import org.elasticsearch.search.fetch.chunk.FetchPhaseResponseChunk;
import org.elasticsearch.search.fetch.chunk.TransportFetchPhaseResponseChunkAction;
import org.elasticsearch.search.internal.InternalScrollSearchRequest;
import org.elasticsearch.search.internal.ShardSearchContextId;
import org.elasticsearch.search.internal.ShardSearchRequest;
Expand Down Expand Up @@ -67,6 +69,8 @@
import java.util.concurrent.Executor;
import java.util.function.BiFunction;

import static org.elasticsearch.search.fetch.chunk.TransportFetchPhaseCoordinationAction.CHUNKED_FETCH_PHASE;

/**
* An encapsulation of {@link SearchService} operations exposed through
* transport.
Expand Down Expand Up @@ -540,8 +544,39 @@ public static void registerRequestHandler(
namedWriteableRegistry
);

final TransportRequestHandler<ShardFetchRequest> shardFetchRequestHandler = (request, channel, task) -> searchService
.executeFetchPhase(request, (SearchShardTask) task, new ChannelActionListener<>(channel));
final TransportRequestHandler<ShardFetchRequest> shardFetchRequestHandler = (request, channel, task) -> {
if (channel.getVersion().supports(CHUNKED_FETCH_PHASE)
&& request instanceof ShardFetchSearchRequest fetchSearchReq
&& fetchSearchReq.getCoordinatingNode() != null) {

// CHUNKED PATH
final FetchPhaseResponseChunk.Writer writer = new FetchPhaseResponseChunk.Writer() {
final Transport.Connection conn = transportService.getConnection(fetchSearchReq.getCoordinatingNode());

@Override
public void writeResponseChunk(FetchPhaseResponseChunk responseChunk, ActionListener<Void> listener) {
transportService.sendChildRequest(
conn,
TransportFetchPhaseResponseChunkAction.TYPE.name(),
new TransportFetchPhaseResponseChunkAction.Request(fetchSearchReq.getCoordinatingTaskId(), responseChunk),
task,
TransportRequestOptions.EMPTY,
new ActionListenerResponseHandler<>(
listener.map(ignored -> null),
in -> ActionResponse.Empty.INSTANCE,
EsExecutors.DIRECT_EXECUTOR_SERVICE
)
);
}
};

searchService.executeFetchPhase(request, (SearchShardTask) task, writer, new ChannelActionListener<>(channel));
} else {
// Normal path
searchService.executeFetchPhase(request, (SearchShardTask) task, new ChannelActionListener<>(channel));
}
};

transportService.registerRequestHandler(
FETCH_ID_SCROLL_ACTION_NAME,
EsExecutors.DIRECT_EXECUTOR_SERVICE,
Expand Down
Loading