Apache Flume Tutorial: Vad är, Architecture & Hadoop Exempel

⚡ Smart sammanfattning

Apache Flume är en distribuerad tjänst för att samla in, aggregera och flytta stora volymer loggdata till HDFS, byggd kring agenter som länkar samman en källa, en kanal och en sink.

  • 🔘 Agentens anatomi: Varje Flume-agent är en JVM-process som innehåller en källa, en eller flera kanaler och en sink.
  • ☑️ Pålitlighet: Leverans enligt bästa förmåga tolererar inga nodfel; leverans från början till slut överlever flera nodfel.
  • Setup: Anpassade källkodeklasser kompileras till en JAR som släpps i Flume lib-katalogen.
  • 🧪 konfiguration: En egenskapsfil namnger källan, kanalen och sinken och anger HDFS-sökvägen och rullningsgränserna.
  • 🛠️ Lansera: Starta pipelinen med flume-ng-agenten, namnge agenten och peka på flume.conf.
  • ⚠️ Daterat exempel: Slutpunkten för Twitter v1.1-strömning stängdes i mars 2023, så behandla övningen som ett anpassat källmönster.

Apache Flume-handledning som täcker agentarkitektur, konfiguration och ett exempel på Hadoop-strömning

Vad är Apache Flume i Hadoop?

Apache Flume är ett pålitligt och distribuerat system för att samla in, aggregera och flytta stora mängder loggdata. Det har en enkel men flexibel arkitektur baserad på strömmande dataflöden. Apache Flume används för att samla in loggdata som finns i loggfiler från webbservrar och aggregera den till HDFS för analys.

Flume i Hadoop stöder flera källor, inklusive:

  • 'tail' (som överför data från en lokal fil och skriver den till HDFS via Flume, liknande Unix-kommandot 'tail')
  • Systemloggar
  • Apache log4j (vilket möjliggör Java applikationer för att skriva händelser till filer i HDFS via Flume).

Den nuvarande utgåvan är Ränna 1.11.0, publicerad den 25 oktober 2022 och tillgänglig från Apache Flume nedladdningssidaDenna genomgång skrevs mot version 1.4.0, så flera steg nedan innehåller en anmärkning där en aktuell version beter sig annorlunda.

Flume Architecture

Ett Flume-medel är ett JVM process med tre komponenter – rännkälla, rännkanal och rännsänka – genom vilka händelser fortplantas efter att ha initierats vid en extern källa. Diagrammet nedan visar hur de är kopplade.

Flumearkitekturdiagram som visar en agent med en källa, en kanal och en sänka som matar HDFS

  1. Händelserna som genereras av den externa källan (en webbserver) förbrukas av Flume-källan. Den externa källan skickar händelser till Flume-källan i ett format som målkällan känner igen.
  2. Flume-källan tar emot en händelse och lagrar den i en eller flera kanaler. Kanalen fungerar som ett minne som lagrar händelsen tills den förbrukas av Flume-sinken. Denna kanal kan använda ett lokalt filsystem för att lagra dessa händelser.
  3. Flume-sinken tar bort händelsen från en kanal och lagrar den i ett externt datalager, till exempel HDFS. Det kan finnas flera Flume-agenter, i vilket fall Flume-sinken vidarebefordrar händelsen till Flume-källan för nästa agent i flödet.

Några viktiga funktioner hos Flume

  • Flume har en flexibel design baserad på strömmande dataflöden. Den är feltolerant och robust, med flera redundans- och återställningsmekanismer. Flume erbjuder olika nivåer av tillförlitlighet, inklusive "leverans på bästa sätt" och 'end-to-end-leverans'. Bästa leverans tolererar inte något fel på Flume-noden, medan end-to-end leverans garanterar leverans även vid flera nodfel.
  • Flume överför data mellan källor och sänkor. Denna datainsamling kan antingen vara schemalagd eller händelsestyrd. Flume har sin egen frågebehandlingsmotor, vilket gör det enkelt att transformera varje ny databatch innan den flyttas till den avsedda sänken.
  • Möjligt Flume sjunker inkluderar HDFS och HBaseFlume kan också transportera händelsedata såsom nätverkstrafikdata, data som genereras av sociala medier och e-postmeddelanden.

