Commit 7e6f8a97 authored by shiyu's avatar shiyu

es索引配置

parent 205cd055
...@@ -13,6 +13,7 @@ import com.wwdz.ch.db.domain.Item; ...@@ -13,6 +13,7 @@ import com.wwdz.ch.db.domain.Item;
import com.wwdz.ch.db.dto.request.CoinRequestDto; import com.wwdz.ch.db.dto.request.CoinRequestDto;
import com.wwdz.ch.db.es.ItemEsDao; import com.wwdz.ch.db.es.ItemEsDao;
import com.xxdxxs.utils.CommonUtils; import com.xxdxxs.utils.CommonUtils;
import com.xxdxxs.utils.DateUtils;
import com.xxdxxs.utils.StringUtils; import com.xxdxxs.utils.StringUtils;
import org.redisson.api.RLock; import org.redisson.api.RLock;
import org.redisson.api.RedissonClient; import org.redisson.api.RedissonClient;
...@@ -59,32 +60,50 @@ public class SynItemToEsJob { ...@@ -59,32 +60,50 @@ public class SynItemToEsJob {
/** /**
* 隔一小时运行一次 * 隔一小时运行一次
*/ */
@Scheduled(fixedDelay = 1000 * 60 * 60) // @Scheduled(fixedDelay = 1000 * 60 * 60)
public void execute() { public void execute() {
RLock lock = redissonClient.getLock(QUANKU_SYN_TO_ES_TASK_KEY); RLock lock = redissonClient.getLock(QUANKU_SYN_TO_ES_TASK_KEY);
if (!lock.tryLock()) { if (!lock.tryLock()) {
logger.warn("当前服务实例获取锁成功: {} 获取锁失败,锁被占用, 不允许执行同步ES作业", Thread.currentThread().getId()); logger.warn("当前服务实例获取锁成功: {} 获取锁失败,锁被占用, 不允许执行同步ES作业", Thread.currentThread().getId());
return; return;
} }
logger.info(">>>>>>>>>>>>>>>>>>>>>>> 同步item数据到ES, 开始执行 <<<<<<<<<<<<<<<<<<<<<"); logger.info(">>>>>>>>>>>>>>>>>>>>>>> 同步item增量数据到ES, 开始执行 <<<<<<<<<<<<<<<<<<<<<");
try { try {
//查询es已有数据的最新修改时间 //查询es已有数据的最新修改时间
Date time = itemEsDao.getMaxTime(); Date time = itemEsDao.getMaxTime();
if (time == null) {
logger.warn(">>>>>>>>>>>> es中没有数据,先同步全量数据 <<<<<<<<<<<<");
return;
}
logger.info("es中最新数据时间为 {}", DateUtils.toString(time));
CoinRequestDto coinRequestDto = new CoinRequestDto(); CoinRequestDto coinRequestDto = new CoinRequestDto();
coinRequestDto.setPage(1); coinRequestDto.setPage(1);
coinRequestDto.setLimit(200); coinRequestDto.setLimit(200);
coinRequestDto.setStartUpdateTime(time); coinRequestDto.setStartUpdateTime(time);
coinRequestDto.setSortColumn("update_time");
List<Item> itemList = itemDao.findListByPage(coinRequestDto); List<Item> itemList = itemDao.findListByPage(coinRequestDto);
PageInfo<Item> pageInfo = new PageInfo<>(itemList); PageInfo<Item> pageInfo = new PageInfo<>(itemList);
logger.info("total : {}, 当前查询 {} 条数据", pageInfo.getTotal(), itemList.size()); logger.info("待同步增量数据总量 : {}, 当前查询 {} 条数据进行同步", pageInfo.getTotal(), itemList.size());
for (Item item : itemList) {
executorService.submit(() -> {
if (itemEsDao.isExised(item.getId())) {
//存在就修改
itemEsDao.update(item);
} else {
itemEsDao.createSingle(item);
}
});
}
itemEsDao.create(itemList); itemEsDao.create(itemList);
logger.info(">>>>>>>>>>>>>>>>>>>>>>> 同步ES成功, 执行完毕 <<<<<<<<<<<<<<<<<<<<<"); logger.info(">>>>>>>>>>>>>>>>>>>>>>> 增量数据同步ES成功, 执行完毕 <<<<<<<<<<<<<<<<<<<<<");
} catch (Exception e) { } catch (Exception e) {
logger.error("同步coins表数据到item表 error {}", e); logger.error("同步item增量数据到ES error {}", e);
} finally { } finally {
if (lock != null && lock.isHeldByCurrentThread()) { if (lock != null && lock.isHeldByCurrentThread()) {
lock.unlock(); lock.unlock();
logger.info("======================== 线程id: {} , 同步ES作业 释放锁成功 ========================", Thread.currentThread().getId()); logger.info("======================== 线程id: {} , 同步增量数据到ES作业 释放锁成功 ========================", Thread.currentThread().getId());
} }
} }
} }
......
package com.wwdz.ch.admin.impl; package com.wwdz.ch.admin.impl;
import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.elasticsearch._types.FieldValue;
import co.elastic.clients.elasticsearch._types.query_dsl.Query;
import co.elastic.clients.elasticsearch._types.query_dsl.QueryBuilders;
import co.elastic.clients.elasticsearch.core.SearchRequest;
import co.elastic.clients.elasticsearch.core.SearchResponse;
import com.wwdz.ch.admin.job.SynCoinJob; import com.wwdz.ch.admin.job.SynCoinJob;
import com.wwdz.ch.core.entity.EsSearchParam; import com.wwdz.ch.core.entity.EsSearchParam;
import com.wwdz.ch.core.util.ElasticsearchUtil; import com.wwdz.ch.core.util.ElasticsearchUtil;
...@@ -17,6 +23,9 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; ...@@ -17,6 +23,9 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
import org.springframework.test.context.web.WebAppConfiguration; import org.springframework.test.context.web.WebAppConfiguration;
import javax.annotation.Resource; import javax.annotation.Resource;
import java.io.IOException;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List; import java.util.List;
@SpringBootTest @SpringBootTest
...@@ -31,6 +40,11 @@ public class ElasticsearchUtilTest { ...@@ -31,6 +40,11 @@ public class ElasticsearchUtilTest {
@Autowired @Autowired
ItemEsDao itemEsDao; ItemEsDao itemEsDao;
@Resource(name="clientByPasswd")
ElasticsearchClient elasticsearchClient;
@Test @Test
public void findDocument() { public void findDocument() {
EsSearchParam esSearchParam = EsSearchParam.of("清代"); EsSearchParam esSearchParam = EsSearchParam.of("清代");
...@@ -42,11 +56,45 @@ public class ElasticsearchUtilTest { ...@@ -42,11 +56,45 @@ public class ElasticsearchUtilTest {
@Test @Test
public void find() { public void find() {
EsSearchRequestDto esSearchRequestDto = new EsSearchRequestDto(); EsSearchRequestDto esSearchRequestDto = new EsSearchRequestDto();
esSearchRequestDto.setContent("清代"); // esSearchRequestDto.setContent("清代");
esSearchRequestDto.setIncludes(Arrays.asList("id", "source", "name","images","updateTime", "createTime"));
esSearchRequestDto.setPage(1); esSearchRequestDto.setPage(1);
esSearchRequestDto.setSize(20); esSearchRequestDto.setSize(20);
EsPageInfo esPageInfo = itemEsDao.findList(esSearchRequestDto); EsPageInfo esPageInfo = itemEsDao.findList(esSearchRequestDto);
logger.info("esPageInfo = {}", esPageInfo); logger.info("esPageInfo = {}", esPageInfo);
} }
/* @Test
public void find2() throws IOException {
EsSearchRequestDto esSearchRequestDto = new EsSearchRequestDto();
// esSearchRequestDto.setContent("清代");
esSearchRequestDto.setIncludes(Arrays.asList("id", "source", "name","images","updateTime", "createTime"));
esSearchRequestDto.setPage(1);
esSearchRequestDto.setSize(20);
EsPageInfo esPageInfo = esDao.findList(esSearchRequestDto);
}
*/
@Test
public void findInTest() throws IOException {
List<FieldValue> list = new ArrayList<>();
list.add(FieldValue.of(258796));
list.add(FieldValue.of(29160));
list.add(FieldValue.of(161073));
list.add(FieldValue.of(325127));
list.add(FieldValue.of(314854));
SearchRequest.Builder builder = new SearchRequest.Builder();
Query terms = QueryBuilders.terms(s->s.field("id").terms(t->t.value(list)));
builder.query(terms);
SearchRequest searchRequest = builder.build();
logger.info("查询语句 : {}", searchRequest.query());
SearchResponse<ItemOfEs> searchResponse = elasticsearchClient.search(searchRequest, ItemOfEs.class);
long total = searchResponse.hits().total().value();
logger.info("总数据量 : {}", total);
logger.info("====== getResponse : {}", searchResponse);
}
} }
\ No newline at end of file
package com.wwdz.ch.admin.impl; package com.wwdz.ch.admin.impl;
import com.wwdz.ch.admin.job.SynCoinJob; import com.wwdz.ch.admin.job.SynCoinJob;
import com.wwdz.ch.admin.job.SynItemToEsJob;
import com.wwdz.ch.db.es.ItemEsDao; import com.wwdz.ch.db.es.ItemEsDao;
import org.junit.Test; import org.junit.Test;
import org.junit.runner.RunWith; import org.junit.runner.RunWith;
...@@ -8,6 +9,8 @@ import org.springframework.beans.factory.annotation.Autowired; ...@@ -8,6 +9,8 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest; import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
import java.util.Date;
import static org.junit.Assert.*; import static org.junit.Assert.*;
@SpringBootTest @SpringBootTest
@RunWith(SpringJUnit4ClassRunner.class) @RunWith(SpringJUnit4ClassRunner.class)
...@@ -19,6 +22,9 @@ public class SynCoinJobTest { ...@@ -19,6 +22,9 @@ public class SynCoinJobTest {
@Autowired @Autowired
ItemEsDao itemEsDao; ItemEsDao itemEsDao;
@Autowired
SynItemToEsJob synItemToEsJob;
@Test @Test
public void execute() { public void execute() {
synCoinJob.execute(); synCoinJob.execute();
...@@ -41,4 +47,10 @@ public class SynCoinJobTest { ...@@ -41,4 +47,10 @@ public class SynCoinJobTest {
public void synToEs() throws Exception { public void synToEs() throws Exception {
synCoinJob.synItemToEs(); synCoinJob.synItemToEs();
} }
@Test
public void getMaxTime() throws Exception {
itemEsDao.getMaxTime();
}
} }
\ No newline at end of file
...@@ -12,6 +12,8 @@ import com.wwdz.ch.db.bean.ItemOfEs; ...@@ -12,6 +12,8 @@ import com.wwdz.ch.db.bean.ItemOfEs;
import com.wwdz.ch.db.domain.Item; import com.wwdz.ch.db.domain.Item;
import com.wwdz.ch.db.dto.request.EsSearchRequestDto; import com.wwdz.ch.db.dto.request.EsSearchRequestDto;
import com.xxdxxs.utils.DateUtils; import com.xxdxxs.utils.DateUtils;
import com.xxdxxs.utils.EntityMapper;
import com.xxdxxs.utils.MapUtils;
import com.xxdxxs.utils.StringUtils; import com.xxdxxs.utils.StringUtils;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
...@@ -162,6 +164,22 @@ public class ItemEsDao { ...@@ -162,6 +164,22 @@ public class ItemEsDao {
} }
public void createSingle(Item item) {
try {
// 创建文档
IndexRequest indexRequest = new IndexRequest.Builder()
.index(INDEX)
.id(String.valueOf(item.getId()))
.document(item)
.build();
IndexResponse indexResponse = elasticsearchClient.index(indexRequest);
logger.info("indexResponse ={}", indexResponse);
logger.info("ES新增数据成功 : {} ", item);
} catch (IOException e) {
logger.info("ES新增数据 content :{} error" , item, e);
}
}
public Boolean isExised(Long id) { public Boolean isExised(Long id) {
try { try {
CountRequest.Builder builder = new CountRequest.Builder(); CountRequest.Builder builder = new CountRequest.Builder();
...@@ -183,8 +201,7 @@ public class ItemEsDao { ...@@ -183,8 +201,7 @@ public class ItemEsDao {
String id = String.valueOf(item.getId()); String id = String.valueOf(item.getId());
try { try {
// 使用map集合封装需要修改的内容 // 使用map集合封装需要修改的内容
Map<String, Object> map = new HashMap<>(); Map<String, Object> map = MapUtils.fromEntity(item);
map.put("images", item.getImages());
UpdateResponse<Item> updateResponse = elasticsearchClient.update(e -> e.index(INDEX).id(id).doc(map), Item.class); UpdateResponse<Item> updateResponse = elasticsearchClient.update(e -> e.index(INDEX).id(id).doc(map), Item.class);
logger.info("updateResponse = {}", updateResponse); logger.info("updateResponse = {}", updateResponse);
logger.info("ES修改数据成功, id:{}, 修改内容: {} ", id, map); logger.info("ES修改数据成功, id:{}, 修改内容: {} ", id, map);
......
Markdown is supported
0%
or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment