[Spark Streams, Scala] - działanie stream i foreachRdd

[Spark Streams, Scala] - działanie stream i foreachRdd
A4
  • Rejestracja:ponad 10 lat
  • Ostatnio:około 5 lat
  • Postów:47
0

Chciałbym się dowiedzieć kilku rzeczy związanych ze streamami w Sparku :

  1. Jak dokładnie pracują streamy w Sparku ? Czy kolejne batche są dołączane do tego samego strumienia, czy to działa w ten sposób że w strumieniu jest tylko ostatni batch (chyba że zastosujemy funkcje okienkowe) ?
  2. Jak pracuje foreachRdd w takim razie ? Wczytuje cały batch (jeśli w strumieniu mamy tylko ostatni batch) czy ze strumienia wczytuje to co aktualnie doszło (jeśli batche są kolejno dołączane w strumieniu) ?
edytowany 1x, ostatnio: ast44
YA
A co u Ciebie jest źródłem danych dla strumienia?
A4
To chyba bez znaczenia, ale dla przykładu odbieram dane z portu przez StreamingContext i socketTextStream
YA
  • Rejestracja:prawie 10 lat
  • Ostatnio:21 minut
  • Postów:2368
0

Ja rozumiem tak:

W StreamingContext określasz "batchDuration", tj. przedział czasu, który będzie dzielił dane wejściowe na "batche", taki batch to RDD (podstawowy koncept sparka), zaś sekwencja RDD to DStream (czyli strumień). Jak źródło coś wyprodukuje, to Reciever (skojarzony ze strumieniem) zapisuje to "coś" w pamięci sparka. Grupa "cosi" wyprodukowanych w okresie "batchDuration" składa się na RDD/batch.

Teraz jak przetwarzasz te dane, to masz 2 opcje przetwarzania:

  1. Bezstanowe - przetwarzasz pojedynczy batch/RDD i nie masz żadnych zależności do innych RDD
  2. "Stanowe" - dla przetwarzanego RDD masz możliwość zajrzenia do wcześniejszych RDD (dla tego przypadku definiujesz okienko, które pozwala zaglądać w przeszłość, "sliding window", które zwraca Ci DStream obejmujący okres ["teraz"-"długość okna"; "teraz"])

Nie wiem jak spark wewnętrznie ogarnia RDD, które przestały być potrzebne, ale zakładam, że jeśli nie ma np. okienek które zaglądają dalej niż 1h, to spark postawiony przed faktem rychłego braku pamięci usuwa RDD "starsze niż" 1h.

Uwaga, mogę to źle rozumieć :-)

edytowany 1x, ostatnio: yarel
A4
  • Rejestracja:ponad 10 lat
  • Ostatnio:około 5 lat
  • Postów:47
0

No właśnie zastanawiam się jak to dokładnie działa, a czy konkretniej w ten sposób (zakładając że nie używam funkcji okienkowej):
a) załóżmy że zbieram dane na batch-a przez 1 s i dodaję do strumienia
b) potem na tym strumieniu wykonuję foreachRdd i w nim operacje które zajmują 2.5 s
c) w czasie przetwarzania pierwszego batcha do strumienia zostały dodane dwa kolejne
d) po przetworzeniu pierwszego batch-a z pierwotnego strumienia jest on usuwany i przetwarzane są następne w kolejce

EDIT: d) ewentualnie kolejne batche nie siedzą w kolejce tylko są wykonywane operacje na nich równolegle, ale potem są i tak usuwane ze strumienia

edytowany 1x, ostatnio: ast44
YA
  • Rejestracja:prawie 10 lat
  • Ostatnio:21 minut
  • Postów:2368
1

Jak spojrzysz na imaplementację DStreama:

  1. RDD trzymane są w mapie [Time, RDD[T]]
    https://github.com/apache/spark/blob/master/streaming/src/main/scala/org/apache/spark/streaming/dstream/DStream.scala#L87

  2. foreachRDD tworzy ForEachDStream:
    https://github.com/apache/spark/blob/master/streaming/src/main/scala/org/apache/spark/streaming/dstream/DStream.scala#L651

  3. ForEachDStream generuje joba na określony punkt w czasie w określony sposób:
    https://github.com/apache/spark/blob/master/streaming/src/main/scala/org/apache/spark/streaming/dstream/ForEachDStream.scala#L47

Kopiuj
 override def generateJob(time: Time): Option[Job] = {
    parent.getOrCompute(time) match {
      case Some(rdd) =>
        val jobFunc = () => createRDDWithLocalProperties(time, displayInnerRDDOps) {
          foreachFunc(rdd, time)
        }
        Some(new Job(time, jobFunc))
      case None => None
    }
  }
}

