-## 一、背景 考虑可扩展性,后台的服务肯定是要能够支持任意的扩展的,这样才能在业务量增长时通过增加机器的方式来应对。这对后台服务提出了一个要求,必须处理好分布式环境和单机环境的不同带来的问题。比如:在文件的分片上传这一场景下,应该负载均衡的问题,一个文件的多个分片请求会分布到不同的服务器上,这导致在将多个分片合并成完整文件时出现问题,而单机情况下则完全不会有这样的问题。
二、难点
这个问题的解决方案有很多种,但是需要根据实际情况尽量选择简洁、易部署和维护的方案进行,并且不能丢掉分布式系统的优点。比如网上的有的方案是使用单独的文件上传服务器,但是这就变成了单机服务了。也有使用NFS挂载的方案,即所有的服务器挂载一个相同的NFS目录,所有上传相关的文件都存放在挂载的目录下,这不仅给运维带来了麻烦,为了保证NFS的高可用,也带来了额外的运维成本。
三、解决方案
3.1 借助负载均衡
借助负载均衡的方案很简单,这是和应用无关的一种方法。即通过负载均衡这一层,将同一个文件的不同分片的请求全部导向到同一个服务器上。比如负载均衡这一块使用的是nginx, 通过url hash的方式来完成,可以使用如下的配置:
upstream backend {
server 0.0.0.0:8080;
server 0.0.0.0:8081;
server 0.0.0.0:8082;
hash $request_uri;
}
server {
location / {
proxy_pass http://backend;
proxy_set_header X-Real-IP $remote_addr;
proxy_set_header Host $host;
}
}
然后写一个简单的服务来验证一下:
func post(resp http.ResponseWriter, req *http.Request) {
v := req.PostFormValue("key")
identity := req.PostFormValue("identity")
log.Printf("identity: %v, value: %v\n", identity, v)
resp.WriteHeader(http.StatusOK)
_, err := resp.Write([]byte("ok"))
if err != nil {
log.Println("error:", err)
}
}
func main() {
var addr string
flag.StringVar(&addr, "addr", "0.0.0.0:8080", "http listen addr:port")
flag.Parse()
http.HandleFunc("/Post", post)
log.Println("listen ", addr)
err := http.ListenAndServe(addr, nil)
if err != nil {
log.Fatal(err)
}
}
我们启动三个服务:
./godemo -addr 0.0.0.0:8080
./godemo -addr 0.0.0.0:8081
./godemo -addr 0.0.0.0:8082
使用curl来post数据过来,请求的地址类似于: http://0.0.0.0/Post?123456。?后面可以当成是文件的唯一标识码,比如文件和用户id的组合的md5信息等,这样对于同一个用户上传的同一个文件的不同分片请求,通过url哈希都会得到同样的结果,这样就会被转发到同一个后端服务器。
curl http://0.0.0.0/Post\?123456 -d "key=value&identity=123456" -X POST
curl http://0.0.0.0/Post\?123456 -d "key=value&identity=123456" -X POST
curl http://0.0.0.0/Post\?123456 -d "key=value&identity=123456" -X POST
curl http://0.0.0.0/Post\?789abc -d "key=vvvvv&identity=789abc" -X POST
curl http://0.0.0.0/Post\?789abc -d "key=vvvvv&identity=789abc" -X POST
curl http://0.0.0.0/Post\?789abc -d "key=vvvvv&identity=789abc" -X POST
结果如下:
端口8080的服务收到了三条请求:
2019/09/26 13:58:52 identity: 123456, value: value
2019/09/26 13:58:56 identity: 123456, value: value
2019/09/26 13:59:03 identity: 123456, value: value
端口8081的服务收到了三条请求:
2019/09/26 13:59:40 identity: 789abc, value: vvvvv
2019/09/26 13:59:43 identity: 789abc, value: vvvvv
2019/09/26 13:59:44 identity: 789abc, value: vvvvv
如果在k8s集群中部署,使用nginx-ingress的话可以使用以下的部署方案:
apiVersion: networking.k8s.io/v1beta1
kind: Ingress
metadata:
name: pipeline-ingress
namespace: default
annotations:
nginx.ingress.kubernetes.io/proxy-body-size: "50m"
nginx.ingress.kubernetes.io/upstream-hash-by: "$request_uri"
spec:
rules:
- host: pipeline.dev.com
http:
paths:
- path: /
backend:
serviceName: pipeline-service
servicePort: 8888
当然这种方式需要注意你的每个请求的URI都需要加上额外的参数(比如用户的userId的md5),这样才能均匀的分布到不同的服务器上。
3.2 后台程序自动proxy请求
这个方案的思路是集群中的每个服务实例都有单独的标识,文件分片在第一次上传时会返回给它一个该请求所属服务器的identity,之后所有的请求会带上这个identity,之后收到请求的服务器会检查这个identity是不是属于自己,如果不属于自己,则把这个请求转发给所属的服务器。 这个方案要求每台服务器都知道其他所有服务器的identity和地址。这里我们可以使用etcd这样的分布式数据库来存储。每台服务器在启动时都向etcd里注册自己的identity和address,之后服务器转发的时候都向etcd里面查找对应的address即可。 下面是一个示例的代码:
package main
import (
"context"
"encoding/json"
"flag"
"io/ioutil"
"log"
"net/http"
"go.etcd.io/etcd/clientv3"
"net/url"
"time"
)
type Server struct {
Etcd string
Addr string
Identify string
}
type ResponseObj struct {
Identity string `json:"identity,omitempty"`
Value string `json:"value"`
}
func getEtcdKV(etcd string) (clientv3.KV, error) {
cfg := clientv3.Config{
Endpoints: []string{etcd},
// set timeout per request to fail fast when the target endpoint is unavailable
DialTimeout: time.Second,
}
cli, err := clientv3.New(cfg)
if err != nil {
return nil, err
}
return clientv3.NewKV(cli), nil
}
func httpProxy(anotherServer string, body map[string]string) (*http.Response, error) {
formData := url.Values{}
for k, v := range body {
formData.Set(k ,v)
}
return http.PostForm(anotherServer, formData)
}
func (s *Server) Register() error {
cli, err := getEtcdKV(s.Etcd)
if err != nil {
return err
}
ctx, cancel := context.WithTimeout(context.Background(), time.Duration(1)*time.Second)
_, err = cli.Put(ctx, s.Identify, "http://" + s.Addr)
cancel()
if err != nil {
return err
}
return nil
}
func (s *Server) Post(resp http.ResponseWriter, req *http.Request) {
v := req.PostFormValue("key")
identity := req.PostFormValue("identity")
log.Printf("identity: %v, value: %v\n", identity, v)
var data ResponseObj
if identity != "" && identity != s.Identify {
// proxy
cli, err := getEtcdKV(s.Etcd)
if err != nil {
log.Println("get kv client error", err)
resp.WriteHeader(http.StatusInternalServerError)
resp.Write([]byte("error"))
return
}
etcdResp, err := cli.Get(context.Background(), identity)
if err != nil {
log.Println("get value error: ", err)
resp.WriteHeader(http.StatusInternalServerError)
resp.Write([]byte("error"))
return
}
if len(etcdResp.Kvs) == 0 {
log.Println("没有值")
resp.WriteHeader(http.StatusInternalServerError)
resp.Write([]byte("error"))
return
}
proxyAddr := string(etcdResp.Kvs[0].Value)
log.Println("proxy to: ", proxyAddr)
text := map[string]string {
"key": v,
"identity": identity,
}
proxyResp, err := httpProxy(proxyAddr+"/Post", text)
if err != nil || proxyResp.Body == nil {
log.Println("http proxy error: ", err)
resp.WriteHeader(http.StatusInternalServerError)
resp.Write([]byte("error"))
return
}
bodyData, err := ioutil.ReadAll(proxyResp.Body)
if err != nil {
log.Println("read proxy body error: ", err)
resp.WriteHeader(http.StatusInternalServerError)
resp.Write([]byte("error"))
return
}
err = json.Unmarshal(bodyData, &data)
if err != nil {
log.Println("json unmarshal error: ", err)
resp.WriteHeader(http.StatusInternalServerError)
resp.Write([]byte("error"))
return
}
} else {
data.Identity = s.Identify
data.Value = v
}
text, err := json.Marshal(data)
if err != nil {
log.Println("marshal json error: ", err)
resp.WriteHeader(http.StatusInternalServerError)
return
}
resp.WriteHeader(http.StatusOK)
_, err = resp.Write(text)
if err != nil {
log.Println("error:", err)
}
}
func main() {
var addr string
var identity string
var etcd string
flag.StringVar(&addr, "addr", "0.0.0.0:8080", "http listen addr:port")
flag.StringVar(&identity, "identity", "", "identify the server")
flag.StringVar(&etcd, "etcd", "http://127.0.0.1:2379", "etcd server url")
flag.Parse()
if identity == "" {
log.Fatal("identity不可为空")
}
var server = &Server{
Addr: addr,
Identify: identity,
Etcd: etcd,
}
err := server.Register()
if err != nil {
log.Fatal("register server error, check your etcd server: ", err)
}
http.HandleFunc("/Post", server.Post)
log.Println("listen ", addr)
err = http.ListenAndServe(addr, nil)
if err != nil {
log.Fatal(err)
}
}
3.3 基于S3存储的方案
兼容S3的存储系统有一个对外的接口叫做: ComposeObject。这个接口可以将多个文件合并成一个文件存储到指定位置。如果我们的底层存储是基于兼容S3的系统,那么就可以利用这个接口轻松的实现。 方案的步骤可以描述成以下:
- 前端实现将文件分片,然后上传
- 上传的请求会因为前端负载均衡的原因分布到不同的服务器,每个上传请求都要带上这个文件、用户id、随机字符串组合的md5值,服务器收到请求后,只管将分片上传到以md5值为名的目录下,一旦上传成功,使用etcd或redis这样作为分片计数器+1.
- 分片总数等于总分片数时调用
ComposeObject组合所有分片存储到指定位置即可。

