rocketMQ design

Motivation

Persistence

和kafka的思想一样,利用顺序IO。

但是,当一个broker上有多个partition的时候,顺序又变成了随机。

RocketMQ为了解决这个问题,采用了单一的日志文件。即把一台机器上的所有的topic的所有queue的消息都存放在同一个文件里面。

先写入Commit log文件里面(单个文件),然后有后台线程异步的同步到ConsumeQueue(也是一个文件),再由Consumer进行消费。这是RocketMQ的方案。

4.2 Persistence

As a result the performance of linear writes on a JBOD configuration with six 7200rpm SATA RAID-5 array is about 600MB/sec but the performance of random writes is only about 100k/sec—a difference of over 6000X.

磁盘的顺序写和随机写,性能相差6000倍,sequential disk access can in some cases be faster than random memory access!

  • OS pagecache

    现代操作系统很乐于使用所有的空闲内存来做disk caching。所有的磁盘读写都会通过这些cache进行。

  • Furthermore, we are building on top of the JVM, and anyone who has spent any time with Java memory usage knows two things:

     The memory overhead of objects is very high, often doubling the size of the data stored (or worse).
    
    Java garbage collection becomes increasingly fiddly and slow as the in-heap data increases.
    

基于上面的两个原因, 得出的结论:

using the filesystem and relying on pagecache is superior to maintaining an in-memory cache

直接使用带pagecache的OS filesystem 甚至性能会比使用内存cache要更好。

This style of pagecache-centric design is described in an article on the design of Varnish

4.3 Efficiency

For more background on the sendfile and zero-copy support in Java, see this article.

4.4 The Producer

任意一个broker都保存着metadate,关于哪些节点是活着的,还有一个topic的partition的leader是谁?所以producer可以找到对应的leader,直接向其发送消息。

4.5 The Consumer

consumer执行‘fetch’操作,同时带着offset,来向leader拉去消息。

offset:

一般的队列都会在broker端记录consumer消费的位置,这样做可以及时的删除消费掉的消息,但是维护这个位置是很费事的,还需要consumer返回ack。

kafka也会记录一个offset,表示下一个可以消费的消息的位置,而且可以‘倒带’,即重现消费之前消费过的消息。

4.6 Message Delivery Semantics

4.7 Replication

BFS and DFS

如何表示一个图

临接列表(adjacency lists)和临接矩阵(adjacency matrix) ,有向图和无向图都可以用这样的方式来表示。

最简单的图的遍历算法之一,而且还是很多重要的其他算法的基础。Prim的最小生成树(minimum-spanning-tree) 和 Dijksra的最短路径算法(single-source shortest-paths)都使用了这个想法。

广度优先搜索在进一步遍历图中顶点之前,先访问当前顶点的所有邻接结点。

a .首先选择一个顶点作为起始结点,并将其染成灰色,其余结点为白色。

b. 将起始结点放入队列中。

c. 从队列首部选出一个顶点,并找出所有与之邻接的结点,将找到的邻接结点放入队列尾部,将已访问过结点涂成黑色,没访问过的结点是白色。如果顶点的颜色是灰色,表示已经发现并且放入了队列,如果顶点的颜色是白色,表示还没有发现

d. 按照同样的方法处理队列中的下一个结点。 基本就是出队的顶点变成黑色,在队列里的是灰色,还没入队的是白色。

算法导论中的伪代码

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
//图G,从s开始广度优先遍历
BFS(G, s)
//初始化s之外的所有其他节点
for s 以外的所有定点 u:
u.color = WHITE
u.d = -1
u.p = NIL //p代表父节点
//初始化s
s.color = GRAY
s.d = 0
s.p = NIL
enqueue(s)
while (u = dequeue()) != NIL
for v in u的所有临接节点:
if v.color == WHITE
v.color = GRAY //入队列前设置为灰
v.d = u.d + 1
v.p = u
enqueue(v)
u.color = BLACK //访问过的设置为黑

因为每个节点都入队一次,而且,每个节点出来之后,都会对他的临接节点进行访问。所有 最后的时间复杂度为O(E+V)

在《算法》中,使用marked[]来代替对节点上色,用edgeTo[]代替父节点

参考:https://www.cnblogs.com/xiehongfeng100/p/4461772.html

深度优先搜索 Depth-first-seatch

深度优先搜索在搜索过程中访问某个顶点后,需要递归地访问此顶点的所有未访问过的相邻顶点。

初始条件下所有节点为白色,选择一个作为起始顶点,按照如下步骤遍历:

a. 选择起始顶点涂成灰色,表示还未访问

b. 从该顶点的邻接顶点中选择一个,继续这个过程(即再寻找邻接结点的邻接结点),一直深入下去,直到一个顶点没有邻接结点了,涂黑它,表示访问过了

c. 回溯到这个涂黑顶点的上一层顶点,再找这个上一层顶点的其余邻接结点,继续如上操作,如果所有邻接结点往下都访问过了,就把自己涂黑,再回溯到更上一层。