Installation av flume, bibliotek och källkod

Innan vi börjar med själva processen, se till att du har Hadoop installerat; om inte, arbeta dig igenom hur man installerar Hadoop först. Ändra användaren till 'hduser' (det ID som användes vid konfiguration av Hadoop; du kan byta till det användar-ID som användes under din egen Hadoop-konfiguration).

Terminalen växlar Linux-användaren till hduser innan Flume-installationen börjar

Steg 1) Skapa en ny katalog med namnet 'FlumeTutorial'.

sudo mkdir FlumeTutorial
  1. Ge läs-, skriv- och körbehörigheter.
    sudo chmod -R 777 FlumeTutorial
  2. Kopiera filerna MyTwitterSource.java och MyTwitterSourceForFlume.java in i den här katalogen.

Ladda ner indatafiler härifrån

Kontrollera filbehörigheterna för alla dessa filer, enligt nedan, och bevilja läsbehörighet om den saknas.

Terminalen listar filbehörigheterna på den nedladdade Java källfiler

Steg 2) Ladda ner 'Apache Flume' från https://flume.apache.org/download.html.

Apache Flume 1.4.0 har använts i denna Flume-handledning.

Apache Flume nedladdningssida som visar länken till den binära tarball-filen för att välja

Klicka sedan vidare till en spegel.

Apache-spegelsidan nådd efter att ha klickat på Flume tarball-länken

Steg 3) Kopiera den nedladdade tarball-filen till valfri katalog och t.ex.tracinnehållet med följande kommando.

sudo tar -xvf apache-flume-1.4.0-bin.tar.gz

Terminal extracLägg till Flume-tarballen med kommandot sudo tar -xvf

Detta skapar en ny katalog med namnet apache-flume-1.4.0-bin och extracts filerna till den. Den katalogen kallas i resten av artikeln.

Steg 4) Installation av Flume-biblioteket. Kopiera twitter4j-core-4.0.1.jar, flume-ng-configuration-1.4.0.jar, flume-ng-core-1.4.0.jar och flume-ng-sdk-1.4.0.jar till

/lib/

En eller alla kopierade JAR-filer kan ha körningsbehörigheten uppsatt, vilket kan orsaka problem med kompileringen av kod, så återkalla den. I mitt fall hade twitter4j-core-4.0.1.jar körningsbehörigheten. Jag återkallade den enligt nedan.

sudo chmod -x twitter4j-core-4.0.1.jar

Terminalen återkallar körbehörigheten på twitter4j-kärnans JAR-fil

Efter detta ger kommandot nedan alla läsbehörighet på twitter4j-core-4.0.1.jar.

sudo chmod +rrr /usr/local/apache-flume-1.4.0-bin/lib/twitter4j-core-4.0.1.jar

Observera att jag laddade ner twitter4j-core-4.0.1.jar från Maven-arkivetoch alla Flume JAR-filer, dvs. flume-ng-*-1.4.0.jar, från org.apache.flume-artefakterna.

Ladda data från Twitter med Flume

Steg 1) Gå till katalogen som innehåller källkodsfilerna.

