临时提交
parent
94436de3e1
commit
0dab86f0b8
|
@ -14,6 +14,7 @@ import org.springframework.beans.factory.annotation.Qualifier;
|
||||||
import org.springframework.data.redis.connection.Message;
|
import org.springframework.data.redis.connection.Message;
|
||||||
import org.springframework.data.redis.connection.MessageListener;
|
import org.springframework.data.redis.connection.MessageListener;
|
||||||
import org.springframework.data.redis.core.RedisTemplate;
|
import org.springframework.data.redis.core.RedisTemplate;
|
||||||
|
import org.springframework.scheduling.annotation.Scheduled;
|
||||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
||||||
import org.springframework.stereotype.Component;
|
import org.springframework.stereotype.Component;
|
||||||
|
|
||||||
|
@ -151,6 +152,7 @@ public class RedisRpcConfig implements MessageListener {
|
||||||
public RedisRpcResponse request(RedisRpcRequest request, int timeOut) {
|
public RedisRpcResponse request(RedisRpcRequest request, int timeOut) {
|
||||||
request.setSn((long) random.nextInt(1000) + 1);
|
request.setSn((long) random.nextInt(1000) + 1);
|
||||||
SynchronousQueue<RedisRpcResponse> subscribe = subscribe(request.getSn());
|
SynchronousQueue<RedisRpcResponse> subscribe = subscribe(request.getSn());
|
||||||
|
|
||||||
try {
|
try {
|
||||||
sendRequest(request);
|
sendRequest(request);
|
||||||
return subscribe.poll(timeOut, TimeUnit.SECONDS);
|
return subscribe.poll(timeOut, TimeUnit.SECONDS);
|
||||||
|
@ -209,4 +211,10 @@ public class RedisRpcConfig implements MessageListener {
|
||||||
public int getCallbackCount(){
|
public int getCallbackCount(){
|
||||||
return callbacks.size();
|
return callbacks.size();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Scheduled(fixedRate = 1000) //每1秒执行一次
|
||||||
|
public void execute(){
|
||||||
|
System.out.println("callbacks的长度: " + callbacks.size());
|
||||||
|
System.out.println("队列的长度: " + topicSubscribers.size());
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
@ -14,14 +14,19 @@ import com.genersoft.iot.vmp.media.zlm.dto.HookSubscribeForStreamChange;
|
||||||
import com.genersoft.iot.vmp.media.zlm.dto.MediaServerItem;
|
import com.genersoft.iot.vmp.media.zlm.dto.MediaServerItem;
|
||||||
import com.genersoft.iot.vmp.media.zlm.dto.hook.HookParam;
|
import com.genersoft.iot.vmp.media.zlm.dto.hook.HookParam;
|
||||||
import com.genersoft.iot.vmp.service.redisMsg.IRedisRpcService;
|
import com.genersoft.iot.vmp.service.redisMsg.IRedisRpcService;
|
||||||
|
import com.genersoft.iot.vmp.utils.SystemInfoUtils;
|
||||||
import com.genersoft.iot.vmp.vmanager.bean.ErrorCode;
|
import com.genersoft.iot.vmp.vmanager.bean.ErrorCode;
|
||||||
import com.genersoft.iot.vmp.vmanager.bean.WVPResult;
|
import com.genersoft.iot.vmp.vmanager.bean.WVPResult;
|
||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
import org.slf4j.LoggerFactory;
|
||||||
import org.springframework.beans.factory.annotation.Autowired;
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
import org.springframework.data.redis.core.RedisTemplate;
|
import org.springframework.data.redis.core.RedisTemplate;
|
||||||
|
import org.springframework.scheduling.annotation.Scheduled;
|
||||||
import org.springframework.stereotype.Service;
|
import org.springframework.stereotype.Service;
|
||||||
|
|
||||||
|
import java.util.List;
|
||||||
|
import java.util.Map;
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
public class RedisRpcServiceImpl implements IRedisRpcService {
|
public class RedisRpcServiceImpl implements IRedisRpcService {
|
||||||
|
|
||||||
|
|
Loading…
Reference in New Issue