腾讯大规模Hadoop集群实践
腾讯数据平台部
翟艳堂(hakeemzhai)
腾讯服务总体框架
Lhotse统一调度
TDW
海量数据存
储与计算
TDBank
数据采
集与分
发
自助提取与分析
数据规范化管理
TRC 实时计算平台
实时采集 流式计算 分布式存储
精准推荐模型
社交广告 电商 视频 其它
数据仓库 数据应用
数
据
分
析
精
准
推
荐
数据开发者平台 + 数据应用门户
专题分析 SNG
IEG
MIG
CDG
ECC
TEG
OMG
腾讯分布式数据仓库(TDW)
任务统一调度 Lhotse IDE
精准推荐 用户画像 数据挖掘
TDBank
数据采集
业务报表
TDW Core
计算引擎
MapReduce
存储引擎
HDFS 查询引擎
Hive
为什么要做大集群
数据共享
计算资源共享
减轻运营负担
同乐微博集群
200台
同乐主集群
SNG/OMG/ECC
1250台
TS4+TS5
福永微博集群
350台
宝安主集群
IEG/MIG/…
450台TS4
SOSO集群
财付通集群
……
宝安公用集群
500台
宝安TA集群
宝安智胜集群
……
枢纽点击流集群
500台
宝安数挖集群
250台
南汇集群
TDW
面临的挑战
400台 4000台
Ø 计算层 Ø 存储层
l NameNode没有容灾
丢失1个小时数据的风险
重启耗时长
不支持灰度变更
l JobTracker调度效率低
集群扩展性不好
高可用
高效
高扩展性
高可用
高效
高扩展性
JobTracker分散化
NameNode高可用
1年
JobTracker分散化
方案选择
TDW基线版本: CDH3u3
Yarn
Corona
版本稳定性 社区开发中,稳定版发布时间未知 facebook发布的版本
代码复杂度 系列代码,完全重构 基于系列代码
HDFS的要求 HDFS
系统HDFS
时间:2012年12月
JobTracker分散化
Task
Tracker
Task
Tracker ...
资源管理
任务调度
任务管理
Cluster
Manager
资源管理
任务调度
JobTracker 任务管理
Job
Tracker
… 任务管理
Task
Tracker
Task
Tracker ...
Task
Tracker
Ø JobTracker分散化平行扩展
Ø 资源管理和任务调度解耦
Ø 更精细地调度
JobTracker分散化
Cluster
Manager
JobClient
Task
Tracker
Task
Tracker ...
Task
Tracker
1. request jobtracker
resource
2. grant jobtracker
resource
3. start jobtracker
JobTracker
4. request
map/reduce
resource
5. grant
map/reduce
resource
6. submit
launch map/
reduce actions
map/
reduce
heartbeat
heartbeat
heartbea
t
JobTracker分散化
Ø 任务调度依赖于心跳
Ø Map/Reduce任务调度耦合
Ø 全排序(nlog(n))
基于心跳模型的拉取调度 独立并发式的下推调度
ClusterManager
Map调度
ClusterManager
Reduce调度
Ø 任务调度即需即用
Ø Map/Reduce任务调度独立、并发
Ø 优先级队列堆排序(log(n))
PULL
Task
Tracker
Task
Tracker ...
Task
Tracker
PUSH PUSH
JobTracker
串行调度
Task
Tracker
NameNode高可用
NameNode高可用
Data
Node
Data
Node ...
Name
Node
Second
NameNode
check
point
blockrepor
t
client
meta ops
ANN BNN BNN
zk1 zk2 …
sync edit
log
learn
meta
Data
Node
Data
Node ...
blockrepo
rt
client
meta ops
Ø 一主两热备
Ø 元数据在主备间实时同步
Ø DataNode同时向3个Master汇报Block
Namenode主备仲裁以及状态转换
Active:IP1 Standby:IP2 client
zookeeper cluster
Standby:IP3
../election/1,2,3 ../repView/ A:ip1,S:ip2,S:ip3
heatbeat X
newbie:IP1 Active:IP2 client Standby:IP3
../election/2,3 ../repView/ A:ip2,S:ip3,n:ip1
heatbeat
zookeeper cluster
重新获取主 更新repview
newbie:IP1 Active:IP2 client Standby:IP3
../election/2,3 ../repView/ A:ip2,S:ip3,n:ip1
heatbeat
重新学习,收集DN状态
standby:IP1 Active:IP2 client Standby:IP3
../election/2,3 4 ../repView/ A:ip2,S:ip3,S:ip1
heatbeat
NameNode分散化
NameNode分散化
Hive Meta
user
Tbl_a, Tbl_b
Tbl_a
namenode 1
Tbl_b namenode 3
...
....
Tbl_a Tbl_b
submit mr
Hive
user
submit mr
获取NN信息
Namenode DN DN ...
计算层
计算层
Ø 按业务分布
Ø 按负载分布资源
HDFS Cluster1
(NameNode1)
HDFS Cluster2
(NameNode2)
HDFS Cluster3
(NameNode3)
其他优化
HDFS兼容
CDH3u3
NameNode
DFSClien
t
DFSClien
t
NameNode
RPC
Server
RPC
Server
DataNode DataNode
FileSyste
m
Abstract
FileSyste
m
FileSyste
m
Abstract
FileSyste
m
routing
table
routing
table
HDFS HDFS
Ø 1个节点慢,整个job慢
监控数据库
CPU利用率最高的/最低的
reduce平均执行时间最大
的
推测执行差异化服务
Ø 一视同仁
• 资源浪费
Ø 关键任务不能执行慢,非关键任务不能卡死
关键任务 非关键任务
推测比例 90% 1%
推测间隔 5s 30m
检测节点短板
Ø 误删除数据将会造成灾难
• NameNode回收站
• 删除黑白名单
• DataNode回收站
防止数据误删除
大Job的困扰
Ø 资源池限制
Ø 生产时段和非生产时段动
态调整
Ø 下手狠一点
Job元数据分散化
Reduce拖取Map数据重试
client
MapReduce
HDFS1 HDFS2
job2
meta job1
meta
Ø 等待时间线性增长
Ø 节点自动下线
集群发展现状
集群容量
–服务器 5600台
–CPU ~10w+核
–内存 ~350TB
–磁盘 ~67200块
–存储容量 ~100PB
单集群支撑规模 400
5600
每日作业数 4万 100万+
每日计算量
存储利用率 85%+ 84%+
CPU利用率 30% 90%+
数据安全性 可能会丢失1个小
时数据
丢数据风险很低
重启暂停线上
服务时间
1小时 秒级自动无缝切换
总存储量 4PB 86PB
文件数+块数 5千万 6亿
未来计划
未来方向——实时化
快
p 内存化
p DAG
Ø 输入数据来源于内存
Ø shuffle数据放内存
Ø 计算结果放内存
Ø 扩充Map->Reduce单一的计算过程
Ø 多个计算过程组成DAG
Spark
Spark介绍
Input
group by
map
union
Output
HDFS read
HDFS write
p DAG模型
p 内存计算
Ø 数据缓存在内存中可以重复使用
Ø 整个计算过程是一个DAG模型
Ø 不需要多次从HDFS输入输出
join
后续计划
p 资源管理精细化(使用节点实际负载来判断是否下发任务)
p 计算本地化和节点负载之间平衡
p Shuffle优化
p 更好的集群扩展性
Ø Spark on GAIA
引入GAIA负责资源管理、
任务调度、资源隔离
GAIA Master
(资源管理、任务调度)
Spark
Client 申请资源
GAIA Slave
(任务执行)
分配Spark Master
GAIA Slave
(任务执行)
分配Spark Task
GAIA
(资源管理、任务调度、资源隔离)
Spark
Spark
Master
Spark
Executor
GAIA Slave
(任务执行)
Spark
Executor
谢谢