Apache Beam ElasticsearchIO.read()多查询处理方案问询
问题描述
我在使用ElasticsearchIO.read()处理多个查询实例时遇到问题。我的查询是基于一组输入值动态构建为PCollection的,希望了解如何通过.withQuery()参数或其他灵活方式实现该功能。
问题核心:ElasticsearchIO.read()方法需要PBegin作为输入来启动管道,但我需要基于已有的PCollection查询列表来触发多次读取操作,而不是从管道起点开始。另外想知道能否将ElasticsearchIO.read()包装在Create转换中,用模拟PBegin的方式实现。
初步尝试代码(不可行)
PCollection<String> queries = ... // 动态构建的查询列表 PCollection<String> queryResults = queries.apply( ParDo.of(new DoFn<String, String>() { @ProcessElement public void processElement(ProcessContext c) { String query = c.element(); PCollection<String> results = c.pipeline() .apply(ElasticsearchIO.read() .withConnectionConfiguration( ElasticsearchIO.ConnectionConfiguration.create(hosts, indexName)) .withQuery(query)); c.output(results); } }) .apply(Flatten.pCollections()));
解决方案
1. 错误原因说明
你在ParDo内部动态创建PCollection的方式不符合Beam的设计规则:所有管道拓扑必须在构建阶段定义完成,运行阶段(ProcessElement执行时)不能修改管道结构或创建新的PCollection。
2. 正确实现:ParDo结合Elasticsearch客户端
放弃使用ElasticsearchIO.read()(它是管道源操作,只能从PBegin启动),直接在ParDo中使用Elasticsearch官方客户端执行每个查询:
import org.elasticsearch.client.RestHighLevelClient; import org.elasticsearch.client.ClientConfiguration; import org.elasticsearch.client.RestClients; import org.elasticsearch.action.search.SearchRequest; import org.elasticsearch.action.search.SearchResponse; import org.elasticsearch.search.builder.SearchSourceBuilder; import org.elasticsearch.index.query.QueryBuilders; // ... PCollection<String> queries = ... // 动态构建的查询列表 PCollection<String> queryResults = queries.apply(ParDo.of(new DoFn<String, String>() { // 客户端实例,在Setup阶段初始化,避免每次处理元素都创建 private transient RestHighLevelClient client; @Setup public void setup() { // 配置Elasticsearch连接 ClientConfiguration clientConfig = ClientConfiguration.builder() .connectedTo(hosts) // hosts为Elasticsearch地址数组,如{"es-host1:9200", "es-host2:9200"} .build(); this.client = RestClients.create(clientConfig).rest(); } @ProcessElement public void processElement(ProcessContext c) throws IOException { String queryJson = c.element(); // 构建搜索请求 SearchRequest searchRequest = new SearchRequest(indexName); SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); // 使用wrapperQuery直接传入JSON格式的查询 sourceBuilder.query(QueryBuilders.wrapperQuery(queryJson)); searchRequest.source(sourceBuilder); // 执行查询并遍历结果输出 SearchResponse searchResponse = client.search(searchRequest, RequestOptions.DEFAULT); for (SearchHit hit : searchResponse.getHits().getHits()) { c.output(hit.getSourceAsString()); } } @Teardown public void teardown() throws IOException { // 关闭客户端释放资源 if (client != null) { client.close(); } } }));
3. 关于PBegin输入IO的集合引入方法
Beam中以PBegin为输入的IO都是源操作(如ElasticsearchIO.read()、BigQueryIO.read()),它们是管道的起始点,无法基于已有PCollection触发。你给出的Create.ofProvider方案适用于将静态/预定义列表转为PCollection,但无法处理动态生成的查询列表场景:
// 示例:静态列表转为PCollection ValueProvider<List<String>> myListProvider = ValueProvider.StaticValueProvider.of(myList); PCollection<String> pcoll = pipeline.apply(Create.ofProvider(myListProvider, ListCoder.of()));
4. 性能优化方案:合并批量查询
如果你的查询可以合并,建议将多个查询合并为一个bool查询的should子句,减少Elasticsearch请求次数:
PCollection<String> queries = ...; // 合并多个查询为一个批量查询 PCollection<String> combinedQuery = queries.apply(Combine.globally(new CombineFn<String, StringBuilder, String>() { @Override public StringBuilder createAccumulator() { return new StringBuilder("{\"bool\": {\"should\": ["); } @Override public StringBuilder addInput(StringBuilder accum, String input) { if (accum.length() > "{\"bool\": {\"should\": [".length()) { accum.append(","); } accum.append(input); return accum; } @Override public StringBuilder mergeAccumulators(Iterable<StringBuilder> accums) { StringBuilder merged = new StringBuilder(); for (StringBuilder accum : accums) { merged.append(accum); } return merged; } @Override public String extractOutput(StringBuilder accum) { return accum.append("]}}").toString(); } })); // 执行合并后的查询(同样使用客户端方式) PCollection<String> results = combinedQuery.apply(ParDo.of(new DoFn<String, String>() { private transient RestHighLevelClient client; @Setup public void setup() { ClientConfiguration clientConfig = ClientConfiguration.builder() .connectedTo(hosts) .build(); this.client = RestClients.create(clientConfig).rest(); } @ProcessElement public void processElement(ProcessContext c) throws IOException { String query = c.element(); SearchRequest searchRequest = new SearchRequest(indexName); SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); sourceBuilder.query(QueryBuilders.wrapperQuery(query)); searchRequest.source(sourceBuilder); SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT); for (SearchHit hit : response.getHits().getHits()) { c.output(hit.getSourceAsString()); } } @Teardown public void teardown() throws IOException { if (client != null) { client.close(); } } }));
内容的提问来源于stack exchange,提问作者user21627820

