0%

大数据流式计算及应用实践

面向感知的数据流处理模型
R树、B+树

CAP(Consistency,Availability,tolerance to network Partitions) 理论说明,分布式系统中的一致性、可用性、分区容错性三者不可兼得。

所以,并行关系数据库必然无法获得较强的扩展性和良好的系统可用性。
面对这些挑战,以Google的BigTable为代表的NoSQL(notonlySQL)数据库发展起来。BigTable是一个多维稀疏排序表,由行和列组成,每个存储单元都有一个时间戳,形成三维结构。同一个数据单元的多个操作形成数据的多个版本,由时间戳来区分。

高性能批数据处理模式、流式数据处理模式和两者混合的模式。

批处理模式就是修改以Hadoop为代表的批处理框架,减少中间结果的写盘次数,增加作业间的流水化程度来提高单位时间的吞吐量。流处理模式,是针对流式数据的一种天然适合的处理方法,将到达的数据在维护的滑动窗口内进行处理,这将在下一章详细说明,典型系统如YahooS4和Storm。而两者混合的模式,主要思路是基于MapReduce模型增加或改变其中的某些处理步骤,以实现流处理。

四个时间点?
契税 多少?
查封 离婚 继承

距离 环境 学校

高性能批数据处理模式、流式数据处理模式和两者混合的模式。其中,批处理模式就是修改以Hadoop为代表的批处理框架,减少中间结果的写盘次数,增加作业间的流水化程度来提高单位时间的吞吐量。流处理模式,是针对流式数据的一种天然适合的处理方法,将到达的数据在维护的滑动窗口内进行处理,这将在下一章详细说明,典型系统如Yahoo S4和Storm。而两者混合的模式,主要思路是基于MapReduce模型增加或改变其中的某些处理步骤,以实现流处理。例如,DEDUCE系统扩展了IBM的流式数据处理System S,使其支持MapReduce。此外,SparkStreaming通过引入离散流(discretized streams)编程模型,改进了批处理模式,大幅提高了处理速度,并在Spark系统上实现了集成。

Storm 作为 一种 分布式 系统 的 体系 结构、 分布式 通信、 作业、 接 入/ 处理 组件 和 若干 进阶 的 概念。

丁维龙; 等. Storm:大数据流式计算及应用实践 (高端云计算与大数据丛书) (Kindle 位置 1278-1279). 电子工业出版社. Kindle 版本.

supervisor

工作进程执行指定topology的子集,而同一个topology可以由多个工作进程完成;一个工作进程由多个工作线程组成,工作线程是spout/bolt的运行时实例,数量由spout/bolt的数目及其配置确定。supervisor是分布式部署的,在Storm中的地位类似于Hadoop中的TaskTracker。

首先是作业在Storm系统中的部署。
① 用户将作业打包(即将topology的代码组织为jar文件),通过Storm的客户端命令或者控制台节点的Web接口,提交至Storm系统的主控节点;
② 主控节点根据系统的全局配置和作业中的局部配置,将接收的代码分发至调度的工作节点;
③ 工作节点下载来自主控节点的代码包,并根据主控节点的调度生成相关的工作进程和线程。

其次是系统节点状态的协调,包括如下几个部分。
①主控节点与协调节点之间:主控节点将系统全局的配置、节点局部的配置、主控节点运行时的状态,通过服务接口交由Zookeeper维护,由协调节点实现配置管理与状态监控;
②工作节点与协调节点之间:工作节点将自身的状态、工作进程和工作线程的状态,通过服务接口交由Zookeeper维护,由协调节点实现状态获取与更新;
③主控节点与控制台节点之间:控制台节点调用主控节点开放相关的接口,可以获取系统、作业、工作进程、工作线程和任务的运行时状态。

最后是工作节点间作业计算结果的数据传输,包括如下几个部分。
①组件间线程级的数据传递:在同一进程内,spout/bolt的工作线程将处理结果向作业的下游组件传递;
②组件间进程级的数据传递:在不同进程间(无论是否在同一台机器中),spout/bolt的工作线程将处理的结果向作业的下游组件传递。进程是操作系统中程序运行时的基本单元,线程是进程的最小调度单位,Storm的分布式数据处理,也是通过进程线程的调度实现的。
注意:考虑到作业的独立性与安全性,Storm不支持跨作业的数据传递,故这里工作节点间的数据传输,一定存在于同一作业(隶属同一topology)的运行时实例进程/线程之间。

Dubbo-分布式服务治理框架

学习dubbo教程 笔记整理

Dubbo官方文档

关键词:分布式治理

Dubbo是阿里巴巴推出的一款分布式服务治理框架(虽然现在阿里内部一些部门已经不再使用)。

Dubbo将分布式服务分为四个角色:服务提供者、服务消费者、注册中心和监控中心。

Dubbo角色

![][dubbo-roles]

  1. 服务提供者注册到服务注册中心
  2. 服务消费者从服务提供者订阅服务
  3. 当服务提供者发生变化(出现新服务提供者、旧的服务提供者死掉等),注册中心通知服务消费者
  4. 当服务消费者调用服务时,首先从注册中心查找服务,注册中心直接将选定的服务提供者的ip和端口等信息返回给服务消费者,服务消费者直接调用服务提供者
  5. 服务提供者和服务消费者每个一分钟将收集的运行信息上报到监控中心