...jak dla mnie foreachRDD wykona się dla RDD z określonego batch intervalu (czyli "aktualnie" przetwarzanego)
I będzie się wykonywać co batch interval, tzn. job będzie tworzony.

W opisanym przypadku nowy RDD co sekundę, przetwarzanie 2.5 sekundy, pewnie ilość jobów będzie rosła...

A4
Zastanawiam się jeszcze jak jest jest z tymi RDD, które zostały już obsłużone. Czy są one usuwane ze strumienia na którym pracuje foreachRdd ?
YA
  • Rejestracja:prawie 10 lat
  • Ostatnio:21 minut
  • Postów:2368
1

Stare RDD usuwane są via DStream.clearMetadata, które dla określenia "stare" uwzględnia coś co się nazywa "rememberDuration" (zależne od slideDuration/checkpointDuration). Jak określisz okno, to strumień będzie pamiętał to co potrzebne do obsługi okna, może nawet więcej, jeśli chceckpointDuration jest większe niż to wynikające z okna.

Nie doszukiwałem się, gdzie wywoływane jest to clearMetadata, podejrzewam, że po zakończeniu przetwarzania joba.

A4
  • Rejestracja:ponad 10 lat
  • Ostatnio:około 5 lat
  • Postów:47
0

A jakby działało coś takiego ?

Kopiuj
val x = dstream.foreach{ x => ...}
val y = dstream.reduceByWindow(...)

slideDuration musiałoby być nie większe niż rememberDuration, ale w takim razie foreach nie mógłby wywalać tych RDD. Czyli Spark sam by musiał ustalić ten warunek zanim wykonał się jeszcze foreach ?

YA
  • Rejestracja:prawie 10 lat
  • Ostatnio:21 minut
  • Postów:2368
1

Ten dstream, którego używasz w przykładzie nie wziął się ot tak z nicości, wcześniej musiał zostać zainicjalizowany i mieć ustawione to rememberDuration.
Jak robisz coś na streamie, to najczęściej dostajesz nowy stream, dla którego "parentem" jest ten wyjściowy, i ten potomny mówi parentowi "hej, potrzebuję RDD za okres X", parent stream sprawdza sobie czy aktualnie obejmuje pamięcią dłuższy okres, czy krótszy i w razie potrzeby "zwiększa" okres zapamiętywania.

Kliknij, aby dodać treść...

Pomoc 1.18.8

Typografia

Edytor obsługuje składnie Markdown, w której pojedynczy akcent *kursywa* oraz _kursywa_ to pochylenie. Z kolei podwójny akcent **pogrubienie** oraz __pogrubienie__ to pogrubienie. Dodanie znaczników ~~strike~~ to przekreślenie.

Możesz dodać formatowanie komendami , , oraz .

Ponieważ dekoracja podkreślenia jest przeznaczona na linki, markdown nie zawiera specjalnej składni dla podkreślenia. Dlatego by dodać podkreślenie, użyj <u>underline</u>.

Komendy formatujące reagują na skróty klawiszowe: Ctrl+B, Ctrl+I, Ctrl+U oraz Ctrl+S.

Linki

By dodać link w edytorze użyj komendy lub użyj składni [title](link). URL umieszczony w linku lub nawet URL umieszczony bezpośrednio w tekście będzie aktywny i klikalny.

Jeżeli chcesz, możesz samodzielnie dodać link: <a href="link">title</a>.

Wewnętrzne odnośniki

Możesz umieścić odnośnik do wewnętrznej podstrony, używając następującej składni: [[Delphi/Kompendium]] lub [[Delphi/Kompendium|kliknij, aby przejść do kompendium]]. Odnośniki mogą prowadzić do Forum 4programmers.net lub np. do Kompendium.

Wspomnienia użytkowników

By wspomnieć użytkownika forum, wpisz w formularzu znak @. Zobaczysz okienko samouzupełniające nazwy użytkowników. Samouzupełnienie dobierze odpowiedni format wspomnienia, zależnie od tego czy w nazwie użytkownika znajduje się spacja.

Znaczniki HTML

Dozwolone jest używanie niektórych znaczników HTML: <a>, <b>, <i>, <kbd>, <del>, <strong>, <dfn>, <pre>, <blockquote>, <hr/>, <sub>, <sup> oraz <img/>.

Skróty klawiszowe

Dodaj kombinację klawiszy komendą notacji klawiszy lub skrótem klawiszowym Alt+K.

Reprezentuj kombinacje klawiszowe używając taga <kbd>. Oddziel od siebie klawisze znakiem plus, np <kbd>Alt+Tab</kbd>.

Indeks górny oraz dolny

