This commit is contained in:
wyl
2022-05-16 23:27:50 +08:00
parent 55d934e294
commit 50a0a6b2f6
9 changed files with 252 additions and 0 deletions
+25
View File
@@ -0,0 +1,25 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<parent>
<artifactId>JavaBasiceDemo</artifactId>
<groupId>com.wyl.example</groupId>
<version>1.0-SNAPSHOT</version>
</parent>
<modelVersion>4.0.0</modelVersion>
<artifactId>DelayQueue</artifactId>
<properties>
<maven.compiler.source>8</maven.compiler.source>
<maven.compiler.target>8</maven.compiler.target>
</properties>
<dependencies>
<dependency>
<groupId>org.redisson</groupId>
<artifactId>redisson</artifactId>
</dependency>
</dependencies>
</project>
@@ -0,0 +1,32 @@
package com.wyl.delayqueue.redis.conf;
import org.redisson.Redisson;
import org.redisson.api.RedissonClient;
import org.redisson.config.Config;
/**
*
* @ClassName: RedissonConf
* @Date: 2022/3/3 22:37
* @author wangyl
* @version V1.0
*/
public class RedissonConf {
private volatile static RedissonClient redissonClient;
public static RedissonClient getRedissonClient() {
if (redissonClient == null) {
synchronized (RedissonConf.class) {
if (redissonClient == null) {
Config config = new Config();
config.useSingleServer()
.setAddress("redis://192.168.123.102:6379")
.setDatabase(0);
return Redisson.create(config);
}
}
}
return redissonClient;
}
}
@@ -0,0 +1,63 @@
package com.wyl.delayqueue.system.entity;
import lombok.Data;
import java.time.Duration;
import java.time.LocalDateTime;
import java.util.concurrent.Delayed;
import java.util.concurrent.TimeUnit;
/**
* @author: wangyl
* @date: 2022/5/14
* @description: 延时消息
*/
@Data
public class DelayMessage<T> implements Delayed {
/**
* 数据
*/
private T data;
/**
* 到期时间
*/
private LocalDateTime delayTime;
/**
* 获取到期的时间
* @param unit
* @return long
* @Date 2022/5/14
* @Author wangyl
*/
@Override
public long getDelay(TimeUnit unit) {
return unit.convert(Duration.between(LocalDateTime.now(),delayTime).toMillis(),TimeUnit.MILLISECONDS);
}
/**
* 队列中的元素是否到期
* @param o
* @return int
* @Date 2022/5/14
* @Author wangyl
*/
@Override
public int compareTo(Delayed o) {
return this.getDelay(TimeUnit.MILLISECONDS)>o.getDelay(TimeUnit.MILLISECONDS)?1:-1;
}
@Override
public String toString() {
String str = "任务失效时间为:"+delayTime.toString() +"任务内容为:"+data.toString();
return str;
}
}
@@ -0,0 +1,18 @@
package com.wyl.delayqueue.system.entity;
import lombok.Data;
/**
* @author: wangyl
* @date: 2022/5/14
* @description: 消息数据
*/
@Data
public class Message {
private Integer id;
private String message;
@Override
public String toString() {
String str = "id:" + id + " message:" + message;
return str;
}
}
@@ -0,0 +1,49 @@
package com.wyl.delayqueue.system;
import com.wyl.delayqueue.system.entity.DelayMessage;
import com.wyl.delayqueue.system.entity.Message;
import org.junit.jupiter.api.Test;
import java.time.LocalDateTime;
import java.util.concurrent.DelayQueue;
public class SystemDayTest {
@Test
public void test() throws InterruptedException {
Message message = new Message();
message.setId(1);
message.setMessage("hello world 1");
Message message1 = new Message();
message1.setId(2);
message1.setMessage("hello world 2");
Message message2 = new Message();
message2.setId(3);
message2.setMessage("hello world 3");
DelayMessage<Message> messageDelayMessage = new DelayMessage<>();
messageDelayMessage.setData(message);
messageDelayMessage.setDelayTime(LocalDateTime.now()
.plusSeconds(10));
DelayMessage<Message> messageDelayMessage1 = new DelayMessage<>();
messageDelayMessage1.setData(message1);
messageDelayMessage1.setDelayTime(LocalDateTime.now()
.plusSeconds(20));
DelayMessage<Message> messageDelayMessage2 = new DelayMessage<>();
messageDelayMessage2.setData(message2);
messageDelayMessage2.setDelayTime(LocalDateTime.now()
.plusSeconds(30));
DelayQueue<DelayMessage<Message>> delayQueue = new DelayQueue<>();
delayQueue.add(messageDelayMessage);
delayQueue.add(messageDelayMessage1);
delayQueue.add(messageDelayMessage2);
while (true) {
DelayMessage<Message> take = delayQueue.take();
System.out.println(take);
if (delayQueue.size()==0){
break;
}
}
}
}
+1
View File
@@ -16,6 +16,7 @@
<module>chain-of-responsibility</module>
<module>spring-boot</module>
<module>redis</module>
<module>DelayQueue</module>
</modules>
<packaging>pom</packaging>
<properties>
+13
View File
@@ -11,6 +11,7 @@
<modules>
<module>spring-boot-drools</module>
<module>spring-boot-xxjob</module>
<module>spring-boot-common</module>
</modules>
<artifactId>spring-boot</artifactId>
@@ -29,6 +30,18 @@
<type>pom</type>
<scope>import</scope>
</dependency>
<dependency>
<groupId>com.baomidou</groupId>
<artifactId>mybatis-plus-boot-starter</artifactId>
<version>3.3.2</version>
</dependency>
</dependencies>
</dependencyManagement>
<dependencies>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<scope>provided</scope>
</dependency>
</dependencies>
</project>
@@ -0,0 +1,18 @@
package com.wyl.springbootmybatis;
import org.junit.runner.RunWith;
import org.mybatis.spring.annotation.MapperScan;
import org.mybatis.spring.annotation.MapperScans;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.context.annotation.ComponentScan;
import org.springframework.test.context.junit4.SpringRunner;
@RunWith(SpringRunner.class)
@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.NONE, classes = BaseTest.Application.class)
public class BaseTest {
@ComponentScan(value = "com.wyl.springbootmybatis")
@MapperScan({"com.wyl.springbootmybatis.mybatis.mapper"})
public static class Application {
}
}
@@ -0,0 +1,33 @@
package com.wyl.springbootmybatis.mybatis;
import com.wyl.springbootmybatis.BaseTest;
import com.wyl.springbootmybatis.mybatis.domain.MybatisTest;
import com.wyl.springbootmybatis.mybatis.service.MybatisTestService;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import java.util.List;
import java.util.UUID;
class SpringBootMybatisApplicationTests extends BaseTest {
@Autowired
MybatisTestService mybatisTestService;
@Test
void getDataTest() {
List<MybatisTest> list = mybatisTestService.list();
System.out.println(list);
}
@Test
void saveDataTest() {
for(int i =0;i<300000;i++)
{
MybatisTest mybatisTest = new MybatisTest();
mybatisTest.setCode(UUID.randomUUID().toString());
mybatisTestService.save(mybatisTest);
}
}
}