d. 上一层继续做如上操作,知道所有顶点都访问过。

算法导论中的实现:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
DFS(G)
//对所有节点做初始化
for u in G.V: // O(V)
u.color = WHITE
u.p = NIL
for u in G.V: // O(V)
if u.color == WHITE
DFS-Visit(G,u)

GFS-Visit(G,u)
//先设置为灰色
u.color = GRAY
//尝试所有的临接节点,
for v in u的临接节点: //adj[v]次
if v.color == WHITE
v.p = u
DFS-Visit(G, v)
//在自己的所有临接节点都访问完成后,设置为黑色
u.color = BLACK

DFS的时间复杂度为O(V+E)。

Kafka design

Kafka官方文档翻译以及扩展

官方文档

中文版本:https://www.cnblogs.com/coprince/p/5893066.html

其他参考:https://colobu.com/2017/11/02/kafka-replication/

4.1 Motivation

for handling all the real-time data feeds

  • high-throughput
  • large data backlogs 大量的数据积压
  • low-latency delivery 满足传统的队列的用法

4.2 Persistence

As a result the performance of linear writes on a JBOD configuration with six 7200rpm SATA RAID-5 array is about 600MB/sec but the performance of random writes is only about 100k/sec—a difference of over 6000X.

磁盘的顺序写和随机写,性能相差6000倍,sequential disk access can in some cases be faster than random memory access!

  • OS pagecache

    现代操作系统很乐于使用所有的空闲内存来做disk caching。所有的磁盘读写都会通过这些cache进行。

  • Furthermore, we are building on top of the JVM, and anyone who has spent any time with Java memory usage knows two things:

     The memory overhead of objects is very high, often doubling the size of the data stored (or worse).
    
    Java garbage collection becomes increasingly fiddly and slow as the in-heap data increases.
    

基于上面的两个原因, 得出的结论:

using the filesystem and relying on pagecache is superior to maintaining an in-memory cache

直接使用带pagecache的OS filesystem 甚至性能会比使用内存cache要更好。

This style of pagecache-centric design is described in an article on the design of Varnish

partition存储细节:

broker的机器上,xxx/message-folder为数据文件存储根目录,在这个目录下, 每个partition一个文件夹。

在partition的文件夹下,有多个segment:

  • segment file组成:由2大部分组成,分别为index file和data file,此2个文件一一对应,成对出现,后缀”.index”和“.log”分别表示为segment索引文件、数据文件.
  • segment文件命名规则:partion全局的第一个segment从0开始,后续每个segment文件名为上一个segment文件最后一条消息的offset值。数值最大为64位long大小,19位数字字符长度,没有数字用0填充。

image-20181105233252954

image-20181105233332064

segment data file由许多message组成,每个message物理结构如下:

image-20181105233407991

4.3 Efficiency

For more background on the sendfile and zero-copy support in Java, see this article.

4.4 The Producer

任意一个broker都保存着metadate,关于哪些节点是活着的,还有一个topic的partition的leader是谁?所以producer可以找到对应的leader,直接向其发送消息。

4.5 The Consumer

consumer执行‘fetch’操作,同时带着offset,来向leader拉去消息。

offset:

一般的队列都会在broker端记录consumer消费的位置,这样做可以及时的删除消费掉的消息,但是维护这个位置是很费事的,还需要consumer返回ack。

kafka也会记录一个offset,表示下一个可以消费的消息的位置,而且可以‘倒带’,即重现消费之前消费过的消息。

4.6 Message Delivery Semantics

4.7 Replication

  • The unit of replication is the topic partition

    relication的最小单位是topic的partition。

  • 支持automatic faiover

    TODO 怎么支持的

  • 不像其他队列,replication的作用就只有做副本,无副作用。例如不会提供读。

  • followers从leader消费数据,就好像是一个kafka的consumer一样。

  • 节点存活的定义

    1 A node must be able to maintain its session with ZooKeeper

    2 If it is a slave it must replicate the writes happening on the leader and not fall “too far” behind

  • leader 维护ISR (in sync replication)

    如果一个follower挂了,阻塞了或者落后了,leader会把他从ISR中删除。

  • kafka不解决拜占庭问题。

commit

定义:

所有的ISR的节点都提交了log之后,才说一个message被committed。只有committed message才会提供给consumer。

对于consumer来说:

consumer不会担心看到过的message 会因为leader切换而丢失。

对于producer:

可以选择是否要等待message被committed,是latency 和 durability之间的tradeoff。

The guarantee that Kafka offers is that a committed message will not be lost, as long as there is at least one in sync replica alive, at all times.

只要有一个ISR存活,就可以保障committed message 不会丢失。

Kafka will remain available in the presence of node failures after a short fail-over period, but may not remain available in the presence of network partitions.

kafka在节点短期失效的情况下可以保持可用,但是在存在网络分区的情况下,是不可用的。

TODO:怎么理解