只有服务消费者调用服务生成者是同步调用,其他都是异步的。

Provider与Registry、Consumer与Register之间都保持着长连接,用于保持信息同步

简洁的项目结构

Dubbo支持的RPC协议

  1. 支持常见的传输协议:RMI、Dubbo、Hessain、WebService、Http等,
    其中Dubbo和RMI协议基于TCP实现,Hessian和WebService基于HTTP实现。
  2. 传输框架:Netty、Mina、以及基于servlet等方式。
  3. 序列化方式:Hessian2、dubbo、JSON( fastjson 实现)、JAVA、SOAP 等。
  4. 注册中心可以选择 zooKeeper Redis Dubbo Multicast

Dubbo 服务降级

Dubbo学习(七):服务的升级和降级

服务降级方式:

  • 服务接口拒绝服务:无用户特定信息,页面能访问,但是添加删除提示服务器繁忙。页面内容也可在Varnish或CDN内获取。
  • 页面拒绝服务:页面提示由于服务繁忙此服务暂停。跳转到varnish或nginx的一个静态页面。
  • 延迟持久化:页面访问照常,但是涉及记录变更,会提示稍晚能看到结果,将数据记录到异步队列或log,服务恢复后执行。
  • 随机拒绝服务:服务接口随机拒绝服务,让用户重试,目前较少有人采用。因为用户体验不佳。

服务降级埋点的地方:

  • 消息中间件:所有API调用可以使用消息中间件进行控制
  • 前端页面:指定网址不可访问(NGINX+LUA)
  • 底层数据驱动:拒绝所有增删改动作,只允许查询

  1. 从网易镜像仓库下载Jenkins

网易镜像仓库

拉取到本地

1
docker pull hub.c.163.com/library/jenkins:latest

查看镜像

1
docker images
1
2
REPOSITORY                      TAG                 IMAGE ID            CREATED             SIZE
hub.c.163.com/library/jenkins latest 88d9d8a30b47 8 months ago 810MB

从官网下载最新的Jenkins

1
wget http://updates.jenkins-ci.org/download/war/2.107.2/jenkins.war

DockerFile

1
2
3
FROM hub.c.163.com/library/jenkins
# 将最新的Jenkins拷贝到`/usr/share/jenkins/`目录
ADD jenkins.war /usr/share/jenkins/

常见镜像

1
docker build -t avery/jenkins .

启动一个Container

1
docker run --name prod_jenkins2 -d -p 9005:8080 -p 9006:50000 -v  /Users/avery/docker/jenkins/backup/jenkins_home:/var/jenkins_home avery/jenkins
  1. 创建本地Jenkins数据目录,将该目录映射到docker中的Jenkins目录,方便备份数据
  2. 8080端口是Jenkins web service的默认端口,将宿主机的端口映射到8080端口,可以直接访问
  3. –name 命名为prod_jenkins2
  4. -d 后台运行
  5. -p 端口映射
  6. -v 文件目录映射

查看

1
2
3
docker container ls -a
CONTAINER ID IMAGE COMMAND CREATED STATUS PORTS NAMES
cb8cdf3be078 avery/jenkins "/bin/tini -- /usr/l…" 3 hours ago Exited (143) 3 hours ago prod_jenkins2

启动后,可以直接通过url ‘http://localhost:9001‘ 访问,Jenkins首次打开,需要获取随机密码

1
2
3
4
5
# 获取jenkins密码
docker container exec cb8cdf3be078 cat /var/jenkins_home/secrets/initialAdminPassword
a7cefff7d2634cc6a9542491f465eb52
# 或者直接进入shell查看
docker exec -it cb8cdf3be078 bash

为了学习Hadoop,尝试使用Vbox搭建环境,不是很方便。后面转向Docker。本文将使用Docker搭建Hadoop集群记录下来,以备后用

基础环境

  • MacOS 10.12
  • Docker 17.03.1-ce

网易镜像 https://c.163yun.com/hub#/m/home/

安装Docker

在MacOS 上安装Docker即为简单,从官网上下载dmg包,拖到应用目录启动即可。

安装完成后 查看docker 版本:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
% docker version

Client:
Version: 17.03.1-ce
API version: 1.27
Go version: go1.7.5
Git commit: c6d412e
Built: Tue Mar 28 00:40:02 2017
OS/Arch: darwin/amd64

Server:
Version: 17.03.1-ce
API version: 1.27 (minimum version 1.12)
Go version: go1.7.5
Git commit: c6d412e
Built: Fri Mar 24 00:00:50 2017
OS/Arch: linux/amd64
Experimental: true

docker 建立镜像

建立三个镜像,存储目录如下:

1
2
3
4
5
6
7
8
9
10
% tree
.
├── centos-ssh-root
│   └── Dockerfile
├── centos-ssh-root-jdk
│   ├── Dockerfile
│   └── jdk-7u79.tar.gz
└── centos-ssh-root-jdk-hadoop
   ├── Dockerfile
   └── hadoop-2.7.3.tar.gz

