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
5c29dae1
Commit
5c29dae1
authored
Aug 23, 2023
by
shiyu
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
jackson版本升级至2.11.2,springboot升级至2.5.3,为了支持es8.9 date类型排序
parent
ec5b3b5d
Changes
9
Show whitespace changes
Inline
Side-by-side
Showing
9 changed files
with
275 additions
and
14 deletions
+275
-14
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
+33
-1
ch-admin-api/src/main/java/com/wwdz/ch/admin/job/SynItemToEsJob.java
...i/src/main/java/com/wwdz/ch/admin/job/SynItemToEsJob.java
+117
-0
ch-admin-api/src/main/resources/application.yml
ch-admin-api/src/main/resources/application.yml
+2
-2
ch-core/pom.xml
ch-core/pom.xml
+10
-0
ch-dao/pom.xml
ch-dao/pom.xml
+48
-1
ch-dao/src/main/java/com/wwdz/ch/db/dto/request/CoinRequestDto.java
.../main/java/com/wwdz/ch/db/dto/request/CoinRequestDto.java
+6
-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
+44
-8
ch-dao/src/main/java/com/wwdz/ch/db/impl/ItemDaoImpl.java
ch-dao/src/main/java/com/wwdz/ch/db/impl/ItemDaoImpl.java
+1
-0
pom.xml
pom.xml
+14
-2
No files found.
ch-admin-api/src/main/java/com/wwdz/ch/admin/job/SynCoinJob.java
View file @
5c29dae1
...
...
@@ -5,6 +5,7 @@ import com.wwdz.ch.core.consts.CommConsts;
import
com.wwdz.ch.core.consts.CommonEnum
;
import
com.wwdz.ch.core.consts.MediaTypeEnum
;
import
com.wwdz.ch.core.storage.QiniuStorage
;
import
com.wwdz.ch.core.util.RedisUtils
;
import
com.wwdz.ch.db.bean.ItemOfEs
;
import
com.wwdz.ch.db.dao.CoinsDao
;
import
com.wwdz.ch.db.dao.ItemDao
;
...
...
@@ -14,11 +15,16 @@ import com.wwdz.ch.db.dto.request.CoinRequestDto;
import
com.wwdz.ch.db.es.ItemEsDao
;
import
com.xxdxxs.utils.CommonUtils
;
import
com.xxdxxs.utils.StringUtils
;
import
org.redisson.api.RLock
;
import
org.redisson.api.RedissonClient
;
import
org.slf4j.Logger
;
import
org.slf4j.LoggerFactory
;
import
org.springframework.beans.factory.InitializingBean
;
import
org.springframework.beans.factory.annotation.Autowired
;
import
org.springframework.scheduling.annotation.Scheduled
;
import
org.springframework.stereotype.Component
;
import
javax.annotation.PostConstruct
;
import
javax.imageio.ImageIO
;
import
java.awt.image.BufferedImage
;
import
java.io.ByteArrayInputStream
;
...
...
@@ -39,6 +45,8 @@ public class SynCoinJob {
private
static
final
Logger
logger
=
LoggerFactory
.
getLogger
(
SynCoinJob
.
class
);
private
String
QUANKU_SYN_TO_ES_TASK_KEY
=
"QUANKU_SYN_TO_ES_TASK_KEY"
;
@Autowired
ItemDao
itemDao
;
...
...
@@ -51,6 +59,12 @@ public class SynCoinJob {
@Autowired
ItemEsDao
itemEsDao
;
@Autowired
RedisUtils
redisUtils
;
@Autowired
RedissonClient
redissonClient
;
public
void
execute
()
{
logger
.
info
(
">>>>>>>>>>>>>>>>>>>>>>> 同步coins表数据到item表, 开始执行 <<<<<<<<<<<<<<<<<<<<<"
);
...
...
@@ -175,7 +189,6 @@ public class SynCoinJob {
}
public
void
synItemToEs
()
{
try
{
ExecutorService
executorService
=
Executors
.
newFixedThreadPool
(
10
);
...
...
@@ -201,4 +214,23 @@ public class SynCoinJob {
logger
.
error
(
"同步数据到es error {}"
,
e
);
}
}
public
void
syn
(){
RLock
lock
=
redissonClient
.
getLock
(
QUANKU_SYN_TO_ES_TASK_KEY
);
if
(!
lock
.
tryLock
())
{
logger
.
warn
(
"当前服务实例获取锁成功: {} 获取锁失败,锁被占用, 不允许执行同步订单作业"
,
Thread
.
currentThread
().
getId
());
return
;
}
try
{
synItemToEs
();
}
catch
(
Exception
e
)
{
logger
.
error
(
"同步数据到es error {}"
,
e
);
}
finally
{
if
(
lock
!=
null
&&
lock
.
isHeldByCurrentThread
())
{
lock
.
unlock
();
logger
.
info
(
"======================== 线程id: {} , 释放锁成功 ========================"
,
Thread
.
currentThread
().
getId
());
}
}
}
}
ch-admin-api/src/main/java/com/wwdz/ch/admin/job/SynItemToEsJob.java
0 → 100644
View file @
5c29dae1
package
com.wwdz.ch.admin.job
;
import
com.github.pagehelper.PageInfo
;
import
com.wwdz.ch.core.consts.CommConsts
;
import
com.wwdz.ch.core.consts.CommonEnum
;
import
com.wwdz.ch.core.consts.MediaTypeEnum
;
import
com.wwdz.ch.core.storage.QiniuStorage
;
import
com.wwdz.ch.core.util.RedisUtils
;
import
com.wwdz.ch.db.dao.CoinsDao
;
import
com.wwdz.ch.db.dao.ItemDao
;
import
com.wwdz.ch.db.domain.Coins
;
import
com.wwdz.ch.db.domain.Item
;
import
com.wwdz.ch.db.dto.request.CoinRequestDto
;
import
com.wwdz.ch.db.es.ItemEsDao
;
import
com.xxdxxs.utils.CommonUtils
;
import
com.xxdxxs.utils.StringUtils
;
import
org.redisson.api.RLock
;
import
org.redisson.api.RedissonClient
;
import
org.slf4j.Logger
;
import
org.slf4j.LoggerFactory
;
import
org.springframework.beans.factory.InitializingBean
;
import
org.springframework.beans.factory.annotation.Autowired
;
import
org.springframework.scheduling.annotation.Scheduled
;
import
org.springframework.stereotype.Component
;
import
javax.annotation.PostConstruct
;
import
javax.imageio.ImageIO
;
import
java.awt.image.BufferedImage
;
import
java.io.ByteArrayInputStream
;
import
java.io.ByteArrayOutputStream
;
import
java.io.IOException
;
import
java.io.InputStream
;
import
java.net.URL
;
import
java.util.Arrays
;
import
java.util.Date
;
import
java.util.List
;
import
java.util.concurrent.ExecutorService
;
import
java.util.concurrent.Executors
;
import
java.util.stream.IntStream
;
@Component
public
class
SynItemToEsJob
{
private
static
final
Logger
logger
=
LoggerFactory
.
getLogger
(
SynItemToEsJob
.
class
);
private
String
QUANKU_SYN_TO_ES_TASK_KEY
=
"QUANKU_SYN_TO_ES_TASK_KEY"
;
ExecutorService
executorService
=
Executors
.
newFixedThreadPool
(
10
);
@Autowired
ItemDao
itemDao
;
@Autowired
ItemEsDao
itemEsDao
;
@Autowired
RedissonClient
redissonClient
;
/**
* 隔一小时运行一次
*/
@Scheduled
(
fixedDelay
=
1000
*
60
*
60
)
public
void
execute
()
{
RLock
lock
=
redissonClient
.
getLock
(
QUANKU_SYN_TO_ES_TASK_KEY
);
if
(!
lock
.
tryLock
())
{
logger
.
warn
(
"当前服务实例获取锁成功: {} 获取锁失败,锁被占用, 不允许执行同步ES作业"
,
Thread
.
currentThread
().
getId
());
return
;
}
logger
.
info
(
">>>>>>>>>>>>>>>>>>>>>>> 同步item表数据到ES, 开始执行 <<<<<<<<<<<<<<<<<<<<<"
);
try
{
//查询es已有数据的最新修改时间
Date
time
=
itemEsDao
.
getMaxTime
();
CoinRequestDto
coinRequestDto
=
new
CoinRequestDto
();
coinRequestDto
.
setPage
(
1
);
coinRequestDto
.
setLimit
(
200
);
coinRequestDto
.
setStartUpdateTime
(
time
);
List
<
Item
>
itemList
=
itemDao
.
findListByPage
(
coinRequestDto
);
PageInfo
<
Item
>
pageInfo
=
new
PageInfo
<>(
itemList
);
logger
.
info
(
"total : {}, 当前查询 {} 条数据"
,
pageInfo
.
getTotal
(),
itemList
.
size
());
itemEsDao
.
create
(
itemList
);
logger
.
info
(
">>>>>>>>>>>>>>>>>>>>>>> 同步ES成功, 执行完毕 <<<<<<<<<<<<<<<<<<<<<"
);
}
catch
(
Exception
e
)
{
logger
.
error
(
"同步coins表数据到item表 error {}"
,
e
);
}
finally
{
if
(
lock
!=
null
&&
lock
.
isHeldByCurrentThread
())
{
lock
.
unlock
();
logger
.
info
(
"======================== 线程id: {} , 同步ES作业 释放锁成功 ========================"
,
Thread
.
currentThread
().
getId
());
}
}
}
public
void
synItemToEs
()
{
try
{
CoinRequestDto
coinRequestDto
=
new
CoinRequestDto
();
coinRequestDto
.
setPage
(
1
);
coinRequestDto
.
setLimit
(
100
);
PageInfo
<
Item
>
pageInfo
=
itemDao
.
findPageExcludeFilter
(
coinRequestDto
);
int
totalPageNum
=
pageInfo
.
getPages
();
logger
.
info
(
"total page : {}, total : {}"
,
totalPageNum
,
pageInfo
.
getTotal
());
IntStream
.
rangeClosed
(
1
,
totalPageNum
).
forEach
(
i
->
{
logger
.
info
(
">>>>>>>>>>>>>>>>> 当前同步数据的页码: {} 更新开始 <<<<<<<<<<<<<<<<<"
,
i
);
CoinRequestDto
dto
=
new
CoinRequestDto
();
dto
.
setPage
(
i
);
dto
.
setLimit
(
100
);
List
<
Item
>
itemList
=
itemDao
.
findListByPage
(
dto
);
executorService
.
submit
(()
->
{
itemEsDao
.
create
(
itemList
);
});
logger
.
info
(
">>>>>>>>>>>>>>>>> 页码 = {}, 同步es 成功 ! <<<<<<<<<<<<<<<<<"
,
i
);
});
logger
.
info
(
">>>>>>>>>>>>>>>>>数据同步es 完成 SUCCESS!!!!! <<<<<<<<<<<<<<<<<"
);
}
catch
(
Exception
e
)
{
logger
.
error
(
"同步数据到es error {}"
,
e
);
}
}
}
ch-admin-api/src/main/resources/application.yml
View file @
5c29dae1
...
...
@@ -5,9 +5,9 @@ spring:
encoding
:
UTF-8
servlet
:
multipart
:
max-file-size
:
50Mb
max-request-size
:
80Mb
max-request-size
:
80MB
location
:
/tmp
max-file-size
:
50MB
jackson
:
date-format
:
yyyy-MM-dd HH:mm:ss
...
...
ch-core/pom.xml
View file @
5c29dae1
...
...
@@ -92,6 +92,10 @@
<artifactId>
druid
</artifactId>
<groupId>
com.alibaba
</groupId>
</exclusion>
<exclusion>
<artifactId>
jakarta.json-api
</artifactId>
<groupId>
jakarta.json
</groupId>
</exclusion>
</exclusions>
</dependency>
...
...
@@ -118,6 +122,12 @@
<artifactId>
spring-boot-starter-data-redis
</artifactId>
</dependency>
<dependency>
<groupId>
jakarta.json
</groupId>
<artifactId>
jakarta.json-api
</artifactId>
<version>
2.0.1
</version>
</dependency>
<dependency>
<groupId>
org.springframework.data
</groupId>
<artifactId>
spring-data-redis
</artifactId>
...
...
ch-dao/pom.xml
View file @
5c29dae1
...
...
@@ -59,7 +59,14 @@
<dependency>
<groupId>
com.fasterxml.jackson.core
</groupId>
<artifactId>
jackson-databind
</artifactId>
<version>
2.11.2
</version>
</dependency>
<dependency>
<groupId>
com.fasterxml.jackson.core
</groupId>
<artifactId>
jackson-core
</artifactId>
<version>
2.11.2
</version>
</dependency>
<dependency>
<groupId>
tk.mybatis
</groupId>
...
...
@@ -101,6 +108,14 @@
<artifactId>
druid
</artifactId>
<groupId>
com.alibaba
</groupId>
</exclusion>
<exclusion>
<artifactId>
jackson-databind
</artifactId>
<groupId>
com.fasterxml.jackson.core
</groupId>
</exclusion>
<exclusion>
<artifactId>
jackson-core
</artifactId>
<groupId>
com.fasterxml.jackson.core
</groupId>
</exclusion>
</exclusions>
</dependency>
...
...
@@ -114,9 +129,41 @@
<dependency>
<groupId>
co.elastic.clients
</groupId>
<artifactId>
elasticsearch-java
</artifactId>
<version>
8.9.0
</version>
<version>
8.8.2
</version>
<exclusions>
<exclusion>
<artifactId>
elasticsearch-rest-client
</artifactId>
<groupId>
org.elasticsearch.client
</groupId>
</exclusion>
</exclusions>
</dependency>
<dependency>
<groupId>
org.elasticsearch.client
</groupId>
<artifactId>
elasticsearch-rest-client
</artifactId>
<version>
8.8.2
</version>
</dependency>
<!-- <dependency>
<groupId>org.elasticsearch.client</groupId>
<artifactId>elasticsearch-rest-high-level-client</artifactId>
<version>7.8.0</version>
<exclusions>
<exclusion>
<artifactId>elasticsearch-rest-client</artifactId>
<groupId>org.elasticsearch.client</groupId>
</exclusion>
<exclusion>
<artifactId>jackson-core</artifactId>
<groupId>com.fasterxml.jackson.core</groupId>
</exclusion>
</exclusions>
</dependency>-->
<dependency>
<groupId>
junit
</groupId>
<artifactId>
junit
</artifactId>
<scope>
test
</scope>
</dependency>
</dependencies>
...
...
ch-dao/src/main/java/com/wwdz/ch/db/dto/request/CoinRequestDto.java
View file @
5c29dae1
...
...
@@ -3,6 +3,7 @@ package com.wwdz.ch.db.dto.request;
import
com.xxdxxs.entity.Entity
;
import
lombok.Data
;
import
java.util.Date
;
import
java.util.List
;
@Data
...
...
@@ -87,4 +88,9 @@ public class CoinRequestDto extends BaseRequestDto implements Entity {
* 商品视频, 分号分隔
*/
private
String
videos
;
/**
* 用于同步数据到es
*/
private
Date
startUpdateTime
;
}
ch-dao/src/main/java/com/wwdz/ch/db/es/ItemEsDao.java
View file @
5c29dae1
package
com.wwdz.ch.db.es
;
import
co.elastic.clients.elasticsearch.ElasticsearchClient
;
import
co.elastic.clients.elasticsearch._types.
FieldValue
;
import
co.elastic.clients.elasticsearch._types.
*
;
import
co.elastic.clients.elasticsearch._types.query_dsl.*
;
import
co.elastic.clients.elasticsearch.core.*
;
import
co.elastic.clients.elasticsearch.core.bulk.BulkOperation
;
import
co.elastic.clients.elasticsearch.core.bulk.CreateOperation
;
import
co.elastic.clients.elasticsearch.core.search.Hit
;
import
co.elastic.clients.elasticsearch.core.search.ScoreMode
;
import
com.wwdz.ch.db.bean.EsPageInfo
;
import
com.wwdz.ch.db.bean.ItemOfEs
;
import
com.wwdz.ch.db.domain.Item
;
import
com.wwdz.ch.db.dto.request.EsSearchRequestDto
;
import
com.xxdxxs.utils.DateUtils
;
import
com.xxdxxs.utils.StringUtils
;
import
org.slf4j.Logger
;
import
org.slf4j.LoggerFactory
;
import
org.springframework.stereotype.Repository
;
import
org.springframework.util.CollectionUtils
;
import
javax.annotation.Resource
;
import
java.io.IOException
;
import
java.util.ArrayList
;
import
java.util.Arrays
;
import
java.util.List
;
import
java.util.Random
;
import
java.util.*
;
@Repository
...
...
@@ -113,7 +110,7 @@ public class ItemEsDao {
build
.
must
(
wildcardQuery
);
builder
.
size
(
size
);
SearchRequest
searchRequest
=
builder
.
build
();
logger
.
info
(
"查询语句 : {}"
,
searchRequest
.
query
());
logger
.
info
(
"查询语句 : {}"
,
searchRequest
.
toString
());
SearchResponse
<
ItemOfEs
>
searchResponse
=
elasticsearchClient
.
search
(
searchRequest
,
ItemOfEs
.
class
);
long
total
=
searchResponse
.
hits
().
total
().
value
();
logger
.
info
(
"总数据量 : {}"
,
total
);
...
...
@@ -163,4 +160,43 @@ public class ItemEsDao {
}
public
void
update
(
Item
item
)
{
String
id
=
String
.
valueOf
(
item
.
getId
());
try
{
// 使用map集合封装需要修改的内容
Map
<
String
,
Object
>
map
=
new
HashMap
<>();
map
.
put
(
"images"
,
item
.
getImages
());
UpdateResponse
<
Item
>
updateResponse
=
elasticsearchClient
.
update
(
e
->
e
.
index
(
INDEX
).
id
(
id
).
doc
(
map
),
Item
.
class
);
logger
.
info
(
"ES修改数据成功, id:{}, 修改内容: {} "
,
id
,
map
);
}
catch
(
Exception
e
)
{
logger
.
error
(
"修改es数据出错 id: {}"
,
id
,
e
);
}
}
public
Date
getMaxTime
()
{
try
{
SearchRequest
.
Builder
builder
=
new
SearchRequest
.
Builder
();
Query
query
=
QueryBuilders
.
matchAll
().
build
().
_toQuery
();
builder
.
index
(
INDEX
);
builder
.
source
(
sourceBuild
->
sourceBuild
.
filter
(
f
->
f
.
includes
(
"id"
,
"updateTime"
,
"createTime"
)));
builder
.
query
(
query
);
builder
.
sort
(
s
->
s
.
field
(
a
->
a
.
field
(
"updateTime"
).
order
(
SortOrder
.
Desc
)));
builder
.
size
(
1
);
SearchRequest
searchRequest
=
builder
.
build
();
logger
.
info
(
"查询语句 : {}"
,
searchRequest
.
toString
());
SearchResponse
<
ItemOfEs
>
searchResponse
=
elasticsearchClient
.
search
(
searchRequest
,
ItemOfEs
.
class
);
logger
.
info
(
">>>>>>>>>>> getMaxTime Response : {}"
,
searchResponse
);
List
<
Hit
<
ItemOfEs
>>
hitList
=
searchResponse
.
hits
().
hits
();
if
(!
CollectionUtils
.
isEmpty
(
hitList
))
{
Date
time
=
hitList
.
get
(
0
).
source
().
getUpdateTime
();
logger
.
info
(
">>>>>>>>>>>> 查询出的最大更新时间: {}"
,
DateUtils
.
toString
(
time
));
return
time
;
}
}
catch
(
Exception
e
)
{
logger
.
error
(
"查询es中数据最大更新时间出错 :{}"
,
e
);
}
return
null
;
}
}
ch-dao/src/main/java/com/wwdz/ch/db/impl/ItemDaoImpl.java
View file @
5c29dae1
...
...
@@ -66,6 +66,7 @@ public class ItemDaoImpl implements ItemDao {
ItemExample
example
=
new
ItemExample
();
ItemExample
.
Criteria
criteria
=
example
.
createCriteria
();
example
.
setOrderByClause
(
"create_time asc"
);
JdbcHelper
.
ifPresent
(
coinRequestDto
.
getStartUpdateTime
(),
criteria:
:
andUpdateTimeGreaterThan
);
PageHelper
.
startPage
(
coinRequestDto
.
getPage
(),
coinRequestDto
.
getLimit
());
return
itemMapper
.
selectByExampleWithBLOBs
(
example
);
}
...
...
pom.xml
View file @
5c29dae1
...
...
@@ -9,7 +9,7 @@
<parent>
<groupId>
org.springframework.boot
</groupId>
<artifactId>
spring-boot-starter-parent
</artifactId>
<version>
2.
0.4.RELEASE
</version>
<version>
2.
5.3
</version>
<relativePath/>
</parent>
...
...
@@ -18,7 +18,8 @@
<java.version>
1.8
</java.version>
<maven.test.skip>
true
</maven.test.skip>
<fastjson.version>
1.2.70
</fastjson.version>
<jackson-databind.version>
2.9.6
</jackson-databind.version>
<jackson-databind.version>
2.11.2
</jackson-databind.version>
<jakartajson.version>
2.0.1
</jakartajson.version>
<mybatis-spring-boot-starter.version>
1.3.2
</mybatis-spring-boot-starter.version>
<pagehelper-spring-boot-starter.version>
1.2.5
</pagehelper-spring-boot-starter.version>
<shiro-spring-boot-web-starter.version>
1.4.0
</shiro-spring-boot-web-starter.version>
...
...
@@ -143,6 +144,17 @@
<groupId>
com.fasterxml.jackson.core
</groupId>
<artifactId>
jackson-databind
</artifactId>
<version>
${jackson-databind.version}
</version>
</dependency>
<dependency>
<groupId>
jakarta.json
</groupId>
<artifactId>
jakarta.json-api
</artifactId>
<version>
${jakartajson.version}
</version>
</dependency>
<dependency>
<groupId>
org.elasticsearch.client
</groupId>
<artifactId>
elasticsearch-rest-client
</artifactId>
<version>
8.8.2
</version>
</dependency>
</dependencies>
</dependencyManagement>
...
...
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