设为首页 收藏本站
查看: 989|回复: 0

[经验分享] Flume-TailFileSource

[复制链接]

尚未签到

发表于 2017-5-21 12:44:02 | 显示全部楼层 |阅读模式
  package org.apache.flume.source;
  import java.io.File;
import java.io.RandomAccessFile;
import java.util.HashSet;
import java.util.Iterator;
import java.util.List;
import java.util.Set;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;
  import org.apache.flume.Context;
import org.apache.flume.Event;
import org.apache.flume.EventDrivenSource;
import org.apache.flume.conf.Configurable;
import org.apache.flume.event.EventBuilder;
import org.apache.flume.instrumentation.SourceCounter;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
  public class TailFileSource extends AbstractSource implements EventDrivenSource,
  Configurable {
   private static final Logger logger = LoggerFactory.getLogger(TailFileSource.class);
 private SourceCounter sourceCounter;
 private String pointerFile;
 private String tailFile;
 private long collectInterval;
 private int batchLine;
 private boolean batch;
 private AtomicBoolean tailRun=new AtomicBoolean(true);
 private AtomicLong cursor=new AtomicLong(0);
 
  
 @Override
 public void configure(Context context) {
  if (sourceCounter == null) {
   sourceCounter = new SourceCounter(getName());
  }
  pointerFile = context.getString("tailfile.pointer","cursor.pt");
  tailFile = context.getString("tailfile.file","data.txt");
  collectInterval=context.getLong("tailfile.collectInterval",3000L);
  batchLine=context.getInteger("tailfile.batchLine",10);
  batch=context.getBoolean("tailfile.batch",false);
  
 }
  @Override
 public synchronized void start() {
  super.start();
  sourceCounter.start();
  TailFile tf=new TailFile();
  tf.addFileTailerListener(new FileTailerListener(){
   @Override
   public void newFileLine(List<Event> events) {
    if(!events.isEmpty()){
     getChannelProcessor().processEventBatch(events);
     sourceCounter.incrementAppendBatchAcceptedCount();
        sourceCounter.addToEventAcceptedCount(events.size());
    }
   }
  @Override
   public void newFileLine(Event event) {
     getChannelProcessor().processEvent(event);
     sourceCounter.incrementAppendAcceptedCount();
   }
  });
  Thread t=new java.lang.Thread(tf);
  t.setDaemon(true);
  t.start();
 }
  @Override
 public synchronized void stop() {
  tailRun.set(false);
  super.stop();
  sourceCounter.stop();
 }
 
 
 protected interface FileTailerListener{
     public void newFileLine(List<Event> events);
     public void newFileLine(Event event);
 }
 
 
 protected class TailFile implements java.lang.Runnable{
  private Set<FileTailerListener> listeners = new HashSet<FileTailerListener>();
  
  
     public void addFileTailerListener(FileTailerListener l) {
         this.listeners.add(l);
     }
  public void removeFileTailerListener(FileTailerListener l) {
         this.listeners.remove(l);
     }
  @Override
  public void run() {
   long[] st=this.readPointerFile();
   RandomAccessFile file=null;
   boolean flag=true;
   while(flag){
    try {
     File tf=new File(tailFile);
     file = new RandomAccessFile(tf, "r");
     if(st[0]==tf.lastModified()){
      cursor.set(st[1]);
     }else{
      st[0]=tf.lastModified();
      cursor.set(0);
     }
     flag=false;
    } catch (Exception e) {
     try {
      logger.error(e.getMessage()+",will retry file:"+tailFile);
      Thread.sleep(5000);
     } catch (Exception e1) {
      
     }
    }
   }
   
   while(tailRun.get()){
    try {
     if(!this.sameTailFile(st[0])) {
      logger.error("file change:"+tailFile);
      File tf=new File(tailFile);
      file = new RandomAccessFile(tf, "r");
      st[0]=tf.lastModified();
      cursor.set(0);
     }
     
     long fileLength =file.length();
     if (fileLength < cursor.get()) {
         cursor.set(fileLength);
     }
     
     if (fileLength > cursor.get()) {
         file.seek(cursor.get());
         String line = file.readLine();
         int i=1;
        
         if(batch){
          java.util.List<Event> batchAl=new java.util.ArrayList<Event>(batchLine);
          while (line != null) {
           batchAl.add(EventBuilder.withBody(line.getBytes()));
           if(i%batchLine==0) {
            fireNewFileLine(batchAl);
            batchAl.clear();
            cursor.set(file.getFilePointer());
            st[1]=cursor.get();
            writePointerFile(st);
           }
              line = file.readLine();
              i++;
          }
         
          if(!batchAl.isEmpty()){
           fireNewFileLine(batchAl);
           batchAl.clear();
           cursor.set(file.getFilePointer());
           st[1]=cursor.get();
           writePointerFile(st);
          }
         }else{
          while(line!=null){
           fireNewFileLine(EventBuilder.withBody(line.getBytes()));
           line = file.readLine();
          }
           cursor.set(file.getFilePointer());
           st[1]=cursor.get();
           writePointerFile(st);
         }
  }
     Thread.sleep(collectInterval);
    } catch (Exception e) {
     
    }
    
   }
   
   try {
    if(file!=null) file.close();
   } catch (Exception e) {
    
   }
  }
  
  
  private long[] readPointerFile(){
   logger.info("read pointerFile:"+pointerFile);
   java.io.ObjectInputStream ois=null;
   long[] temp={0L,0L};
   try {
    ois=new java.io.ObjectInputStream(new java.io.FileInputStream(new File(pointerFile)));
    temp[0]=ois.readLong();
    temp[1]=ois.readLong();
   } catch (Exception e) {
    logger.error("can't read pointerFile:"+pointerFile);
   } finally{
    try {
     if(ois!=null)ois.close();
    } catch (Exception e) {
     
    }
   }
   return temp;
  }
  
  private void writePointerFile(long[] temp){
   logger.debug("write pointerFile:"+pointerFile);
   java.io.ObjectOutputStream oos=null;
   try {
    oos=new java.io.ObjectOutputStream(new java.io.FileOutputStream(pointerFile));
    oos.writeLong(temp[0]);
    oos.writeLong(temp[1]);
   } catch (Exception e) {
    logger.error("can't write pointerFile:"+pointerFile);
   }finally{
    try {
     if(oos!=null)oos.close();
    } catch (Exception e) {
     
    }
   }
  
  }
  
  private boolean sameTailFile(long time){
   return new File(tailFile).lastModified()==time?true:false;
  }
  
     
     protected void fireNewFileLine(List<Event> events) {
         for (Iterator<FileTailerListener> i = this.listeners.iterator(); i.hasNext();) {
          FileTailerListener l =  i.next();
             l.newFileLine(events);
         }
     }
    
    
     protected void fireNewFileLine(Event event) {
         for (Iterator<FileTailerListener> i = this.listeners.iterator(); i.hasNext();) {
          FileTailerListener l =  i.next();
             l.newFileLine(event);
         }
     }
 }
}