建立CentOS-SSH-root基础镜像

建立一个带有ssh和root账户的CentOS镜像

  1. 新建Dockerfile文件
1
2
3
4
% mkdir  centos-ssh-root
% cd centos-ssh-root
# 新建 Dockerfile 文件
% vi Dockerfile

Dockerfile的内容为:

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
# 选择一个已有的os镜像作为基础  国内首选仓库 c.163.com
FROM hub.c.163.com/public/centos:7.2.1511
RUN yum clean all
RUN yum install -y yum-plugin-ovl || true
# 安装基础的工具包
RUN yum install -y vim tar wget curl rsync bzip2 iptables tcpdump less telnet net-tools lsof sysstat cronie python-setuptools
RUN yum clean all
RUN easy_install supervisor
RUN cp -f /usr/share/zoneinfo/Asia/Shanghai /etc/localtime
EXPOSE 22
RUN mkdir -p /etc/supervisor/conf.d/
RUN /usr/bin/echo_supervisord_conf > /etc/supervisord.conf
RUN echo [include] >> /etc/supervisord.conf
RUN echo 'files = /etc/supervisor/conf.d/*.conf' >> /etc/supervisord.conf
#COPY sshd.conf /etc/supervisor/conf.d/sshd.conf
CMD ["/usr/bin/supervisord"]



# 镜像的作者
MAINTAINER avery

# 安装openssh-server和sudo软件包,并且将sshd的UsePAM参数设置成no
RUN yum install -y openssh-server sudo
RUN sed -i 's/UsePAM yes/UsePAM no/g' /etc/ssh/sshd_config
#安装openssh-clients
RUN yum install -y openssh-clients

# 添加测试用户root,密码abc.123,并且将此用户添加到sudoers里
RUN echo "root:abc.123" | chpasswd
RUN echo "root ALL=(ALL) ALL" >> /etc/sudoers
# 下面这两句比较特殊,
# 在centos6上必须要有,否则创建出来的容器sshd不能登录

# 为了避免文件已存在报错,首先删掉私钥文件
RUN rm -rf /etc/ssh/ssh_host_rsa_key
RUN rm -rf /etc/ssh/ssh_host_dsa_key

RUN ssh-keygen -t dsa -f /etc/ssh/ssh_host_dsa_key
RUN ssh-keygen -t rsa -f /etc/ssh/ssh_host_rsa_key

# 启动sshd服务并且暴露22端口
# RUN mkdir /var/run/sshd
EXPOSE 22
CMD ["/usr/sbin/sshd", "-D"]

  1. 创建镜像
1
2
# 在Dockerfile同级目录下执行下面的命令,最后的.必须有
docker build -t="avery/centos-ssh-root" .

查看刚刚创建成功的镜像

1
2
3
% docker images
REPOSITORY TAG IMAGE ID CREATED SIZE
avery/centos-ssh-root latest ba72ccadf508 35 hours ago 776 MB

创建带有JDK的镜像

  1. 准备 下载JDK

此处使用JDK7u79

1
2
3
4
% mkdir centos-ssh-root-jdk
% cd centos-ssh-root-jdk
% cp ~/Download/jdk-7u79.tar.gz .
% vi Dockerfile
  1. Dockerfile文件
1
2
3
4
5
6
7
8
9
10
11
FROM avery/centos-ssh-root
ADD jdk-7u79.tar.gz /usr/local/
RUN mv /usr/local/jdk1.7.0_79 /usr/local/jdk1.7
# 添加环境变量
ENV JAVA_HOME /usr/local/jdk1.7
ENV PATH $JAVA_HOME/bin:$PATH
# 可能上面的设置不会生效
RUN echo 'JAVA_HOME=/usr/local/jdk1.7/' >> .bash_profile
RUN echo 'PATH=$JAVA_HOME/bin:$PATH' >> .bash_profile
RUN echo 'export JAVA_HOME' >> .bash_profile
RUN echo 'export PATH' >> .bash_profile
  1. 创建镜像
1
docker build -t="avery/centos-ssh-root" .
  1. 查看镜像
1
2
3
4
% docker images
REPOSITORY TAG IMAGE ID CREATED SIZE
avery/centos-ssh-root-jdk latest e773519d0452 34 hours ago 1.39 GB
avery/centos-ssh-root latest ba72ccadf508 35 hours ago 776 MB

构建hadoop镜像

  1. 准备

官网下载 hadoop-2.7.3.tar.gz

1
2
3
4
% mkdir centos-ssh-root-jdk-hadoop 
% cd centos-ssh-root-jdk-hadoop
% cp ../hadoop-2.7.3.tar.gz .
% vi Dockerfile
  1. Dockerfile
1
2
3
4
5
6
7
8
9
10
11
FROM avery/centos-ssh-root-jdk
ADD hadoop-2.7.3.tar.gz /usr/local
RUN mv /usr/local/hadoop-2.7.3 usr/local/hadoop
ENV HADOOP_HOME /usr/local/hadoop
ENV PATH $HADOOP_HOME/bin:$PATH

