Elasticsearch 8.1.3插件开发:基于前置查询结果实现关联查询的技术求助
Hey there, let's work through this two-step join query plugin for Elasticsearch 8.1.3. I've messed around with similar plugin development challenges before, so here's what I've learned that should help you out:
Problem with Scheme 1 (Internal HTTP Calls)
Elasticsearch's security sandbox is really strict, especially around thread permissions. The async HTTP client you tried uses thread pools that trigger the ThreadPermission error you saw, and AccessController.doPrivileged() won't save you here—this is a fundamental restriction of how Elasticsearch isolates plugin code. This approach is a dead end, so let's move on.
Problem with Scheme 2 (Direct RestSearchAction Calls)
RestSearchAction is meant to handle incoming HTTP requests from external clients, not to be invoked directly from plugin code. It's part of the REST layer, not the internal core API. You were looking in the wrong place—instead of REST handlers, you need to use Elasticsearch's internal client to execute search requests programmatically.
This is the standard way to run queries from within a plugin. You'll use the internal Client (or NodeClient in 8.x) to execute two sequential SearchRequest operations: first the geo query on device-location, then the terms query on device-logs.
Step 1: Set Up Your Rest Handler with Dependency Injection
First, create a custom BaseRestHandler that injects the internal client. This client has the necessary permissions to run internal searches without hitting security roadblocks.
import org.elasticsearch.action.ActionListener; import org.elasticsearch.action.search.SearchRequest; import org.elasticsearch.action.search.SearchResponse; import org.elasticsearch.client.Client; import org.elasticsearch.client.node.NodeClient; import org.elasticsearch.common.unit.DistanceUnit; import org.elasticsearch.common.xcontent.XContentBuilder; import org.elasticsearch.common.xcontent.XContentFactory; import org.elasticsearch.index.query.QueryBuilders; import org.elasticsearch.rest.BaseRestHandler; import org.elasticsearch.rest.BytesRestResponse; import org.elasticsearch.rest.RestChannel; import org.elasticsearch.rest.RestRequest; import org.elasticsearch.rest.RestStatus; import org.elasticsearch.search.builder.SearchSourceBuilder; import java.io.IOException; import java.util.Arrays; import java.util.List; import java.util.stream.Collectors; public class TwoStepJoinRestHandler extends BaseRestHandler { private final Client client; // Inject the internal client via constructor public TwoStepJoinRestHandler(Client client) { this.client = client; } @Override public String getName() { return "two_step_join_query_handler"; } @Override public RestChannelConsumer prepareRequest(RestRequest request, NodeClient client) throws IOException { // Extract input parameters from the incoming request (adjust as needed) double targetLat = request.paramAsDouble("lat", 37.7749); double targetLon = request.paramAsDouble("lon", -122.4194); double searchDistance = request.paramAsDouble("distance", 5.0); // 1. Build and execute the geo query on device-location SearchRequest geoSearchRequest = new SearchRequest("device-location"); SearchSourceBuilder geoSource = new SearchSourceBuilder(); geoSource.query(QueryBuilders.geoDistanceQuery("location") .point(targetLat, targetLon) .distance(searchDistance, DistanceUnit.KILOMETERS)); geoSource.fetchSource(new String[]{"deviceId"}, null); // Only retrieve deviceId to save bandwidth geoSearchRequest.source(geoSource); // Return a channel consumer to handle async execution return channel -> client.search(geoSearchRequest, new ActionListener<SearchResponse>() { @Override public void onResponse(SearchResponse geoResponse) { // Extract device IDs from the geo query results List<String> deviceIds = Arrays.stream(geoResponse.getHits().getHits()) .map(hit -> (String) hit.getSourceAsMap().get("deviceId")) .collect(Collectors.toList()); if (deviceIds.isEmpty()) { // No devices found—return empty response try { channel.sendResponse(new BytesRestResponse(RestStatus.OK, "[]")); } catch (IOException e) { onFailure(e); } return; } // 2. Build and execute the logs query on device-logs SearchRequest logSearchRequest = new SearchRequest("device-logs"); SearchSourceBuilder logSource = new SearchSourceBuilder(); logSource.query(QueryBuilders.termsQuery("deviceId", deviceIds)); logSearchRequest.source(logSource); client.search(logSearchRequest, new ActionListener<SearchResponse>() { @Override public void onResponse(SearchResponse logResponse) { // Format the final response try { XContentBuilder responseBuilder = XContentFactory.jsonBuilder(); responseBuilder.startObject(); responseBuilder.field("matching_devices_count", deviceIds.size()); responseBuilder.field("logs", logResponse.getHits().getHits()); responseBuilder.endObject(); channel.sendResponse(new BytesRestResponse(RestStatus.OK, responseBuilder)); } catch (IOException e) { onFailure(e); } } @Override public void onFailure(Exception e) { channel.sendResponse(new BytesRestResponse(e)); } }); } @Override public void onFailure(Exception e) { channel.sendResponse(new BytesRestResponse(e)); } }); } }
Step 2: Register Your Plugin and Handler
Next, create a plugin class to register your rest handler with Elasticsearch:
import org.elasticsearch.common.settings.Settings; import org.elasticsearch.plugins.Plugin; import org.elasticsearch.plugins.RestPlugin; import org.elasticsearch.rest.RestController; import org.elasticsearch.rest.RestHandler; import java.util.List; public class TwoStepJoinPlugin extends Plugin implements RestPlugin { @Override public List<RestHandler> getRestHandlers(Settings settings, RestController restController, org.elasticsearch.cluster.ClusterSettings clusterSettings, org.elasticsearch.index.IndexScopedSettings indexScopedSettings, org.elasticsearch.common.settings.SettingsFilter settingsFilter, org.elasticsearch.cluster.metadata.IndexNameExpressionResolver indexNameExpressionResolver, java.util.function.Supplier<org.elasticsearch.cluster.node.DiscoveryNodes> nodesInCluster) { // Register your custom handler with the internal client return List.of(new TwoStepJoinRestHandler(restController.getClient())); } }
Step 3: Update Security Policy
Your plugin-security.policy needs additional permissions to let the plugin execute internal searches and handle threads. Update it to:
grant { permission java.lang.RuntimePermission "setFactory"; // Allow thread operations needed for async execution permission org.elasticsearch.secure_sm.ThreadPermission "modifyArbitraryThreadGroup"; // Allow access to REST actions permission org.elasticsearch.rest.RestPermission "*"; // Allow search operations on the target indices permission org.elasticsearch.action.ActionPermission "indices:data/read/search*"; };
Key Notes for Success
- Use Async Execution: Elasticsearch is designed to be non-blocking, so always use
ActionListenerinstead of blocking calls. Blocking the thread pool can cause performance issues. - Limit Fetched Fields: In the geo query, we only fetch
deviceIdto reduce data transfer and processing time. - Handle Edge Cases: Make sure to handle empty device ID lists gracefully to avoid unnecessary log queries.
- Make sure your Maven pom.xml includes the correct Elasticsearch dependencies for version 8.1.3 (use
providedscope since Elasticsearch provides these at runtime). - Package the plugin as a zip with the correct structure (including the security policy and plugin descriptor).
内容的提问来源于stack exchange,提问作者WISERDIVISOR

