|
|
package com.yoho.search.recall.scene.component;
|
|
|
|
|
|
|
|
|
import com.alibaba.fastjson.JSON;
|
|
|
import com.alibaba.fastjson.JSONObject;
|
|
|
import com.yoho.core.redis.cluster.operations.serializer.RedisKeyBuilder;
|
|
|
import com.yoho.search.base.utils.ISearchConstants;
|
|
|
import com.yoho.search.base.utils.ProductIndexEsField;
|
|
|
import com.yoho.search.base.utils.Transfer;
|
|
|
import com.yoho.search.common.cache.impls.SearchRedis;
|
|
|
import com.yoho.search.core.es.model.SearchParam;
|
|
|
import com.yoho.search.core.es.model.SearchResult;
|
|
|
import com.yoho.search.recall.scene.models.CacheRequestResponse;
|
|
|
import com.yoho.search.service.base.SearchCommonService;
|
|
|
import com.yoho.search.service.base.index.ProductIndexBaseService;
|
|
|
import org.apache.commons.collections.MapUtils;
|
|
|
import org.apache.commons.lang.StringUtils;
|
|
|
import org.elasticsearch.index.query.QueryBuilders;
|
|
|
import org.slf4j.Logger;
|
|
|
import org.slf4j.LoggerFactory;
|
|
|
import org.springframework.beans.factory.annotation.Autowired;
|
|
|
import org.springframework.stereotype.Component;
|
|
|
|
|
|
import java.util.*;
|
|
|
|
|
|
@Component
|
|
|
public class BacthQueryBySknComponent {
|
|
|
|
|
|
private static final Logger logger = LoggerFactory.getLogger(BacthQueryBySknComponent.class);
|
|
|
|
|
|
@Autowired
|
|
|
private SearchRedis searchRedis;
|
|
|
@Autowired
|
|
|
private SearchCommonService searchCommonService;
|
|
|
@Autowired
|
|
|
private ProductIndexBaseService productIndexBaseService;
|
|
|
|
|
|
/**
|
|
|
* 按skn查询并按顺序返回
|
|
|
*
|
|
|
* @param productSkns
|
|
|
* @return
|
|
|
*/
|
|
|
public List<Map<String, Object>> queryProductListBySkn(List<Integer> productSkns,int size){
|
|
|
//1、构建请求
|
|
|
List<CacheRequestResponse<Integer, Map<String, Object>>> sknInfoRequestResponses = new ArrayList<>();
|
|
|
for (Integer productSkn : productSkns) {
|
|
|
sknInfoRequestResponses.add(new CacheRequestResponse(productSkn,redisKeyBuilderTransfer(),toValueTransfer(),fromValueTransfer()));
|
|
|
}
|
|
|
|
|
|
//2、批量从缓存中获取
|
|
|
sknInfoRequestResponses = this.batchQueryFromCache(sknInfoRequestResponses);
|
|
|
|
|
|
//3、获取未命中缓存的skn
|
|
|
List<Integer> notCacheSkns = new ArrayList<>();
|
|
|
for (CacheRequestResponse<Integer, Map<String, Object>> cacheRequestResponse : sknInfoRequestResponses) {
|
|
|
if(cacheRequestResponse.getValueObject()==null){
|
|
|
notCacheSkns.add(cacheRequestResponse.getKeyObject());
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
//4、执行批量查询
|
|
|
Map<Integer,Map<String, Object>> queryResults = this.batchQuery(notCacheSkns);
|
|
|
|
|
|
//5、填充查询结果
|
|
|
for (CacheRequestResponse<Integer, Map<String, Object>> cacheRequestResponse : sknInfoRequestResponses) {
|
|
|
if(cacheRequestResponse.getValueObject()==null){
|
|
|
cacheRequestResponse.setValueObject(queryResults.get(cacheRequestResponse.getKeyObject()));
|
|
|
cacheRequestResponse.setNeedRecache(true);
|
|
|
}
|
|
|
}
|
|
|
|
|
|
//6、将CacheRequestResponse中需要缓存的key加入缓存
|
|
|
this.batchAddToCache(sknInfoRequestResponses,10*60);//缓存10分钟
|
|
|
|
|
|
//7、构造返回结果
|
|
|
Map<Integer,Map<String, Object>> productInfoMap = new HashMap<>();
|
|
|
for (CacheRequestResponse<Integer, Map<String, Object>> sknInfoRequestResponse: sknInfoRequestResponses) {
|
|
|
productInfoMap.put(sknInfoRequestResponse.getKeyObject(),sknInfoRequestResponse.getValueObject());
|
|
|
}
|
|
|
List<Map<String, Object>> finalResults = new ArrayList<>();
|
|
|
for (Integer productSkn : productSkns) {
|
|
|
if(productInfoMap.get(productSkn)!=null){
|
|
|
finalResults.add(productInfoMap.get(productSkn));
|
|
|
}
|
|
|
if(finalResults.size()>=size){
|
|
|
break;
|
|
|
}
|
|
|
}
|
|
|
return finalResults;
|
|
|
}
|
|
|
|
|
|
|
|
|
private <K, V> List<CacheRequestResponse<K, V>> batchQueryFromCache(List<CacheRequestResponse<K, V>> cacheRequestResponses) {
|
|
|
try {
|
|
|
Collection<RedisKeyBuilder> keys = new ArrayList<>();
|
|
|
for (CacheRequestResponse<K, V> requestResponse : cacheRequestResponses) {
|
|
|
RedisKeyBuilder redisKeyBuilder = requestResponse.getRedisKeyBuilderTransfer().transfer(requestResponse.getKeyObject());
|
|
|
keys.add(redisKeyBuilder);
|
|
|
}
|
|
|
List<String> cachedValues = searchRedis.searchValueOperations.multiGet(keys);
|
|
|
for (int i = 0; i < cacheRequestResponses.size(); i++) {
|
|
|
CacheRequestResponse<K, V> requestResponse = cacheRequestResponses.get(i);
|
|
|
String cachedValue = cachedValues.get(i);
|
|
|
if (!StringUtils.isBlank(cachedValue)) {
|
|
|
requestResponse.setValueObject(requestResponse.getToValueTransfer().transfer(cachedValue));
|
|
|
requestResponse.setNeedRecache(false);
|
|
|
}
|
|
|
}
|
|
|
return cacheRequestResponses;
|
|
|
}catch (Exception e){
|
|
|
logger.error(e.getMessage(),e);
|
|
|
return cacheRequestResponses;
|
|
|
}
|
|
|
}
|
|
|
|
|
|
private <K, V> void batchAddToCache(List<CacheRequestResponse<K, V>> cacheRequestResponses,long timeOutInSecond){
|
|
|
try {
|
|
|
Map<RedisKeyBuilder, String> toCacheMap = new HashMap<>();
|
|
|
for (CacheRequestResponse<K, V> requestResponse : cacheRequestResponses) {
|
|
|
if(requestResponse.isNeedRecache() && requestResponse.getValueObject()!=null){
|
|
|
RedisKeyBuilder redisKeyBuilder = requestResponse.getRedisKeyBuilderTransfer().transfer(requestResponse.getKeyObject());
|
|
|
V value = requestResponse.getValueObject();
|
|
|
String cacheValue = requestResponse.getFromValueTransfer().transfer(value);
|
|
|
toCacheMap.put(redisKeyBuilder,cacheValue);
|
|
|
}
|
|
|
}
|
|
|
if(!toCacheMap.isEmpty()){
|
|
|
searchRedis.searchRedisTemplate.mset(toCacheMap,timeOutInSecond);
|
|
|
}
|
|
|
}catch (Exception e){
|
|
|
logger.error(e.getMessage(),e);
|
|
|
}
|
|
|
}
|
|
|
|
|
|
|
|
|
private Transfer<Integer, RedisKeyBuilder> redisKeyBuilderTransfer() {
|
|
|
return new Transfer<Integer, RedisKeyBuilder>() {
|
|
|
@Override
|
|
|
public RedisKeyBuilder transfer(Integer productSkn) {
|
|
|
return RedisKeyBuilder.newInstance().appendFixed("YOHOSEARCH:").appendFixed("SKN:").appendVar(productSkn);
|
|
|
}
|
|
|
};
|
|
|
}
|
|
|
|
|
|
private Transfer<String, Map<String, Object>> toValueTransfer() {
|
|
|
return new Transfer<String, Map<String, Object>>() {
|
|
|
@Override
|
|
|
public Map<String, Object> transfer(String value) {
|
|
|
Map<String, Object> product = new HashMap<>();
|
|
|
product.putAll(JSONObject.parseObject(value));
|
|
|
return product;
|
|
|
}
|
|
|
};
|
|
|
}
|
|
|
|
|
|
private Transfer<Map<String, Object>, String> fromValueTransfer() {
|
|
|
return new Transfer<Map<String, Object>, String>() {
|
|
|
@Override
|
|
|
public String transfer( Map<String, Object> product) {
|
|
|
return JSON.toJSONString(product);
|
|
|
}
|
|
|
};
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
private Map<Integer,Map<String, Object>> batchQuery(List<Integer> productSkns) {
|
|
|
SearchParam searchParam = new SearchParam();
|
|
|
searchParam.setOffset(0);
|
|
|
searchParam.setSize(productSkns.size());
|
|
|
searchParam.setFiter(QueryBuilders.termsQuery(ProductIndexEsField.productSkn, productSkns));
|
|
|
searchParam.setIncludeFields(productIndexBaseService.getProductIndexIncludeFields());
|
|
|
SearchResult searchResult = searchCommonService.doSearch(ISearchConstants.INDEX_NAME_PRODUCT_INDEX, searchParam);
|
|
|
List<Map<String, Object>> productList = productIndexBaseService.getProductListWithPricePlan(searchResult.getResultList());
|
|
|
Map<Integer,Map<String, Object>> results = new HashMap<>();
|
|
|
for (Map<String, Object> product: productList) {
|
|
|
results.put(MapUtils.getIntValue(product,"product_skn",0),product);
|
|
|
}
|
|
|
return results;
|
|
|
}
|
|
|
|
|
|
} |
...
|
...
|
|