Quorum algorithms( majority vote)

简单来说,一半以上的follower commit了才能保障选举出来的leader有最新的commit。

好处:latency只依赖于最快的node。因为只要有一半ack了,就是commit成功,最快的最先完成。

坏处:写的增多,整体吞吐量下降,使得他不适合于大量的写的情况。举个例子:如果为了可以容忍2个节点的实效,就必须要5份数据的copy。 这个坏处也就可以解释为什么 quorum算法被经常用于类似zk这样的配置共享集群,而很少用于data storage。

kafka没有使用quorum算法。

Instead of majority vote, Kafka dynamically maintains a set of in-sync replicas (ISR) that are caught-up to the leader.

作为替代,kafka维护了一个ISR,其中每个replication都和leader保持一致。这个ISR是被leader维护的,且会保存到zk上。

只有ISR的成员,才会是leader的候选人。

This ISR set is persisted to ZooKeeper whenever it changes.

这个IRS集合是被持久化到zk。

Another important design distinction is that Kafka does not require that crashed nodes recover with all their data intact

另外一个kafka的特点是,crashed的nodes恢复后,不用恢复所有的数据。

Unclean leader election

如果所有的ISR都挂了,怎么办?

  1. Wait for a replica in the ISR to come back to life and choose this replica as the leader (hopefully it still has all its data).
  2. Choose the first replica (not necessarily in the ISR) that comes back to life as the leader.

Replica Management

We attempt to balance partitions within a cluster in a round-robin fashion to avoid clustering all partitions for high-volume topics on a small number of nodes

平衡的将partition分布到不同的broker中,用round-robin的方式。

we elect one of the brokers as the “controller”. This controller detects failures at the broker level and is responsible for changing the leader of all affected partitions in a failed broker. If the controller fails, one of the surviving brokers will become the new controller.

为每个partition进行leader选举是非常低效的,因为部署单元是broker,所以选一个borker当controller,它会去在broker层面去检测失效,并且还会去负责为失效的broker中受影响的partitions去做leader选举。如果controller broker挂了,就会选出一个新的。

实际上,Kafka选举一个broker作为controller,这个controller通过 watch Zookeeper检测所有的broker failure,并负责为所有受影响的parition选举leader,再将相应的leader调整命令发送至受影响的broker。如果 controller失败了,幸存的所有broker都会尝试在Zookeeper中创建/controller->{this broker id},如果创建成功(只可能有一个创建成功),则该broker会成为controller,若创建不成功,则该broker会等待新 controller的命令。

Red Black Tree

定义:

1 二分查找树

2 每个node有一个额外的存储:color (红,黑)

3 root 是黑

4 原本在BST中指向NULL的pointer,在RBT中,全部指向了NIL,每个NIL是黑

5 如果一个node是红,它的孩子都是黑。

6 任意node开始,到NIL的所有路径中,包含了相同的黑node。

image-20180820170413020

###

Fix Up

什麼情況需要對InsertRBT()做修正? 當新增node接在紅色的node的child pointer,形成紅色與紅色相連時。

  • node(X)為其parent,顏色為紅色;
  • node(Y)為其uncle,其顏色可能為紅色或黑色
  • node(Z)為其grandparent,顏色必定為黑色(因為node(X)是紅色)

image-20180820193242724

###

根據uncle的顏色是紅色或者黑色,可以將修正(FixUp)分成三種情形(case):

  1. Case1:uncle是紅色,不論新增的node是node(X)的leftchildrightchild
  2. Case2:uncle是黑色,而且新增的node為node(X)的rightchild
  3. Case3:uncle是黑色,而且新增的node為node(X)的leftchild

旋转(Rotation)

搜索树的操作:Tree-Insert 和 Tree-Delete,会改变树的结构,通过旋转进行恢复。

image-20180815004944884

有左和右两种方式。

插入(Insertion)

We can insert a node into an n-node red-black tree in O(lg n) time.

TODO

参考

http://alrightchiu.github.io/SecondRound/red-black-tree-introjian-jie.html

http://alrightchiu.github.io/SecondRound/red-black-tree-insertxin-zeng-zi-liao-yu-fixupxiu-zheng.html

JVM最大线程数

JVM最大线程数

至于操作系统栈大小(ulimit -s):这个配置只影响进程的初始线程;后续用pthread_create创建的线程都可以指定栈大小。HotSpot VM为了能精确控制Java线程的栈大小,特意不使用进程的初始线程(primordial thread)作为Java线程。

不显式设置-Xss或-XX:ThreadStackSize时,在Linux x64上ThreadStackSize的默认值就是1024KB

  • StackOverflowError (the stack size is greater than the limit), increase the value
  • OutOfMemoryError: unable to create new native thread (too many threads, each thread has a large stack), decrease it.

JVM中可以生成的最大数量由JVM的堆内存大小、Thread的Stack内存大小、系统最大可创建的线程数量(Java线程的实现是基于底层系统的线程机制来实现的,Windows下_beginthreadex,Linux下pthread_create)三个方面影响。

