flume安裝與使用

簡(jiǎn)介

安裝

java開(kāi)發(fā)

依賴(lài)

<dependency>
        <groupId>org.apache.flume</groupId>
        <artifactId>flume-ng-core</artifactId>
        <version>1.5.2</version>
</dependency>
<dependency>
        <groupId>org.apache.flume</groupId>
        <artifactId>flume-ng-configuration</artifactId>
        <version>1.5.2</version>
</dependency>

說(shuō)明

192.168.1.100是hadoop主服務(wù)器
192.168.60.8是安裝flume的服務(wù)器

配置文件
example.conf(192.168.60.8)

agent1.channels = c1
agent1.sources = r1
agent1.sinks = sink1

agent1.channels.c1.type = memory

agent1.sources.r1.channels = c1
agent1.sources.r1.type = avro
agent1.sources.r1.bind = 0.0.0.0
agent1.sources.r1.port = 41414

agent1.sinks.sink1.type=hdfs
agent1.sinks.sink1.hdfs.path=hdfs://192.168.1.100:50040/flume
agent1.sinks.sink1.hdfs.fileType=DataStream
agent1.sinks.sink1.hdfs.writeFormat=TEXT
agent1.sinks.sink1.hdfs.rollInterval=4
agent1.sinks.sink1.channel=c1

運(yùn)行配置

./bin/flume-ng agent -n agent1 -c conf -f example.conf -Dflume.root.logger=DEBUG,console

java程序

import java.nio.charset.Charset;
import org.apache.flume.Event;
import org.apache.flume.EventDeliveryException;
import org.apache.flume.api.RpcClient;
import org.apache.flume.api.RpcClientFactory;
import org.apache.flume.event.EventBuilder;

public class MyApp {
      public static void main(String[] args) {
        MyRpcClientFacade client = new MyRpcClientFacade();
        //配置文件機(jī)器地址、端口
        client.init("192.168.60.8", 41414);
        String sampleData = "Hello Flume!";
        for (int i = 0; i < 10; i++) {
          client.sendDataToFlume(sampleData);
        }
        client.cleanUp();
      }
}

    class MyRpcClientFacade {
      private RpcClient client;
      private String hostname;
      private int port;

      public void init(String hostname, int port) {
        this.hostname = hostname;
        this.port = port;
        this.client = RpcClientFactory.getDefaultInstance(hostname, port);
      }

      public void sendDataToFlume(String data) {
        Event event = EventBuilder.withBody(data, Charset.forName("UTF-8"));
        try {
          client.append(event);
        } catch (EventDeliveryException e) {
          client.close();
          client = null;
          client = RpcClientFactory.getDefaultInstance(hostname, port);
        }
      }

      public void cleanUp() {
        client.close();
        System.out.println("程序執(zhí)行完畢。。。");
      }

    }

查看結(jié)果(192.168.1.100)

/opt/hadoop/bin/hadoop fs -ls /flume
Paste_Image.png
/opt/hadoop/bin/hadoop fs -cat /flume/FlumeData.1467803628159
Paste_Image.png

參考文章

flume開(kāi)發(fā)詳解
開(kāi)發(fā)者API

最后編輯于
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請(qǐng)聯(lián)系作者
【社區(qū)內(nèi)容提示】社區(qū)部分內(nèi)容疑似由AI輔助生成,瀏覽時(shí)請(qǐng)結(jié)合常識(shí)與多方信息審慎甄別。
平臺(tái)聲明:文章內(nèi)容(如有圖片或視頻亦包括在內(nèi))由作者上傳并發(fā)布,文章內(nèi)容僅代表作者本人觀點(diǎn),簡(jiǎn)書(shū)系信息發(fā)布平臺(tái),僅提供信息存儲(chǔ)服務(wù)。

相關(guān)閱讀更多精彩內(nèi)容

  • 1.安裝 Java 安裝Flume配置文件解析Flume Agent啟動(dòng)(1)Flume 配置文件(Agent的配...
    lmem閱讀 1,726評(píng)論 0 0
  • Spring Cloud為開(kāi)發(fā)人員提供了快速構(gòu)建分布式系統(tǒng)中一些常見(jiàn)模式的工具(例如配置管理,服務(wù)發(fā)現(xiàn),斷路器,智...
    卡卡羅2017閱讀 136,553評(píng)論 19 139
  • Linux-Server-Notes PMS /home/softwareluke/圖片/2017-09-11 0...
    燕京博士閱讀 637評(píng)論 0 1
  • Type 初始化變量,需要明確值 References const 當(dāng)變量沒(méi)有改變引用(reassign)時(shí),優(yōu)先...
    我_巨可愛(ài)閱讀 508評(píng)論 0 0
  • 首先,向簡(jiǎn)書(shū)說(shuō)聲抱歉。 昨天為了讓更多的人知道"瓦工學(xué)社",我多次在與創(chuàng)業(yè)相關(guān)的文章下面評(píng)論,做宣傳,違反了簡(jiǎn)書(shū)的...
    Rafi閱讀 157評(píng)論 0 1

友情鏈接更多精彩內(nèi)容