1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143
| public Boolean isExist(String index, String id,String routing) { RestHighLevelClient client = ESClient.instance(esConfig); GetRequest request = new GetRequest(index).id(id); request.fetchSourceContext(new FetchSourceContext(false)); request.storedFields("_none_"); request.routing(routing); try { return client.exists(request, RequestOptions.DEFAULT); } catch (IOException e) { return false; } }
public Map<String, Object> get(String index, String id,String routing) throws IOException { RestHighLevelClient client = ESClient.instance(esConfig); GetRequest request = new GetRequest(esConfig.getIndex(), id); request.routing(routing); GetResponse getResponse = client.get(request, RequestOptions.DEFAULT); Map<String, Object> doc = new HashMap<String, Object>(); if(getResponse.isExists()) { doc = getResponse.getSourceAsMap(); doc.put("id",getResponse.getId()); } return doc; }
public SearchResponse getByIds(List<String> ids) throws IOException { RestHighLevelClient client = ESClient.instance(esConfig); SearchRequest searchRequest = new SearchRequest(esConfig.getIndex()); SearchSourceBuilder searchSourceBuilder = new SearchSourceBuilder(); BoolQueryBuilder booleanQueryBuilder = QueryBuilders.boolQuery(); IdsQueryBuilder idsQueryBuilder = QueryBuilders.idsQuery(); idsQueryBuilder.ids().addAll(ids); booleanQueryBuilder.must(idsQueryBuilder); searchSourceBuilder.from(0); searchSourceBuilder.size(ids.size()); searchSourceBuilder.query(booleanQueryBuilder); searchRequest.source(searchSourceBuilder); SearchResponse searchResponse = client.search(searchRequest, RequestOptions.DEFAULT); return searchResponse; }
public void createOrUpdate(String index, String id,String routing,JSONObject item,Boolean bulk) throws IOException { RestHighLevelClient client = ESClient.instance(esConfig); GetRequest request = new GetRequest(esConfig.getIndex(), id); request.routing(routing); JSONObject doc = new JSONObject(); GetResponse getResponse = client.get(request, RequestOptions.DEFAULT); if(getResponse.isExists()) { Map<String, Object> source = getResponse.getSourceAsMap(); JSONObject content = JSONObject.parseObject(source.get("content").toString()); item.keySet().forEach(key->{ content.put(key,item.get(key)); }); doc = this.createDocSource(content,routing,true); } else doc = this.createDocSource(item,routing,false); if(bulk) this.bulkUpdate(index, id, doc, true, routing); else this.update(index, id, doc, true,routing); }
public void delete(String index, String id,String routing) { RestHighLevelClient client = ESClient.instance(esConfig); DeleteRequest request = new DeleteRequest(index, id); request.routing(routing); client.deleteAsync(request, RequestOptions.DEFAULT, new ActionListener<DeleteResponse>() { @Override public void onResponse(DeleteResponse response) { if (response.getResult() == DocWriteResponse.Result.DELETED) { logger.info("delete>>>"+"{\"id\":\""+id+"\",\"type\":\""+routing+"\"}"); } } @Override public void onFailure(Exception e) { logger.info("delete fail>>>"+"{\"id\":\""+id+"\",\"type\":\""+routing+"\"}"); } }); }
public void update(String index, String id, JSONObject item,Boolean upsert,String routing) { if(item == null) { return; } RestHighLevelClient client = ESClient.instance(esConfig); UpdateRequest request = new UpdateRequest(index,id).doc(item); request.docAsUpsert(upsert); request.retryOnConflict(3); request.routing(routing); client.updateAsync(request, RequestOptions.DEFAULT, new ActionListener<UpdateResponse>() { @Override public void onResponse(UpdateResponse response) { if (response.getResult() == DocWriteResponse.Result.CREATED) { logger.info("created>>>"+"{\"id\":\""+id+"\",\"type\":\""+routing+"\"}"); } else if (response.getResult() == DocWriteResponse.Result.UPDATED) { logger.info("updated>>>"+"{\"id\":\""+id+"\",\"type\":\""+routing+"\"}"); } } @Override public void onFailure(Exception e) { logger.warn("not_store_record("+e.getMessage()+")>>>"+item.toJSONString()); } }); } public void bulkUpdate(String index, String id, JSONObject item,Boolean upsert,String routing) { if(item == null) { return; } UpdateRequest request = new UpdateRequest(index,id).doc(item); request.docAsUpsert(upsert); request.retryOnConflict(3); request.routing(routing); queue.offer(request); }
public void bulk(ConcurrentLinkedQueue<UpdateRequest> queue, RestHighLevelClient client) { BulkRequest bulk = new BulkRequest(); for (int i = 0; i < esConfig.getBulk(); i++) { if(queue.isEmpty()) break; bulk.add(queue.poll()); } if(bulk.numberOfActions() == 0) return; logger.info(DateUtils.getNowDate()+":starting handle doc queue,actions:"+bulk.numberOfActions()); client.bulkAsync(bulk, RequestOptions.DEFAULT, new ActionListener<BulkResponse>() { @Override public void onResponse(BulkResponse response) { if(response.hasFailures()) logger.info(response.buildFailureMessage()); } @Override public void onFailure(Exception e) { logger.warn("bulk handle docs error:"+e.getMessage()); } }); }
|