具体数量可以根据Java进程可以访问的最大内存(32位系统上一般2G)、堆内存、Thread的Stack内存来估算。

1
(MaxProcessMemory - JVMMemory – ReservedOsMemory) / (ThreadStackSize) = Number of threads
  • MaxProcessMemory : 进程的最大寻址空间
  • JVMMemory : JVM内存
  • ReservedOsMemory : 保留的操作系统内存,如Native heap,JNI之类,一般100多M
  • ThreadStackSize : 线程栈的大小,jvm启动时由Xss指定

综上所述:

jvm的最大线程数受一下jvm参数限制

-Xms 最小堆内存 -Xmx 最大堆内存 -Xss 设置每个线程的堆栈大小。JDK5.0以后每个线程堆栈大小为1M

线程数量 = (机器本身可用内存 - JVM分配的堆内存) / Xss的值。

补充下: jvm线程数还和系统的最大线程限制有关:

操作系统限制 系统最大可开线程数,主要受以下几个参数影响

​ /proc/sys/kernel/pid_max

​ /proc/sys/kernel/thread-max 系统可以生成最大线程数量

​ /proc/sys/vm/max_map_count 文件包含限制一个进程可以拥有的VMA(虚拟内存区域)的数量 。单个JVM能开启的最大线程数是/proc/sys/vm/max_map_count的设置数的一半

https://segmentfault.com/a/1190000004694232

https://www.oschina.net/translate/understanding-virtual-memory?print

Mysql explain

#Mysql Explain type

对mysql explain后的type解释:

原文:https://dev.mysql.com/doc/refman/5.6/en/explain-output.html#explain-join-types

  • system

    表只有一行的情况。

  • const

    只有一行匹配,并且在查询的开头就命中

    当使用主键或者唯一索引只匹配一条数据的时候。

  • eq_ref

    多表关联查询的时候,只匹配到一行。

  • ref

    根据索引匹配到多行。 “等于”或者”不等于”判断。

  • unique_subquery

    使用了In的子查询, 并且子查询是主键或者唯一索引

  • index_subquery

    和unique_subquery类似,但是子查询是非唯一索引

  • range

    根据索引进行范围查找的

  • index

    对索引的扫描

    只select索引中的列,不加where条件,就是index类型。

  • all

    对全表的扫描

FeignClient

从源码分析FeignClient

1
public abstract class Feign

Feign这个类是为了简化对http apis的请求。

通过 newInstance 方法,产生Target,来代表对应的Http Apis。

1
public abstract <T> T newInstance(Target<T> target);

产生一个 http api的实例。

1
2
3
4
5
6
7
8
9
10
11
12
13
public interface Target<T> {
/* The type of the interface this target applies to. */
Class<T> type();

/* configuration key associated with this target. */
String name();

/* base HTTP URL of the target. */
String url();

//从 template中产生Request
public Request apply(RequestTemplate input);
}

apply被用来产生request,基于传入的template input,加入header和参数,产生一个不可变的Request。

例如:

1
2
3
4
5
public Request apply(RequestTemplate input) {
input.insert(0, url());
input.replaceHeader(&quot;X-Auth&quot;, currentToken);
return input.asRequest();
}

target的主要作用是产生Request, 作为Http请求的参数。

Feign的一个默认实现:ReflectiveFeign,newInstance最终返回了一个动态代理:

1
2
3
Proxy.newProxyInstance(ClassLoader loader,
Class<?>[] interfaces,
InvocationHandler h)

最重要的是要看InvocationHandler的实现。

在初始化 Feign的时候,在Feign.Builder中,选择的FeignInvocationHandler作为InvocationHandler的实现。

FeignInvocationHandler中,包含多个MethodHandler, 默认实现为:SynchronousMethodHandler。

在SynchronousMethodHandler中,包含着Client和Retryer,这个类中,实现了最终的Http的请求。

MethodHandler

SynchronousMethodHandler , 在这个类里面最终实现了对 http client 的调用。

1
2
3
4
class SynchronousMethodHandler implements MethodHandler {
private final Client client;
private final Retryer retryer;
}

InvocationHandler

