赞
踩
深秋初冬的一个晚上,突然间收到业务一个需求,要在老系统上使用新系统Elasticsearch
库的数据。
目前项目情况,新、老系统并行运行,根据产品、渠道路由,但是老系统未使用Elasticsearch
新系统基础框架Spring Cloud Alibaba
version
-> 2.2.1RELEASE
,老系统基础框架 Spring Framework
version
-> 4.3.24RELEASE
。
为了满足业务需求,经过技术分析讨论有两种实现方案:
方案一:
在新系统中暴露
HTTP
服务接口,让老系统直接调用新系统,完成数据获取;
方案二:
在老系统以最小侵入单元的形式集成
Elasticsearch
,完成数据获取;
经过系统交互分析,从系统架构设计角度考虑,为减少系统耦合,采用方案二完成数据接入。
<dependency>
<groupId>org.elasticsearch</groupId>
<artifactId>elasticsearch</artifactId>
<version>6.8.6</version>
</dependency>
<dependency>
<groupId>org.elasticsearch.client</groupId>
<artifactId>elasticsearch-rest-client</artifactId>
<version>6.8.6</version>
</dependency>
<dependency>
<groupId>org.elasticsearch.client</groupId>
<artifactId>elasticsearch-rest-high-level-client</artifactId>
<version>6.8.6</version>
</dependency>
注释:本文采用version
-> 6.8.6
客户端完成接入,原因:与生产Elasticsearch
版本保持一致。
import lombok.extern.slf4j.Slf4j; import org.apache.http.Header; import org.apache.http.HttpHost; import org.apache.http.auth.AuthScope; import org.apache.http.auth.UsernamePasswordCredentials; import org.apache.http.client.CredentialsProvider; import org.apache.http.impl.client.BasicCredentialsProvider; import org.apache.http.message.BasicHeader; import org.elasticsearch.client.RestClient; import org.elasticsearch.client.RestClientBuilder; import org.elasticsearch.client.RestHighLevelClient; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; /** * elasticsearch 配置类 */ @Slf4j @Configuration public class ElasticsearchConfig { @Value("${elasticsearch.cluster.address}") private String clusterAddress; @Value("${elasticsearch.username}") private String username; @Value("${elasticsearch.password}") private String password; @Value("${elasticsearch.shards}") private Integer numberOfShards; @Value("${elasticsearch.replicas}") private Integer numberOfReplicas; @Value("${elasticsearch.connect_timeout}") private Long connectTimeout; @Value("${elasticsearch.socket_timeout}") private Long socketTimeout; public static RestHighLevelClient client = null; public Integer getNumberOfShards() { return numberOfShards; } public Integer getNumberOfReplicas() { return numberOfReplicas; } /** * RestHighLevelClient bean创建 */ @Bean public RestHighLevelClient restClient() { final CredentialsProvider credentialsProvider = new BasicCredentialsProvider(); credentialsProvider.setCredentials(AuthScope.ANY, new UsernamePasswordCredentials(username, password)); Header[] defaultHeaders = {new BasicHeader("content-type", "application/json")}; RestClientBuilder restClientBuilder = RestClient.builder(HttpHost.create(clusterAddress)); restClientBuilder .setHttpClientConfigCallback(httpAsyncClientBuilder -> httpAsyncClientBuilder.setDefaultCredentialsProvider(credentialsProvider)) .setDefaultHeaders(defaultHeaders) .setRequestConfigCallback(requestConfigBuilder -> { // 连接5秒超时,套接字连接60s超时 return requestConfigBuilder.setConnectTimeout(connectTimeout.intValue()).setSocketTimeout(socketTimeout.intValue()); }) .setHttpClientConfigCallback(httpClientBuilder -> { httpClientBuilder.disableAuthCaching(); return httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider); }); client = new RestHighLevelClient(restClientBuilder); return client; } }
Elasticsearch
持久层接入import cn.hutool.core.map.MapUtil; import lombok.extern.slf4j.Slf4j; import org.elasticsearch.action.admin.indices.delete.DeleteIndexRequest; import org.elasticsearch.action.bulk.BulkRequest; import org.elasticsearch.action.delete.DeleteRequest; import org.elasticsearch.action.index.IndexRequest; import org.elasticsearch.action.search.SearchRequest; import org.elasticsearch.action.search.SearchResponse; import org.elasticsearch.action.support.WriteRequest; import org.elasticsearch.action.support.master.AcknowledgedResponse; import org.elasticsearch.action.update.UpdateRequest; import org.elasticsearch.client.RequestOptions; import org.elasticsearch.client.RestHighLevelClient; import org.elasticsearch.client.core.CountRequest; import org.elasticsearch.client.core.CountResponse; import org.elasticsearch.client.indices.CreateIndexRequest; import org.elasticsearch.client.indices.GetIndexRequest; import org.elasticsearch.index.query.BoolQueryBuilder; import org.elasticsearch.index.query.QueryBuilders; import org.elasticsearch.index.query.RangeQueryBuilder; import org.elasticsearch.search.SearchHit; import org.elasticsearch.search.builder.SearchSourceBuilder; import org.elasticsearch.search.sort.SortOrder; import org.elasticsearch.xcontent.XContentType; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import java.io.IOException; import java.util.*; /** * elasticsearch 持久层 */ @Slf4j @Service public class ElasticsearchRepository { @Autowired private RestHighLevelClient client ; private final RequestOptions options = RequestOptions.DEFAULT; /** * 写入数据 * @param indexName * @param dataMap 数据实体 * @return */ public boolean insert (String indexName, Map<String,Object> dataMap){ try { BulkRequest request = new BulkRequest(); request.add(new IndexRequest(indexName,"record").id(dataMap.remove("id").toString()) .opType("create").source(dataMap, XContentType.JSON)); client.bulk(request, options); return Boolean.TRUE ; } catch (Exception e){ log.error("ElasticsearchRepository#insert, 索引名称:{}, 执行异常:{}", indexName, e); } return Boolean.FALSE; } /** * 批量写入数据 * @param indexName * @param userIndexList * @return */ public boolean batchInsert (String indexName, List<Map<String,Object>> userIndexList){ try { BulkRequest request = new BulkRequest(); for (Map<String,Object> dataMap:userIndexList){ request.add(new IndexRequest(indexName,"record").id(dataMap.remove("id").toString()) .opType("create").source(dataMap,XContentType.JSON)); } client.bulk(request, options); return Boolean.TRUE ; } catch (Exception e){ log.error("ElasticsearchRepository#batchInsert, 索引名称:{}, 执行异常:{}", indexName, e); } return Boolean.FALSE; } /** * 更新数据 * 可以直接修改索引结构 * @param indexName * @param dataMap * @return */ public boolean update (String indexName, Map<String,Object> dataMap){ try { UpdateRequest updateRequest = new UpdateRequest(indexName,"record", dataMap.remove("id").toString()); updateRequest.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE); updateRequest.doc(dataMap) ; client.update(updateRequest, options); return Boolean.TRUE ; } catch (Exception e){ log.error("ElasticsearchRepository#update, 索引名称:{}, 执行异常:{}", indexName, e); } return Boolean.FALSE; } /** * 根据 id 及索引删除数据 * @param indexName * @param id * @return */ public boolean delete (String indexName, String id){ try { DeleteRequest deleteRequest = new DeleteRequest(indexName,"record", id); client.delete(deleteRequest, options); return Boolean.TRUE ; } catch (Exception e){ log.error("ElasticsearchRepository#delete, 索引名称:{}, 执行异常:{}", indexName, e); } return Boolean.FALSE; } /** * 判断索引是否存在 * @param indexName * @return */ public boolean checkIndex (String indexName) { try { return client.indices().exists(new GetIndexRequest(indexName), options); } catch (IOException e) { log.error("ElasticsearchRepository#checkIndex, 索引名称:{}, 执行异常:{}", indexName, e); } return Boolean.FALSE ; } /** * 创建索引 * @param indexName * @param columnMap * @return */ public boolean createIndex (String indexName ,Map<String, Object> columnMap){ try { if(!checkIndex(indexName)){ CreateIndexRequest request = new CreateIndexRequest(indexName); if (columnMap != null && columnMap.size()>0) { Map<String, Object> source = new HashMap<>(); source.put("properties", columnMap); request.mapping(source); } client.indices().create(request, options); return Boolean.TRUE ; } } catch (IOException e) { log.error("ElasticsearchRepository#createIndex, 索引名称:{}, 执行异常:{}", indexName, e); } return Boolean.FALSE; } /** * 删除索引 * @param indexName * @return */ public boolean deleteIndex(String indexName) { try { if(checkIndex(indexName)){ DeleteIndexRequest request = new DeleteIndexRequest(indexName); AcknowledgedResponse response = client.indices().delete(request, options); return response.isAcknowledged(); } } catch (Exception e) { log.error("ElasticsearchRepository#deleteIndex, 索引名称:{}, 执行异常:{}", indexName, e); } return Boolean.FALSE; } /** * 查询满足条件的数据条数 * @param indexName * @param matchMap * @return */ public Long count (String indexName, LinkedHashMap<String, Object> matchMap){ // 查询器构造 BoolQueryBuilder queryBuilder = QueryBuilders.boolQuery(); SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); // 查询时间范围内的数据 if (MapUtil.isNotEmpty(matchMap)){ if (matchMap.containsKey("startTime") && matchMap.containsKey("endTime")){ RangeQueryBuilder rangequerybuilder = QueryBuilders .rangeQuery("createTime") .from(Long.parseLong(String.valueOf(matchMap.get("startTime")))) .to(Long.parseLong(String.valueOf(matchMap.get("endTime")))); queryBuilder.must(rangequerybuilder); } // 移除时间查询条件 matchMap.remove("startTime"); matchMap.remove("endTime"); // 时间查询条件外的参数拼接 if (MapUtil.isNotEmpty(matchMap)){ matchMap.forEach((k, v) -> {queryBuilder.must(QueryBuilders.termQuery(k, v));}); sourceBuilder.query(queryBuilder); } } CountRequest countRequest = new CountRequest(indexName); countRequest.source(sourceBuilder); try { CountResponse countResponse = client.count(countRequest, options); return countResponse.getCount(); } catch (Exception e) { log.error("ElasticsearchRepository#count, 索引名称:{}, 执行异常:{}", indexName, e); } return 0L; } /** * 查询满足条件的数据集合 * 适用于满足条件的数据条数可控的全量查询 PS:单次查询条数不超过 10000条 * @param indexName * @param matchMap * @return */ public List<Map<String,Object>> list (String indexName, LinkedHashMap<String, Object> matchMap) { // 查询条件,指定时间并过滤指定字段值 SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); BoolQueryBuilder queryBuilder = QueryBuilders.boolQuery(); // 查询参数拼接 if (MapUtil.isNotEmpty(matchMap)){ matchMap.forEach((k, v) -> {queryBuilder.must(QueryBuilders.termQuery(k, v));}); } sourceBuilder.query(queryBuilder); SearchRequest searchRequest = new SearchRequest(indexName); searchRequest.source(sourceBuilder); try { SearchResponse searchResp = client.search(searchRequest, options); List<Map<String,Object>> data = new ArrayList<>() ; SearchHit[] searchHitArr = searchResp.getHits().getHits(); for (SearchHit searchHit:searchHitArr){ Map<String,Object> temp = searchHit.getSourceAsMap(); temp.put("id",searchHit.getId()) ; data.add(temp); } return data; } catch (Exception e) { log.error("ElasticsearchRepository#list, 索引名称:{}, 执行异常:{}", indexName, e); } return null ; } /** * 根据查询条件,分页查询 * 适用于满足条件的数据总量较大的循环查询场景 * @param indexName * @param offset 偏移量 * @param size 条数 * @param matchMap * @return */ public List<Map<String,Object>> page (String indexName, Integer offset, Integer size, LinkedHashMap<String, Object> matchMap) { // 添加分页参数 SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); sourceBuilder.from(offset); sourceBuilder.size(size); sourceBuilder.sort("createTime", SortOrder.DESC); BoolQueryBuilder queryBuilder = QueryBuilders.boolQuery(); // 查询时间范围内的数据 if (MapUtil.isNotEmpty(matchMap)){ if (matchMap.containsKey("startTime") && matchMap.containsKey("endTime")){ RangeQueryBuilder rangequerybuilder = QueryBuilders .rangeQuery("createTime") .from(Long.parseLong(String.valueOf(matchMap.get("startTime")))) .to(Long.parseLong(String.valueOf(matchMap.get("endTime")))); queryBuilder.must(rangequerybuilder); } // 移除时间查询条件 matchMap.remove("startTime"); matchMap.remove("endTime"); // 时间查询条件外的参数拼接 if (MapUtil.isNotEmpty(matchMap)){ matchMap.forEach((k, v) -> {queryBuilder.must(QueryBuilders.termQuery(k, v));}); sourceBuilder.query(queryBuilder); } } // 查询请求 SearchRequest searchRequest = new SearchRequest(indexName); searchRequest.source(sourceBuilder); try { SearchResponse searchResp = client.search(searchRequest, options); List<Map<String,Object>> data = new ArrayList<>() ; SearchHit[] searchHitArr = searchResp.getHits().getHits(); for (SearchHit searchHit:searchHitArr){ Map<String,Object> temp = searchHit.getSourceAsMap(); temp.put("id",searchHit.getId()) ; data.add(temp); } return data; } catch (Exception e) { log.error("ElasticsearchRepository#page, 索引名称:{}, 执行异常:{}", indexName, e); } return null ; } /** * 根据条件查询,按照创建时间进行降序排列 * 可扩展,根据更新时间、 id、证件号等 * @param indexName * @param matchMap * @return */ public List<Map<String,Object>> sort (String indexName, LinkedHashMap<String, Object> matchMap) { // 先升序时间,在倒序年龄 SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); sourceBuilder.sort("createTime", SortOrder.ASC); // 查询器构造 BoolQueryBuilder queryBuilder = QueryBuilders.boolQuery(); // 查询时间范围内的数据 if (MapUtil.isNotEmpty(matchMap)){ if (matchMap.containsKey("startTime") && matchMap.containsKey("endTime")){ RangeQueryBuilder rangequerybuilder = QueryBuilders .rangeQuery("createTime") .from(Long.parseLong(String.valueOf(matchMap.get("startTime")))) .to(Long.parseLong(String.valueOf(matchMap.get("endTime")))); queryBuilder.must(rangequerybuilder); } // 移除时间查询条件 matchMap.remove("startTime"); matchMap.remove("endTime"); // 时间查询条件外的参数拼接 if (MapUtil.isNotEmpty(matchMap)){ matchMap.forEach((k, v) -> {queryBuilder.must(QueryBuilders.termQuery(k, v));}); sourceBuilder.query(queryBuilder); } } SearchRequest searchRequest = new SearchRequest(indexName); searchRequest.source(sourceBuilder); try { SearchResponse searchResp = client.search(searchRequest, options); List<Map<String,Object>> data = new ArrayList<>() ; SearchHit[] searchHitArr = searchResp.getHits().getHits(); for (SearchHit searchHit:searchHitArr){ Map<String,Object> temp = searchHit.getSourceAsMap(); temp.put("id",searchHit.getId()) ; data.add(temp); } return data; } catch (Exception e) { log.error("ElasticsearchRepository#sort, 索引名称:{}, 执行异常:{}", indexName, e); } return null ; } }
采用 Junit
实现
1.在承接业务需求时,首先要结合功能实现的复杂度,考虑架构的合理性,在相对更合理的系统设计背景下进行功能设计、开发;
2.进行技术开发时,首先要考虑功能对模块的侵入性,在最小侵入性的前提下,采用与基础框架融合的方式,完成开发任务。
Copyright © 2003-2013 www.wpsshop.cn 版权所有,并保留所有权利。