# 为了防止环境变量失效
RUN echo 'HADOOP_HOME=/usr/local/hadoop/' >>.bash_profile
RUN echo 'PATH=$HADOOP_HOME/bin:$PATH' >> .bash_profile
RUN echo 'export HADOOP_HOME' >> .bash_profile
RUN echo 'export PATH' >> .bash_profile
  1. 创建镜像
1
docker build -t="avery/centos-ssh-root-jdk-hadoop" .
  1. 查看
1
2
3
4
5
% docker images
REPOSITORY TAG IMAGE ID CREATED SIZE
avery/centos-ssh-root-jdk-hadoop latest c93aadf2b6c5 1 hours ago 2.05 GB
avery/centos-ssh-root-jdk latest e773519d0452 1 hours ago 1.39 GB
avery/centos-ssh-root latest ba72ccadf508 1 hours ago 776 MB

docker 搭建 Hadoop集群

此处搭建的集群只在本机使用,各个Docker container没有独立的ip,如需要通过IP访问可以参考使用docker搭建hadoop分布式集群

在MacOS上修改系统文件,需要关闭Rootless:

重启按住 Command+R,进入恢复模式,打开Terminal。

1
csrutil disable

重启即可。如果要恢复默认,那么

1
csrutil enable

创建3台container

1
2
3
4
5
docker run --name hadoop0 --hostname hadoop0 -d -P -p 50070:50070 -p 8088:8088 avery/centos-ssh-root-jdk-hadoop

docker run --name hadoop1 --hostname hadoop1 -d -P avery/centos-ssh-root-jdk-hadoop

docker run --name hadoop2 --hostname hadoop2 -d -P avery/centos-ssh-root-jdk-hadoop

查看三个container的信息

1
2
3
4
5
docker container ls
CONTAINER ID IMAGE COMMAND CREATED STATUS PORTS NAMES
78e5e35adfa4 avery/centos-ssh-root-jdk-hadoop "/usr/sbin/sshd -D" 3 minutes ago Up 3 minutes 0.0.0.0:32776->22/tcp hadoop2
5f27fd91c7ba avery/centos-ssh-root-jdk-hadoop "/usr/sbin/sshd -D" 3 minutes ago Up 3 minutes 0.0.0.0:32775->22/tcp hadoop1
1d7fd37371cc avery/centos-ssh-root-jdk-hadoop "/usr/sbin/sshd -D" 3 minutes ago Up 3 minutes 0.0.0.0:8088->8088/tcp, 0.0.0.0:50070->50070/tcp, 0.0.0.0:32774->22/tcp hadoop0

查看三个container的ip:

1
2
3
4
5
6
% docker inspect --format='{{.NetworkSettings.IPAddress}}' hadoop0
172.17.0.2
% docker inspect --format='{{.NetworkSettings.IPAddress}}' hadoop1
172.17.0.3
% docker inspect --format='{{.NetworkSettings.IPAddress}}' hadoop2
172.17.0.4

集群规划

准备搭建一个具有三个节点的集群,一主两从

  • 主节点:hadoop0 ip:172.17.0.2
  • 从节点1:hadoop1 ip:172.17.0.3
  • 从节点2:hadoop2 ip:172.17.0.4

连接container

验证ssh连接

1
2
3
4
5
6
% ssh root@localhost -p 32774
# 输入密码 abc.123
% ssh root@localhost -p 32775
# 输入密码 abc.123
% ssh root@localhost -p 32776
# 输入密码 abc.123

使用exec

1
docker exec -it hadoop0 /bin/bash

修改container的主机名

分别修改三个container的hosts

1
vi /etc/hosts 

添加下面配置

1
2
3
172.17.0.2    hadoop0
172.17.0.3 hadoop1
172.17.0.4 hadoop2

ssh 免密

设置ssh免密码登录
在hadoop0上执行下面操作

1
2
3
4
5
6
7
8
9
cd  ~
mkdir .ssh
cd .ssh
ssh-keygen -t rsa
# (一直按回车即可)
ssh-copy-id -i localhost
ssh-copy-id -i hadoop0
ssh-copy-id -i hadoop1
ssh-copy-id -i hadoop2

在hadoop1上执行下面操作

1
2
3
4
5
6
cd  ~
cd .ssh
ssh-keygen -t rsa
# (一直按回车即可)
ssh-copy-id -i localhost
ssh-copy-id -i hadoop1

在hadoop2上执行下面操作

1
2
3
4
5
6
cd  ~
cd .ssh
ssh-keygen -t rsa
# (一直按回车即可)
ssh-copy-id -i localhost
ssh-copy-id -i hadoop2

至此,Docker搭建Hadoop集群的准备工作

build debian-conajdk8-hadoop

1
docker pull hub.c.163.com/library/centos:latest

docker pull caiserkaiser/hadoop:2.7.2

1
docker network create -d bridge --subnet "172.173.16.0/24" --gateway "172.173.16.1"  datastore_net

基于一个image构建

centos7_jdk8

hub.c.163.com/housan993/centos7_jdk8:latest

yum install -y openssh-server openssh-clients

docker pull hub.c.163.com/public/centos:7.0

基于caiserkaiser的快速搭建

