记得上下班打卡 | git大法好,push需谨慎
Skip to content
Projects
Groups
Snippets
Help
Loading...
Help
Submit feedback
Contribute to GitLab
Sign in
Toggle navigation
L
liquidnet-bus-v1
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
董敬伟
liquidnet-bus-v1
Commits
86fc3936
Commit
86fc3936
authored
Apr 13, 2022
by
anjiabin
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
修改队列消费pending相关
parent
4ea9f164
Changes
1
Hide whitespace changes
Inline
Side-by-side
Showing
1 changed file
with
16 additions
and
3 deletions
+16
-3
RedisStreamConfig.java
...iquidnet.common.cache/redis/config/RedisStreamConfig.java
+16
-3
No files found.
liquidnet-bus-common/liquidnet-common-cache/liquidnet-common-cache-redis/src/main/java/com.liquidnet.common.cache/redis/config/RedisStreamConfig.java
View file @
86fc3936
...
@@ -2,8 +2,6 @@ package com.liquidnet.common.cache.redis.config;
...
@@ -2,8 +2,6 @@ package com.liquidnet.common.cache.redis.config;
import
lombok.extern.slf4j.Slf4j
;
import
lombok.extern.slf4j.Slf4j
;
import
lombok.var
;
import
lombok.var
;
import
org.springframework.beans.factory.annotation.Autowired
;
import
org.springframework.context.annotation.Configuration
;
import
org.springframework.data.redis.connection.RedisConnectionFactory
;
import
org.springframework.data.redis.connection.RedisConnectionFactory
;
import
org.springframework.data.redis.connection.stream.MapRecord
;
import
org.springframework.data.redis.connection.stream.MapRecord
;
import
org.springframework.data.redis.connection.stream.RecordId
;
import
org.springframework.data.redis.connection.stream.RecordId
;
...
@@ -11,7 +9,6 @@ import org.springframework.data.redis.connection.stream.StreamRecords;
...
@@ -11,7 +9,6 @@ import org.springframework.data.redis.connection.stream.StreamRecords;
import
org.springframework.data.redis.core.StreamOperations
;
import
org.springframework.data.redis.core.StreamOperations
;
import
org.springframework.data.redis.core.StringRedisTemplate
;
import
org.springframework.data.redis.core.StringRedisTemplate
;
import
org.springframework.data.redis.stream.StreamMessageListenerContainer
;
import
org.springframework.data.redis.stream.StreamMessageListenerContainer
;
import
org.springframework.stereotype.Component
;
import
java.net.InetAddress
;
import
java.net.InetAddress
;
import
java.net.UnknownHostException
;
import
java.net.UnknownHostException
;
...
@@ -54,6 +51,22 @@ public class RedisStreamConfig {
...
@@ -54,6 +51,22 @@ public class RedisStreamConfig {
.
builder
()
.
builder
()
.
pollTimeout
(
Duration
.
ofMillis
(
1
))
.
pollTimeout
(
Duration
.
ofMillis
(
1
))
.
build
();
.
build
();
// 创建配置对象
StreamMessageListenerContainer
.
StreamMessageListenerContainerOptions
<
String
,
MapRecord
<
String
,
String
,
String
>>
streamMessageListenerContainerOptions
=
StreamMessageListenerContainer
.
StreamMessageListenerContainerOptions
.
builder
()
// 一次性最多拉取多少条消息
.
batchSize
(
10
)
// 执行消息轮询的执行器
// .executor(this.threadPoolTaskExecutor)
// 消息消费异常的handler
.
errorHandler
(
t
->
{
// StreamListener中的异常会在这里抛出
System
.
out
.
println
(
"errorHandler: "
+
t
.
getMessage
());
})
// 超时时间,设置为0,表示不超时(超时后会抛出异常)
.
pollTimeout
(
Duration
.
ZERO
)
.
build
();
return
StreamMessageListenerContainer
.
create
(
factory
,
options
);
return
StreamMessageListenerContainer
.
create
(
factory
,
options
);
}
}
...
...
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