运维网声明 1、欢迎大家加入本站运维交流群:群②:261659950 群⑤:202807635 群⑦870801961 群⑧679858003
2、本站所有主题由该帖子作者发表,该帖子作者与运维网享有帖子相关版权
3、所有作品的著作权均归原作者享有,请您和我们一样尊重他人的著作权等合法权益。如果您对作品感到满意,请购买正版
4、禁止制作、复制、发布和传播具有反动、淫秽、色情、暴力、凶杀等内容的信息,一经发现立即删除。若您因此触犯法律,一切后果自负,我们对此不承担任何责任
5、所有资源均系网友上传或者通过网络收集,我们仅提供一个展示、介绍、观摩学习的平台,我们不对其内容的准确性、可靠性、正当性、安全性、合法性等负责,亦不承担任何法律责任
6、所有作品仅供您个人学习、研究或欣赏,不得用于商业或者其他用途,否则,一切后果均由您自己承担,我们对此不承担任何法律责任
7、如涉及侵犯版权等问题,请您及时通知我们,我们将立即采取措施予以解决
8、联系人Email:admin@iyunv.com 网址:www.yunweiku.com

所有资源均系网友上传或者通过网络收集,我们仅提供一个展示、介绍、观摩学习的平台,我们不对其承担任何法律责任,如涉及侵犯版权等问题,请您及时通知我们,我们将立即处理,联系人Email:kefu@iyunv.com,QQ:1061981298 本贴地址:https://www.yunweiku.com/thread-379587-1-1.html 上篇帖子: flume-ng安装 下篇帖子: flume iterceptor
您需要登录后才可以回帖 登录 | 立即注册

本版积分规则

扫码加入运维网微信交流群X

扫码加入运维网微信交流群

扫描二维码加入运维网微信交流群,最新一手资源尽在官方微信交流群!快快加入我们吧...

扫描微信二维码查看详情

客服E-mail:kefu@iyunv.com 客服QQ:1061981298


QQ群⑦:运维网交流群⑦ QQ群⑧:运维网交流群⑧ k8s群:运维网kubernetes交流群


提醒:禁止发布任何违反国家法律、法规的言论与图片等内容;本站内容均来自个人观点与网络等信息,非本站认同之观点.


本站大部分资源是网友从网上搜集分享而来,其版权均归原作者及其网站所有,我们尊重他人的合法权益,如有内容侵犯您的合法权益,请及时与我们联系进行核实删除!



合作伙伴: 青云cloud

快速回复 返回顶部 返回列表