Docker 构建 Hadoop 2.7.2 集群
Docker 构建 Flink 1.10.2 集群(ON YARN)

1
docker network create -d bridge --subnet "172.173.16.0/24" --gateway "172.173.16.1"  datastore_net
1
2
3
docker run -it -d --network datastore_net --ip 172.173.16.10 --name hadoop01 caiser/hadoop:2.7.2 bin/bash
docker run -it -d --network datastore_net --ip 172.173.16.11 --name hadoop02 caiser/hadoop:2.7.2 bin/bash
docker run -it -d --network datastore_net --ip 172.173.16.12 --name hadoop03 caiser/hadoop:2.7.2 bin/bash

DockerFile免密登录

1
2
3
4
5
6
7
8
9
10
11
FROM ubuntu:14.04
MAINTAINER yjt xxx
RUN sudo apt-get update && \
sudo apt-get install -y net-tools openssh-server psmisc iproute wget vim
RUN ssh-keygen -t rsa -f ~/.ssh/id_rsa -P '' && cat /root/.ssh/id_rsa.pub >> /root/.ssh/authorized_keys && \
sed -i 's/PermitEmptyPasswords yes/PermitEmptyPasswords no /' /etc/ssh/sshd_config && \
sed -i 's/PermitRootLogin without-password/PermitRootLogin yes /' /etc/ssh/sshd_config && \
echo " StrictHostKeyChecking no" >> /etc/ssh/ssh_config && \
echo " UserKnownHostsFile /dev/null" >> /etc/ssh/ssh_config && \
echo "root:1234" | chpasswd
CMD [ "sh", "-c", "sudo service ssh start; bash"]

接下来,命令行执行:

1
docker build -t ubuntu-ssh .

通过刚刚生成的镜像启动两台容器:

1
2
docker run -it --rm --name=yjt1 --net mynetwork --ip 172.20.1.1 --privileged ubuntu-ssh
docker run -it --rm --name=yjt2 --net mynetwork --ip 172.20.1.2 --privileged ubuntu-ssh

【参考文献】

  1. 使用docker搭建hadoop分布式集群

MySQL 会对 sql 语句做优化,

  1. in 后面的条件不超过一定数量仍然会使用索引。mysql 会根据索引长度和 in 后面条件数量判断是否使用索引。
  2. 如果是 in 后面是子查询,则不会使用索引。此时采用join来替换
  3. 使用union all代替inor

使用union all优化的样例

一个文章库,里面有两个表:category 和 article。category 里面有 10 条分类数据。article 里面有 20 万条。article 里面有一个”article_category”字段是与 category 里的”category_id”字段相对应的。 article 表里面已经把 article_category 字义为了索引。数据库大小为 1.3G。

问题描述:

执行一个很普通的查询:

1
Select * FROM `article` Where article_category=11 orDER BY article_id DESC LIMIT 5

执行时间大约要 5 秒左右

解决方案:
建一个索引:

1
create index idx_u on article (article_category,article_id);
1
Select * FROM `article` Where article_category=11 orDER BY article_id DESC LIMIT 5

减少到 0.0027 秒

继续问题:

1
Select * FROM `article` Where article_category IN (2,3) orDER BY article_id DESC LIMIT 5

执行时间要 11.2850 秒。

使用 OR:

1
2
3
4
5
select * from article
where article_category=2
or article_category=3
order by article_id desc
limit 5

执行时间:11.0777

解决方案:
避免使用 in 或者 or (or 会导致扫表),使用 union all

使用 UNION ALL:

1
2
3
4
(select * from article where article_category=2 order by article_id desc limit 5)
UNION ALL (select * from article where article_category=3 order by article_id desc limit 5)
orDER BY article_id desc
limit 5

执行时间:0.0261

mysql

启动mysql

brew services start mysql

abc.ABC.123

集群搭建

  1. mac osx 搭建hadoop开发环境
  2. spark on yarn
  3. hive环境搭建

按照上面文章的方式搭建,启动hdfs、yarn

必须使用这个命令

1
2
3
4
5
6
## 启动yarn
hadoop/sbin/./yarn-daemon.sh start resourcemanager
hadoop/sbin./yarn-daemon.sh start nodemanager
## 启动hdfs
hadoop/sbin/start-hdfs.sh

hdfs-web
yarn-web

hive

create user ‘hadoop‘@’%’ identified by ‘hadoop%HADOOP%123’;

grant all privileges on . to ‘hadoop‘@’%’ with grant option;

flush privileges;

flink on yarn部署

在flink lib下增加flink-hadoop兼容包Pre-bundled Hadoop 2.7.5

1
2
# 启动
bin/start-cluster.sh

搭建AthenaX

Building AthenaX and Flink

mac install kafka

1
2
brew install zookeeper
brew install kafka

修改 /usr/local/etc/kafka/server.properties, 找到 listeners=PLAINTEXT://:9092 那一行,把注释取消掉。

1
2
$ brew services start zookeeper
$ brew services start kafka

创建topic

1
/usr/local/bin/kafka-topics --create --zookeeper localhost:2181 --replication-factor 1 --partitions 1 --topic order-streaming-test

清理topic