Steg 2) Ställ in CLASSPATH till att innehålla /lib/* och ~/FlumeTutorial/flume/mytwittersource/*.

export CLASSPATH="/usr/local/apache-flume-1.4.0-bin/lib/*:~/FlumeTutorial/flume/mytwittersource/*"

Terminal exporterar CLASSPATH som pekar på Flume-biblioteket och källkodskatalogerna

Steg 3) Kompilera källkoden med kommandot nedan.

javac -d . MyTwitterSourceForFlume.java MyTwitterSource.java

Terminal som kompilerar de två Java källfiler med javac

Steg 4) Skapa en JAR-fil. Skapa först en Manifest.txt-fil med en textredigerare du väljer och lägg till raden nedan i den.

Main-Class: flume.mytwittersource.MyTwitterSourceForFlume

Här är flume.mytwittersource.MyTwitterSourceForFlume namnet på huvudklassen. Observera att du måste trycka på Enter-tangenten i slutet av den här raden, som visas nedan.

Manifest.txt öppnas i en textredigerare med Main-Class-posten

Skapa nu JAR-filen 'MyTwitterSourceForFlume.jar' enligt följande.

jar cfm MyTwitterSourceForFlume.jar Manifest.txt flume/mytwittersource/*.class

Terminalpaketering av de kompilerade klasserna i MyTwitterSourceForFlume.jar

Steg 5) Kopiera denna JAR till /lib/.

sudo cp MyTwitterSourceForFlume.jar <Flume Installation Directory>/lib/

Terminalen kopierar den anpassade källkoden JAR till Flume lib-katalogen

Steg 6) Gå till Flumes konfigurationskatalog, /konf.

Om flume.conf inte finns, kopiera flume-conf.properties.template och byt namn på den till flume.conf.

sudo cp flume-conf.properties.template flume.conf

Terminalen kopierar flume-conf.properties.template till flume.conf

Om flume-env.sh inte finns, kopiera flume-env.sh.template och byt namn på den till flume-env.sh.

sudo cp flume-env.sh.template flume-env.sh

Terminalkopiering av flume-env.sh.template till flume-env.sh

Skapa en Twitter-applikation

Läs detta först. V1.1-strömningen statuses/filter slutpunkten som twitter4j 4.0.1 behöver togs bort den 9 mars 2023, och API v2-filtrerade strömmen som ersatte den, nu på utvecklare.x.com, ligger bakom en betald nivå. Behandla skärmarna nedan som ett anpassat källkodsmönster och peka sedan samma agent mot en fil, exec eller Kafka-källkod.

Steg 1) Skapa en Twitter-applikation genom att logga in på utvecklarportalen.

Twitter-utvecklarens inloggningssida som används för att nå applistan

Hemsida för Twitter-utvecklarkontot visas efter inloggning

Steg 2) Gå till "Mina program" (det här alternativet visas när du klickar på "Ägget"-knappen i det övre högra hörnet).

Sidan Mina applikationer på Twitter-utvecklarportalen

Steg 3) Skapa en ny applikation genom att klicka på "Skapa ny app".

Steg 4) Fyll i ansökningsuppgifterna genom att ange ansökans namn, en beskrivning och en webbplats. Du kan läsa informationen under varje inmatningsruta.

Formulär för att skapa en Twitter-applikation med fält för namn, beskrivning och webbplats

Steg 5) Scrolla ner på sidan, acceptera villkoren genom att markera "Ja, jag godkänner" och klicka på knappen "Skapa din Twitter-applikation".

Kryssrutan för villkor och knappen "skapa ansökan" längst ner i Twitter-formuläret

Steg 6) I fönstret för den nyskapade applikationen, gå till fliken "API-nycklar", skrolla ner på sidan och klicka på knappen "Skapa min åtkomsttoken".

Fliken API-nycklar i den nya Twitter-applikationen innan en åtkomsttoken finns

Information om åtkomsttoken som visas efter att knappen Skapa min åtkomsttoken har använts

Steg 7) Uppdatera sidan.

Steg 8) Klicka på "Testa OAuth". Detta visar programmets "OAuth"-inställningar.

Testskärmen för OAuth som visar programmets OAuth-inställningar

Steg 9) Ändra 'flume.conf' med dessa OAuth-inställningar. Stegen för att ändra 'flume.conf' ges nedan.

OAuth-inställningar som listar konsumentnyckeln, konsumenthemligheten och åtkomsttokenvärdena

Vi måste kopiera konsumentnyckeln, konsumenthemligheten, åtkomsttoken och åtkomsttokenhemligheten för att kunna uppdatera 'flume.conf'.

Obs: Dessa värden tillhör användaren och är därför konfidentiella, så de bör inte delas.

Ändra 'flume.conf'-filen

Steg 1) Öppna 'flume.conf' i skrivläge och ange värden för parametrarna nedan.

sudo gedit flume.conf

Kopiera innehållet nedan.

MyTwitAgent.sources = Twitter
MyTwitAgent.channels = MemChannel
MyTwitAgent.sinks = HDFS
MyTwitAgent.sources.Twitter.type = flume.mytwittersource.MyTwitterSourceForFlume
MyTwitAgent.sources.Twitter.channels = MemChannel
MyTwitAgent.sources.Twitter.consumerKey = <Copy consumer key value from Twitter App>
MyTwitAgent.sources.Twitter.consumerSecret = <Copy consumer secret value from Twitter App>
MyTwitAgent.sources.Twitter.accessToken = <Copy access token value from Twitter App>
MyTwitAgent.sources.Twitter.accessTokenSecret = <Copy access token secret value from Twitter App>
MyTwitAgent.sources.Twitter.keywords = guru99
MyTwitAgent.sinks.HDFS.channel = MemChannel
MyTwitAgent.sinks.HDFS.type = hdfs
MyTwitAgent.sinks.HDFS.hdfs.path = hdfs://localhost:54310/user/hduser/flume/tweets/
MyTwitAgent.sinks.HDFS.hdfs.fileType = DataStream
MyTwitAgent.sinks.HDFS.hdfs.writeFormat = Text
MyTwitAgent.sinks.HDFS.hdfs.batchSize = 1000
MyTwitAgent.sinks.HDFS.hdfs.rollSize = 0
MyTwitAgent.sinks.HDFS.hdfs.rollCount = 10000
MyTwitAgent.channels.MemChannel.type = memory
MyTwitAgent.channels.MemChannel.capacity = 10000
MyTwitAgent.channels.MemChannel.transactionCapacity = 1000

Filen flume.conf öppnas i en editor med MyTwitAgent-käll-, kanal- och sink-egenskaperna

Steg 2) Ställ också in TwitterAgent.sinks.HDFS.hdfs.path enligt nedan.

TwitterAgent.sinks.HDFS.hdfs.path = hdfs:// : / /flume/tweets/

HDFS-sinkens hdfs.path-egenskap är inställd på ett värdnamn, portnummer och HDFS-hemkatalog.

Att hitta , och , se värdet för parametern 'fs.defaultFS' som är inställd i $HADOOP_HOME/etc/hadoop/core-site.xml, som visas nedan.

Egenskapen fs.defaultFS i core-site.xml, som anger värdnamnet och porten

Steg 3) För att spola data till HDFS när de kommer, radera posten nedan om den finns.

TwitterAgent.sinks.HDFS.hdfs.rollInterval = 600

Exempel: Streama Twitter-data med Flume

Steg 1) Öppna 'flume-env.sh' i skrivläge och ange värden för parametrarna nedan.

JAVA_HOME=<Installation directory of Java>
FLUME_CLASSPATH="<Flume Installation Directory>/lib/MyTwitterSourceForFlume.jar"

flume-env.sh öppnas i en editor med JAVA_HOME och FLUME_CLASSPATH uppsatta

Steg 2) Starta Hadoop.

$HADOOP_HOME/sbin/start-dfs.sh
$HADOOP_HOME/sbin/start-yarn.sh

Steg 3) Två av JAR-filerna från Flume-tarballen är inte kompatibla med Hadoop 2.2.0, så i detta Apache Flume-exempel följer vi stegen nedan för att göra Flume kompatibelt med Hadoop 2.2.0. Denna JAR-växling är en fix från 1.4.0-eran; Flume 1.11.0 levererar redan aktuella protobuf- och Guava-byggen, så en modern tarball behöver normalt inget av det.

a. Flytta protobuf-java-2.4.1.jar från ' /lib'. Gå till den katalogen först.

CD /lib

sudo mv protobuf-java-2.4.1.jar ~/

Terminalen flyttar protobuf-java-2.4.1.jar från Flume lib-katalogen

b. Hitta JAR-filen 'guava' enligt nedan.

find . -name "guava*"

Terminal find-kommandot lokaliserar den medföljande guava-JAR-filen

Flytta guava-10.0.1.jar ut ur ' /lib'.

sudo mv guava-10.0.1.jar ~/

Terminalen flyttar guava-10.0.1.jar från Flume lib-katalogen

c. Ladda ner guava-17.0.jar från Maven-arkivet, visas nedan.

Maven Repository-sida för guava 17.0, den nya JAR-filen att ladda ner

Kopiera nu den här nedladdade JAR-filen till ' /lib'.

Steg 4) Gå till ' /bin' och starta Flume enligt följande.

./flume-ng agent -n MyTwitAgent -c conf -f <Flume Installation Directory>/conf/flume.conf

Terminal som startar Flume-agenten MyTwitAgent med kommandot flume-ng

Kommandotolksfönstret där Flume hämtar tweets ser ut så här.

Kommandotolken som visar Flume-agenten som hämtar tweets och skriver dem till HDFS

Från kommandofönstrets meddelande kan vi se att utdata skrivs till katalogen /user/hduser/flume/tweets/. Öppna nu den här katalogen med en webbläsare.

Steg 5) För att se resultatet av datainläsningen, öppna http://localhost:50070/ i en webbläsare, bläddra i filsystemet och gå sedan till katalogen där data har laddats, det vill säga

/flume/tweets/

Port 50070 är NameNode-webbgränssnittet på Hadoop 2; Hadoop 3 flyttade samma sida till port 9870.

HDFS-webbläsaren som visar flume/tweets-katalogen med de laddade tweet-filerna

Flume är ena halvan av intag: Sqoop importerar tabeller i batchar, Flume strömmar händelser, sedan Pig or Bikupa forma filerna och Oozie schemalägger kedjan. Se även stora dataanalysverktyg, MapReduce-kopplingar och räknare och Talang.

Vanliga frågor

Inte som skrivet. Slutpunkten för streamingstatusar/filter i v1.1 togs bort den 9 mars 2023 och ersättningen för API v2 kräver en betald nivå. Flume-mekaniken gäller fortfarande som en övning i anpassad källkod.

Modellerar baslinjenormal loggvolym och meddelandeform, och flaggar sedan avvikelser som fasta tröskelvärden missar. De klustrar även upprepade stackar. traces till en enda incident och utarbeta en sannolik orsak, vilket förkortar triage.

Copilot utarbetar källkods-, kanal- och sink-block snabbt, men den hittar på egenskapsnamn och blandar utgåvor. Kontrollera varje nyckel mot Flume användarhandbok för din version innan du startar agenten.

Flume strömmar händelsedata som loggar kontinuerligt till HDFS. Sqoop flyttar strukturerade tabeller mellan relationsdatabaser och Hadoop i schemalagda batchar. De täcker olika halvor av inmatningen och parar ihop väl.

En minneskanal är snabbast men förlorar buffrade händelser om agenten dör. En filkanal skriver till disk och överlever omstarter, med lägre dataflöde. Föredra hållbarhet för allt du inte kan skicka om.

Kafka är nu den vanliga standardversionen eftersom den lagrar data och betjänar många konsumenter. Flume 1.11.0, från oktober 2022, fungerar fortfarande med enkel envägslogginsamling till HDFS.

Nästan alltid en JAR-krock: Flume-tarballen paketerar sina egna Guava- och Protobuf-versioner, vilka konflikterar med de som Hadoop laddar. Att ta bort den äldre paketerade JAR-filen rensar den vanligtvis.

De avgör när sinken stänger en fil och öppnar en ny: rollSize per byte, rollCount per händelseantal, rollInterval per sekunder. Noll inaktiverar just den utlösaren.

Sammanfatta detta inlägg med: