Commit c8e1496e authored by shiyu's avatar shiyu

商品同步es

parent 6684964a
...@@ -120,6 +120,7 @@ public class SynCoinJob { ...@@ -120,6 +120,7 @@ public class SynCoinJob {
CoinRequestDto coinRequestDto = new CoinRequestDto(); CoinRequestDto coinRequestDto = new CoinRequestDto();
coinRequestDto.setPage(1); coinRequestDto.setPage(1);
coinRequestDto.setLimit(100); coinRequestDto.setLimit(100);
coinRequestDto.setKind(3);
PageInfo<Item> pageInfo = itemDao.findPageExcludeFilter(coinRequestDto); PageInfo<Item> pageInfo = itemDao.findPageExcludeFilter(coinRequestDto);
int totalPageNum = pageInfo.getPages(); int totalPageNum = pageInfo.getPages();
logger.info("total page : {}, total : {}", totalPageNum, pageInfo.getTotal()); logger.info("total page : {}, total : {}", totalPageNum, pageInfo.getTotal());
...@@ -128,6 +129,7 @@ public class SynCoinJob { ...@@ -128,6 +129,7 @@ public class SynCoinJob {
CoinRequestDto dto = new CoinRequestDto(); CoinRequestDto dto = new CoinRequestDto();
dto.setPage(i); dto.setPage(i);
dto.setLimit(100); dto.setLimit(100);
dto.setKind(3);
List<Item> itemList = itemDao.findListByPage(dto); List<Item> itemList = itemDao.findListByPage(dto);
itemList.parallelStream().forEach(item -> { itemList.parallelStream().forEach(item -> {
String imageUrl = item.getImages(); String imageUrl = item.getImages();
......
...@@ -8,8 +8,10 @@ import com.wwdz.ch.core.storage.QiniuStorage; ...@@ -8,8 +8,10 @@ import com.wwdz.ch.core.storage.QiniuStorage;
import com.wwdz.ch.core.util.RedisUtils; import com.wwdz.ch.core.util.RedisUtils;
import com.wwdz.ch.db.dao.CoinsDao; import com.wwdz.ch.db.dao.CoinsDao;
import com.wwdz.ch.db.dao.ItemDao; import com.wwdz.ch.db.dao.ItemDao;
import com.wwdz.ch.db.dao.SynDataRecordDao;
import com.wwdz.ch.db.domain.Coins; import com.wwdz.ch.db.domain.Coins;
import com.wwdz.ch.db.domain.Item; import com.wwdz.ch.db.domain.Item;
import com.wwdz.ch.db.domain.SynDataRecord;
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;
...@@ -44,7 +46,9 @@ public class SynItemToEsJob { ...@@ -44,7 +46,9 @@ public class SynItemToEsJob {
private static final Logger logger = LoggerFactory.getLogger(SynItemToEsJob.class); private static final Logger logger = LoggerFactory.getLogger(SynItemToEsJob.class);
private String QUANKU_SYN_TO_ES_TASK_KEY = "QUANKU_SYN_TO_ES_TASK_KEY"; private final static String QUANKU_SYN_TO_ES_TASK_KEY = "QUANKU_SYN_TO_ES_TASK_KEY";
private final int SYN_ITEM_TO_ES_TYPE = 1;
ExecutorService executorService = Executors.newFixedThreadPool(10); ExecutorService executorService = Executors.newFixedThreadPool(10);
...@@ -57,10 +61,13 @@ public class SynItemToEsJob { ...@@ -57,10 +61,13 @@ public class SynItemToEsJob {
@Autowired @Autowired
RedissonClient redissonClient; RedissonClient redissonClient;
@Autowired
SynDataRecordDao synDataRecordDao;
/** /**
* 隔2分钟运行一次 * 隔2分钟运行一次
*/ */
@Scheduled(fixedDelay = 1000 * 60 * 2) @Scheduled(fixedDelay = 1000 * 60 * 1)
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()) {
...@@ -69,17 +76,34 @@ public class SynItemToEsJob { ...@@ -69,17 +76,34 @@ public class SynItemToEsJob {
} }
logger.info(">>>>>>>>>>>>>>>>>>>>>>> 同步item增量数据到ES, 开始执行 <<<<<<<<<<<<<<<<<<<<<"); logger.info(">>>>>>>>>>>>>>>>>>>>>>> 同步item增量数据到ES, 开始执行 <<<<<<<<<<<<<<<<<<<<<");
try { try {
//查询es已有数据的最新修改时间 //查询同步数据记录表中数据,判断上次的同步操作是否完成,没有则继续上次的查询条件继续执行,完成则查询es中最新时间作为本次查询开始时间
Date time = itemEsDao.getMaxTime(); Date startTime;
if (time == null) { Date endTime;
logger.warn(">>>>>>>>>>>> es中没有数据,先同步全量数据 <<<<<<<<<<<<"); int page = 1;
return; int limit = 100;
SynDataRecord synDataRecord = synDataRecordDao.find(SYN_ITEM_TO_ES_TYPE);
if (synDataRecord.getIsEnd().intValue() != 1) {
startTime = synDataRecord.getQueryStartTime();
endTime = synDataRecord.getQueryEndTime();
page = synDataRecord.getPageIndex() + 1;
limit = synDataRecord.getPageSize();
logger.info("同步item增量数据到ES, 承接上次同步的时间范围执行, 时间范围 {} to {}", DateUtils.toString(startTime), DateUtils.toString(endTime));
} else {
//查询es已有数据的最新修改时间
startTime = itemEsDao.getMaxTime();
if (startTime == null) {
logger.warn(">>>>>>>>>>>> es中没有数据,先同步全量数据 <<<<<<<<<<<<");
return;
}
endTime = new Date();
logger.info("上次同步任务执行完毕,开始新一轮时间范围的数据同步, 时间范围 {} to {}", DateUtils.toString(startTime), DateUtils.toString(endTime));
} }
logger.info("es中最新数据时间为 {}", DateUtils.toString(time)); logger.info("es中最新数据时间为 {}", DateUtils.toString(startTime));
CoinRequestDto coinRequestDto = new CoinRequestDto(); CoinRequestDto coinRequestDto = new CoinRequestDto();
coinRequestDto.setPage(1); coinRequestDto.setPage(page);
coinRequestDto.setLimit(100); coinRequestDto.setLimit(limit);
coinRequestDto.setStartUpdateTime(time); coinRequestDto.setStartUpdateTime(startTime);
coinRequestDto.setEndUpdateTime(endTime);
coinRequestDto.setSortColumn("update_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);
...@@ -89,6 +113,26 @@ public class SynItemToEsJob { ...@@ -89,6 +113,26 @@ public class SynItemToEsJob {
itemEsDao.upsert(item); itemEsDao.upsert(item);
}); });
} }
//如果有下一页则记录下来,下一次任务执行接着同步该更新时间范围的
if (pageInfo.isHasNextPage()) {
SynDataRecord updateRecord = new SynDataRecord();
updateRecord.setPageIndex(page);
updateRecord.setPageSize(limit);
updateRecord.setQueryStartTime(startTime);
updateRecord.setQueryEndTime(endTime);
updateRecord.setIsEnd(0);
updateRecord.setType(SYN_ITEM_TO_ES_TYPE);
synDataRecordDao.update(updateRecord);
} else {
SynDataRecord updateRecord = new SynDataRecord();
updateRecord.setPageIndex(page);
updateRecord.setPageSize(limit);
updateRecord.setQueryStartTime(startTime);
updateRecord.setQueryEndTime(endTime);
updateRecord.setIsEnd(1);
updateRecord.setType(SYN_ITEM_TO_ES_TYPE);
synDataRecordDao.update(updateRecord);
}
logger.info(">>>>>>>>>>>>>>>>>>>>>>> 增量数据同步ES成功, 执行完毕 <<<<<<<<<<<<<<<<<<<<<"); logger.info(">>>>>>>>>>>>>>>>>>>>>>> 增量数据同步ES成功, 执行完毕 <<<<<<<<<<<<<<<<<<<<<");
} catch (Exception e) { } catch (Exception e) {
logger.error("同步item增量数据到ES error {}", e); logger.error("同步item增量数据到ES error {}", e);
......
package com.wwdz.ch.db.dao;
import com.wwdz.ch.db.domain.SynDataRecord;
public interface SynDataRecordDao {
boolean create(SynDataRecord record);
boolean update(SynDataRecord record);
SynDataRecord find(int type);
}
This diff is collapsed.
...@@ -99,6 +99,11 @@ public class CoinRequestDto extends BaseRequestDto implements Entity { ...@@ -99,6 +99,11 @@ public class CoinRequestDto extends BaseRequestDto implements Entity {
*/ */
private Date startUpdateTime; private Date startUpdateTime;
/**
* 用于同步数据到es
*/
private Date endUpdateTime;
private String detail; private String detail;
private Double price; private Double price;
......
...@@ -315,12 +315,12 @@ public class ItemEsDao { ...@@ -315,12 +315,12 @@ public class ItemEsDao {
IndexResponse indexResponse = elasticsearchClient.index(indexRequest); IndexResponse indexResponse = elasticsearchClient.index(indexRequest);
logger.info("indexResponse ={}", indexResponse); logger.info("indexResponse ={}", indexResponse);
if (indexResponse.shards().successful().intValue() == 1 ) { if (indexResponse.shards().successful().intValue() == 1 ) {
logger.info("ES upsert 数据成功 : {} ", item); logger.info("ES upsert 数据成功 : {} ", item.getId());
} else { } else {
logger.error("ES upsert 数据失败 : {} ", item); logger.error("ES upsert 数据失败 : {} ", item.getId());
} }
} catch (IOException e) { } catch (IOException e) {
logger.info("ES upsert 数据 content :{} error" , item, e); logger.info("ES upsert 数据 content :{} error" , item.getId(), e);
} }
} }
......
...@@ -94,10 +94,11 @@ public class ItemDaoImpl implements ItemDao { ...@@ -94,10 +94,11 @@ public class ItemDaoImpl implements ItemDao {
ItemExample example = new ItemExample(); ItemExample example = new ItemExample();
ItemExample.Criteria criteria = example.createCriteria(); ItemExample.Criteria criteria = example.createCriteria();
example.setOrderByClause("create_time asc"); example.setOrderByClause("create_time asc");
JdbcHelper.ifPresent(coinRequestDto.getStartUpdateTime(), criteria::andUpdateTimeGreaterThan);
JdbcHelper.ifPresent(coinRequestDto.getEndUpdateTime(), criteria::andUpdateTimeLessThanOrEqualTo);
JdbcHelper.ifPresent(coinRequestDto.getKind(), criteria :: andKindEqualTo); JdbcHelper.ifPresent(coinRequestDto.getKind(), criteria :: andKindEqualTo);
JdbcHelper.ifPresent(coinRequestDto.getIsDeleted(), criteria::andIsDeletedEqualTo); JdbcHelper.ifPresent(coinRequestDto.getIsDeleted(), criteria::andIsDeletedEqualTo);
JdbcHelper.ifPresent(coinRequestDto.getCid(), criteria::andCidEqualTo); JdbcHelper.ifPresent(coinRequestDto.getCid(), criteria::andCidEqualTo);
JdbcHelper.ifPresent(coinRequestDto.getStartUpdateTime(), criteria::andUpdateTimeGreaterThan);
PageHelper.startPage(coinRequestDto.getPage(), coinRequestDto.getLimit()); PageHelper.startPage(coinRequestDto.getPage(), coinRequestDto.getLimit());
return itemMapper.selectByExampleWithBLOBs(example); return itemMapper.selectByExampleWithBLOBs(example);
} }
......
package com.wwdz.ch.db.impl;
import com.wwdz.ch.db.dao.SynDataRecordDao;
import com.wwdz.ch.db.domain.SynDataRecord;
import com.wwdz.ch.db.domain.SynDataRecordExample;
import com.wwdz.ch.db.mapper.SynDataRecordMapper;
import lombok.Data;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Repository;
import java.util.Date;
/**
* 同步数据记录
*/
@Repository
public class SynDataRecordDaoImpl implements SynDataRecordDao {
@Autowired
SynDataRecordMapper synDataRecordMapper;
@Override
public boolean create(SynDataRecord record) {
Date now = new Date();
record.setCreateTime(now);
record.setUpdateTime(now);
return synDataRecordMapper.insert(record) > 0;
}
@Override
public boolean update(SynDataRecord record) {
SynDataRecordExample synDataRecordExample = new SynDataRecordExample();
SynDataRecordExample.Criteria criteria = synDataRecordExample.createCriteria();
criteria.andTypeEqualTo(record.getType());
Date now = new Date();
record.setUpdateTime(now);
return synDataRecordMapper.updateByExampleSelective(record, synDataRecordExample) > 0;
}
@Override
public SynDataRecord find(int type) {
SynDataRecordExample synDataRecordExample = new SynDataRecordExample();
SynDataRecordExample.Criteria criteria = synDataRecordExample.createCriteria();
criteria.andTypeEqualTo(type);
return synDataRecordMapper.selectOneByExample(synDataRecordExample);
}
}
package com.wwdz.ch.db.mapper;
import com.wwdz.ch.db.domain.SynDataRecord;
import com.wwdz.ch.db.domain.SynDataRecordExample;
import java.util.List;
import org.apache.ibatis.annotations.Mapper;
import org.apache.ibatis.annotations.Param;
@Mapper
public interface SynDataRecordMapper {
long countByExample(SynDataRecordExample example);
int deleteByExample(SynDataRecordExample example);
int deleteByPrimaryKey(Integer id);
int insert(SynDataRecord record);
int insertSelective(SynDataRecord record);
/**
* This method was generated by MyBatis Generator.
* This method corresponds to the database table syn_data_record
*
* @mbg.generated
* @project https://github.com/itfsw/mybatis-generator-plugin
*/
SynDataRecord selectOneByExample(SynDataRecordExample example);
/**
* This method was generated by MyBatis Generator.
* This method corresponds to the database table syn_data_record
*
* @mbg.generated
* @project https://github.com/itfsw/mybatis-generator-plugin
*/
SynDataRecord selectOneByExampleSelective(@Param("example") SynDataRecordExample example, @Param("selective") SynDataRecord.Column ... selective);
/**
* This method was generated by MyBatis Generator.
* This method corresponds to the database table syn_data_record
*
* @mbg.generated
* @project https://github.com/itfsw/mybatis-generator-plugin
*/
List<SynDataRecord> selectByExampleSelective(@Param("example") SynDataRecordExample example, @Param("selective") SynDataRecord.Column ... selective);
List<SynDataRecord> selectByExample(SynDataRecordExample example);
/**
* This method was generated by MyBatis Generator.
* This method corresponds to the database table syn_data_record
*
* @mbg.generated
* @project https://github.com/itfsw/mybatis-generator-plugin
*/
SynDataRecord selectByPrimaryKeySelective(@Param("id") Integer id, @Param("selective") SynDataRecord.Column ... selective);
SynDataRecord selectByPrimaryKey(Integer id);
int updateByExampleSelective(@Param("record") SynDataRecord record, @Param("example") SynDataRecordExample example);
int updateByExample(@Param("record") SynDataRecord record, @Param("example") SynDataRecordExample example);
int updateByPrimaryKeySelective(SynDataRecord record);
int updateByPrimaryKey(SynDataRecord record);
}
\ No newline at end of file
...@@ -72,7 +72,7 @@ ...@@ -72,7 +72,7 @@
<javaClientGenerator type="XMLMAPPER" targetPackage="com.wwdz.ch.db.mapper" <javaClientGenerator type="XMLMAPPER" targetPackage="com.wwdz.ch.db.mapper"
targetProject="ch-dao/src/main/java"/> targetProject="ch-dao/src/main/java"/>
<table tableName="item" enableCountByExample="true" enableUpdateByExample="true" enableDeleteByExample="true" enableSelectByExample="true" selectByExampleQueryId="true"> <table tableName="syn_data_record" enableCountByExample="true" enableUpdateByExample="true" enableDeleteByExample="true" enableSelectByExample="true" selectByExampleQueryId="true">
<generatedKey column="id" sqlStatement="Mysql" identity="true" /> <generatedKey column="id" sqlStatement="Mysql" identity="true" />
</table> </table>
......
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