1
2
3
4
class FeignInvocationHandler implements InvocationHandler {

private final Target target;
private final Map<Method, MethodHandler> dispatch;

同时,在Feign.Builder中,

Client使用的是java默认的java.net.HttpUrlConnection

Retryer,的默认值为重试5次。

Feign-hystrix

对Feign的扩展:

1 容许Feign的接口返回HystrixCommand 或者 rx.Observable (通过HystrixDelegatingContract 来实现)

2 对接口调用进行包装,加入了断路器。(通过HystrixInvocationHandler实现)

在HystrixInvocationHandler中,invoke方法,对httpclient的调用,使用HystrixCommand的方式:

1
2
3
4
5
6
7
8
9
10
11
12
HystrixCommand<Object> hystrixCommand = new HystrixCommand<Object>(setterMethodMap.get(method)) {
@Override
protected Object run() throws Exception {
try {
return HystrixInvocationHandler.this.dispatch.get(method).invoke(args);
} catch (Exception e) {
throw e;
} catch (Throwable t) {
throw (Error) t;
}
}
}

与spring的集成

在spring-cloud-netflix-core中,有spring和feign集成的代码。

在FeignClientFactoryBean中,重新定义了 Feign.Builder:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
class FeignClientFactoryBean implements FactoryBean<Object>, InitializingBean,
ApplicationContextAware {

@Override
public Object getObject() throws Exception {
FeignContext context = applicationContext.getBean(FeignContext.class);
//从FeignContext中生成Feign.Builder
Feign.Builder builder = feign(context);
if (!StringUtils.hasText(this.url)) {
String url;//处理url
//loadbalance ribbon?
return loadBalance(builder, context,
new HardCodedTarget<>(this.type,
this.name, url));
}

}

在spring中,Feign的初始化依赖于FeignContext。

1
2
3
4
5
6
7
public class FeignContext extends NamedContextFactory<FeignClientSpecification> {

public FeignContext() {
super(FeignClientsConfiguration.class, "feign", "feign.client.name");
}

}

NamedContextFactory可以创建一组子上下文, 每个子上下文中可以使用一组的Specification来定义bean

Ribbon

一个客户端负载均衡器,运行在客户端上。

###LoadBalancerClient

最重要的一个类: LoadBalancerClient 负载均衡器的客户端。

继承关系: ServiceInstanceChooser <- LoadBalancerClient <- RibbonLoadBalancerClient

1
2
3
4
5
6
7
8
9
10
11
/**
* Represents a client side load balancer
*/
public interface LoadBalancerClient extends ServiceInstanceChooser {

/**
* execute request using a ServiceInstance from the LoadBalancer for the specified
* service
*/
<T> T execute(String serviceId, LoadBalancerRequest<T> request) throws IOException;
//......
1
2
3
4
5
6
7
8
9
10
11
/**
* Implemented by classes which use a load balancer to choose a server to
* send a request to.
*/
public interface ServiceInstanceChooser {

/**
* Choose a ServiceInstance from the LoadBalancer for the specified service
*/
ServiceInstance choose(String serviceId);
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
 class RibbonLoadBalancerClient implements LoadBalancerClient {
@Override
public <T> T execute(String serviceId, LoadBalancerRequest<T> request) throws IOException {
//最终选择service的任务还是交与ILoadBalancer来做。
ILoadBalancer loadBalancer = getLoadBalancer(serviceId);
Server server = getServer(loadBalancer);
if (server == null) {
throw new IllegalStateException("No instances available for " + serviceId);
}
RibbonServer ribbonServer = new RibbonServer(serviceId, server, isSecure(server,
serviceId), serverIntrospector(serviceId).getMetadata(server));

return execute(serviceId, ribbonServer, request);
}

###ILoadBalancer

最终选择service的任务还是交与ILoadBalancer来做。

ILoadBalancer <- BaseLoadBalancer <- DynamicServerListLoadBalancer

在DynamicServerListLoadBalancer中,是如何获取和刷新服务列表的?

首先在构造函数中->initWithNiwsConfig() -> restOfInit() -> updateListOfServers()

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
class DynamicServerListLoadBalance{
@VisibleForTesting
public void updateListOfServers() {
List<T> servers = new ArrayList<T>();
if (serverListImpl != null) {
servers = serverListImpl.getUpdatedListOfServers();
LOGGER.debug("List of Servers for {} obtained from Discovery client: {}",
getIdentifier(), servers);

if (filter != null) {
servers = filter.getFilteredListOfServers(servers);
LOGGER.debug("Filtered List of Servers for {} obtained from Discovery client: {}",
getIdentifier(), servers);
}
}
updateAllServerList(servers);
}

ServerList 的具体实现类: serverListImpl 来负责最终的刷新。ServerList 有各种实现,比方说用Consul的话,实现就是ConsulServerList。

参考:https://blog.csdn.net/forezp/article/details/74820899

FeignLoadBalancer

CachingSpringLoadBalancerFactory 会返回一个 FeignLoadBalancer

Robbin会retry:

在CachingSpringLoadBalancerFactory中,创建FeignLoadBalancer的时候

1
2
3
4
5
6
7
8
9
10
11
12
13
public FeignLoadBalancer create(String clientName) {
if (this.cache.containsKey(clientName)) {
return this.cache.get(clientName);
}
IClientConfig config = this.factory.getClientConfig(clientName);
ILoadBalancer lb = this.factory.getLoadBalancer(clientName);
ServerIntrospector serverIntrospector = this.factory.getInstance(clientName, ServerIntrospector.class);
//通过enableRetry来控制是否重试
FeignLoadBalancer client = enableRetry ? new RetryableFeignLoadBalancer(lb, config, serverIntrospector,
loadBalancedRetryPolicyFactory) : new FeignLoadBalancer(lb, config, serverIntrospector);
this.cache.put(clientName, client);
return client;
}

Feign and Ribbon

当结合使用Feign和Ribbon的时候, SynchronousMethodHandler 中的client 类型为:LoadBalancerFeignClient。 在execute的时候,会创建一个FeignLoadBalancer 来执行executeWithLoadBalancer(),这个函数的作用是把请求交给选中的服务来处理,而不是指定一个服务。FeignLoadBalancer包含一个ILoadBalancer

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
class LoadBalancerFeignClient implements Client{
@Override
public Response execute(Request request, Request.Options options) throws IOException {
try {
URI asUri = URI.create(request.url());
String clientName = asUri.getHost();
URI uriWithoutHost = cleanUrl(request.url(), clientName);
FeignLoadBalancer.RibbonRequest ribbonRequest = new FeignLoadBalancer.RibbonRequest(
this.delegate, request, uriWithoutHost);

IClientConfig requestConfig = getClientConfig(options, clientName);
return lbClient(clientName).executeWithLoadBalancer(ribbonRequest,
requestConfig).toResponse();
}
catch (ClientException e) {
IOException io = findIOException(e);
if (io != null) {
throw io;
}
throw new RuntimeException(e);
}
}

}
1
2
3
4
5
6
7
8
9
10
11
public T executeWithLoadBalancer(final S request, final IClientConfig requestConfig) throws ClientException {
RequestSpecificRetryHandler handler = getRequestSpecificRetryHandler(request, requestConfig);
LoadBalancerCommand<T> command = LoadBalancerCommand.<T>builder()
.withLoadBalancerContext(this)
.withRetryHandler(handler)
.withLoadBalancerURI(request.getUri())
.build();

try {
//在这个里面选择一个服务进行调用。
return command.submit(

HTTP 2

http2协议不是为了推翻http1,重新设计,它的目标是为了性能的提升等等,一个最主要的目标是容许用户使用一个连接来访问网站。

http2的官方网站:https://http2.github.io/

RFC 7540

https://httpwg.org/specs/rfc7540.html

Hypertext Transfer Protocol Version 2 (HTTP/2)

之前:

In particular, HTTP/1.0 allowed only one request to be outstanding at a time on a given TCP connection. HTTP/1.1 added request pipelining, but this only partially addressed request concurrency and still suffers from head-of-line blocking.

Therefore, HTTP/1.0 and HTTP/1.1 clients that need to make many requests use multiple connections to a server in order to achieve concurrency and thereby reduce latency

Http 头 :

1 无法压缩

2 每次请求都重复传递

HTTP2 Overview

http2支持1的所有核心特性,并且更加高效。Frame是http2中的最基本的单位。

Stream

stream 是一个独立的,双向的Frame的队列,这些frame在server 和client之间通过Http/2 传送。

stream有以下特性:

1 一个HTTP/2链接可以包含多个stream

2 stream可以被client和server创建

3 stream可以被client和server关闭

4 frame的顺序很重要

5 streams 被一个数字唯一标识

stream的状态

idle

所有的stream开始于这个状态。

  • 接受或者发出一个 HEADERS frame 会导致这个steam 变成”open”
  • 向其他的stream发出一个 PUSH_PROMISE frame,本stream会变为”reserved(local)”状态
  • 从其他的stream收到一个 PUSH_PROMISE frame,本stream会变为”reserved(remote)”状态

reserved(local/remote)

用于实现push,暂时不讨论

open

处于这个状态的stream,可以被两方使用,发送各种frames。

  • 双方都可以发送一个标记了 END_STREAM的frame,来使stream进入”half-closed”状态。发送 END_STREAM的一方,stream进入”half-closed(local)”状态,收到 END_STREAM的一方,stream进入”half-closed(remote)”状态。
  • 双方也可以发送一个RST_STREAM frame,导致双方进入”closed”状态

half-closed(local)

half-closed(remote) –END_STREAM flag –> half-close(local) 会是两边都进入“close”状态

或者收到一个 RST_STREAM ,也会直接进入“close”状态。

half-closed 有什么用处?

Stream Identifiers

stream 的标识符是一个31bit的无符号int。client端使用奇数,server端使用偶数。id为0的stream被用来作为控制信息的连接。

新创建的stream的标号必须比当前endpoint上现有的要大。

Stream concurrency

TODO

Flow Control

HTTP/2通过WINDOW_UPDATE frame来支持流量控制。

HTTP/2“流”的流量控制的目标是:在不改变协议的情况下允许使用多种流量控制算法。

HTTP/2的流量控制有一下特征:

  • 流量控制是特定于一个连接的。
  • 流量控制是基于WINDOW_UPDATE帧的。接收方公布自己打算在每个流以及整个连接上分别接收多少字节。这是一个以信用为基础的方案。
  • 流量控制是有方向的,由接收者全面控制。
  • 无论是新流还是整个连接,流量控制窗口的初始值是65535字节
  • 帧的类型决定了流量控制是否适用于帧。目前,只有DATA帧服从流量控制,所有其它类型的帧并不消耗流量控制窗口的空间。这保证了重要的控制帧不会被流量控制阻塞。
  • 流量控制不能被禁用

DBs replication

replication的定义:在不同的联网的机器上,保存相同的数据。

需要replication的理由

  • 使数据地理上离用户更近
  • 保障系统的可用性
  • 应付大量的读请求

replication的主要难点在于处理数据的变更

在实现replication的时候,有许多可以权衡的点:同步或者异步,怎么处理失效的replicas。

主从模型

被好多数据库产品采用:mysql, PostgresSQL,MongoDB,RethinkDB。

同步异步

其中的一个重要的细节:备份是同步还是异步发生的,在关系数据库中,这个选项大多是可配置的,其他的数据库,一般会选择其中一种情况。所有的从节点的同步,无论是全同步,或是全异步都有问题:全同步的情况,如果一个节点断了,就会导致整个系统的写入变慢;全异步的情况,数据又容易丢失。一般采用混合的方式:一个从同步写,其他的采用异步,如果同步写的从节点断线了,可以立刻用其他的节点来补上,这样的方式,通常被叫做:semi-synchronous 。在持久性上的弱化,莫阿斯

新增从节点

不管是扩容还是替换掉废节点,都需要做这个操作。怎么保证新增的从节点和主节点的同步?

在不停服务的情况下,一般这样操作:

1 获取主节点某个时刻的snapshot。

2 复制这个snapshot到新增的从节点。

3 新从节点和主节点获取时间点之后的所有更新。这个需要一个准确的位置,在PostgresSQL中叫’log sequence number’,在mysql中叫‘binlog coordinates’。

4 当从snapshot开始的所有数据变更都被同步了之后,这个从节点可以像正常的从节点一样了。

处理节点掉线

为了保证整个系统长时间运行,需要有这样的能力,保证任何节点 (主从)挂掉都不会影响整个系统的可用性。

从节点挂掉:Catch-up recovery

被重新启动的从节点会向主节点请求这期间的所有数据改动。

主节点挂掉:Failover

失效转移可以人工或者自动执行。但都包含以下几步:

1 主节点的失效检测

2 选举出新的主节点

3 修改配置,使用新的主节点

但是也存在很多难以解决的问题:

1 在异步同步的情况下,如果保障选出的新leader有最新的更新。

2 丢弃的写操作很危险。

3 脑裂

4 主节点失效检测的timeout选取

复制日志(replication log)的实现

1 Statement-based replication

最简单的方式,主节点记录所有收到的请求,并转发给从节点。mysql 5.1 之前都在用这种方式。

2 Write-ahead log (WAL) shipping

几乎都在使用WAL日志,不管是log-structured 存储引擎(SSTable 和 LSM-tree)还是B-tree存储引擎。

但还是有他的坏处:太底层了。一条WAL日志记录了一个磁盘block上的某个byte的变更。这个坏处导致了relication和存储引擎的实现耦合太紧。

3Logical (row-based) log replication

也就是logical log,和底层的存储引擎解耦。 Mysql的binlog(when configured to use row-based replication) 就是这一类。

4 Trigger-based replication

这种方式,一般是依赖外部工具来实现,更容易带来bug,但是更灵活。

read-after-write consistency

写入后,由主从同步的延迟,可能读不到写入的变更。 read-after-write consistency可以防止这样的情况 .

怎么做到这种一致性:有这样几种思路, 1 任何可能被当前用户修改的数据的读取,都从主节点进行。举个例子:用户的profile,我们就可以永远从主库读取自己的profile, 从从库读取他人的profile。2 如果上面的方法不行,可以记录updatetime,然后在每次写入后的一段时间内都读取主库。3 客户端可以记录一个它最近一次写操作的“时间戳”,然后,每次读取的时候都带着这个“时间戳”。系统来保障,服务读请求的从节点,自身同步过的写操作的“时间戳”要大于读请求中带着的。4 如果涉及到多数据中心,会更加复杂。

Monotonic Read(单调读)

当一个客户端多次进行读时,他不会感觉到时间倒流,每次读到的结果一定是相同的或者更新的。

一个实现的方法是:保障每个客户端被绑定在一个replica上,

Consistent Prefix Reads

确保的是:无论任何一个客户端,他读到的写操作的顺序是和实际的写入顺序是一致的。

多主模型

之前说的都是单主模型:只有一个主节点接受写操作。但是如果因为什么原因无法连接上主节点,导致的结果就是,无法向数据库写入。自然的,产生了多主模型。

使用场景

1 多数据中心:每个中心一个leader。这样比单leader有不少好处,但是相应的坏处是:数据可能被并发的修改,在不同的数据中心之间,这个“冲突”必须被解决。

2 Clients with offline operation :一部分datacenter离线一段时间,然后又回到集群。离开的那段时间,也可以正常运行。比方说 CouchDB

3 Collaborative editing :一般我们不把协同编辑功能,当做一个数据库的replication的问题来看,但是他们有好多相似的地方。

解决写冲突

  • 冲突检测:多主的情况下,每个写入都是成功的,冲突只能在后续的一个时间点被检测出来,那个时候,已经无法让用户去解决冲突。 如果要同步的冲突检测,就只能等待所有的replica都同步完,这个和单leader没什么区别。
  • 避免冲突:最简单的解决冲突的办法就是避免冲突。举个🌰:用户可以修改自己的信息,可以让这个用户一直路由到同一个dc。但是当用户旅游到了其他地方,就近接入了另外一个dc, 这个方法就失效了, 冲突就会出现。
  • Converging toward a consistent state: 如果是单个leader,我们可以知道更新操作的顺序,如果有对同一个field的更新操作,最后一个更新决定了这个field的值。如果是多个leader,这个顺序很难知道。

收敛的方法:

1 每一个写操作都带一个唯一Id。冲突的几个版本中,id大的胜利。

2 在一个额外的地方记录冲突,后续解决。可能通知用户来选择。

  • 自定义解决冲突的逻辑:让应用去解决冲突是大部分数据库的一般做法,大部分多主的数据库工具可以自定义冲突解决逻辑。一般有2个时机可以执行冲突解决逻辑:当冲突被检测到或者当读请求到来的时候。

自动冲突解决:

1
2
3
• Conflict-free replicated datatypes (CRDTs) [32, 38] are a family of data structures for sets, maps, ordered lists, counters, etc. that can be concurrently edited by multiple users, and which automatically resolve conflicts in sensible ways. Some CRDTs have been implemented in Riak 2.0 [39, 40].
• Mergeable persistent data structures [41] track history explicitly, similarly to the Git version control system, and use a three-way merge function (whereas CRDTs use two-way merges).
• Operational transformation [42] is the conflict resolution algorithm behind col‐ laborative editing applications such as Etherpad [30] and Google Docs [31]. It was designed particularly for concurrent editing of an ordered list of items, such as the list of characters that constitute a text document.

无主

废除主节点,任何节点都可以接受写操作。比较有代表性的此类数据库:Dynamo,Riak,Cassandra。

有节点挂掉时的写入

假设有个3个节点,其中一个挂掉。 写入的时候,是能更新其余2个节点。当挂掉的节点又启动起来,不能立刻与另外2个同步。这个时候如果读数据,可能会得到未更新的结果(读请求落到刚刚恢复的节点)。过一会,数据都同步了,读也会全部为新的了。

怎么解决这个问题?读请求也并行的发到多个节点,如果得到了多个结果,选择“版本号”最大的结果。

Read repair and anti-entropy

副本集群应该保证最终集群中的数据是一致的。一个离开的节点,又回到了集群,怎么保证他的同步?和Dynamo类似的数据库,一般会采用两种方法,来实现:

  • Read repair

    读取的时候,会读多个节点,把其中落后的节点,通知他们补上。

  • Anti-entropy process

​ 有一个后台的任务不断的去发现节点中的不同,进行修复。

Quorums for reading and writing (读写的法定人数)

假设有n个replicas,写入的时候,必须有至少w个节点写入成功,读的时候,至少从r个节点中读取。只要 w + r > n ,就可以保证读的时候,总可以读到最新的结果,因为读写至少有一个节点是重合的。在类Dynamo数据库中,w,r,n都是可配置的。

Sloppy Quorums and Hinted Handoff

无主的集群,可以忍受任意一个节点的挂掉, 因为不需要fallover, 也可以忍受个别节点变慢,因为因为只需要w或r个节点返回就可以。所以这个导致了,它十分适合需要 高可用低延迟,但是可以容忍偶发的 stale reads

如果因为网络原因,导致一个clinet只能看到一部分的节点,这样会导致他的读写无法达到法定人数。

这个时候,如果选择?

1 失败。

2

Hinted Handoff:当需要向一个已知宕机的节点写入数据的时候,Cassandra会将数据作为提示数据写入另一个正常的复本节点。稍后会对恢复的宕机节点进行数据重放。如果所有的复本节点都不能访问,则将提示数据写入到协调者本地。(能够理解提示的含义么) 当一个节点,通过Gossip发现它所存储的提示数据的所属节点恢复正常,他就会将提示数据发送到其所属目标节点。

缩短了暂时失败的节点恢复一致的时间。尤其是在一些奇怪的网络问题下,更加有用。(网络短时间不可用之类)

Hinted Handoff并不是修复机制的代替品。

上面的方式也叫做 Sloppy quorums,它可以有效的增加写入的可用性,但是它不是真正的法定人数,不能保证读到最新的数据。

Limitations of Quorum Consistency (Quorum Consistency 的限制)

就算 w + r > n, 也是有一些边界case,会返回不是最新的结果。

1 sloppy quorum