Skip to content
Projects
Groups
Snippets
Help
Loading...
Help
Submit feedback
Sign in / Register
Toggle navigation
Y
yzg-util
Project
Project
Details
Activity
Releases
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
YZG
yzg-util
Commits
534c2e5c
Commit
534c2e5c
authored
Aug 10, 2021
by
yanzg
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
修改实例化关系
parent
e3f38cad
Changes
3
Show whitespace changes
Inline
Side-by-side
Showing
3 changed files
with
15 additions
and
18 deletions
+15
-18
MqConfigurable.java
...src/main/java/com/yanzuoguang/mq/base/MqConfigurable.java
+1
-10
MessageSendServiceImpl.java
...m/yanzuoguang/mq/service/impl/MessageSendServiceImpl.java
+11
-5
MqServiceImpl.java
...n/java/com/yanzuoguang/mq/service/impl/MqServiceImpl.java
+3
-3
No files found.
yzg-util-mq/src/main/java/com/yanzuoguang/mq/base/MqConfigurable.java
View file @
534c2e5c
...
@@ -63,11 +63,7 @@ public class MqConfigurable implements RabbitTemplate.ConfirmCallback, RabbitTem
...
@@ -63,11 +63,7 @@ public class MqConfigurable implements RabbitTemplate.ConfirmCallback, RabbitTem
try
{
try
{
if
(
ack
&&
correlationData
!=
null
if
(
ack
&&
correlationData
!=
null
&&
!
StringHelper
.
isEmpty
(
correlationData
.
getId
()))
{
&&
!
StringHelper
.
isEmpty
(
correlationData
.
getId
()))
{
String
toId
=
getId
(
correlationData
.
getId
());
messageSendService
.
onSuccess
(
correlationData
.
getId
());
// 不是临时数据
if
(
toId
.
equals
(
correlationData
.
getId
()))
{
messageSendService
.
onSuccess
(
toId
);
}
}
else
if
(!
ack
)
{
}
else
if
(!
ack
)
{
System
.
out
.
println
(
"丢失消息:"
+
ack
+
" msg:"
+
cause
);
System
.
out
.
println
(
"丢失消息:"
+
ack
+
" msg:"
+
cause
);
}
}
...
@@ -94,15 +90,10 @@ public class MqConfigurable implements RabbitTemplate.ConfirmCallback, RabbitTem
...
@@ -94,15 +90,10 @@ public class MqConfigurable implements RabbitTemplate.ConfirmCallback, RabbitTem
String
content
=
new
String
(
message
.
getBody
(),
charset
);
String
content
=
new
String
(
message
.
getBody
(),
charset
);
// 组成消息
// 组成消息
MessageVo
messageVo
=
new
MessageVo
(
exchange
,
routingKey
,
content
);
MessageVo
messageVo
=
new
MessageVo
(
exchange
,
routingKey
,
content
);
messageVo
.
setMessageId
(
getId
(
messageProperties
.
getMessageId
()));
// 写入数据库
// 写入数据库
messageSendService
.
onError
(
messageVo
);
messageSendService
.
onError
(
messageVo
);
}
catch
(
Exception
ex
)
{
}
catch
(
Exception
ex
)
{
Log
.
error
(
MqConfigurable
.
class
,
ex
);
Log
.
error
(
MqConfigurable
.
class
,
ex
);
}
}
}
}
private
String
getId
(
String
from
)
{
return
from
.
replace
(
"temp:"
,
""
);
}
}
}
yzg-util-mq/src/main/java/com/yanzuoguang/mq/service/impl/MessageSendServiceImpl.java
View file @
534c2e5c
...
@@ -32,6 +32,8 @@ import java.util.List;
...
@@ -32,6 +32,8 @@ import java.util.List;
@Component
@Component
public
class
MessageSendServiceImpl
implements
MessageSendService
{
public
class
MessageSendServiceImpl
implements
MessageSendService
{
public
static
final
String
TEMP_ID
=
"temp"
;
@Autowired
@Autowired
private
MyRabbitTemplate
rabbitTemplate
;
private
MyRabbitTemplate
rabbitTemplate
;
...
@@ -78,7 +80,7 @@ public class MessageSendServiceImpl implements MessageSendService {
...
@@ -78,7 +80,7 @@ public class MessageSendServiceImpl implements MessageSendService {
public
String
send
(
MessageVo
req
)
{
public
String
send
(
MessageVo
req
)
{
req
.
check
();
req
.
check
();
// 获取消息临时Id,消息Id为空时标识为第一次发送,并设置默认消息Id
// 获取消息临时Id,消息Id为空时标识为第一次发送,并设置默认消息Id
String
finalMessageId
=
StringHelper
.
getFirst
(
req
.
getMessageId
(),
StringHelper
.
getId
(
"temp"
,
StringHelper
.
getNewID
()));
String
finalMessageId
=
StringHelper
.
getFirst
(
req
.
getMessageId
(),
StringHelper
.
getId
(
TEMP_ID
,
StringHelper
.
getNewID
()));
// 设置编号
// 设置编号
CorrelationData
correlationData
=
new
CorrelationData
();
CorrelationData
correlationData
=
new
CorrelationData
();
correlationData
.
setId
(
finalMessageId
);
correlationData
.
setId
(
finalMessageId
);
...
@@ -90,7 +92,7 @@ public class MessageSendServiceImpl implements MessageSendService {
...
@@ -90,7 +92,7 @@ public class MessageSendServiceImpl implements MessageSendService {
// 设置持久化
// 设置持久化
properties
.
setDeliveryMode
(
MessageDeliveryMode
.
PERSISTENT
);
properties
.
setDeliveryMode
(
MessageDeliveryMode
.
PERSISTENT
);
// 设置消息编号
// 设置消息编号
properties
.
setMessageId
(
StringHelper
.
getIdShort
(
finalMessageId
,
"temp"
));
properties
.
setMessageId
(
StringHelper
.
getIdShort
(
finalMessageId
,
TEMP_ID
));
if
(
req
.
getDedTime
()
>
0
)
{
if
(
req
.
getDedTime
()
>
0
)
{
properties
.
setExpiration
(
req
.
getDedTime
()
+
""
);
properties
.
setExpiration
(
req
.
getDedTime
()
+
""
);
}
}
...
@@ -107,10 +109,13 @@ public class MessageSendServiceImpl implements MessageSendService {
...
@@ -107,10 +109,13 @@ public class MessageSendServiceImpl implements MessageSendService {
*/
*/
@Override
@Override
public
String
onSuccess
(
String
messageId
)
{
public
String
onSuccess
(
String
messageId
)
{
if
(!
StringHelper
.
isEmpty
(
messageId
))
{
String
toId
=
StringHelper
.
getIdShort
(
messageId
,
TEMP_ID
);
messageDao
.
remove
(
messageId
);
// 不是临时数据
if
(!
toId
.
equals
(
messageId
)
||
StringHelper
.
isEmpty
(
toId
))
{
return
StringHelper
.
EMPTY
;
}
}
return
messageId
;
messageDao
.
remove
(
toId
);
return
toId
;
}
}
/**
/**
...
@@ -120,6 +125,7 @@ public class MessageSendServiceImpl implements MessageSendService {
...
@@ -120,6 +125,7 @@ public class MessageSendServiceImpl implements MessageSendService {
*/
*/
@Override
@Override
public
String
onError
(
MessageVo
messageVo
)
{
public
String
onError
(
MessageVo
messageVo
)
{
messageVo
.
setMessageId
(
StringHelper
.
getIdShort
(
messageVo
.
getMessageId
(),
TEMP_ID
));
messageVo
.
check
();
messageVo
.
check
();
// 设置处理次数
// 设置处理次数
messageVo
.
setHandleCount
(
messageVo
.
getHandleCount
()
+
1
);
messageVo
.
setHandleCount
(
messageVo
.
getHandleCount
()
+
1
);
...
...
yzg-util-mq/src/main/java/com/yanzuoguang/mq/service/impl/MqServiceImpl.java
View file @
534c2e5c
...
@@ -71,9 +71,9 @@ public class MqServiceImpl implements MqService {
...
@@ -71,9 +71,9 @@ public class MqServiceImpl implements MqService {
// 设置默认消息Id
// 设置默认消息Id
String
defaultId
=
StringHelper
.
getFirst
(
req
.
getMessageId
(),
StringHelper
.
getNewID
());
String
defaultId
=
StringHelper
.
getFirst
(
req
.
getMessageId
(),
StringHelper
.
getNewID
());
// 将Id去掉temp:
// 将Id去掉temp:
String
simpleId
=
StringHelper
.
getId
Short
(
defaultId
,
"temp"
);
String
simpleId
=
StringHelper
.
getId
(
MessageSendServiceImpl
.
TEMP_ID
,
StringHelper
.
getIdShort
(
defaultId
,
MessageSendServiceImpl
.
TEMP_ID
)
);
// 增加temp标识第一次发送
// 增加temp标识第一次发送
req
.
setMessageId
(
StringHelper
.
getId
(
"temp"
,
simpleId
)
);
req
.
setMessageId
(
simpleId
);
return
yzgMqProcedure
.
send
(
req
,
now
);
return
yzgMqProcedure
.
send
(
req
,
now
);
}
}
...
...
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