通过上图可以看出,ceph的底层核心是
Ceph Monitor维护了一个关于集群信息映射的主副本。Ceph monitors集群保证了在一个monitor服务停止时的高可用性。存储集群客户端从Ceph Monitor获取一份集群信息映射的拷贝。 Ceph OSD服务负责检查自身的状态以及其他OSD的状态,然后汇报给monitors。 存储集群客户端和每一个Ceph OSD服务使用
Ceph OSD服务在一个扁平的命名空间中(没有目录结构)存储所有的数据,每个数据都作为一个对象。一个对象有唯一标识符,二进制数据,以及有一系列键值对的元数据。这完全取决于Ceph客户端的实现。举例来说,CephFS使用元数据来存储文件属性,比如文件拥有者,创建日期,最后修改日期等等。 
为了通过监控服务的认证,客户端把用户名发送给监控服务,监控服务生成一个会话钥匙,使用秘钥附加上用户名来加密。之后,监控服务把加密后的ticket返回给客户端。客户端使用共享的秘钥解密,获取到会话钥匙。会话钥匙为当前的会话提供认证。客户端之后请求由这个会话钥匙签名的ticket。监控服务生成ticket,使用用户秘钥加密并返回给客户端。客户端解密ticket,使用它来签名给OSD和元数据服务器的请求。
cephx协议认证每个在客户端机器和Ceph服务器之间的会话。每一个在客户端和服务器之间发送的信息,都会接连通过初始化认证,使用ticket签名,这个签名可以由监控服务,OSD和元数据服务器用他们共享的秘钥来验证。
由这种认证提供的保护是在Ceph客户端和Ceph服务端之间的。在Ceph客户端之外这个认证是不管用的。如果用户通过远程连接到Ceph客户端上,Ceph的认证不会应用到这个远程连接上。
因为有这样的数据复制的能力,Ceph OSD服务使得Ceph客户端不需要负责这件事,同时也能保证高的数据可用性和数据安全。
pool至少会设置以下参数:
客户端在有集群映射表的拷贝和CRUSH算法的情况下,可以准确的计算出在读取或写入对象时使用哪一个OSD。
当对象NYAN从纠删码pool中读取时,解码功能读取三个块: 块1包含ABC,块3包含GHI,块4包含YXY。然后它重建了对象的原始内容ABCDEFGHI。这个解码功能被通知块2和块5是缺失的(它们被称作’擦除’)。块5不会被读取是因为OSD4从集群中退出了。一旦3个块被读取解码功能就会被调用,而此时OSD2因为最慢所以根本没有被考虑进来。 
OSD1是主服务,接受客户端的完整的写入,这意味着写入的数据要完整的替换对象而不是覆盖一部分。版本2(v2)的对象被创建来覆盖版本1(v1)。OSD1将数据编码到3块: D1v2在OSD1上,D2v2在OSD2上,C1v2在OSD3上。每一个块都被发送到目标OSD上。当OSD接受到消息要写入数据块时,它同样创建一个新的PG日志来反映这个改变。举例来说,只要OSD3存储了C1V2,它添加1,2到日志中。因为OSD是异步工作的,当其他块已经存储好了(比如C1v1和D1v1)一些数据块可能还在处理(比如D2v2)。
如果一切顺利,数据块都会被存储,日志中的
最后,之前的版本都可以被移除了: OSD1上的D1v1,OSD2上的D2v1, OSD3上的C1v1
但是如果发生了异常,OSD1在D2v2仍然写入的时候停止了,对象的版本2只是部分写入:OSD3有一个块但是不足以恢复。它丢失了两个块: D1v2和D2v2,但是纠删码参数K=2,M=1要求至少2个块是可用的,才能恢复第三个块。OSD4成为新的主服务,找到last_complete日志是1,1,这应该是最新的日志了。
日志1,2被发现在OSD3上,这和在OSD4上存储的最新版本1,1不一样、因此1,2被取消,C1v2块被移除。D1v1块被解码功能重建出来存在OSD4上。 



