Commit 5486ed29 authored by shiyu's avatar shiyu

查询排除部分数据

parent 50618e8e
......@@ -4,6 +4,7 @@ import com.wwdz.ch.admin.entity.vo.CoinVo;
import com.wwdz.ch.admin.service.ItemService;
import com.wwdz.ch.admin.service.OperateLogService;
import com.wwdz.ch.admin.util.AuthSupport;
import com.wwdz.ch.core.consts.CommConsts;
import com.wwdz.ch.core.consts.CommonEnum;
import com.wwdz.ch.core.storage.QiniuStorage;
import com.wwdz.ch.core.type.Result;
......@@ -134,14 +135,14 @@ public class ItemServiceImpl implements ItemService {
String extraInfo = JsonUtils.fromMap(map);
item.setExtra(extraInfo);
String imgBaseStr = coinRequestDto.getImgBaseStr();
String keyName = String.valueOf(coinRequestDto.getId()) + System.currentTimeMillis() + IMG_SUFFIX ;
String keyName = CommConsts.UPLOAD_PRE_DIRECTORY + coinRequestDto.getId() + System.currentTimeMillis() ;
Map<String, Object> paramMap = new HashMap<>();
if (StringUtils.hasLength(item.getExtra())) {
paramMap = JsonUtils.toMap(item.getExtra());
}
if (StringUtils.hasLength(imgBaseStr)) {
qiniuStorage.storeBase64Str(imgBaseStr, IMG_TYPE, keyName);
String url = qiniuStorage.generateUrl(keyName);
String uploadKeyName = qiniuStorage.storeBase64Str(imgBaseStr, keyName);
String url = qiniuStorage.generateUrl(uploadKeyName);
paramMap.put("topImage", url);
item.setTopImage(url);
}
......
package com.wwdz.ch.admin.impl;
import com.wwdz.ch.admin.service.UploadService;
import com.wwdz.ch.core.consts.CommConsts;
import com.wwdz.ch.core.storage.QiniuStorage;
import com.wwdz.ch.core.type.Result;
import com.xxdxxs.utils.StringUtils;
......@@ -13,8 +14,6 @@ public class UploadServiceImpl implements UploadService {
private final static String IMG_TYPE = "image/png";
private final static String IMG_SUFFIX = ".png";
@Autowired
QiniuStorage qiniuStorage;
......@@ -25,9 +24,9 @@ public class UploadServiceImpl implements UploadService {
if (StringUtils.isEmpty(base64Str) || ObjectUtils.isEmpty(id) ) {
return Result.failed("参数不能为空");
}
String keyName = String.valueOf(id) + System.currentTimeMillis() + IMG_SUFFIX;
qiniuStorage.storeBase64Str(base64Str, IMG_TYPE, keyName);
String url = qiniuStorage.generateUrl(keyName);
String keyName = CommConsts.UPLOAD_PRE_DIRECTORY + id + System.currentTimeMillis();
String uploadKeyName = qiniuStorage.storeBase64Str(base64Str, keyName);
String url = qiniuStorage.generateUrl(uploadKeyName);
return Result.success(url);
} catch (Exception e) {
e.printStackTrace();
......@@ -35,6 +34,4 @@ public class UploadServiceImpl implements UploadService {
return Result.failed();
}
}
......@@ -33,7 +33,57 @@ public class SynCoinJob {
CoinRequestDto coinRequestDto = new CoinRequestDto();
coinRequestDto.setPage(1);
coinRequestDto.setLimit(100);
PageInfo<Coins> pageInfo = coinsDao.findByPage(coinRequestDto);
PageInfo<Item> pageInfo = itemDao.findByPage(coinRequestDto);
int totalPageNum = pageInfo.getPages();
logger.info("total page : {}, total : {}", totalPageNum, pageInfo.getTotal());
IntStream.rangeClosed(1, totalPageNum).parallel().forEach(i -> {
logger.info(">>>>>>>>>>>>>>>>> 当前同步数据的页码: {} 同步开始 <<<<<<<<<<<<<<<<<", i);
CoinRequestDto dto = new CoinRequestDto();
dto.setPage(i);
dto.setLimit(100);
List<Coins> coinsList = coinsDao.findList(dto);
coinsList.parallelStream().forEach(coins -> {
Item item = new Item();
item.setId(coins.getId());
item.setName(coins.getName());
item.setCid(coins.getCid());
item.setIsOnSale(false);
item.setIsStd(false);
item.setSort(1);
item.setImages(coins.getImages());
item.setShareImage(null);
item.setTopImage(null);
item.setPrice(coins.getPrice());
item.setDetail(coins.getDetail());
item.setCreateTime(coins.getCreated());
item.setSourceType(CommonEnum.SourceTypeEnum.SYSTEM_GET.getCode());
item.setUpdateTime(coins.getUpdated());
item.setIsDeleted(false);
item.setBidTime("");
item.setRemark(null);
item.setIsSetted(0);
item.setOriginalExtra(coins.getExtra());
item.setOriginalSourceUrl(coins.getSourceUrl());
item.setSource(coins.getSource());
itemDao.insert(item);
});
logger.info(">>>>>>>>>>>>>>>>> 同步数据的页码: {}, 同步成功 <<<<<<<<<<<<<<<<<", i);
});
logger.info(">>>>>>>>>>>>>>>>>>>>>>> 同步coins表数据到item表, 执行完毕 <<<<<<<<<<<<<<<<<<<<<");
} catch (Exception e) {
logger.error("同步coins表数据到item表 error {}", e);
}
}
public void updateImageUrlOfsSize() {
logger.info(">>>>>>>>>>>>>>>>>>>>>>> 更新item表, 开始执行 <<<<<<<<<<<<<<<<<<<<<");
try {
CoinRequestDto coinRequestDto = new CoinRequestDto();
coinRequestDto.setPage(1);
coinRequestDto.setLimit(100);
PageInfo<Item> pageInfo = itemDao.findByPage(coinRequestDto);
int totalPageNum = pageInfo.getPages();
logger.info("total page : {}, total : {}", totalPageNum, pageInfo.getTotal());
IntStream.rangeClosed(1, totalPageNum).parallel().forEach(i -> {
......
......@@ -124,11 +124,7 @@
<version>2.0</version>
</dependency>
<dependency>
<groupId>co.elastic.clients</groupId>
<artifactId>elasticsearch-java</artifactId>
<version>8.8.2</version>
</dependency>
</dependencies>
</project>
......@@ -16,4 +16,10 @@ public interface CommConsts {
public static final String MISS_PARAMS = "缺少必要参数";
public static final Long CACHE_EXPIRE_MINUTES = 60 * 6L; // 设置缓存 6小时,有效期单位 : 分钟
/**
* 上传文件目录
*/
public static final String UPLOAD_PRE_DIRECTORY = "quanku/";
}
package com.wwdz.ch.core.storage;
import java.awt.image.BufferedImage;
import java.io.ByteArrayInputStream;
import java.io.InputStream;
import java.net.MalformedURLException;
......@@ -8,6 +9,7 @@ import java.nio.file.Path;
import java.util.Base64;
import java.util.stream.Stream;
import com.xxdxxs.utils.JsonUtils;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.core.io.Resource;
......@@ -20,10 +22,14 @@ import com.qiniu.storage.Configuration;
import com.qiniu.storage.UploadManager;
import com.qiniu.util.Auth;
import javax.imageio.ImageIO;
public class QiniuStorage implements Storage {
private final Log logger = LogFactory.getLog(QiniuStorage.class);
private final static String BASE64_IMG_SUFFIX = ".png";
private String endpoint;
private String accessKey;
private String secretKey;
......@@ -79,6 +85,8 @@ public class QiniuStorage implements Storage {
try {
String upToken = auth.uploadToken(bucketName);
Response response = uploadManager.put(inputStream, keyName, upToken, null, contentType);
logger.info("七牛存储结果 response:{}" + response);
logger.info("七牛存储结果 response:{}" + JsonUtils.from(response));
logger.info("七牛存储结果:" + response.statusCode);
} catch (QiniuException ex) {
ex.printStackTrace();
......@@ -86,7 +94,27 @@ public class QiniuStorage implements Storage {
}
public void storeBase64Str(String baseStr, String contentType, String keyName) throws Exception {
public void store(byte[] data, String keyName) {
if (uploadManager == null) {
if (auth == null) {
auth = Auth.create(accessKey, secretKey);
}
uploadManager = new UploadManager(new Configuration());
}
try {
String upToken = auth.uploadToken(bucketName);
Response response = uploadManager.put(data, keyName, upToken);
logger.info("七牛存储结果 response:{}" + response);
logger.info("七牛存储结果 response:{}" + JsonUtils.from(response));
logger.info("七牛存储结果:" + response.statusCode);
} catch (QiniuException ex) {
ex.printStackTrace();
}
}
/*public void storeBase64Str(String baseStr, String contentType, String keyName) throws Exception {
if (uploadManager == null) {
if (auth == null) {
auth = Auth.create(accessKey, secretKey);
......@@ -95,8 +123,31 @@ public class QiniuStorage implements Storage {
}
Base64.Decoder decoder = Base64.getDecoder();
byte[] imageByte = decoder.decode(baseStr.split(",")[1]);
InputStream inputStream = new ByteArrayInputStream(imageByte);
BufferedImage image = ImageIO.read(inputStream);
keyName = keyName + "_" + image.getWidth() + "x" + image.getHeight()+ ".jpg";
store(inputStream, 1, contentType, keyName);
}*/
public String storeBase64Str(String baseStr, String keyName) throws Exception {
if (uploadManager == null) {
if (auth == null) {
auth = Auth.create(accessKey, secretKey);
}
uploadManager = new UploadManager(new Configuration());
}
Base64.Decoder decoder = Base64.getDecoder();
byte[] imageByte = decoder.decode(baseStr.split(",")[1]);
InputStream inputStream = new ByteArrayInputStream(imageByte);
BufferedImage image = ImageIO.read(inputStream);
keyName = keyName + "_" + image.getWidth() + "x" + image.getHeight() + BASE64_IMG_SUFFIX;
store(imageByte, keyName);
String url =generateUrl(keyName);
System.out.println("url = " + url);
return keyName;
}
......
package com.wwdz.ch.core.util;
import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.elasticsearch._types.FieldValue;
import co.elastic.clients.elasticsearch._types.SortOrder;
import co.elastic.clients.elasticsearch._types.query_dsl.MatchQuery;
import co.elastic.clients.elasticsearch._types.query_dsl.Query;
import co.elastic.clients.elasticsearch._types.query_dsl.TermQuery;
import co.elastic.clients.elasticsearch._types.query_dsl.*;
import co.elastic.clients.elasticsearch.core.GetResponse;
import co.elastic.clients.elasticsearch.core.SearchRequest;
import co.elastic.clients.elasticsearch.core.SearchResponse;
......@@ -20,6 +19,7 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.function.Function;
......@@ -29,7 +29,7 @@ public class ElasticsearchUtil {
private static final Logger logger = LoggerFactory.getLogger(ElasticsearchUtil.class);
private final static String INDEX = "quanku";
private final static String INDEX = "quanku_new";
@Resource(name="clientByPasswd")
ElasticsearchClient elasticsearchClient;
......@@ -40,40 +40,33 @@ public class ElasticsearchUtil {
int size = esSearchParam.getSize();
int from = (page - 1) * size;
Query termQuery = TermQuery.of(m -> m
.field("name")
.value(FieldValue.of(esSearchParam.getContent())))._toQuery();
Query wildcardQuery = WildcardQuery.of(w->w.field("name.keyword").value("中国*"))._toQuery();
Query matchQuery = MatchQuery.of(m->m
.field("name").query(esSearchParam.getContent())
)._toQuery();
SearchRequest.Builder builder = new SearchRequest.Builder();
builder.index(INDEX);
builder.query(queryBuild ->
queryBuild.term(termQueryBuilder ->
termQueryBuilder.field("name").value(esSearchParam.getContent())));
builder.source(sourceBuild->sourceBuild.filter(f->f.includes("coin_id", "name", "price", "top_image_path")));
builder.query(matchQuery);
builder.from(from);
builder.size(size);
builder.sort(f -> f.field(o -> o.field("coin_id")
.order(SortOrder.Desc)));
Query nameQuery = MatchQuery.of(m -> m
.field("name")
.query(esSearchParam.getContent())
)._toQuery();
// SearchResponse<ItemOfEs> searchResponse = elasticsearchClient.search((Function<SearchRequest.Builder, ObjectBuilder<SearchRequest>>) builder, ItemOfEs.class);
// SearchRequest searchRequest = (SearchRequest)((ObjectBuilder)fn.apply(new co.elastic.clients.elasticsearch.core.SearchRequest.Builder())).build();
SearchResponse<ItemOfEs> searchResponse = elasticsearchClient.search(s -> s
.index(INDEX)
.source(sourceBuild->sourceBuild.filter(f->f.includes("coin_id", "name", "price", "top_image_path")))
.query(nameQuery)
.from(page).size(size)
.sort(f -> f.field(o -> o.field("coin_id").order(SortOrder.Desc))),
ItemOfEs.class);
logger.info("getResponse:{}", searchResponse);
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);
List<Hit<ItemOfEs>> hitList = searchResponse.hits().hits();
logger.info("====== getResponse : {}", searchResponse);
for (Hit<ItemOfEs> itemOfEsHit : hitList) {
ItemOfEs itemOfEs = itemOfEsHit.source();
logger.info("id : {}, item : {}", itemOfEs.getCoin_id(), itemOfEs);
logger.info("[data] id : {}, item : {}", itemOfEs.getCoin_id(), itemOfEs);
}
} catch (Exception e) {
e.printStackTrace();
......
......@@ -110,6 +110,12 @@
<version>1.2.12</version>
</dependency>
<dependency>
<groupId>co.elastic.clients</groupId>
<artifactId>elasticsearch-java</artifactId>
<version>8.8.2</version>
</dependency>
</dependencies>
......
package com.wwdz.ch.db.bean;
import com.xxdxxs.entity.Entity;
import lombok.Data;
import java.util.List;
@Data
public class EsPageInfo implements Entity {
/**
* 数据总条数
*/
private Long total;
private List<? extends Entity> dataList;
/**
* 当前页码
*/
private Integer currentPage;
private Integer size;
public EsPageInfo(Long total, List<? extends Entity> dataList, Integer currentPage, Integer size) {
this.total = total;
this.dataList = dataList;
this.currentPage = currentPage;
this.size = size;
}
public EsPageInfo(List<? extends Entity> dataList) {
this.dataList = dataList;
}
public static EsPageInfo of(List<? extends Entity> dataList) {
return new EsPageInfo(dataList);
}
public static EsPageInfo of(Long total, List<? extends Entity> dataList, Integer currentPage, Integer size) {
return new EsPageInfo(total, dataList, currentPage, size);
}
}
package com.wwdz.ch.db.bean;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.xxdxxs.entity.Entity;
import lombok.Data;
@Data
@JsonIgnoreProperties(ignoreUnknown = true)
public class ItemOfEs implements Entity {
/**
......
......@@ -71,4 +71,10 @@ public class CoinRequestDto extends BaseRequestDto implements Entity {
* 置顶图片
*/
private String topImage;
/**
* 是否排除部分数据
* 0:否;1是
*/
private Integer isExcluded;
}
package com.wwdz.ch.db.dto.request;
import com.xxdxxs.entity.Entity;
import lombok.Data;
import java.util.Map;
@Data
public class EsSearchRequestDto implements Entity {
private Map<String, String> paramMap;
private String content;
private String sortField;
private Boolean isDesc = true;
private int page = 1;
private int size = 20;
public EsSearchRequestDto() {
}
public EsSearchRequestDto(Map<String, String> paramMap, String content, String sortField, Boolean isDesc, int page, int size) {
this.paramMap = paramMap;
this.content = content;
this.sortField = sortField;
this.isDesc = isDesc;
this.page = page;
this.size = size;
}
public EsSearchRequestDto(String content, String sortField, Boolean isDesc, int page, int size) {
this.content = content;
this.sortField = sortField;
this.isDesc = isDesc;
this.page = page;
this.size = size;
}
public EsSearchRequestDto(String content) {
this.content = content;
}
public EsSearchRequestDto(String content, int page, int size) {
this.content = content;
this.page = page;
this.size = size;
}
public static EsSearchRequestDto of(String content, int page, int size) {
return new EsSearchRequestDto(content, page, size);
}
public static EsSearchRequestDto of(String content) {
return new EsSearchRequestDto(content);
}
}
package com.wwdz.ch.db.es;
import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.elasticsearch._types.SortOrder;
import co.elastic.clients.elasticsearch._types.query_dsl.MatchQuery;
import co.elastic.clients.elasticsearch._types.query_dsl.Query;
import co.elastic.clients.elasticsearch.core.SearchRequest;
import co.elastic.clients.elasticsearch.core.SearchResponse;
import co.elastic.clients.elasticsearch.core.search.Hit;
import com.wwdz.ch.db.bean.EsPageInfo;
import com.wwdz.ch.db.bean.ItemOfEs;
import com.wwdz.ch.db.dto.request.EsSearchRequestDto;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Repository;
import javax.annotation.Resource;
import java.util.ArrayList;
import java.util.List;
@Repository
public class ItemEsDao {
private static final Logger logger = LoggerFactory.getLogger(ItemEsDao.class);
private final static String INDEX = "quanku_new";
@Resource(name="clientByPasswd")
ElasticsearchClient elasticsearchClient;
public EsPageInfo findList(EsSearchRequestDto esSearchRequestDto) {
EsPageInfo esPageInfo = null;
try {
List<ItemOfEs> list = new ArrayList<>();
int page = esSearchRequestDto.getPage();
int size = esSearchRequestDto.getSize();
int from = (page - 1) * size;
Query matchQuery = MatchQuery.of(m -> m
.field("name")
.query(esSearchRequestDto.getContent())
)._toQuery();
SearchRequest.Builder builder = new SearchRequest.Builder();
builder.index(INDEX);
builder.source(sourceBuild->sourceBuild.filter(f->f.includes("coin_id", "name", "price", "top_image_path")));
builder.query(matchQuery);
builder.from(from);
builder.size(size);
builder.sort(f -> f.field(o -> o.field("coin_id")
.order(SortOrder.Desc)));
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);
List<Hit<ItemOfEs>> hitList = searchResponse.hits().hits();
logger.info("====== getResponse : {}", searchResponse);
for (Hit<ItemOfEs> itemOfEsHit : hitList) {
ItemOfEs itemOfEs = itemOfEsHit.source();
list.add(itemOfEs);
}
esPageInfo = EsPageInfo.of(total, list, page, size);
} catch (Exception e) {
logger.error("ES 查询商品列表 error : {}", e);
}
return esPageInfo;
}
}
......@@ -11,6 +11,7 @@ import com.xxdxxs.db.component.JdbcHelper;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Repository;
import org.springframework.util.CollectionUtils;
import org.springframework.util.ObjectUtils;
import org.springframework.util.StringUtils;
import java.util.Arrays;
......@@ -54,6 +55,11 @@ public class ItemDaoImpl implements ItemDao {
ItemExample example = new ItemExample();
ItemExample.Criteria criteria = example.createCriteria();
criteria.andIsDeletedEqualTo(false);
//暂时去除华夏古泉和天眷堂的钱币, 对应编码1和7
if (!ObjectUtils.isEmpty(coinRequestDto.getIsExcluded())) {
criteria.andSourceNotEqualTo(1);
criteria.andSourceNotEqualTo(7);
}
example.orderBy("update_time desc");
JdbcHelper.ifPresent(coinRequestDto.getId(), criteria :: andIdEqualTo);
JdbcHelper.ifPresent(coinRequestDto.getName(), criteria :: andNameLike);
......@@ -63,6 +69,7 @@ public class ItemDaoImpl implements ItemDao {
return new PageInfo<>(list);
}
@Override
public Item findDetails(Long id) {
ItemExample example = new ItemExample();
......
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