Skip to content
GitLab
Projects
Groups
Snippets
Help
Loading...
Help
Help
Support
Keyboard shortcuts
?
Submit feedback
Sign in
Toggle navigation
Q
qk_backend
Project overview
Project overview
Details
Activity
Releases
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Issues
0
Issues
0
List
Boards
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Analytics
Analytics
CI / CD
Repository
Value Stream
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
quanku
qk_backend
Commits
9615e793
Commit
9615e793
authored
Aug 18, 2023
by
shiyu
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
es
parent
7279c75e
Changes
3
Hide whitespace changes
Inline
Side-by-side
Showing
3 changed files
with
82 additions
and
18 deletions
+82
-18
ch-admin-api/src/main/java/com/wwdz/ch/admin/job/SynCoinJob.java
...n-api/src/main/java/com/wwdz/ch/admin/job/SynCoinJob.java
+1
-6
ch-dao/src/main/java/com/wwdz/ch/db/bean/ItemOfEs.java
ch-dao/src/main/java/com/wwdz/ch/db/bean/ItemOfEs.java
+3
-0
ch-dao/src/main/java/com/wwdz/ch/db/es/ItemEsDao.java
ch-dao/src/main/java/com/wwdz/ch/db/es/ItemEsDao.java
+78
-12
No files found.
ch-admin-api/src/main/java/com/wwdz/ch/admin/job/SynCoinJob.java
View file @
9615e793
...
@@ -184,7 +184,7 @@ public class SynCoinJob {
...
@@ -184,7 +184,7 @@ public class SynCoinJob {
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
());
IntStream
.
rangeClosed
(
1
,
totalPageNum
).
parallel
().
forEach
(
i
->
{
IntStream
.
rangeClosed
(
1
,
totalPageNum
).
forEach
(
i
->
{
logger
.
info
(
">>>>>>>>>>>>>>>>> 当前同步数据的页码: {} 更新开始 <<<<<<<<<<<<<<<<<"
,
i
);
logger
.
info
(
">>>>>>>>>>>>>>>>> 当前同步数据的页码: {} 更新开始 <<<<<<<<<<<<<<<<<"
,
i
);
CoinRequestDto
dto
=
new
CoinRequestDto
();
CoinRequestDto
dto
=
new
CoinRequestDto
();
dto
.
setPage
(
i
);
dto
.
setPage
(
i
);
...
@@ -192,11 +192,6 @@ public class SynCoinJob {
...
@@ -192,11 +192,6 @@ public class SynCoinJob {
List
<
Item
>
itemList
=
itemDao
.
findListByPage
(
dto
);
List
<
Item
>
itemList
=
itemDao
.
findListByPage
(
dto
);
itemEsDao
.
create
(
itemList
);
itemEsDao
.
create
(
itemList
);
logger
.
info
(
">>>>>>>>>>>>>>>>> 页码 = {}, 同步es 成功 ! <<<<<<<<<<<<<<<<<"
,
i
);
logger
.
info
(
">>>>>>>>>>>>>>>>> 页码 = {}, 同步es 成功 ! <<<<<<<<<<<<<<<<<"
,
i
);
try
{
Thread
.
sleep
(
100
);
}
catch
(
InterruptedException
e
)
{
throw
new
RuntimeException
(
e
);
}
});
});
logger
.
info
(
">>>>>>>>>>>>>>>>>数据同步es 完成 SUCCESS!!!!! <<<<<<<<<<<<<<<<<"
);
logger
.
info
(
">>>>>>>>>>>>>>>>>数据同步es 完成 SUCCESS!!!!! <<<<<<<<<<<<<<<<<"
);
}
catch
(
Exception
e
)
{
}
catch
(
Exception
e
)
{
...
...
ch-dao/src/main/java/com/wwdz/ch/db/bean/ItemOfEs.java
View file @
9615e793
...
@@ -8,6 +8,9 @@ import lombok.Data;
...
@@ -8,6 +8,9 @@ import lombok.Data;
@JsonIgnoreProperties
(
ignoreUnknown
=
true
)
@JsonIgnoreProperties
(
ignoreUnknown
=
true
)
public
class
ItemOfEs
implements
Entity
{
public
class
ItemOfEs
implements
Entity
{
private
Long
id
;
/**
/**
* 商品id
* 商品id
*/
*/
...
...
ch-dao/src/main/java/com/wwdz/ch/db/es/ItemEsDao.java
View file @
9615e793
package
com.wwdz.ch.db.es
;
package
com.wwdz.ch.db.es
;
import
co.elastic.clients.elasticsearch.ElasticsearchClient
;
import
co.elastic.clients.elasticsearch.ElasticsearchClient
;
import
co.elastic.clients.elasticsearch._types.SortOrder
;
import
co.elastic.clients.elasticsearch._types.*
;
import
co.elastic.clients.elasticsearch._types.query_dsl.MatchQuery
;
import
co.elastic.clients.elasticsearch._types.aggregations.AggregationBuilders
;
import
co.elastic.clients.elasticsearch._types.query_dsl.Query
;
import
co.elastic.clients.elasticsearch._types.analysis.Language
;
import
co.elastic.clients.elasticsearch._types.query_dsl.*
;
import
co.elastic.clients.elasticsearch.core.*
;
import
co.elastic.clients.elasticsearch.core.*
;
import
co.elastic.clients.elasticsearch.core.bulk.BulkOperation
;
import
co.elastic.clients.elasticsearch.core.bulk.BulkOperation
;
import
co.elastic.clients.elasticsearch.core.bulk.CreateOperation
;
import
co.elastic.clients.elasticsearch.core.bulk.CreateOperation
;
import
co.elastic.clients.elasticsearch.core.search.Hit
;
import
co.elastic.clients.elasticsearch.core.search.Hit
;
import
co.elastic.clients.elasticsearch.indices.CreateIndexResponse
;
import
co.elastic.clients.elasticsearch.indices.CreateIndexResponse
;
import
co.elastic.clients.elasticsearch.indices.DeleteIndexResponse
;
import
co.elastic.clients.elasticsearch.indices.DeleteIndexResponse
;
import
co.elastic.clients.json.JsonData
;
import
co.elastic.clients.util.ObjectBuilder
;
import
com.fasterxml.jackson.databind.util.ArrayBuilders
;
import
com.wwdz.ch.db.bean.EsPageInfo
;
import
com.wwdz.ch.db.bean.EsPageInfo
;
import
com.wwdz.ch.db.bean.ItemOfEs
;
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.StringUtils
;
import
org.slf4j.Logger
;
import
org.slf4j.Logger
;
import
org.slf4j.LoggerFactory
;
import
org.slf4j.LoggerFactory
;
import
org.springframework.stereotype.Repository
;
import
org.springframework.stereotype.Repository
;
import
javax.annotation.Resource
;
import
javax.annotation.Resource
;
import
java.io.ByteArrayInputStream
;
import
java.io.IOException
;
import
java.io.IOException
;
import
java.io.InputStream
;
import
java.io.SequenceInputStream
;
import
java.util.ArrayList
;
import
java.util.ArrayList
;
import
java.util.Arrays
;
import
java.util.List
;
import
java.util.List
;
import
java.util.Random
;
import
java.util.concurrent.TimeUnit
;
import
java.util.function.Function
;
import
java.util.stream.IntStream
;
import
java.util.stream.LongStream
;
@Repository
@Repository
public
class
ItemEsDao
{
public
class
ItemEsDao
{
...
@@ -29,6 +44,9 @@ public class ItemEsDao {
...
@@ -29,6 +44,9 @@ public class ItemEsDao {
private
final
static
String
INDEX
=
"quanku_temp"
;
private
final
static
String
INDEX
=
"quanku_temp"
;
//商品来源是中国,所对应的编码
private
final
static
List
<
String
>
CHINESE_SOURCE_LIST
=
Arrays
.
asList
(
"1"
,
"2"
,
"3"
,
"7"
,
"8"
,
"10"
,
"11"
,
"1000"
);
@Resource
(
name
=
"clientByPasswd"
)
@Resource
(
name
=
"clientByPasswd"
)
ElasticsearchClient
elasticsearchClient
;
ElasticsearchClient
elasticsearchClient
;
...
@@ -66,32 +84,79 @@ public class ItemEsDao {
...
@@ -66,32 +84,79 @@ public class ItemEsDao {
EsPageInfo
esPageInfo
=
null
;
EsPageInfo
esPageInfo
=
null
;
try
{
try
{
List
<
ItemOfEs
>
list
=
new
ArrayList
<>();
List
<
ItemOfEs
>
list
=
new
ArrayList
<>();
/* //查询总条数
CountResponse countResponse = elasticsearchClient.count();
long total = countResponse.count();
logger.info(">>>>>>>>>> 数据总量 :{} <<<<<<<<<<", total);*/
int
page
=
esSearchRequestDto
.
getPage
();
int
page
=
esSearchRequestDto
.
getPage
();
int
size
=
esSearchRequestDto
.
getSize
();
int
size
=
esSearchRequestDto
.
getSize
();
int
from
=
(
page
-
1
)
*
size
;
int
from
=
(
page
-
1
)
*
size
;
Query
matchQuery
=
MatchQuery
.
of
(
m
->
m
Query
sourceTermQuery
=
TermQuery
.
of
(
t
->
.
field
(
"name"
)
t
.
field
(
"source"
)
.
query
(
esSearchRequestDto
.
getContent
())
.
value
(
2
)).
_toQuery
();
).
_toQuery
();
Query
sourceTermQuery2
=
TermQuery
.
of
(
t
->
t
.
field
(
"source"
)
.
value
(
7
)).
_toQuery
();
SearchRequest
.
Builder
builder
=
new
SearchRequest
.
Builder
();
SearchRequest
.
Builder
builder
=
new
SearchRequest
.
Builder
();
builder
.
index
(
INDEX
);
builder
.
index
(
INDEX
);
builder
.
source
(
sourceBuild
->
sourceBuild
.
filter
(
f
->
f
.
includes
(
"coin_id"
,
"name"
,
"price"
,
"top_image_path"
)));
builder
.
source
(
sourceBuild
->
sourceBuild
.
filter
(
f
->
f
.
includes
(
"id"
,
"images"
,
"name"
,
"price"
)));
builder
.
query
(
matchQuery
);
List
<
Query
>
sourceTypeQueryList
=
new
ArrayList
<>();
builder
.
from
(
from
);
CHINESE_SOURCE_LIST
.
forEach
(
i
->
{
Query
query
=
TermQuery
.
of
(
t
->
t
.
field
(
"sourceType"
).
value
(
i
)).
_toQuery
();
sourceTypeQueryList
.
add
(
query
);
});
Query
shouldQuery
=
QueryBuilders
.
bool
(
q
->
q
.
should
(
sourceTypeQueryList
));
BoolQuery
.
Builder
build
=
QueryBuilders
.
bool
().
mustNot
(
sourceTermQuery
).
mustNot
(
sourceTermQuery2
).
must
(
shouldQuery
);
if
(
StringUtils
.
hasLength
(
esSearchRequestDto
.
getContent
()))
{
Query
matchQuery
=
MatchQuery
.
of
(
m
->
m
.
field
(
"name"
)
.
query
(
esSearchRequestDto
.
getContent
())
).
_toQuery
();
build
.
must
(
matchQuery
);
builder
.
from
(
from
);
}
else
{
int
num
=
10000
-
size
;
Random
random
=
new
Random
();
from
=
random
.
nextInt
(
num
);
builder
.
from
(
from
);
/* Script script = new Script.Builder(ScriptType.INLINE, "painless", "Math.random()", null).build();
ScriptSortBuilder scriptSortBuilder = new ScriptSortBuilder(script).type(ScriptSortType.NUMBER).order(SortOrder.DESC);
SearchRequest searchRequest = new SearchRequest().source(new SearchSourceBuilder()
.query(QueryBuilders.functionScoreQuery(QueryBuilders.matchAllQuery()).randomScore(script))
.sort(scriptSortBuilder)
.size(20));
builder.sort(ScriptSort.of(s -> s.type(ScriptSortType.Number).script(new Script("ss"))).order(SortOrder.Desc))));
logger.info(">>>>>>>>>>>>> script = {}", scriptSort.script().toString());*/
}
builder
.
query
(
build
.
build
().
_toQuery
());
builder
.
size
(
size
);
builder
.
size
(
size
);
builder
.
sort
(
f
->
f
.
field
(
o
->
o
.
field
(
"coin_id"
)
/*
builder.sort(f -> f.field(o -> o.field("coin_id")
.
order
(
SortOrder
.
Desc
)));
.order(SortOrder.Desc)));
*/
SearchRequest
searchRequest
=
builder
.
build
();
SearchRequest
searchRequest
=
builder
.
build
();
logger
.
info
(
"查询语句 : {}"
,
searchRequest
.
query
());
logger
.
info
(
"查询语句 : {}"
,
searchRequest
.
query
());
SearchResponse
<
ItemOfEs
>
searchResponse
=
elasticsearchClient
.
search
(
searchRequest
,
ItemOfEs
.
class
);
SearchResponse
<
ItemOfEs
>
searchResponse
=
elasticsearchClient
.
search
(
searchRequest
,
ItemOfEs
.
class
);
long
total
=
searchResponse
.
hits
().
total
().
value
();
long
total
=
searchResponse
.
hits
().
total
().
value
();
logger
.
info
(
"总数据量 : {}"
,
total
);
logger
.
info
(
"总数据量 : {}"
,
total
);
List
<
Hit
<
ItemOfEs
>>
hitList
=
searchResponse
.
hits
().
hits
();
List
<
Hit
<
ItemOfEs
>>
hitList
=
searchResponse
.
hits
().
hits
();
logger
.
info
(
"====== getResponse : {}"
,
searchResponse
);
logger
.
info
(
"====== getResponse : {}"
,
searchResponse
);
List
<
Long
>
idList
=
new
ArrayList
<>();
for
(
Hit
<
ItemOfEs
>
itemOfEsHit
:
hitList
)
{
for
(
Hit
<
ItemOfEs
>
itemOfEsHit
:
hitList
)
{
ItemOfEs
itemOfEs
=
itemOfEsHit
.
source
();
ItemOfEs
itemOfEs
=
itemOfEsHit
.
source
();
list
.
add
(
itemOfEs
);
list
.
add
(
itemOfEs
);
idList
.
add
(
itemOfEs
.
getId
());
}
}
logger
.
info
(
"======== 随机返回的数据id : {}"
,
idList
);
esPageInfo
=
EsPageInfo
.
of
(
total
,
list
,
page
,
size
);
esPageInfo
=
EsPageInfo
.
of
(
total
,
list
,
page
,
size
);
}
catch
(
Exception
e
)
{
}
catch
(
Exception
e
)
{
logger
.
error
(
"ES 查询商品列表 error : {}"
,
e
);
logger
.
error
(
"ES 查询商品列表 error : {}"
,
e
);
...
@@ -128,4 +193,5 @@ public class ItemEsDao {
...
@@ -128,4 +193,5 @@ public class ItemEsDao {
}
}
}
}
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment