Skip to content
Projects
Groups
Snippets
Help
Loading...
Sign in
Toggle navigation
M
mq-man
Project
Project
Details
Activity
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
M. A. Rased
mq-man
Commits
b2301c44
Commit
b2301c44
authored
Sep 06, 2022
by
M. A. Rased
Browse files
Options
Browse Files
Download
Plain Diff
Merge branch 'feature/remove-dl' into 'master'
remove dlq and dlx feature See merge request
!1
parents
59758a73
2fa3d1bb
Hide whitespace changes
Inline
Side-by-side
Showing
3 changed files
with
202 additions
and
32 deletions
+202
-32
DeadLetterApplication.java
...in/java/com/example/deadletter/DeadLetterApplication.java
+0
-1
DeadLetterQueueService.java
...n/java/com/example/deadletter/DeadLetterQueueService.java
+175
-31
PolicyData.java
src/main/java/com/example/deadletter/dto/PolicyData.java
+27
-0
No files found.
src/main/java/com/example/deadletter/DeadLetterApplication.java
View file @
b2301c44
...
...
@@ -15,5 +15,4 @@ public class DeadLetterApplication {
DeadLetterQueueService
deadLetterQueueService
=
context
.
getBean
(
"deadLetterQueueService"
,
DeadLetterQueueService
.
class
);
deadLetterQueueService
.
setup
(
args
);
}
}
src/main/java/com/example/deadletter/DeadLetterQueueService.java
View file @
b2301c44
package
com
.
example
.
deadletter
;
import
com.example.deadletter.dto.PolicyData
;
import
com.example.deadletter.dto.QueueData
;
import
org.json.JSONObject
;
import
org.springframework.beans.factory.annotation.Value
;
...
...
@@ -12,6 +13,7 @@ import org.springframework.web.client.RestTemplate;
import
java.util.ArrayList
;
import
java.util.Base64
;
import
java.util.List
;
import
java.util.Scanner
;
@Service
public
class
DeadLetterQueueService
{
...
...
@@ -34,14 +36,18 @@ public class DeadLetterQueueService {
private
final
RestTemplate
restTemplate
=
new
RestTemplate
();
public
void
setup
(
String
[]
args
)
{
if
(
args
.
length
==
0
)
{
args
=
getAllQueue
();
}
if
(
args
.
length
==
1
&&
args
[
0
].
equals
(
"remove"
))
{
removeDeadLetterExchange
();
removeDeadLetterQueue
();
removeDeadLetterPolicies
();
return
;
}
String
deadLetterExchange
=
"DLX"
;
createDeadLetterExchange
(
deadLetterExchange
);
for
(
String
queue
:
args
)
{
System
.
out
.
println
(
"\n\nQUEUE NAME : "
+
queue
);
...
...
@@ -52,14 +58,180 @@ public class DeadLetterQueueService {
createDeadLetterQueue
(
queue
,
deadLetterExchange
);
setDLXPolicy
(
queue
,
deadLetterExchange
);
}
}
private
void
removeDeadLetterPolicies
()
{
List
<
String
>
policyList
=
getAllPolicies
();
if
(
policyList
.
isEmpty
())
return
;
System
.
out
.
println
(
"\n\n\nPOLICIES TO DELETE \n"
+
policyList
);
System
.
out
.
println
(
"\npress y to delete"
);
Scanner
scanner
=
new
Scanner
(
System
.
in
);
String
input
=
scanner
.
next
();
if
(!
input
.
equalsIgnoreCase
(
"y"
))
return
;
for
(
String
policy
:
policyList
)
{
removePolicy
(
policy
);
}
}
private
void
createDeadLetterExchange
(
String
deadLetterExchange
)
{
private
void
removePolicy
(
String
policyName
)
{
ResponseEntity
<
String
>
responseEntity
=
null
;
String
apiUrl
=
"http://"
+
host
+
":"
+
managementPort
+
"/api/policies"
;
String
wholeUrl
=
apiUrl
+
"/"
+
vhost
+
"/"
+
policyName
;
try
{
responseEntity
=
restTemplate
.
exchange
(
wholeUrl
,
HttpMethod
.
DELETE
,
new
HttpEntity
<>(
getHeader
()),
String
.
class
);
System
.
out
.
print
(
" try to delete policy: "
+
policyName
+
"\n"
);
}
catch
(
Exception
ex
)
{
throw
new
RuntimeException
(
ex
);
}
}
private
List
<
String
>
getAllPolicies
()
{
List
<
String
>
policyList
=
new
ArrayList
<>();
String
apiUrl
=
"http://"
+
host
+
":"
+
managementPort
+
"/api/policies"
;
String
wholeUrl
=
apiUrl
+
"/"
+
vhost
;
ResponseEntity
<
List
<
PolicyData
>>
responseEntity
=
null
;
try
{
responseEntity
=
restTemplate
.
exchange
(
wholeUrl
,
HttpMethod
.
GET
,
new
HttpEntity
<>(
getHeader
()),
new
ParameterizedTypeReference
<
List
<
PolicyData
>>()
{
}
);
List
<
PolicyData
>
policyDataList
=
responseEntity
.
getBody
();
for
(
PolicyData
policyData
:
policyDataList
)
{
if
(!
policyData
.
getName
().
startsWith
(
"dlx"
))
{
continue
;
}
policyList
.
add
(
policyData
.
getName
());
}
return
policyList
;
}
catch
(
Exception
ex
)
{
throw
new
RuntimeException
(
ex
);
}
}
private
void
removeDeadLetterQueue
()
{
List
<
String
>
deadLetterQueueList
=
getDeadLetterQueue
();
if
(
deadLetterQueueList
.
isEmpty
())
return
;
System
.
out
.
println
(
"\n\nDEAD LETTER QUEUES TO DELETE \n"
+
deadLetterQueueList
);
System
.
out
.
println
(
"\npress y to delete"
);
Scanner
scanner
=
new
Scanner
(
System
.
in
);
String
input
=
scanner
.
next
();
if
(!
input
.
equalsIgnoreCase
(
"y"
))
return
;
for
(
String
deadLetterQueue
:
deadLetterQueueList
)
{
removeQueue
(
deadLetterQueue
);
}
}
private
void
removeQueue
(
String
deadLetterQueue
)
{
ResponseEntity
<
String
>
responseEntity
=
null
;
String
apiUrl
=
"http://"
+
host
+
":"
+
managementPort
+
"/api/queues"
;
String
wholeUrl
=
apiUrl
+
"/"
+
vhost
+
"/"
+
deadLetterQueue
+
"?if-empty=true&if-unused=true"
;
System
.
out
.
printf
(
"%n try to delete dead letter queue: %s "
,
deadLetterQueue
);
try
{
responseEntity
=
restTemplate
.
exchange
(
wholeUrl
,
HttpMethod
.
DELETE
,
new
HttpEntity
<>(
getHeader
()),
String
.
class
);
}
catch
(
Exception
ex
)
{
System
.
out
.
println
(
ex
.
getMessage
());
return
;
}
}
private
List
<
String
>
getDeadLetterQueue
()
{
String
apiUrl
=
"http://"
+
host
+
":"
+
managementPort
+
"/api/queues"
;
String
wholeUrl
=
apiUrl
+
"/"
+
vhost
;
ResponseEntity
<
List
<
QueueData
>>
responseEntity
=
null
;
try
{
responseEntity
=
restTemplate
.
exchange
(
wholeUrl
,
HttpMethod
.
GET
,
new
HttpEntity
<>(
getHeader
()),
new
ParameterizedTypeReference
<
List
<
QueueData
>>()
{
});
if
(
responseEntity
.
getStatusCode
()
==
HttpStatus
.
OK
)
{
List
<
String
>
queueList
=
new
ArrayList
<>();
for
(
QueueData
queueData
:
responseEntity
.
getBody
())
{
//reject non dl queues
if
(!
queueData
.
getName
().
startsWith
(
"dl-"
))
{
continue
;
}
queueList
.
add
(
queueData
.
getName
());
}
return
queueList
;
}
throw
new
RuntimeException
(
String
.
valueOf
(
responseEntity
.
getBody
()));
}
catch
(
Exception
ex
)
{
throw
new
RuntimeException
(
ex
);
}
}
private
void
removeDeadLetterExchange
()
{
String
deadLetterExchange
=
"DLX"
;
String
apiUrl
=
"http://"
+
host
+
":"
+
managementPort
+
"/api/exchanges"
;
String
wholeUrl
=
apiUrl
+
"/"
+
vhost
+
"/"
+
deadLetterExchange
;
ResponseEntity
<
String
>
responseEntity
=
null
;
try
{
responseEntity
=
restTemplate
.
exchange
(
wholeUrl
,
HttpMethod
.
GET
,
new
HttpEntity
<>(
getHeader
()),
String
.
class
);
if
(
responseEntity
.
getStatusCode
()
==
HttpStatus
.
OK
)
{
System
.
out
.
printf
(
"%n dead letter exchange %s exists %n"
,
deadLetterExchange
);
//delete it
try
{
responseEntity
=
restTemplate
.
exchange
(
wholeUrl
,
HttpMethod
.
DELETE
,
new
HttpEntity
<>(
getHeader
()),
String
.
class
);
System
.
out
.
printf
(
"%n deleted dead letter exchange: %s "
,
deadLetterExchange
);
System
.
out
.
print
(
responseEntity
.
getStatusCode
()
+
"\n"
);
}
catch
(
Exception
ex
)
{
throw
new
RuntimeException
(
ex
);
}
}
else
{
System
.
out
.
print
(
responseEntity
.
getStatusCode
()
+
"\n"
);
return
;
}
}
catch
(
HttpClientErrorException
ex
)
{
System
.
out
.
printf
(
"%n dead letter exchange %s does not exist %n"
,
deadLetterExchange
);
return
;
}
catch
(
Exception
ex
)
{
throw
new
RuntimeException
(
ex
);
}
}
private
void
createDeadLetterExchange
(
String
deadLetterExchange
)
{
String
apiUrl
=
"http://"
+
host
+
":"
+
managementPort
+
"/api/exchanges"
;
String
wholeUrl
=
apiUrl
+
"/"
+
vhost
+
"/"
+
deadLetterExchange
;
JSONObject
jsonObject
=
new
JSONObject
();
jsonObject
.
put
(
"type"
,
"topic"
);
...
...
@@ -67,7 +239,6 @@ public class DeadLetterQueueService {
ResponseEntity
<
String
>
responseEntity
=
null
;
try
{
responseEntity
=
restTemplate
.
exchange
(
wholeUrl
,
HttpMethod
.
GET
,
...
...
@@ -84,31 +255,24 @@ public class DeadLetterQueueService {
}
}
catch
(
HttpClientErrorException
ex
)
{
System
.
out
.
printf
(
"%n dead letter exchange %s does not exist %n"
,
deadLetterExchange
);
}
catch
(
Exception
ex
)
{
throw
new
RuntimeException
(
ex
);
}
try
{
responseEntity
=
restTemplate
.
exchange
(
wholeUrl
,
HttpMethod
.
PUT
,
new
HttpEntity
<>(
jsonObject
.
toString
(),
getHeader
()),
String
.
class
);
System
.
out
.
printf
(
"%n created dead letter exchange: %s "
,
deadLetterExchange
);
System
.
out
.
print
(
responseEntity
.
getStatusCode
()
+
"\n"
);
}
catch
(
Exception
ex
)
{
throw
new
RuntimeException
(
ex
);
}
}
private
void
setDLXPolicy
(
String
queueName
,
String
deadLetterExchange
)
{
JSONObject
jsonObject2
=
new
JSONObject
();
jsonObject2
.
put
(
"dead-letter-exchange"
,
deadLetterExchange
);
jsonObject2
.
put
(
"dead-letter-routing-key"
,
queueName
);
...
...
@@ -125,7 +289,6 @@ public class DeadLetterQueueService {
String
apiUrl
=
"http://"
+
host
+
":"
+
managementPort
+
"/api/policies"
;
String
wholeUrl
=
apiUrl
+
"/"
+
vhost
+
"/"
+
policyName
;
try
{
responseEntity
=
restTemplate
.
exchange
(
wholeUrl
,
HttpMethod
.
PUT
,
...
...
@@ -138,9 +301,7 @@ public class DeadLetterQueueService {
}
catch
(
Exception
ex
)
{
throw
new
RuntimeException
(
ex
);
}
}
void
createDeadLetterQueue
(
String
queueName
,
String
deadLetterExchange
)
{
...
...
@@ -151,7 +312,6 @@ public class DeadLetterQueueService {
if
(
checkQueuePresent
(
deadLetterQueue
))
{
System
.
out
.
printf
(
"%n dead letter queue already present : %s "
,
deadLetterQueue
);
System
.
out
.
printf
(
"%n try to delete dead letter queue: %s "
,
deadLetterQueue
);
try
{
...
...
@@ -160,14 +320,11 @@ public class DeadLetterQueueService {
new
HttpEntity
<>(
getHeader
()),
String
.
class
);
System
.
out
.
print
(
responseEntity
.
getStatusCode
()
+
"\n"
);
}
catch
(
Exception
ex
)
{
System
.
out
.
println
(
ex
.
getMessage
());
return
;
}
}
JSONObject
jsonObject
=
new
JSONObject
();
...
...
@@ -186,10 +343,8 @@ public class DeadLetterQueueService {
}
catch
(
Exception
ex
)
{
throw
new
RuntimeException
(
ex
);
}
apiUrl
=
"http://"
+
host
+
":"
+
managementPort
+
"/api/bindings"
;
wholeUrl
=
apiUrl
+
"/"
+
vhost
+
"/e/"
+
deadLetterExchange
+
"/q/"
+
deadLetterQueue
;
...
...
@@ -203,25 +358,20 @@ public class DeadLetterQueueService {
String
.
class
);
System
.
out
.
printf
(
"%n created binding for exchange %s and queue %s "
,
deadLetterExchange
,
deadLetterQueue
);
System
.
out
.
print
(
responseEntity
.
getStatusCode
()
+
"\n"
);
}
catch
(
Exception
ex
)
{
throw
new
RuntimeException
(
ex
);
}
}
boolean
checkQueuePresent
(
String
queueName
)
{
String
apiUrl
=
"http://"
+
host
+
":"
+
managementPort
+
"/api/queues"
;
String
wholeUrl
=
apiUrl
+
"/"
+
vhost
+
"/"
+
queueName
;
ResponseEntity
<
String
>
responseEntity
=
null
;
try
{
responseEntity
=
restTemplate
.
exchange
(
wholeUrl
,
HttpMethod
.
GET
,
...
...
@@ -235,13 +385,11 @@ public class DeadLetterQueueService {
System
.
out
.
print
(
responseEntity
.
getStatusCode
()
+
"\n"
);
return
false
;
}
catch
(
HttpClientErrorException
ex
)
{
System
.
out
.
printf
(
"%n dead letter queue %s does not exist %n"
,
queueName
);
return
false
;
}
catch
(
Exception
ex
)
{
throw
new
RuntimeException
(
ex
);
}
}
HttpHeaders
getHeader
()
{
...
...
@@ -258,10 +406,8 @@ public class DeadLetterQueueService {
String
apiUrl
=
"http://"
+
host
+
":"
+
managementPort
+
"/api/queues"
;
String
wholeUrl
=
apiUrl
+
"/"
+
vhost
;
ResponseEntity
<
List
<
QueueData
>>
responseEntity
=
null
;
try
{
responseEntity
=
restTemplate
.
exchange
(
wholeUrl
,
HttpMethod
.
GET
,
...
...
@@ -269,7 +415,6 @@ public class DeadLetterQueueService {
new
ParameterizedTypeReference
<
List
<
QueueData
>>()
{
});
if
(
responseEntity
.
getStatusCode
()
==
HttpStatus
.
OK
)
{
List
<
String
>
queueList
=
new
ArrayList
<>();
...
...
@@ -283,7 +428,6 @@ public class DeadLetterQueueService {
return
queueList
.
toArray
(
new
String
[
0
]);
}
throw
new
RuntimeException
(
String
.
valueOf
(
responseEntity
.
getBody
()));
}
catch
(
Exception
ex
)
{
throw
new
RuntimeException
(
ex
);
}
...
...
src/main/java/com/example/deadletter/dto/PolicyData.java
0 → 100644
View file @
b2301c44
package
com
.
example
.
deadletter
.
dto
;
public
class
PolicyData
{
String
name
;
public
PolicyData
(
String
name
)
{
this
.
name
=
name
;
}
public
PolicyData
()
{
}
public
String
getName
()
{
return
name
;
}
public
void
setName
(
String
name
)
{
this
.
name
=
name
;
}
@Override
public
String
toString
()
{
return
"PolicyData{"
+
"name='"
+
name
+
'\''
+
'}'
;
}
}
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