Samouczek Apache Flume: Co to jest, Architecture i przykład Hadoopa

⚡ Inteligentne podsumowanie

Apache Flume to rozproszona usługa służąca do zbierania, agregowania i przesyłania dużych ilości danych dziennika do systemu HDFS, zbudowana na agentach, które łączą w łańcuch źródło, kanał i odbiornik.

  • 🔘 Anatomia agenta: Każdy agent Flume jest procesem JVM przechowującym źródło, jeden lub więcej kanałów i odbiornik.
  • Niezawodność: Dostarczanie z najlepszym wysiłkiem nie toleruje awarii żadnego węzła; dostarczanie kompleksowe jest w stanie przetrwać awarie wielu węzłów.
  • Konfiguracja: Niestandardowe klasy źródłowe kompilują się do pliku JAR, który jest umieszczany w katalogu bibliotek Flume.
  • 🧪 Konfiguracja: Jeden plik właściwości określa nazwę źródła, kanału i odbiornika oraz ustawia ścieżkę HDFS i limity przewijania.
  • 🛠️. Uruchomić: Uruchom potok za pomocą agenta flume-ng, nadając agentowi nazwę i wskazując na flume.conf.
  • ⚠️ Przykład z datą: Obsługa strumieniowania Twitter v1.1 została zamknięta w marcu 2023 r., dlatego należy traktować to ćwiczenie jako wzorzec niestandardowego źródła.

Samouczek Apache Flume obejmujący architekturę agentów, konfigurację i przykład przesyłania strumieniowego Hadoop

Co to jest Apache Flume w Hadoop?

Apache Flume to niezawodny i rozproszony system do gromadzenia, agregowania i przesyłania ogromnych ilości danych logów. Posiada prostą, a zarazem elastyczną architekturę opartą na strumieniowym przepływie danych. Apache Flume służy do gromadzenia danych logów obecnych w plikach logów z serwerów WWW i agregowania ich w celu… HDFS Do analizy.

Flume w Hadoop obsługuje wiele źródeł, w tym:

  • „tail” (który przesyła dane z pliku lokalnego i zapisuje je do HDFS za pomocą Flume, podobnie do polecenia „tail” w systemie Unix)
  • Dzienniki systemowe
  • Log Apache4j (co umożliwia Java aplikacje do zapisywania zdarzeń do plików w HDFS poprzez Flume).

Obecna wersja to Flume 1.11.0, opublikowano 25 października 2022 r. i jest dostępne na stronie Strona pobierania Apache FlumeTen poradnik został napisany dla wersji 1.4.0, dlatego poniżej zamieszczono kilka wskazówek, w których bieżąca wersja zachowuje się inaczej.

Przepływ Architektura

Agent Flume to FMV Proces składający się z trzech komponentów – źródła Flume, kanału Flume i odbiornika Flume – przez które zdarzenia rozprzestrzeniają się po zainicjowaniu ich w źródle zewnętrznym. Poniższy diagram pokazuje, jak się one łączą.

Diagram architektury Flume przedstawiający agenta ze źródłem, kanałem i odbiorem zasilającym HDFS

  1. Zdarzenia generowane przez źródło zewnętrzne (serwer WWW) są przetwarzane przez źródło Flume. Źródło zewnętrzne wysyła zdarzenia do źródła Flume w formacie rozpoznawanym przez źródło docelowe.
  2. Źródło Flume odbiera zdarzenie i zapisuje je w jednym lub kilku kanałach. Kanał działa jak magazyn, który przechowuje zdarzenie do momentu jego wykorzystania przez odbiornik Flume. Kanał ten może wykorzystywać lokalny system plików do przechowywania tych zdarzeń.
  3. Odbiornik Flume usuwa zdarzenie z kanału i zapisuje je w zewnętrznym repozytorium, takim jak HDFS. Może istnieć wielu agentów Flume, w takim przypadku odbiornik Flume przekazuje zdarzenie do źródła Flume kolejnego agenta w przepływie.

