前言

最近因为工作需要在学习Dubbo的各种机制。其中深入学习了一下AbstractRegistry的实现机制。在此根据Dubbo源码对其实现进行一个总结。

Registry是干啥的

首先看一下dubbo最简单的架构图。架构图中一共有五个元素,而Registry类就是对注册中心的抽象。AbstractRegistry是对注册中心的一个抽象的实现。

可以看到它主要实现了RegistryServiceNode接口。这两个接口分别定义了节点属性如Url地址,是否可用,以及注册中心服务的属性如注册,注销,订阅,通知等等。

当服务启动时,会调用注册中心的register方法将自己的服务通过URL的方式发布到注册中心,而订阅其它服务时,会将订阅的服务通过URL发送给注册中心(URL中通常包含各种配置)。当服务需要优雅关闭时,会将自己从注册中心上解除注册。当服务出现变更时,会调用notify方法触发所有的监听器。

阅读全文 »

前言

最近在参与一个识别热点数据的需求开发。其中涉及了限流算法相关的内容。所以这里记录一下自己了解的各种限流算法,以及各个限流算法的实现。

限流算法的应用场景非常广泛,比如通过限流来确保下游配置较差的应用不会被上游应用的大量请求击穿,无论是HTTP请求还是RPC请求,从而使得服务保持稳定。限流也同样可以用于客户端,比如当我们需要从微博上爬取数据时,我们需要在请求中携带token从而通过微博的网关验证。但是微博为了防止服务被单个客户端大量访问,往往会在服务端进行限流,比如可能是一个token一个小时只能发起1000次请求。但是爬虫发出的请求通常远远不止这个量级。所以在客户端进行限流可以确保我们的token不会失效或是查封。

限流算法可以从多种角度分类,比如按照处理方式分为两种,一种是在超出限定流量之后会拒绝多余的访问,另一种是超出限定流量之后,只是报警或者是记录日志,访问仍然正常进行。

目前比较常见的限流算法有以下几种:

  • 固定窗口
  • 滑动窗口
  • 令牌桶算法
  • 漏桶算法

本文主要记录一下固定窗口和滑动窗口。令牌桶算法在谷歌的开源guava包中有实现,下次再开一篇文章分享一下。文中错误的地方欢迎指出!如果guava中实现了滑动窗口算法也请告诉我,急需,目前没有找到orz。

固定窗口

这是限流算法中最暴力的一种想法。既然我们希望某个API在一分钟内只能固定被访问N次(可能是出于安全考虑,也可能是出于服务器资源的考虑),那么我们就可以直接统计这一分钟开始对API的访问次数,如果访问次数超过了限定值,则抛弃后续的访问。直到下一分钟开始,再开放对API的访问。

所有的暴力算法的共同点都是容易实现,而固定窗口限流的缺点也同样很明显。假设现在有一个恶意用户在上一分钟的最后一秒和下一分钟的第一秒疯狂的冲击API。按照固定窗口的限流规则,这些请求都能够访问成功,但是在这一秒内,服务将承受超过规定值的访问冲击(这个规定值很可能是服务器能够承受的最大负载),从而导致服务无法稳定提供。而且因为用户在这一秒内耗光了上一分钟和下一分钟的访问定额,从而导致别的用户无法享受正常的服务,对于服务提供方来说是完全不能接收的。

这里自己做了一个简单的实现:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81

import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;