cd /usr/local/var/lib/kafka-logs

install kafka on mac

RPC - Remote Procedure Call 远程过程调用
是一种进程间通信方式。它允许程序调用另一个地址空间(通常是共享网络的另一台机器上)的过程或函数,而不用程序员显式编码这个远程调用的细节。即程序员无论是调用本地的还是远程的,本质上编写的调用代码基本相同。

RPC 的鼻祖Bruce Jay Nelson 在1980s提出了基本的实现结构
RPC基本实现结构

图中的user即为client

1 常见的RPC协议

CORBAR: 为了解决异构平台的 RPC,使用了 IDL(Interface Definition Language)来定义远程接口,并将其映射到特定的平台语言中。后来大部分的跨语言平台 RPC 基本都采用了此类方式,比如我们熟悉的 Web Service(SOAP),近年开源的Thrift 等。他们大部分都通过 IDL 定义,并提供工具来映射生成不同语言平台的 user-stub 和 server-stub,并通过框架库来提供 RPCRuntime 的支持。不过貌似每个不同的 RPC 框架都定义了各自不同的 IDL 格式,导致程序员的学习成本进一步上升(苦逼啊),Web Service 尝试建立业界标准,无赖标准规范复杂而效率偏低,否则 Thrift 等更高效的 RPC 框架就没必要出现了。

IDL 是为了跨平台语言实现 RPC 不得已的选择,要解决更广泛的问题自然导致了更复杂的方案。而对于同一平台内的 RPC 而言显然没必要搞个中间语言出来,例如Java原生的 RMI,这样对于 java 程序员而言显得更直接简单,降低使用的学习成本。目前市面上提供的 RPC 框架已经可算是五花八门,百家争鸣了。需要根据实际使用场景谨慎选型,需要考虑的选型因素我觉得至少包括下面几点:

  1. 性能指标
  2. 是否需要跨语言平台
  3. 内网开放还是公网开放
  4. 开源 RPC 框架本身的质量、社区活跃度

1.1 CORBAR

1.2 RMI

1.3 Web Service(SOAP)

1.4 Thrift

2 异步调用与同步调用

  1. 同步调用
    客户方等待调用执行完成并返回结果。
  2. 异步调用
    客户方调用后不用等待执行结果返回,但依然可以通过回调通知等方式获取返回结果。
    若客户方不关心调用返回结果,则变成单向异步调用,单向调用不用返回结果。

2.1 RPC 结构拆解

img

RPC 服务方通过 RpcServer 去导出(export)远程接口方法,而客户方通过 RpcClient 去引入(import)远程接口方法。客户方像调用本地方法一样去调用远程接口方法,RPC 框架提供接口的代理实现,实际的调用将委托给代理RpcProxy 。代理封装调用信息并将调用转交给RpcInvoker 去实际执行。在客户端的RpcInvoker 通过连接器RpcConnector 去维持与服务端的通道RpcChannel,并使用RpcProtocol 执行协议编码(encode)并将编码后的请求消息通过通道发送给服务方。

RPC 服务端接收器 RpcAcceptor 接收客户端的调用请求,同样使用RpcProtocol 执行协议解码(decode)。解码后的调用信息传递给RpcProcessor 去控制处理调用过程,最后再委托调用给RpcInvoker 去实际执行并返回调用结果。

2.2 RPC 组件职责

上面我们进一步拆解了 RPC 实现结构的各个组件组成部分,下面我们详细说明下每个组件的职责划分。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
1. RpcServer  
   负责导出(export)远程接口  
2. RpcClient  
   负责导入(import)远程接口的代理实现  
3. RpcProxy  
   远程接口的代理实现  
4. RpcInvoker  
   客户方实现:负责编码调用信息和发送调用请求到服务方并等待调用结果返回  
   服务方实现:负责调用服务端接口的具体实现并返回调用结果  
5. RpcProtocol  
   负责协议编/解码  
6. RpcConnector  
   负责维持客户方和服务方的连接通道和发送数据到服务方  
7. RpcAcceptor  
   负责接收客户方请求并返回请求结果  
8. RpcProcessor  
   负责在服务方控制调用过程,包括管理调用线程池、超时时间等  
9. RpcChannel  
   数据传输通道  

2.3 RPC 实现分析

在进一步拆解了组件并划分了职责之后,这里以在  Java 平台实现该 RPC 框架概念模型为例,详细分析下实现中需要考虑的因素。

2.3.1 导出远程接口

导出远程接口的意思是指只有导出的接口可以供远程调用,而未导出的接口则不能。在 java 中导出接口的代码片段可能如下:

1
2
3
DemoService demo   = new ...;  
RpcServer   server = new ...;  
server.export(DemoService.class, demo, options);  

我们可以导出整个接口,也可以更细粒度一点只导出接口中的某些方法,如:

1
2
// 只导出 DemoService 中签名为 hi(String s) 的方法  
server.export(DemoService.class, demo, "hi"new Class<?>[] { String.class }, options);  

java 中还有一种比较特殊的调用就是多态,也就是一个接口可能有多个实现,那么远程调用时到底调用哪个?这个本地调用的语义是通过 jvm 提供的引用多态性隐式实现的,那么对于 RPC 来说跨进程的调用就没法隐式实现了。如果前面DemoService 接口有 2 个实现,那么在导出接口时就需要特殊标记不同的实现,如:

1
2
3
4
5
DemoService demo   = new ...;  
DemoService demo2  = new ...;  
RpcServer   server = new ...;  
server.export(DemoService.class, demo, options);  
server.export("demo2", DemoService.class, demo2, options);  

上面 demo2 是另一个实现,我们标记为 “demo2” 来导出,那么远程调用时也需要传递该标记才能调用到正确的实现类,这样就解决了多态调用的语义。

2.3.2 导入远程接口与客户端代理

导入相对于导出远程接口,客户端代码为了能够发起调用必须要获得远程接口的方法或过程定义。目前,大部分跨语言平台 RPC 框架采用根据 IDL 定义通过 code generator 去生成 stub 代码,这种方式下实际导入的过程就是通过代码生成器在编译期完成的。我所使用过的一些跨语言平台 RPC 框架如 CORBAR、WebService、ICE、Thrift 均是此类方式。

代码生成的方式对跨语言平台 RPC 框架而言是必然的选择,而对于同一语言平台的 RPC 则可以通过共享接口定义来实现。在 java 中导入接口的代码片段可能如下:

1
2
3
RpcClient client = new ...;  
DemoService demo = client.refer(DemoService.class);  
demo.hi("how are you?");  

在 java 中 ‘import’ 是关键字,所以代码片段中我们用 refer 来表达导入接口的意思。这里的导入方式本质也是一种代码生成技术,只不过是在运行时生成,比静态编译期的代码生成看起来更简洁些。java 里至少提供了两种技术来提供动态代码生成,一种是 jdk 动态代理,另外一种是字节码生成。动态代理相比字节码生成使用起来更方便,但动态代理方式在性能上是要逊色于直接的字节码生成的,而字节码生成在代码可读性上要差很多。两者权衡起来,个人认为牺牲一些性能来获得代码可读性和可维护性显得更重要。

2.3.3 协议编解码

客户端代理在发起调用前需要对调用信息进行编码,这就要考虑需要编码些什么信息并以什么格式传输到服务端才能让服务端完成调用。出于效率考虑,编码的信息越少越好(传输数据少),编码的规则越简单越好(执行效率高)。我们先看下需要编码些什么信息:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
-- 调用编码 --  
1. 接口方法  
   包括接口名、方法名  
2. 方法参数  
   包括参数类型、参数值  
3. 调用属性  
   包括调用属性信息,例如调用附件隐式参数、调用超时时间等  
  
-- 返回编码 --  
1. 返回结果  
   接口方法中定义的返回值  
2. 返回码  
   异常返回码  
3. 返回异常信息  
   调用异常信息  

除了以上这些必须的调用信息,我们可能还需要一些元信息以方便程序编解码以及未来可能的扩展。这样我们的编码消息里面就分成了两部分,一部分是元信息、另一部分是调用的必要信息。如果设计一种 RPC 协议消息的话,元信息我们把它放在协议消息头中,而必要信息放在协议消息体中。下面给出一种概念上的 RPC 协议消息设计格式:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
-- 消息头 --  
magic      : 协议魔数,为解码设计  
header size: 协议头长度,为扩展设计  
version    : 协议版本,为兼容设计  
st         : 消息体序列化类型  
hb         : 心跳消息标记,为长连接传输层心跳设计  
ow         : 单向消息标记,  
rp         : 响应消息标记,不置位默认是请求消息  
status code: 响应消息状态码  
reserved   : 为字节对齐保留  
message id : 消息 id  
body size  : 消息体长度  
  
-- 消息体 --  
采用序列化编码,常见有以下格式  
xml   : 如 webservie soap  
json  : 如 JSON-RPC  
binary: 如 thrift; hession; kryo 等  

格式确定后编解码就简单了,由于头长度一定所以我们比较关心的就是消息体的序列化方式。序列化我们关心三个方面:
\1. 序列化和反序列化的效率,越快越好。 
\2. 序列化后的字节长度,越小越好。 
\3. 序列化和反序列化的兼容性,接口参数对象若增加了字段,是否兼容。
上面这三点有时是鱼与熊掌不可兼得,这里面涉及到具体的序列化库实现细节,就不在本文进一步展开分析了。

2.3.4 传输服务

协议编码之后,自然就是需要将编码后的 RPC 请求消息传输到服务方,服务方执行后返回结果消息或确认消息给客户方。RPC 的应用场景实质是一种可靠的请求应答消息流,和 HTTP 类似。因此选择长连接方式的 TCP 协议会更高效,与 HTTP 不同的是在协议层面我们定义了每个消息的唯一 id,因此可以更容易的复用连接。

既然使用长连接,那么第一个问题是到底 client 和 server 之间需要多少根连接?实际上单连接和多连接在使用上没有区别,对于数据传输量较小的应用类型,单连接基本足够。单连接和多连接最大的区别在于,每根连接都有自己私有的发送和接收缓冲区,因此大数据量传输时分散在不同的连接缓冲区会得到更好的吞吐效率。所以,如果你的数据传输量不足以让单连接的缓冲区一直处于饱和状态的话,那么使用多连接并不会产生任何明显的提升,反而会增加连接管理的开销。