如果你想要更大的镜像尺寸,更大的S3或Swift对象,或者更大的CephFS目录,你应该考虑通过将数据分片分布到一个对象集合中的多个对象上来提升读写性能。并行操作会显著的提升写入性能。因为对象会被映射到不同的PG中,然后被映射到不同的OSD中,每一个写入操作都会以最大的写入速度来并行操作。单硬盘的写入会因为柱头的移动和设备的最大带宽(100M/s)而受到限制(比如每次寻道都要6ms)。 通过将写入分布到多个对象上,Ceph可以减少每个硬盘的寻道时间,然后将多个硬盘的吞吐量结合以达到更快的写入(或读取)速度。
下面是三个重要的变量决定了Ceph如何对数据分片:


代码的github地址是:
对于一个24-bit的图片来说,通常来说颜色空间是连续的,因为颜色之间的最小差异是几乎不可能察觉的。现在,这个连续的颜色空间要被映射到256种离散的颜色上。将一个连续的变量映射到离散的集合上,就叫做
在递归情况下,B1将一直不会被处理。 2.为了解决1中的问题。我们可以设定一个最大深度(level)。比如我们要得到4个输出颜色。那么log 2 4=2。最大深度应该是2。当切割到最大深度时,我们就不再往下切割,转而切割那些还未被处理的更大的区域。