Niektóre ważne cechy Flume

  • Flume charakteryzuje się elastyczną konstrukcją opartą na strumieniowym przepływie danych. Jest odporny na błędy i solidny, z wieloma mechanizmami przełączania awaryjnego i odzyskiwania. Flume oferuje różne poziomy niezawodności, w tym: „dostawa z najwyższym wysiłkiem” oraz „dostawa od końca do końca”. Dostawa z najwyższą starannością nie toleruje żadnej awarii węzła Flume, podczas gdy dostawa od początku do końca gwarantuje dostawę nawet w przypadku awarii wielu węzłów.
  • Flume przesyła dane między źródłami i odbiorcami. Gromadzenie danych może odbywać się w sposób zaplanowany lub sterowany zdarzeniami. Flume posiada własny moduł przetwarzania zapytań, który ułatwia transformację każdej nowej partii danych przed jej przeniesieniem do docelowego odbiorcy.
  • Możliwy Flume tonie obejmują HDFS i HBaseFlume może również przesyłać dane o zdarzeniach, takie jak dane o ruchu sieciowym, dane generowane przez witryny mediów społecznościowych i wiadomości e-mail.

Konfiguracja Flume, biblioteki i kodu źródłowego

Zanim rozpoczniemy właściwy proces, upewnij się, że masz zainstalowany Hadoop; jeśli nie, przejrzyj jak zainstalować Hadoop Najpierw zmień użytkownika na „hduser” (identyfikator używany podczas konfiguracji Hadoop; możesz zmienić na identyfikator używany podczas własnej konfiguracji Hadoop).

Przełączanie użytkownika Linux na hduser w terminalu przed rozpoczęciem konfiguracji Flume

Krok 1) Utwórz nowy katalog o nazwie „FlumeTutorial”.

sudo mkdir FlumeTutorial
  1. Udzielaj uprawnień do odczytu, zapisu i wykonywania.
    sudo chmod -R 777 FlumeTutorial
  2. Skopiuj pliki Moje źródło Twittera.java oraz MójTwitterSourceForFlume.java do tego katalogu.

Pobierz stąd pliki wejściowe

Sprawdź uprawnienia wszystkich plików, jak pokazano poniżej, i jeśli ich brakuje, przyznaj uprawnienie „odczyt”.

Terminal wyświetlający uprawnienia do pobranego pliku Java pliki źródłowe

Krok 2) Pobierz „Apache Flume” z https://flume.apache.org/download.html.

W tym samouczku Flume użyto Apache Flume 1.4.0.

Strona pobierania Apache Flume pokazująca łącze do pliku tarball z binarnym plikiem do wyboru

Następnie kliknij, aby przejść do lustra.

Strona lustrzana Apache'a dostępna po kliknięciu łącza tarball Flume

Krok 3) Skopiuj pobrany plik tarball do wybranego katalogu i wyeksportujtracPrzejrzyj zawartość za pomocą następującego polecenia.

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

Terminal extractworzenie archiwum Flume za pomocą polecenia sudo tar -xvf

Tworzy nowy katalog o nazwie apache-flume-1.4.0-bin i np.tracPrzenosi do niego pliki. Ten katalog jest nazywany w dalszej części artykułu.

Krok 4) Konfiguracja biblioteki Flume. Skopiuj pliki twitter4j-core-4.0.1.jar, flume-ng-configuration-1.4.0.jar, flume-ng-core-1.4.0.jar i flume-ng-sdk-1.4.0.jar do

/lib/

Jeden lub wszystkie skopiowane pliki JAR mogą mieć ustawione uprawnienia do wykonywania, co może powodować problemy z kompilacją kodu, dlatego należy je cofnąć. W moim przypadku plik twitter4j-core-4.0.1.jar miał uprawnienia do wykonywania. Cofnąłem je w następujący sposób.

sudo chmod -x twitter4j-core-4.0.1.jar

Terminal cofający uprawnienie do wykonywania pliku JAR rdzenia Twitter4j

