Commit dcb78f84 authored by shiyu's avatar shiyu

es增量数据同步

parent 538ab83a
...@@ -58,9 +58,9 @@ public class SynItemToEsJob { ...@@ -58,9 +58,9 @@ public class SynItemToEsJob {
RedissonClient redissonClient; RedissonClient redissonClient;
/** /**
* 隔一小时运行一次 * 隔2分钟运行一次
*/ */
// @Scheduled(fixedDelay = 1000 * 60 * 60) @Scheduled(fixedDelay = 1000 * 60 * 2)
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()) {
...@@ -78,7 +78,7 @@ public class SynItemToEsJob { ...@@ -78,7 +78,7 @@ public class SynItemToEsJob {
logger.info("es中最新数据时间为 {}", DateUtils.toString(time)); 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(10);
coinRequestDto.setStartUpdateTime(time); coinRequestDto.setStartUpdateTime(time);
coinRequestDto.setSortColumn("update_time"); coinRequestDto.setSortColumn("update_time");
List<Item> itemList = itemDao.findListByPage(coinRequestDto); List<Item> itemList = itemDao.findListByPage(coinRequestDto);
...@@ -86,17 +86,9 @@ public class SynItemToEsJob { ...@@ -86,17 +86,9 @@ public class SynItemToEsJob {
logger.info("待同步增量数据总量 : {}, 当前查询 {} 条数据进行同步", pageInfo.getTotal(), itemList.size()); logger.info("待同步增量数据总量 : {}, 当前查询 {} 条数据进行同步", pageInfo.getTotal(), itemList.size());
for (Item item : itemList) { for (Item item : itemList) {
executorService.submit(() -> { executorService.submit(() -> {
if (itemEsDao.isExised(item.getId())) { itemEsDao.upsert(item);
//存在就修改
itemEsDao.update(item);
} else {
itemEsDao.createSingle(item);
}
}); });
} }
itemEsDao.create(itemList);
logger.info(">>>>>>>>>>>>>>>>>>>>>>> 增量数据同步ES成功, 执行完毕 <<<<<<<<<<<<<<<<<<<<<"); logger.info(">>>>>>>>>>>>>>>>>>>>>>> 增量数据同步ES成功, 执行完毕 <<<<<<<<<<<<<<<<<<<<<");
} catch (Exception e) { } catch (Exception e) {
logger.error("同步item增量数据到ES error {}", e); logger.error("同步item增量数据到ES error {}", e);
......
...@@ -57,7 +57,7 @@ public class ElasticsearchUtilTest { ...@@ -57,7 +57,7 @@ public class ElasticsearchUtilTest {
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.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);
......
...@@ -2,6 +2,8 @@ package com.wwdz.ch.admin.impl; ...@@ -2,6 +2,8 @@ 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.admin.job.SynItemToEsJob;
import com.wwdz.ch.db.dao.ItemDao;
import com.wwdz.ch.db.domain.Item;
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;
...@@ -25,6 +27,9 @@ public class SynCoinJobTest { ...@@ -25,6 +27,9 @@ public class SynCoinJobTest {
@Autowired @Autowired
SynItemToEsJob synItemToEsJob; SynItemToEsJob synItemToEsJob;
@Autowired
ItemDao itemDao;
@Test @Test
public void execute() { public void execute() {
synCoinJob.execute(); synCoinJob.execute();
...@@ -42,6 +47,18 @@ public class SynCoinJobTest { ...@@ -42,6 +47,18 @@ public class SynCoinJobTest {
itemEsDao.createIndex(); itemEsDao.createIndex();
} }
@Test
public void upsert() throws Exception {
Item item = itemDao.findById(767537L);
itemEsDao.upsert(item);
}
@Test
public void updateDocumnet() throws Exception {
Item item = itemDao.findById(767540L);
itemEsDao.update(item);
}
@Test @Test
public void synToEs() throws Exception { public void synToEs() throws Exception {
......
...@@ -169,7 +169,7 @@ spring: ...@@ -169,7 +169,7 @@ spring:
uris: 172.16.122.190:9200 uris: 172.16.122.190:9200
username: elastic username: elastic
password: changeme password: changeme
index: quanku_new index: quanku_prod
server: server:
servlet: servlet:
session: session:
......
...@@ -163,10 +163,12 @@ public class ItemEsDao { ...@@ -163,10 +163,12 @@ public class ItemEsDao {
} }
/**
public void createSingle(Item item) { * 更新或创建
* @param item
*/
public void upsert(Item item) {
try { try {
// 创建文档
IndexRequest indexRequest = new IndexRequest.Builder() IndexRequest indexRequest = new IndexRequest.Builder()
.index(INDEX) .index(INDEX)
.id(String.valueOf(item.getId())) .id(String.valueOf(item.getId()))
...@@ -174,9 +176,13 @@ public class ItemEsDao { ...@@ -174,9 +176,13 @@ public class ItemEsDao {
.build(); .build();
IndexResponse indexResponse = elasticsearchClient.index(indexRequest); IndexResponse indexResponse = elasticsearchClient.index(indexRequest);
logger.info("indexResponse ={}", indexResponse); logger.info("indexResponse ={}", indexResponse);
logger.info("ES新增数据成功 : {} ", item); if (indexResponse.shards().successful().intValue() == 1 ) {
logger.info("ES upsert 数据成功 : {} ", item);
} else {
logger.error("ES upsert 数据失败 : {} ", item);
}
} catch (IOException e) { } catch (IOException e) {
logger.info("ES新增数据 content :{} error" , item, e); logger.info("ES upsert 数据 content :{} error" , item, e);
} }
} }
......
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