org.springframework.boot
spring-boot-starter-data-elasticsearch
注解
import java.lang.annotation.*;
/**
* 该注解类用于注解映射到ES中的数据字段,类似Hibernate中实体字段映射数据库中的表字段
* @version 1.0
* @date 2022/5/19 9:46
* @since : JDK 11
*/
@Retention(RetentionPolicy.RUNTIME)
@Target(ElementType.FIELD)
@Documented
@Inherited
public @interface ESMappingField {
/**
* 字段名称
* @return
*/
String fieldName() default "";
/**
* 数据类型
* @return
*/
String dataType() default "";
/**
* 使用哪种分词器
* @return
*/
String setAnalyzer() default "";
/**
* 是否使用分词功能
* @return
*/
boolean isAnalyze() default false;
}
ES索引
controller
import io.swagger.annotations.Api;
import io.swagger.annotations.ApiImplicitParam;
import io.swagger.annotations.ApiImplicitParams;
import io.swagger.annotations.ApiOperation;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;
import java.io.IOException;
/**
* @version 1.0
* @date 2022/5/19 9:48
* @since : JDK 11
*/
@RestController
@RequestMapping("/es/index")
@Api(value = "/es/index", tags = "ES索引操作")
public class EsIndexController {
@Autowired
private IElasticsearchIndexService elasticsearchIndexService;
@GetMapping("/isIndexExists")
@ApiOperation(value = "判断索引是否存在")
@ApiImplicitParam(value = "索引名称", name = "indexName", dataTypeClass = String.class, required = true, paramType = "query")
public R
响应主体R
import java.io.Serializable;
/**
* 响应信息主体
*/
public class R implements Serializable
{
private static final long serialVersionUID = 1L;
/** 成功 */
public static final int SUCCESS = 200;
/** 失败 */
public static final int FAIL = 500;
private int code;
private String msg;
private T data;
public static R ok()
{
return restResult(null, SUCCESS, null);
}
public static R ok(T data)
{
return restResult(data, SUCCESS, null);
}
public static R ok(T data, String msg)
{
return restResult(data, SUCCESS, msg);
}
public static R fail()
{
return restResult(null, FAIL, null);
}
public static R fail(String msg)
{
return restResult(null, FAIL, msg);
}
public static R fail(T data)
{
return restResult(data, FAIL, null);
}
public static R fail(T data, String msg)
{
return restResult(data, FAIL, msg);
}
public static R fail(int code, String msg)
{
return restResult(null, code, msg);
}
private static R restResult(T data, int code, String msg)
{
R apiResult = new R<>();
apiResult.setCode(code);
apiResult.setData(data);
apiResult.setMsg(msg);
return apiResult;
}
public int getCode()
{
return code;
}
public void setCode(int code)
{
this.code = code;
}
public String getMsg()
{
return msg;
}
public void setMsg(String msg)
{
this.msg = msg;
}
public T getData()
{
return data;
}
public void setData(T data)
{
this.data = data;
}
}
service
import java.io.IOException;
/**
* @version 1.0
* @date 2022/5/19 9:41
* @since : JDK 11
*/
public interface IElasticsearchIndexService {
/**
* 判断索引是否存在
* @param indexName
* @return
*/
Boolean isIndexExists(String indexName) throws IOException;
/**
* 创建索引
*
* @param indexName
* @param numberOfShards
* @param numberOfReplicas
* @return
*/
Boolean createIndexWithShards(String indexName, Integer numberOfShards, Integer numberOfReplicas) throws IOException;
/**
* 创建索引(默认1个分片,1个副本)
*
* @param indexName
* @return
*/
Boolean createIndex(String indexName) throws IOException;
/**
* 删除索引
* @param indexName
* @return
* @throws IOException
*/
Boolean deleteIndex(String indexName) throws IOException;
/**
* 判断索引别名是否存在
* @param aliasName
* @return
*/
Boolean isAliasExists(String aliasName) throws IOException;
/**
* 创建索引同时给索引创建别名
*
* @param indexName
* @param aliasName
* @return
*/
Boolean createIndexWithAlias(String indexName, String aliasName) throws IOException;
/**
* 给索引添加别名
* @param indexName
* @param aliasName
* @return
* @throws IOException
*/
Boolean addAlias(String indexName, String aliasName) throws IOException;
/**
* 删除某个索引的别名
* @param aliasName
* @return
* @throws IOException
*/
Boolean deleteAlias(String indexName,String aliasName) throws IOException;
/**
* 重建索引,拷贝数据
*
* @param oldIndexname
* @param newIndexname
*/
void reindex(String oldIndexname, String newIndexname) throws IOException;
/**
* 重建索引后修改别名
*
* @param aliasname
* @param oldIndexname
* @param newIndexname
* @return
*/
Boolean changeAliasAfterReindex(String aliasname, String oldIndexname, String newIndexname) throws IOException ;
/**
* 添加mapping
* @param indexName
* @param clazz
* @return
*/
Boolean addMapping(String indexName, Class<?> clazz) throws IOException;
}
serviceImpl
import com.google.common.collect.Maps;
import org.apache.commons.lang3.StringUtils;
import org.elasticsearch.action.admin.indices.alias.Alias;
import org.elasticsearch.action.admin.indices.alias.IndicesAliasesRequest;
import org.elasticsearch.action.admin.indices.alias.get.GetAliasesRequest;
import org.elasticsearch.action.admin.indices.delete.DeleteIndexRequest;
import org.elasticsearch.action.support.IndicesOptions;
import org.elasticsearch.action.support.master.AcknowledgedResponse;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.client.indices.*;
import org.elasticsearch.common.settings.Settings;
import org.elasticsearch.index.reindex.ReindexRequest;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import java.io.IOException;
import java.lang.reflect.Field;
import java.util.Map;
/**
* @version 1.0
* @date 2022/5/19 9:42
* @since : JDK 11
*/
@Service
public class ElasticsearchIndexServiceImpl implements IElasticsearchIndexService {
@Autowired
private RestHighLevelClient restHighLevelClient;
@Override
public Boolean isIndexExists(String indexName) throws IOException {
GetIndexRequest getIndexRequest = new GetIndexRequest(indexName);
getIndexRequest.humanReadable(true);
return restHighLevelClient.indices().exists(getIndexRequest, RequestOptions.DEFAULT);
}
@Override
public Boolean createIndexWithShards(String indexName, Integer numberOfShards, Integer numberOfReplicas) throws IOException {
CreateIndexRequest createIndexRequest = new CreateIndexRequest(indexName);
//设置分片信息
numberOfShards = numberOfShards == null ? 1 : numberOfShards;
numberOfReplicas = numberOfReplicas == null ? 1 : numberOfReplicas;
createIndexRequest.settings(Settings.builder().
put("index.number_of_shards", numberOfShards)
.put("index.number_of_replicas", numberOfReplicas));
//创建索引
CreateIndexResponse createIndexResponse = restHighLevelClient.indices().create(createIndexRequest, RequestOptions.DEFAULT);
return createIndexResponse.isAcknowledged();
}
@Override
public Boolean createIndex(String indexName) throws IOException {
CreateIndexRequest createIndexRequest = new CreateIndexRequest(indexName);
CreateIndexResponse createIndexResponse = restHighLevelClient.indices().create(createIndexRequest, RequestOptions.DEFAULT);
return createIndexResponse.isAcknowledged();
}
@Override
public Boolean deleteIndex(String indexName) throws IOException {
DeleteIndexRequest deleteIndexRequest = new DeleteIndexRequest(indexName);
deleteIndexRequest.indicesOptions(IndicesOptions.LENIENT_EXPAND_OPEN);
AcknowledgedResponse delete = restHighLevelClient.indices().delete(deleteIndexRequest, RequestOptions.DEFAULT);
return delete.isAcknowledged();
}
@Override
public Boolean isAliasExists(String aliasName) throws IOException {
GetAliasesRequest getAliasesRequest = new GetAliasesRequest(aliasName);
return restHighLevelClient.indices().existsAlias(getAliasesRequest, RequestOptions.DEFAULT);
}
@Override
public Boolean createIndexWithAlias(String indexName, String aliasName) throws IOException {
CreateIndexRequest request = new CreateIndexRequest(indexName);
if (StringUtils.isNotEmpty(aliasName)) {
request.alias(new Alias(aliasName));
}
CreateIndexResponse createIndexResponse = restHighLevelClient.indices().create(request, RequestOptions.DEFAULT);
return createIndexResponse.isAcknowledged();
}
@Override
public Boolean addAlias(String indexName, String aliasName) throws IOException {
IndicesAliasesRequest aliasesRequest = new IndicesAliasesRequest();
IndicesAliasesRequest.AliasActions aliasAction =
new IndicesAliasesRequest.AliasActions(IndicesAliasesRequest.AliasActions.Type.ADD)
.index(indexName)
.alias(aliasName);
aliasesRequest.addAliasAction(aliasAction);
AcknowledgedResponse acknowledgedResponse = restHighLevelClient.indices().updateAliases(aliasesRequest,RequestOptions.DEFAULT);
return acknowledgedResponse.isAcknowledged();
}
@Override
public void reindex(String oldIndexname, String newIndexname) throws IOException {
ReindexRequest request = new ReindexRequest();
request.setSourceIndices(oldIndexname);
request.setDestIndex(newIndexname);
request.setSourceBatchSize(1000);
request.setDestOpType("create");
request.setConflicts("proceed");
// request.setScroll(TimeValue.timeValueMinutes(10));
// request.setTimeout(TimeValue.timeValueMinutes(20));
request.setRefresh(true);
restHighLevelClient.reindex(request, RequestOptions.DEFAULT);
}
@Override
public Boolean deleteAlias(String indexName,String aliasName) throws IOException {
DeleteAliasRequest deleteAliasRequest = new DeleteAliasRequest(indexName,aliasName);
org.elasticsearch.client.core.AcknowledgedResponse acknowledgedResponse = restHighLevelClient.indices().deleteAlias(deleteAliasRequest, RequestOptions.DEFAULT);
return acknowledgedResponse.isAcknowledged();
}
@Override
public Boolean changeAliasAfterReindex(String aliasname, String oldIndexname, String newIndexname) throws IOException {
IndicesAliasesRequest.AliasActions addIndexAction = new IndicesAliasesRequest.AliasActions(
IndicesAliasesRequest.AliasActions.Type.ADD).index(newIndexname).alias(aliasname);
IndicesAliasesRequest.AliasActions removeAction = new IndicesAliasesRequest.AliasActions(
IndicesAliasesRequest.AliasActions.Type.REMOVE).index(oldIndexname).alias(aliasname);
IndicesAliasesRequest indicesAliasesRequest = new IndicesAliasesRequest();
indicesAliasesRequest.addAliasAction(addIndexAction);
indicesAliasesRequest.addAliasAction(removeAction);
AcknowledgedResponse indicesAliasesResponse = restHighLevelClient.indices().updateAliases(indicesAliasesRequest,
RequestOptions.DEFAULT);
return indicesAliasesResponse.isAcknowledged();
}
@Override
public Boolean addMapping(String indexName, Class<?> clazz) throws IOException {
PutMappingRequest putMappingRequest = new PutMappingRequest(indexName);
Map jsonMap = Maps.newHashMap();
Map properties = Maps.newHashMap();
Field[] fields = clazz.getDeclaredFields();
for (Field field : fields) {
ESMappingField esMappingField = field.getDeclaredAnnotation(ESMappingField.class);
if (esMappingField != null) {
String fieldname = esMappingField.fieldName();
String datatype = esMappingField.dataType();
String analyzer = esMappingField.setAnalyzer();
boolean isanalye = esMappingField.isAnalyze();
Map m = Maps.newHashMap();
m.put("type", datatype);
if (isanalye && StringUtils.isNotEmpty(analyzer)) {
m.put("analyzer", analyzer);
m.put("search_analyzer", analyzer);
}
properties.put(fieldname, m);
}
}
jsonMap.put("properties", properties);
putMappingRequest.source(jsonMap);
AcknowledgedResponse putMappingResponse = restHighLevelClient.indices().putMapping(putMappingRequest,
RequestOptions.DEFAULT);
return putMappingResponse.isAcknowledged();
}
}
ES数据访问
import cn.hutool.json.JSON;
import org.elasticsearch.action.search.SearchRequest;
import java.util.List;
import java.util.Map;
/**
* @version 1.0
* @date 2022/5/23 16:35
* @since : JDK 11
*/
public interface IElasticsearchDataService {
boolean saveOne(String indexName, Map data);
boolean saveList(String indexName, List
import cn.hutool.core.collection.CollectionUtil;
import cn.hutool.core.util.ArrayUtil;
import cn.hutool.core.util.StrUtil;
import cn.hutool.http.HttpStatus;
import cn.hutool.json.JSON;
import cn.hutool.json.JSONArray;
import org.elasticsearch.action.bulk.BulkRequest;
import org.elasticsearch.action.bulk.BulkResponse;
import org.elasticsearch.action.index.IndexRequest;
import org.elasticsearch.action.index.IndexResponse;
import org.elasticsearch.action.search.SearchRequest;
import org.elasticsearch.action.search.SearchResponse;
import org.elasticsearch.action.update.UpdateRequest;
import org.elasticsearch.action.update.UpdateResponse;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.core.TimeValue;
import org.elasticsearch.index.query.BoolQueryBuilder;
import org.elasticsearch.index.query.QueryBuilders;
import org.elasticsearch.index.query.TermQueryBuilder;
import org.elasticsearch.index.reindex.BulkByScrollResponse;
import org.elasticsearch.index.reindex.DeleteByQueryRequest;
import org.elasticsearch.search.SearchHit;
import org.elasticsearch.search.SearchHits;
import org.elasticsearch.search.builder.SearchSourceBuilder;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.springframework.util.CollectionUtils;
import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
/**
* @version 1.0
* @date 2022/5/23 16:39
* @since : JDK 11
*/
@Service
public class IElasticsearchDataServiceImpl implements IElasticsearchDataService {
@Autowired
private RestHighLevelClient restHighLevelClient;
@Override
public boolean saveOne(String indexName, Map data) {
if (StrUtil.isEmpty(indexName) || null == data) {
return false;
}
IndexRequest request = new IndexRequest(indexName);
request.source(data);
IndexResponse response = null;
try {
response = restHighLevelClient.index(request, RequestOptions.DEFAULT);
if (response.getResult().name().equalsIgnoreCase("created")) {
return true;
}
} catch (IOException e) {
e.printStackTrace();
return false;
}
return false;
}
@Override
public boolean saveList(String indexName, List> dataList) {
if (!StrUtil.isLowerCase(indexName)) {
// 索引名称必须全部小写
return false;
}
if (CollectionUtils.isEmpty(dataList) || StrUtil.isEmpty(indexName)) {
return false;
}
BulkRequest request = new BulkRequest();
for (Map item : dataList) {
request.add(new IndexRequest(indexName).source(item));
}
try {
BulkResponse bulk = restHighLevelClient.bulk(request, RequestOptions.DEFAULT);
if (bulk.status().getStatus() == HttpStatus.HTTP_OK) {
return true;
}
} catch (IOException e) {
e.printStackTrace();
return false;
}
return false;
}
@Override
public boolean saveObj(String indexName, Object data) {
if (StrUtil.isEmpty(indexName) || null == data) {
return false;
}
Map map = BeanMapTool.beanToMap(data);
return this.saveOne(indexName, map);
}
@Override
public boolean saveObjs(String indexName, List dataList) {
if (StrUtil.isEmpty(indexName) || CollectionUtils.isEmpty(dataList)) {
return false;
}
List> maps = BeanMapTool.objectsToMaps(dataList);
return this.saveList(indexName, maps);
}
@Override
public JSON searchBySearchRequest(SearchRequest request) {
if (null == request) {
return null;
}
JSONArray res = new JSONArray();
try {
SearchHits hits = restHighLevelClient.search(request, RequestOptions.DEFAULT).getHits();
if (hits.getTotalHits().value > 0) {
for (SearchHit hit : hits) {
Map map = hit.getSourceAsMap();
res.add(map);
}
}
} catch (IOException e) {
e.printStackTrace();
return null;
}
return res;
}
@Override
public List> searchByIds(List indexName, String... ids) {
if (CollectionUtil.isEmpty(indexName) || ArrayUtil.isEmpty(ids)) {
return null;
}
SearchRequest request = new SearchRequest();
request.indices(ArrayUtil.toArray(indexName, String.class));
SearchSourceBuilder builder = new SearchSourceBuilder();
BoolQueryBuilder boolQueryBuilder = new BoolQueryBuilder();
boolQueryBuilder.filter(QueryBuilders.idsQuery().addIds(ids));
builder.query(boolQueryBuilder);
request.source(builder);
List> list = new ArrayList<>();
try {
SearchResponse response = restHighLevelClient.search(request, RequestOptions.DEFAULT);
SearchHit[] searchHits = response.getHits().getHits();
for (SearchHit hit : searchHits) {
Map map = hit.getSourceAsMap();
list.add(map);
}
return list;
} catch (IOException e) {
e.printStackTrace();
return null;
}
}
@Override
public boolean updateById(String indexName, Map data, String idMap) {
if (StrUtil.isEmpty(indexName) || null == data || StrUtil.isEmpty(idMap)) {
return false;
}
UpdateRequest request = new UpdateRequest(indexName, data.get(idMap).toString()).doc(data);
try {
UpdateResponse update = restHighLevelClient.update(request, RequestOptions.DEFAULT);
if (update.status().getStatus() == HttpStatus.HTTP_OK) {
return true;
} else {
return false;
}
} catch (IOException e) {
e.printStackTrace();
return false;
}
}
@Override
public boolean deleteById(List indexList, String id) {
if (CollectionUtil.isEmpty(indexList) || StrUtil.isEmpty(id)) {
return false;
}
DeleteByQueryRequest request = new DeleteByQueryRequest();
request.indices(ArrayUtil.toArray(indexList, String.class));
request.setQuery(new TermQueryBuilder("id.keyword", id));
// 更新最大文档数
request.setSize(10);
// 批次大小
request.setBatchSize(1000);
// 并行
request.setSlices(2);
// 使用滚动参数来控制“搜索上下文”存活的时间
request.setScroll(TimeValue.timeValueMinutes(10));
// 超时
request.setTimeout(TimeValue.timeValueMinutes(2));
// 刷新索引
request.setRefresh(true);
BulkByScrollResponse response = null;
try {
response = restHighLevelClient.deleteByQuery(request, RequestOptions.DEFAULT);
return response.getStatus().getUpdated() == 0;
} catch (IOException e) {
e.printStackTrace();
return false;
}
}
}