public class FixedWindowRateLimiter implements RateLimiter, Runnable {

private static final int DEFAULT_ALLOWED_VISIT_PER_SECOND = 5;

private final int maxVisitPerSecond;

private AtomicInteger count;

FixedWindowRateLimiter(){
this.maxVisitPerSecond = DEFAULT_ALLOWED_VISIT_PER_SECOND;
this.count = new AtomicInteger();
}

FixedWindowRateLimiter(int maxVisitPerSecond) {
this.maxVisitPerSecond = maxVisitPerSecond;
this.count = new AtomicInteger();
}

@Override
public boolean isOverLimit() {
return currentQPS() > maxVisitPerSecond;
}

@Override
public int currentQPS() {
return count.get();
}

@Override
public boolean visit() {
count.incrementAndGet();
System.out.print(isOverLimit());
return isOverLimit();
}

@Override
public void run() {
System.out.println(this.currentQPS());
count.set(0);
}

public static void main(String[] args) {
ScheduledExecutorService scheduledExecutorService = Executors.newSingleThreadScheduledExecutor();
FixedWindowRateLimiter rateLimiter = new FixedWindowRateLimiter();
scheduledExecutorService.scheduleAtFixedRate(rateLimiter, 0, 1, TimeUnit.SECONDS);
new Thread(new Runnable() {
@Override
public void run() {
while(true) {
rateLimiter.visit();
try {
Thread.sleep(100);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}).start();

new Thread(new Runnable() {
@Override
public void run() {
while(true) {
rateLimiter.visit();
try {
Thread.sleep(100);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}).start();

}
}

其中RateLimiter是一个通用的接口,后面的其它限流算法也会实现该接口:

1
2
3
4
5
6
7
8
public interface RateLimiter {

boolean isOverLimit();

int currentQPS();

boolean visit();
}

也可以不使用多线程的方式实现,更加简单高效:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
public class FixedWindowRateLimiterWithoutMultiThread implements RateLimiter {
private Long lastVisitAt = System.currentTimeMillis();
private static final int DEFAULT_ALLOWED_VISIT_PER_SECOND = 5;

private final int maxVisitPerSecond;

private AtomicInteger count;

public FixedWindowRateLimiterWithoutMultiThread(int maxVisitPerSecond){
this.maxVisitPerSecond = maxVisitPerSecond;
this.count = new AtomicInteger();
}

public FixedWindowRateLimiterWithoutMultiThread() {
this(DEFAULT_ALLOWED_VISIT_PER_SECOND);
}
@Override
public boolean isOverLimit() {
return count.get() > maxVisitPerSecond;
}

@Override
public int currentQPS() {
return count.get();
}

@Override
public boolean visit() {
long now = System.currentTimeMillis();
synchronized (lastVisitAt) {
if (now - lastVisitAt > 1000) {
lastVisitAt = now;
System.out.println(currentQPS());
count.set(1);
}
}
count.incrementAndGet();
return isOverLimit();
}

public static void main(String[] args) {
RateLimiter rateLimiter = new FixedWindowRateLimiterWithoutMultiThread();
new Thread(new Runnable() {
@Override
public void run() {
while(true) {
rateLimiter.visit();
try {
Thread.sleep(100);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}).start();

new Thread(new Runnable() {
@Override
public void run() {
while(true) {
rateLimiter.visit();
try {
Thread.sleep(100);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}).start();

}
}

滑动窗口

固定窗口就像是滑动窗口的一个特例。滑动窗口将固定窗口再等分为多个小的窗口,每一次对一个小的窗口进行流量控制。这种方法可以很好的解决之前的临界问题。

这里找的网上一个图,假设我们将1s划分为4个窗口,则每个窗口对应250ms。假设恶意用户还是在上一秒的最后一刻和下一秒的第一刻冲击服务,按照滑动窗口的原理,此时统计上一秒的最后750毫秒和下一秒的前250毫秒,这种方式能够判断出用户的访问依旧超过了1s的访问数量,因此依然会阻拦用户的访问。

使用定时任务实现的滑动窗口代码如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
public class SlidingWindowRateLimiter implements RateLimiter, Runnable{
private final long maxVisitPerSecond;

private static final int DEFAULT_BLOCK = 10;
private final int block;
private final AtomicLong[] countPerBlock;

private AtomicLong count;
private volatile int index;

public SlidingWindowRateLimiter(int block, long maxVisitPerSecond) {
this.block = block;
this.maxVisitPerSecond = maxVisitPerSecond;
countPerBlock = new AtomicLong[block];
for (int i = 0 ; i< block ; i++) {
countPerBlock[i] = new AtomicLong();
}
count = new AtomicLong(0);
}

public SlidingWindowRateLimiter() {
this(DEFAULT_BLOCK, DEFAULT_ALLOWED_VISIT_PER_SECOND);
}
@Override
public boolean isOverLimit() {
return currentQPS() > maxVisitPerSecond;
}

@Override
public long currentQPS() {
return count.get();
}

@Override
public boolean visit() {
countPerBlock[index].incrementAndGet();
count.incrementAndGet();
return isOverLimit();
}

@Override
public void run() {
System.out.println(isOverLimit());
System.out.println(currentQPS());
System.out.println("index:" + index);
index = (index + 1) % block;
long val = countPerBlock[index].getAndSet(0);
count.addAndGet(-val);
}

public static void main(String[] args) {
SlidingWindowRateLimiter slidingWindowRateLimiter = new SlidingWindowRateLimiter(10, 1000);
ScheduledExecutorService scheduledExecutorService = Executors.newSingleThreadScheduledExecutor();
scheduledExecutorService.scheduleAtFixedRate(slidingWindowRateLimiter, 100, 100, TimeUnit.MILLISECONDS);

new Thread(new Runnable() {
@Override
public void run() {
while (true) {
slidingWindowRateLimiter.visit();
try {
Thread.sleep(10);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}).start();

new Thread(new Runnable() {
@Override
public void run() {
while (true) {
slidingWindowRateLimiter.visit();
try {
Thread.sleep(10);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}).start();
}
}

参考文章

Protect Your API Resources with Rate Limiting

项目地址:
项目仓库地址

前言

临时性存储是容器的一个很大的买点。“根据一个镜像启动容器,随意变更,然后停止变更重启一个容器。你看,一个全新的文件系统又诞生了。”

在docker的语境下:

1
2
3
4
5
6
7
# docker run -it centos
[root@d42876f95c6a /]# echo "Hello world" > /hello-file
[root@d42876f95c6a /]# exit
exit
# docker run -it centos
[root@a0a93816fcfe /]# cat /hello-file
cat: /hello-file: No such file or directory

当我们围绕容器构建应用程序时,这个临时性存储非常有用。它便于水平扩展:我们只是从同一个镜像创建多个容器实例,每个实例都有自己独立的文件系统。它也易于升级:我们只是创建了一个新版本的映像,我们不必担心从现有容器实例中保留任何内容。它可以轻松地从单个系统移动到群集,或从内部部署移动到云:我们只需要确保集群或云可以访问registry中的镜像。而且它易于恢复:无论我们的程序崩溃对文件系统造成了什么损坏,我们只需要从镜像重新启动一个容器实例,之后就像从未发生过故障一样。

因此,我们希望容器引擎依然提供临时存储。但是从教程示例转换到实际应用程序时,我们确实会遇到问题。真实的应用必修在某个地方存储数据。通常,我们将状态保存到某个数据存储中(SQL或是NOSQL)。这也引来了同样的问题。数据存储也是位于容器中吗?理想情况下,答案是肯定的,这样我们可以利用和应用层相同的滚动升级,冗余和故障转移机制。但是,要在容器中运行我们的数据存储,我们再也不能满足于临时存储。容器实例需要能够访问持久存储。

阅读全文 »

1. 在所有用于where,order bygroup by的列上添加索引

索引除了能够确保唯一的标记一条记录,还能是MySQL服务器更快的从数据库中获取结果。索引在排序中的作用也非常大。

Mysql的索引可能会占据额外的空间,并且会一定程度上降低插入,删除和更新的性能。但是,如果你的表格有超过10行数据,那么索引就能极大的降低查找的执行时间。

强烈建议使用“最坏情况的数据样本”来测试MySql查询,从而更清晰的了解查询在生产中的行为方式。

假设你正在一个超过500行的数据库表中执行如下的查询语句:

1
mysql>select customer_id, customer_name from customers where customer_id='345546'

上述查询会迫使Mysql服务器执行一个全表扫描来获得所查找的数据。

型号,Mysql提供了一个特别的Explain语句,用来分析你的查询语句的性能。当你将查询语句添加到该关键词后面时,MySql会显示优化器对该语句的所有信息。

如果我们用explain语句分析一下上面的查询,会得到如下的分析结果:

1
2
3
4
5
6
mysql> explain select customer_id, customer_name from customers where customer_id='140385';
+----+-------------+-----------+------------+------+---------------+------+---------+------+------+----------+-------------+
| id | select_type | table | partitions | type | possible_keys | key | key_len | ref | rows | filtered | Extra |
+----+-------------+-----------+------------+------+---------------+------+---------+------+------+----------+-------------+
| 1 | SIMPLE | customers | NULL | ALL | NULL | NULL | NULL | NULL | 500 | 10.00 | Using where |
+----+-------------+-----------+------------+------+---------------+------+---------+------+------+----------+-------------+

可以看到,优化器展示出了非常重要的信息,这些信息可以帮助我们微调数据库表。首先,MySql会执行一个全表扫描,因为key列为Null。其次,MySql服务器已经明确表示它将要扫描500行的数据来完成这次查询。

为了优化上述查询,我们只需要在customer_id这一列上添加一个索引m即可:

1
2
3
mysql> Create index customer_id ON customers (customer_Id);
Query OK, 0 rows affected (0.02 sec)
Records: 0 Duplicates: 0 Warnings: 0

如果我们再次执行explain语句,会得到如下结果:

1
2
3
4
5
6
mysql> Explain select customer_id, customer_name from customers where customer_id='140385';
+----+-------------+-----------+------------+------+---------------+-------------+---------+-------+------+----------+-------+
| id | select_type | table | partitions | type | possible_keys | key | key_len | ref | rows | filtered | Extra |
+----+-------------+-----------+------------+------+---------------+-------------+---------+-------+------+----------+-------+
| 1 | SIMPLE | customers | NULL | ref | customer_id | customer_id | 13 | const | 1 | 100.00 | NULL |
+----+-------------+-----------+------------+------+---------------+-------------+---------+-------+------+----------+-------+

从上述的输出结果,显然MySQL服务器会使用索引customer_id来查询表格。可以看需要扫描的行数为1。虽然我只是在一个行数为500的表格中执行这条查询语句,索引在检索一个更大的数据集的时候优化程度更加明显。

阅读全文 »

前言

网上有非常多的关于红黑树理论的描述,本文的重点将不在于此,但是会在文中给出优秀文章的链接。对红黑树不了解的建议先阅读文章再看实现。本红黑树实现不支持多线程环境。因为删除操作灰常复杂,所以后续更新。源码在文末可以查看。

参考链接

https://www.geeksforgeeks.org/red-black-tree-set-3-delete-2/
https://www.geeksforgeeks.org/red-black-tree-set-2-insert/
http://www.codebytes.in/2014/10/red-black-tree-java-implementation.html
https://blog.csdn.net/eson_15/article/details/51144079

阅读全文 »

前言

本文为学习整理和参考文章,不具有教程的功能。其次,后面将会陆续更新各种应用的容器化部署的实践,如MySQL容器化,Jenkins容器化,以供读者参考。

阅读全文 »

前言

今天,我们将介绍一个比较新的设计模式(也就是没有在GoF那本书中出现过的一种设计模式),这个设计模式就是Event Bus设计模式。

起源

假设一个大型应用中,有大量的组件彼此间存在交互。而你希望能够在组件通信的同时能够满足低耦合和关注点分离原则。Event Bus设计模式是一个很好的解决方案。

Event Bus的概念和网络中的总线拓扑概念类似。即存在某种管道,而所有的电脑都连接在这条管道之上。其中的任何一台电脑发送的消息都将分发给总线上所有其它的主机。然后,每台主机决定是否接收还是抛弃掉这条消息。

在组件的层面上也是类似的:主机对应着应用的组件,消息对应于事件(event)或者数据。而管道是Event Bus对象。

阅读全文 »

前言

网上找过很多文章,关于通过docker构建mysql容器并将应用容器和docker容器关联起来的文章不多。本文将给出具体的范例。此处为项目的源码

前置条件

该教程要求在宿主机上配置了:

前言

在多线程中web应用很常见,尤其当你需要开发长期任务。

在Spring中,我们可以额外注意并使用框架已经提供的工具,而不是创造我们自己的线程。

Spring提供了TaskExecutor作为Executors的抽象。这个接口类似于java.util.concurrent.Executor接口。在spring中有许多预先开发好的该接口的实现,可以在官方文档中详细查看。

通过在Spring上下文中配置一个TaskExecutor的实现,你可以将你的TaskExecutor的实现注入到bean中,并可以在bean中访问到该线程池。

在bean中使用线程池的方式如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
package com.gkatzioura.service;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.ApplicationContext;
import org.springframework.core.task.TaskExecutor;
import org.springframework.stereotype.Service;
import java.util.List;
/**
* Created by gkatzioura on 4/26/17.
*/
@Service
public class AsynchronousService {
@Autowired
private ApplicationContext applicationContext;
@Autowired
private TaskExecutor taskExecutor;
public void executeAsynchronously() {
taskExecutor.execute(new Runnable() {
@Override
public void run() {
//TODO add long running task
}
});
}
}
阅读全文 »
0%