Przykład: wpisując H<sub>2</sub>O i m<sup>2</sup> otrzymasz: H2O i m2.

Składnia Tex

By precyzyjnie wyrazić działanie matematyczne, użyj składni Tex.

<tex>arcctg(x) = argtan(\frac{1}{x}) = arcsin(\frac{1}{\sqrt{1+x^2}})</tex>

Kod źródłowy

Krótkie fragmenty kodu

Wszelkie jednolinijkowe instrukcje języka programowania powinny być zawarte pomiędzy obróconymi apostrofami: `kod instrukcji` lub ``console.log(`string`);``.

Kod wielolinijkowy

Dodaj fragment kodu komendą . Fragmenty kodu zajmujące całą lub więcej linijek powinny być umieszczone w wielolinijkowym fragmencie kodu. Znaczniki ``` lub ~~~ umożliwiają kolorowanie różnych języków programowania. Możemy nadać nazwę języka programowania używając auto-uzupełnienia, kod został pokolorowany używając konkretnych ustawień kolorowania składni:

```javascript
document.write('Hello World');
```

Możesz zaznaczyć również już wklejony kod w edytorze, i użyć komendy  by zamienić go w kod. Użyj kombinacji Ctrl+`, by dodać fragment kodu bez oznaczników języka.

Tabelki

Dodaj przykładową tabelkę używając komendy . Przykładowa tabelka składa się z dwóch kolumn, nagłówka i jednego wiersza.

Wygeneruj tabelkę na podstawie szablonu. Oddziel komórki separatorem ; lub |, a następnie zaznacz szablonu.

nazwisko;dziedzina;odkrycie
Pitagoras;mathematics;Pythagorean Theorem
Albert Einstein;physics;General Relativity
Marie Curie, Pierre Curie;chemistry;Radium, Polonium

Użyj komendy by zamienić zaznaczony szablon na tabelkę Markdown.

Lista uporządkowana i nieuporządkowana

Możliwe jest tworzenie listy numerowanych oraz wypunktowanych. Wystarczy, że pierwszym znakiem linii będzie * lub - dla listy nieuporządkowanej oraz 1. dla listy uporządkowanej.

Użyj komendy by dodać listę uporządkowaną.

1. Lista numerowana
2. Lista numerowana

Użyj komendy by dodać listę nieuporządkowaną.

* Lista wypunktowana
* Lista wypunktowana
** Lista wypunktowana (drugi poziom)

Składnia Markdown

Edytor obsługuje składnię Markdown, która składa się ze znaków specjalnych. Dostępne komendy, jak formatowanie , dodanie tabelki lub fragmentu kodu są w pewnym sensie świadome otaczającej jej składni, i postarają się unikać uszkodzenia jej.

Dla przykładu, używając tylko dostępnych komend, nie możemy dodać formatowania pogrubienia do kodu wielolinijkowego, albo dodać listy do tabelki - mogłoby to doprowadzić do uszkodzenia składni.

W pewnych odosobnionych przypadkach brak nowej linii przed elementami markdown również mógłby uszkodzić składnie, dlatego edytor dodaje brakujące nowe linie. Dla przykładu, dodanie formatowania pochylenia zaraz po tabelce, mogłoby zostać błędne zinterpretowane, więc edytor doda oddzielającą nową linię pomiędzy tabelką, a pochyleniem.

Skróty klawiszowe

Skróty formatujące, kiedy w edytorze znajduje się pojedynczy kursor, wstawiają sformatowany tekst przykładowy. Jeśli w edytorze znajduje się zaznaczenie (słowo, linijka, paragraf), wtedy zaznaczenie zostaje sformatowane.

  • Ctrl+B - dodaj pogrubienie lub pogrub zaznaczenie
  • Ctrl+I - dodaj pochylenie lub pochyl zaznaczenie
  • Ctrl+U - dodaj podkreślenie lub podkreśl zaznaczenie
  • Ctrl+S - dodaj przekreślenie lub przekreśl zaznaczenie

Notacja Klawiszy

  • Alt+K - dodaj notację klawiszy

Fragment kodu bez oznacznika

  • Alt+C - dodaj pusty fragment kodu

Skróty operujące na kodzie i linijkach:

  • Alt+L - zaznaczenie całej linii
  • Alt+, Alt+ - przeniesienie linijki w której znajduje się kursor w górę/dół.
  • Tab/⌘+] - dodaj wcięcie (wcięcie w prawo)
  • Shit+Tab/⌘+[ - usunięcie wcięcia (wycięcie w lewo)

Dodawanie postów:

  • Ctrl+Enter - dodaj post
  • ⌘+Enter - dodaj post (MacOS)