Skip to content
Projects
Groups
Snippets
Help
This project
Loading...
Sign in / Register
Toggle navigation
B
beyond-clouds
Overview
Overview
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
4
Issues
4
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
段启岩
beyond-clouds
Commits
cb508aec
Commit
cb508aec
authored
Feb 09, 2020
by
段启岩
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
架构关系-建立父类-TopicMessageListener
parent
6d7a97a6
Show whitespace changes
Inline
Side-by-side
Showing
3 changed files
with
54 additions
and
41 deletions
+54
-41
src/main/java/cn/meteor/beyondclouds/core/listener/TopicMessageListener.java
+47
-0
src/main/java/cn/meteor/beyondclouds/modules/search/listener/ISearchItemUpdateListener.java
+0
-16
src/main/java/cn/meteor/beyondclouds/modules/search/listener/SearchItemUpdateListener.java
+7
-25
No files found.
src/main/java/cn/meteor/beyondclouds/core/listener/TopicMessageListener.java
0 → 100644
View file @
cb508aec
package
cn
.
meteor
.
beyondclouds
.
core
.
listener
;
import
cn.meteor.beyondclouds.modules.queue.message.DataItemUpdateMessage
;
import
cn.meteor.beyondclouds.modules.queue.message.SearchItemUpdateType
;
import
cn.meteor.beyondclouds.modules.search.enums.SearchItemType
;
import
cn.meteor.beyondclouds.util.JsonUtils
;
import
lombok.extern.slf4j.Slf4j
;
import
org.apache.kafka.clients.consumer.ConsumerRecord
;
import
org.springframework.kafka.annotation.KafkaListener
;
import
java.io.Serializable
;
import
java.util.Optional
;
/**
* @author meteor
*/
@Slf4j
public
class
TopicMessageListener
{
@KafkaListener
(
topics
=
"${beyondclouds.kafka.topics.search-item-update}"
)
public
final
void
itemUpdate
(
ConsumerRecord
<?,
String
>
record
)
{
Optional
<
String
>
kafkaMessage
=
Optional
.
ofNullable
(
record
.
value
());
if
(
kafkaMessage
.
isPresent
())
{
DataItemUpdateMessage
itemUpdateMessage
;
try
{
itemUpdateMessage
=
JsonUtils
.
toBean
(
kafkaMessage
.
get
(),
DataItemUpdateMessage
.
class
);
log
.
debug
(
"接收到kafka消息:{}"
,
itemUpdateMessage
.
toString
());
// 调用消息处理函数
onItemUpdate
(
itemUpdateMessage
);
}
catch
(
Exception
e
)
{
e
.
printStackTrace
();
log
.
error
(
"DataItemUpdateMessage consume failed:{}"
,
e
.
getMessage
());
}
}
}
/**
* 数据更新事件
* @param dataItemUpdateMessage
*/
public
void
onItemUpdate
(
DataItemUpdateMessage
dataItemUpdateMessage
)
throws
Exception
{
}
}
src/main/java/cn/meteor/beyondclouds/modules/search/listener/ISearchItemUpdateListener.java
deleted
100644 → 0
View file @
6d7a97a6
package
cn
.
meteor
.
beyondclouds
.
modules
.
search
.
listener
;
import
org.apache.kafka.clients.consumer.ConsumerRecord
;
/**
* @author meteor
*/
public
interface
ISearchItemUpdateListener
{
/**
* 搜索条目更新
* @param record
*/
void
onSearchItemUpdate
(
ConsumerRecord
<?,
String
>
record
);
}
src/main/java/cn/meteor/beyondclouds/modules/search/listener/
impl/
SearchItemUpdateListener.java
→
src/main/java/cn/meteor/beyondclouds/modules/search/listener/SearchItemUpdateListener.java
View file @
cb508aec
package
cn
.
meteor
.
beyondclouds
.
modules
.
search
.
listener
.
impl
;
package
cn
.
meteor
.
beyondclouds
.
modules
.
search
.
listener
;
import
cn.meteor.beyondclouds.core.listener.TopicMessageListener
;
import
cn.meteor.beyondclouds.modules.queue.message.DataItemUpdateMessage
;
import
cn.meteor.beyondclouds.modules.queue.message.DataItemUpdateMessage
;
import
cn.meteor.beyondclouds.modules.queue.message.SearchItemUpdateType
;
import
cn.meteor.beyondclouds.modules.queue.message.SearchItemUpdateType
;
import
cn.meteor.beyondclouds.modules.search.enums.SearchItemType
;
import
cn.meteor.beyondclouds.modules.search.enums.SearchItemType
;
import
cn.meteor.beyondclouds.modules.search.listener.ISearchItemUpdateListener
;
import
cn.meteor.beyondclouds.modules.search.service.ISearchService
;
import
cn.meteor.beyondclouds.modules.search.service.ISearchService
;
import
cn.meteor.beyondclouds.util.JsonUtils
;
import
lombok.extern.slf4j.Slf4j
;
import
lombok.extern.slf4j.Slf4j
;
import
org.apache.kafka.clients.consumer.ConsumerRecord
;
import
org.springframework.beans.factory.annotation.Autowired
;
import
org.springframework.beans.factory.annotation.Autowired
;
import
org.springframework.kafka.annotation.KafkaListener
;
import
org.springframework.stereotype.Component
;
import
org.springframework.stereotype.Component
;
import
java.io.Serializable
;
import
java.io.Serializable
;
import
java.util.Optional
;
/**
/**
* @author meteor
* @author meteor
*/
*/
@Component
@Component
@Slf4j
@Slf4j
public
class
SearchItemUpdateListener
implements
ISearchItemUpdat
eListener
{
public
class
SearchItemUpdateListener
extends
TopicMessag
eListener
{
private
ISearchService
searchService
;
private
ISearchService
searchService
;
...
@@ -30,20 +26,12 @@ public class SearchItemUpdateListener implements ISearchItemUpdateListener {
...
@@ -30,20 +26,12 @@ public class SearchItemUpdateListener implements ISearchItemUpdateListener {
}
}
@Override
@Override
@KafkaListener
(
topics
=
"${beyondclouds.kafka.topics.search-item-update}"
)
public
void
onItemUpdate
(
DataItemUpdateMessage
dataItemUpdateMessage
)
throws
Exception
{
public
void
onSearchItemUpdate
(
ConsumerRecord
<?,
String
>
record
)
{
Optional
<
String
>
kafkaMessage
=
Optional
.
ofNullable
(
record
.
value
());
if
(
kafkaMessage
.
isPresent
())
{
DataItemUpdateMessage
itemUpdateMessage
;
try
{
itemUpdateMessage
=
JsonUtils
.
toBean
(
kafkaMessage
.
get
(),
DataItemUpdateMessage
.
class
);
log
.
debug
(
"接收到kafka消息:{}"
,
itemUpdateMessage
.
toString
());
// 处理搜索条目更新
// 处理搜索条目更新
SearchItemUpdateType
updateType
=
i
temUpdateMessage
.
getUpdateType
();
SearchItemUpdateType
updateType
=
dataI
temUpdateMessage
.
getUpdateType
();
Serializable
itemId
=
i
temUpdateMessage
.
getItemId
();
Serializable
itemId
=
dataI
temUpdateMessage
.
getItemId
();
SearchItemType
searchItemType
=
i
temUpdateMessage
.
getItemType
();
SearchItemType
searchItemType
=
dataI
temUpdateMessage
.
getItemType
();
// 根据不同的更新类型调用对应的方法
// 根据不同的更新类型调用对应的方法
switch
(
updateType
)
{
switch
(
updateType
)
{
...
@@ -59,11 +47,5 @@ public class SearchItemUpdateListener implements ISearchItemUpdateListener {
...
@@ -59,11 +47,5 @@ public class SearchItemUpdateListener implements ISearchItemUpdateListener {
default
:
default
:
break
;
break
;
}
}
}
catch
(
Exception
e
)
{
e
.
printStackTrace
();
log
.
error
(
"searchItemUpdateMessage covert failed:{}"
,
e
.
getMessage
());
}
}
}
}
}
}
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