连接是由 client 端发起建立并维持。如果 client 和 server 之间是直连的,那么连接一般不会中断(当然物理链路故障除外)。如果 client 和 server 连接经过一些负载中转设备,有可能连接一段时间不活跃时会被这些中间设备中断。为了保持连接有必要定时为每个连接发送心跳数据以维持连接不中断。心跳消息是 RPC 框架库使用的内部消息,在前文协议头结构中也有一个专门的心跳位,就是用来标记心跳消息的,它对业务应用透明。

2.3.5 执行调用

client stub 所做的事情仅仅是编码消息并传输给服务方,而真正调用过程发生在服务方。server stub 从前文的结构拆解中我们细分了 RpcProcessor 和RpcInvoker 两个组件,一个负责控制调用过程,一个负责真正调用。这里我们还是以 java 中实现这两个组件为例来分析下它们到底需要做什么?

java 中实现代码的动态接口调用目前一般通过反射调用。除了原生的 jdk 自带的反射,一些第三方库也提供了性能更优的反射调用,因此 RpcInvoker 就是封装了反射调用的实现细节。

调用过程的控制需要考虑哪些因素,RpcProcessor 需要提供什么样地调用控制服务呢?下面提出几点以启发思考:

1
2
3
4
5
6
1. 效率提升  
   每个请求应该尽快被执行,因此我们不能每请求来再创建线程去执行,需要提供线程池服务。  
2. 资源隔离  
   当我们导出多个远程接口时,如何避免单一接口调用占据所有线程资源,而引发其他接口执行阻塞。  
3. 超时控制  
   当某个接口执行缓慢,而 client 端已经超时放弃等待后,server 端的线程继续执行此时显得毫无意义。  

2.4 RPC 异常处理

无论 RPC 怎样努力把远程调用伪装的像本地调用,但它们依然有很大的不同点,而且有一些异常情况是在本地调用时绝对不会碰到的。在说异常处理之前,我们先比较下本地调用和 RPC 调用的一些差异:

  1. 本地调用一定会执行,而远程调用则不一定,调用消息可能因为网络原因并未发送到服务方。
  2. 本地调用只会抛出接口声明的异常,而远程调用还会跑出 RPC 框架运行时的其他异常。
  3. 本地调用和远程调用的性能可能差距很大,这取决于 RPC 固有消耗所占的比重。
    正是这些区别决定了使用 RPC 时需要更多考量。当调用远程接口抛出异常时,异常可能是一个业务异常,也可能是 RPC 框架抛出的运行时异常(如:网络中断等)。业务异常表明服务方已经执行了调用,可能因为某些原因导致未能正常执行,而 RPC 运行时异常则有可能服务方根本没有执行,对调用方而言的异常处理策略自然需要区分。

由于 RPC 固有的消耗相对本地调用高出几个数量级,本地调用的固有消耗是纳秒级,而 RPC 的固有消耗是在毫秒级。那么对于过于轻量的计算任务就并不合适导出远程接口由独立的进程提供服务,只有花在计算任务上时间远远高于 RPC 的固有消耗才值得导出为远程接口提供服务。

2.5 总结

至此我们提出了一个 RPC 实现的概念框架,并详细分析了需要考虑的一些实现细节。无论 RPC 的概念是如何优雅,但是“草丛中依然有几条蛇隐藏着”,只有深刻理解了 RPC 的本质,才能更好地应用。

1
2
3
4
5
6
7
8
9
10
docker serach superset
docker pull amancevice/superset
docker images
mkdir ~/apps/data/superset/
docker run -d -p 8088:8088 -v ~/apps/data/superset/:/home/superset amancevice/superset
docker ps

CONTAINER ID IMAGE COMMAND CREATED STATUS PORTS NAMES
71a1400d1e98 amancevice/superset "gunicorn superset.a…" 10 seconds ago Up 9 seconds (health: starting) 0.0.0.0:8088->8088/tcp quizzical_wright
docker exec -it 71a1400d1e98 superset-init

1 mysql

启动mysql

brew services start mysql

abc.ABC.123

2 集群搭建

  1. mac osx 搭建hadoop开发环境
  2. spark on yarn
  3. hive环境搭建

按照上面文章的方式搭建,启动hdfs、yarn

必须使用这个命令

1
2
3
4
5
6
## 启动yarn
hadoop/sbin/./yarn-daemon.sh start resourcemanager
hadoop/sbin./yarn-daemon.sh start nodemanager
## 启动hdfs
hadoop/sbin/start-hdfs.sh

hdfs-web
yarn-web

2.1 hive

create user ‘hadoop‘@’%’ identified by ‘hadoop%HADOOP%123’;

grant all privileges on . to ‘hadoop‘@’%’ with grant option;

flush privileges;

flink on yarn部署

在flink lib下增加flink-hadoop兼容包Pre-bundled Hadoop 2.7.5

1
2
# 启动
bin/start-cluster.sh

2.3 搭建AthenaX

Building AthenaX and Flink