全部博文(51)
分类: 大数据
2017-11-09 10:08:39
Plumber的设计可以与Flume进行类比。
Plumber的设计和开发思路是基于Flume的,但是实际上只要可以执行注册/注销,并按照格式上报心跳,任何组件都可以作为Plumber的Souce/Sink使用。
目前Plumber使用Flume作为Source,使用Kafka2HDFS作为Sink。Source作为Plumber Agent Source的一个实现例子。
二,监控与心跳
Plumber的数据采集监控主要目标:
目前考虑的准确性主要是分时段对比,每个小时一条汇总数据
心跳数据通过Kafka进行收集,这样做有以下几个好处:
每一次心跳消息中,Agent上报当前节点的采集状态(每个文件采集了多少record,多少byte等)
优点心跳使用Kafka KeyedMessage发送到Kafka。
key的数据要保证同一个agent被发送到Kafka的同一个topic的同一个partition里面去。Key使用ip:port的格式,example:
127.0.0.1:10086
未维护或者不适用的字段上报-1
考虑序列化压力不大, value采用Json格式,便于直接消费检查。
{
"timestamp" : 1470123010 , //时间戳,精度到毫秒
"type" : "source" , //类型,source/sink
"data" : [ //心跳数据
{
"topic" : "app-test2" , //处理的topic
"recordCounter" : 1238432, //启动开始到现在处理的条数
"items":[
{
"timeMap" : 1470123000, //时间段,通常截取到了小时,精确到毫秒
"fileNum" : 5 , //文件数量, 如果不适用,此字段可以上报-1
"fileSize" : 65535, //文件实际大小, 如果不适用,此字段可以上报-1
"bytes" : 6423, //已经处理的字节, 如果不适用,此字段可以上报-1
"records" : 230 //已经处理的record数量
},
{
"timeMap" : 1470123000, //时间段,通常截取到了小时,精确到毫秒
"fileNum" : 5 , //文件数量, 如果不适用,此字段可以上报-1
"fileSize" : 65535, //文件实际大小, 如果不适用,此字段可以上报-1
"bytes" : 6423, //已经处理的字节, 如果不适用,此字段可以上报-1
"records" : 230 //已经处理的record数量
}
]
} //第一条数据
]
}
topic 心跳默认使用的Kafka topic名称为 plumber