flume拦截器

拦截器做用:拦截器是简单的插件式组件,设置在source和channel之间。source接收到的事件,在写入channel以前,拦截器均可以进行转换或者删除这些事件。每一个拦截器只处理同一个source接收到的事件。能够自定义拦截器。git

flume修改时间戳的插件见 https://github.com/haebin/flume-timestamp-interceptorgithub

 

有一个缺陷是,DateUtils.parseDate(timestamp, dateFormat)里面的dateFormat不支持unix时间戳,只能本身手动添加了apache

原来是:app

  1. String timestamp = get(index, data);
  2. now = DateUtils.parseDate(timestamp, dateFormat).getTime();
  3. headers.put(TIMESTAMP, Long.toString(now));

修改后ui

  1. String timestamp = get(index, data);
  2. if (dateFormat[0].equals("tsecond")){
  3. now = Long.parseLong(timestamp)*1000;
  4. }
  5. else if(dateFormat[0].equals("tmillisecond")){
  6. now = Long.parseLong(timestamp);
  7. }
  8. else if(dateFormat[0].equals("tnanosecond")){
  9. now = Long.parseLong(timestamp)/1000000;
  10. }
  11. else {
  12. now = DateUtils.parseDate(timestamp, dateFormat).getTime();
  13. }
  14. headers.put(TIMESTAMP, Long.toString(now));

flume配置:spa

  1. kafka_sn_hive.sources.s1.interceptors = timestamp
  2. kafka_sn_hive.sources.s1.interceptors.timestamp.type = org.apache.flume.interceptor.EventTimestampInterceptor$Builder
  3. kafka_sn_hive.sources.s1.interceptors.timestamp.preserveExisting = false
  4. kafka_sn_hive.sources.s1.interceptors.timestamp.delimiter = ,
  5. kafka_sn_hive.sources.s1.interceptors.timestamp.dateIndex = 4
  6. kafka_sn_hive.sources.s1.interceptors.timestamp.dateFormat = tsecond

表示按逗号作分隔符的第四个(从0开始)字段是一个秒单位的时间戳。插件

在flume里面,时间戳是毫秒级别的,因此要判断这个字段是秒仍是毫秒纳秒unix

 

见http://lisux.me/lishuai/?p=867orm

相关文章
相关标签/搜索