Następnie polecenie poniżej nadaje wszystkim uprawnienie do odczytu pliku twitter4j-core-4.0.1.jar.

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

Proszę zwrócić uwagę, że pobrałem plik twitter4j-core-4.0.1.jar z Repozytorium Maveni wszystkie pliki JAR Flume, tj. flume-ng-*-1.4.0.jar, z artefakty org.apache.flume.

Załaduj dane z Twittera za pomocą Flume

Krok 1) Przejdź do katalogu zawierającego pliki kodu źródłowego.

Krok 2) Ustaw CLASSPATH tak, aby zawierał /lib/* i ~/FlumeTutorial/flume/mytwittersource/*.

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

Terminal eksportujący CLASSPATH wskazujący na katalogi biblioteki Flume i źródła

Krok 3) Skompiluj kod źródłowy za pomocą poniższego polecenia.

javac -d . MyTwitterSourceForFlume.java MyTwitterSource.java

Terminal kompilujący dwa Java pliki źródłowe z javac

Krok 4) Utwórz plik JAR. Najpierw utwórz plik Manifest.txt za pomocą dowolnego edytora tekstu i dodaj do niego poniższy wiersz.

Main-Class: flume.mytwittersource.MyTwitterSourceForFlume

Tutaj flume.mytwittersource.MyTwitterSourceForFlume to nazwa klasy głównej. Pamiętaj, że musisz nacisnąć klawisz Enter na końcu tego wiersza, jak pokazano poniżej.

Plik Manifest.txt otwiera się w edytorze tekstu z wpisem klasy głównej

Teraz utwórz plik JAR „MyTwitterSourceForFlume.jar” w następujący sposób.

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

Pakowanie skompilowanych klas w formacie MyTwitterSourceForFlume.jar

Krok 5) Skopiuj ten plik JAR do /biblioteka/.

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

Terminal kopiuje niestandardowy plik JAR źródłowy do katalogu biblioteki Flume

Krok 6) Przejdź do katalogu konfiguracyjnego Flume, /konf.

Jeśli plik flume.conf nie istnieje, skopiuj plik flume-conf.properties.template i zmień jego nazwę na flume.conf.

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

Kopiowanie flume-conf.properties.template przez terminal do flume.conf

Jeśli plik flume-env.sh nie istnieje, skopiuj plik flume-env.sh.template i zmień jego nazwę na flume-env.sh.

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

Kopiowanie pliku flume-env.sh.template do flume-env.sh w terminalu

Tworzenie aplikacji na Twitterze

Przeczytaj to najpierw. Transmisja strumieniowa v1.1 statuses/filter punkt końcowy, którego potrzebuje twitter4j 4.0.1, został wycofany 9 marca 2023 r., a filtrowany strumień API v2, który go zastąpił, jest teraz dostępny pod adresem developer.x.com, znajduje się za płatną wersją. Potraktuj poniższe ekrany jako niestandardowy wzorzec źródłowy, a następnie skieruj tego samego agenta na plik, plik wykonywalny lub źródło Kafki.

Krok 1) Utwórz aplikację Twitter, logując się do portalu dla programistów.

Strona logowania dla programistów Twittera służąca do dostępu do listy aplikacji

Strona główna konta programisty Twittera wyświetlana po zalogowaniu

Krok 2) Przejdź do „Moich aplikacji” (opcja ta rozwija się po kliknięciu przycisku „Jajko” w prawym górnym rogu).

Strona Moje aplikacje w portalu dla programistów Twittera

Krok 3) Utwórz nową aplikację klikając „Utwórz nową aplikację”.

Krok 4) Uzupełnij dane aplikacji, podając jej nazwę, opis i adres strony internetowej. Możesz skorzystać z uwag podanych pod każdym polem wprowadzania danych.

Formularz tworzenia aplikacji na Twitterze z polami nazwy, opisu i witryny internetowej

Krok 5) Przewiń stronę w dół, zaakceptuj warunki, zaznaczając „Tak, zgadzam się”, i kliknij przycisk „Utwórz aplikację na Twitterze”.

Pole wyboru warunków i przycisk „Utwórz aplikację” na dole formularza Twittera

Krok 6) W oknie nowo utworzonej aplikacji przejdź do zakładki „Klucze API”, przewiń stronę w dół i kliknij przycisk „Utwórz mój token dostępu”.

Karta kluczy API nowej aplikacji Twittera przed utworzeniem tokena dostępu

Szczegóły tokena dostępu wyświetlane po użyciu przycisku Utwórz mój token dostępu

Krok 7) Odśwież stronę.

Krok 8) Kliknij „Testuj OAuth”. Wyświetlą się ustawienia OAuth aplikacji.

Ekran testowy OAuth wyświetlający ustawienia OAuth aplikacji

Krok 9) Zmodyfikuj plik „flume.conf” za pomocą tych ustawień OAuth. Poniżej przedstawiono kroki modyfikacji pliku „flume.conf”.

Ustawienia OAuth zawierające listę wartości klucza konsumenta, tajnego klucza konsumenta i tokena dostępu

Musimy skopiować klucz konsumenta, tajny klucz konsumenta, token dostępu i tajny token dostępu, aby zaktualizować plik „flume.conf”.

Uwaga: Wartości te należą do użytkownika i są poufne, dlatego nie należy ich udostępniać.

Zmodyfikuj plik „flume.conf”.

Krok 1) Otwórz „flume.conf” w trybie zapisu i ustaw wartości poniższych parametrów.

sudo gedit flume.conf

Skopiuj treść poniżej.

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

Plik flume.conf otwiera się w edytorze z właściwościami źródła, kanału i odbiornika MyTwitAgent

Krok 2) Ponadto ustaw TwitterAgent.sinks.HDFS.hdfs.path jak poniżej.

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

Właściwość hdfs.path odbiornika HDFS ustawiona na nazwę hosta, numer portu i katalog domowy HDFS

Znaleźć , I , zobacz wartość parametru 'fs.defaultFS' ustawionego w $HADOOP_HOME/etc/hadoop/core-site.xml, pokazaną poniżej.

Właściwość fs.defaultFS w pliku core-site.xml, która dostarcza nazwę hosta i port

Krok 3) Aby móc przesyłać dane do systemu HDFS w miarę ich pojawiania się, usuń poniższy wpis, jeśli taki istnieje.

TwitterAgent.sinks.HDFS.hdfs.rollInterval = 600

Przykład: przesyłanie strumieniowe danych z Twittera za pomocą Flume

Krok 1) Otwórz „flume-env.sh” w trybie zapisu i ustaw wartości poniższych parametrów.

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

flume-env.sh otwarty w edytorze z ustawionymi zmiennymi JAVA_HOME i FLUME_CLASSPATH

Krok 2) Uruchom Hadoop.

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

Krok 3) Dwa pliki JAR z archiwum tarball Flume nie są kompatybilne z Hadoop 2.2.0, dlatego w tym przykładzie Apache Flume wykonujemy poniższe kroki, aby zapewnić zgodność Flume z Hadoop 2.2.0. Ta zamiana plików JAR to poprawka z ery 1.4.0; Flume 1.11.0 zawiera już aktualne kompilacje protobuf i Guava, więc współczesne archiwum tarball zazwyczaj ich nie potrzebuje.

a. Przenieś protobuf-java-2.4.1.jar z ' /lib'. Przejdź najpierw do tego katalogu.

płyta CD /lib

sudo mv protobuf-java-2.4.1.jar ~/

Terminal przenoszący plik protobuf-java-2.4.1.jar z katalogu biblioteki Flume

b. Znajdź plik JAR „guava”, jak pokazano poniżej.

find . -name "guava*"

Polecenie terminala find lokalizujące dołączony plik JAR z guava

Przenieś guava-10.0.1.jar z ' /biblioteka'.

sudo mv guava-10.0.1.jar ~/

Terminal przenosi guava-10.0.1.jar z katalogu biblioteki Flume

c. Pobierz guava-17.0.jar z Repozytorium Maven, pokazane poniżej.

Strona repozytorium Maven dla guava 17.0, zamienny plik JAR do pobrania

Teraz skopiuj pobrany plik JAR do ' /biblioteka'.

Krok 4) Przejdź do ' /bin' i uruchom Flume w następujący sposób.

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

Terminal uruchamiający agenta Flume o nazwie MyTwitAgent za pomocą polecenia flume-ng

Okno wiersza poleceń, z którego Flume pobiera tweety, wygląda następująco.

Wiersz poleceń pokazujący agenta Flume pobierającego tweety i zapisującego je w systemie HDFS

Z komunikatu w oknie poleceń wynika, że ​​dane wyjściowe są zapisywane w katalogu /user/hduser/flume/tweets/. Teraz otwórz ten katalog za pomocą przeglądarki internetowej.

Krok 5) Aby zobaczyć wynik ładowania danych, otwórz http://localhost:50070/ w przeglądarce, przejrzyj system plików, a następnie przejdź do katalogu, do którego załadowano dane, czyli

/flume/tweety/

Port 50070 to internetowy interfejs użytkownika NameNode w Hadoop 2; Hadoop 3 przeniósł tę samą stronę na port 9870.

Przeglądarka HDFS pokazująca katalog flume/tweets z załadowanymi plikami tweetów

Flume jest połową spożycia: Łyżka importuje tabele w partiach, Flume przesyła strumieniowo zdarzenia, a następnie Świnia or Ul kształtować pliki i Oozie harmonogramuje łańcuch. Zobacz także narzędzia do analizy dużych zbiorów danych, MapReduce łączy i licznikuje oraz Taland.

FAQ

Nie tak, jak napisano. Punkt końcowy statusów/filtrów strumieniowych v1.1 został wycofany 9 marca 2023 r., a zastąpienie API v2 wymaga płatnego poziomu. Mechanika Flume nadal obowiązuje jako ćwiczenie ze źródła niestandardowego.

Modele bazują na normalnej objętości logarytmu i kształcie wiadomości, a następnie sygnalizują odchylenia, których nie spełniają ustalone progi. Klastrują również powtarzane stosy. traces w pojedynczy incydent i sporządź prawdopodobną przyczynę, skracając tym samym triaż.

Copilot szybko tworzy bloki źródłowe, kanałowe i ujściowe, ale wymyśla nazwy właściwości i miksuje wydania. Przed uruchomieniem agenta sprawdź każdy klucz z podręcznikiem użytkownika Flume dla swojej wersji.

Flume przesyła strumieniowo dane zdarzeń, takie jak logi, do systemu HDFS. Sqoop przenosi tabele strukturalne między relacyjnymi bazami danych a Hadoop w zaplanowanych partiach. Obejmują one różne etapy przetwarzania i dobrze się ze sobą łączą.

Kanał pamięci jest najszybszy, ale traci buforowane zdarzenia w przypadku awarii agenta. Kanał pliku zapisuje dane na dysk i zachowuje je po ponownym uruchomieniu, ale z niższą przepustowością. Preferuj trwałość dla wszystkiego, czego nie można ponownie wysłać.

Kafka jest obecnie standardowym rozwiązaniem domyślnym, ponieważ przechowuje dane i obsługuje wielu użytkowników. Flume 1.11.0, od października 2022 r., nadal obsługuje proste, jednokierunkowe gromadzenie logów w systemie HDFS.

Prawie zawsze występuje konflikt JAR: archiwum Flume zawiera własne wersje Guava i protobuf, które kolidują z wersjami ładowanymi przez Hadoop. Usunięcie starszego pliku JAR zazwyczaj usuwa problem.

Decydują, kiedy odbiornik zamyka plik i otwiera nowy: rollSize według bajtów, rollCount według liczby zdarzeń, rollInterval według sekund. Zero wyłącza dany wyzwalacz.

Podsumuj ten post następująco: