DataWeave pod obciążeniem: streaming, deferred i kiedy lepiej użyć Batch
DataWeave w Mule 4 może czytać dane na trzy sposoby: in-memory (cały dokument w RAM), indexed (dysk + indeks, random access) albo streaming (sekwencyjnie, jednostka = rekord CSV / element tablicy JSON / kolekcja XML). Połączenie streamingu na źródle z deferred=true na wyjściu pozwala przepchnąć dane end-to-end bez pełnego indeksu i bez zapisu outputu na dysk — szybciej i z mniejszym zużyciem zasobów niż domyślne ścieżki odczytu/zapisu.
Czego ten tekst nie jest: listą „Top 10 funkcji DataWeave”, zamiennikiem Batch Job przy ETL z resume i błędami per rekord, ani obietnicą, że streaming=true „zawsze przyspieszy map/filter”. Streaming nie jest włączony domyślnie. Działa tylko przy sekwencyjnym dostępie do jednostek strumienia — i psuje się cicho, gdy skrypt wymaga random access do całego dokumentu.
Dla kogo: developerzy i leadzi integracji z GB-owymi plikami, memory pressure na CloudHub / workerach oraz transformacjami w hot path. Poniżej: trzy strategie odczytu, karty pułapek (objaw → przyczyna → docs → wzorzec), decision tree DW vs Batch vs Cache, checklista oraz FAQ pod AEO.
Trzy strategie odczytu (zanim włączysz streaming=true)
Oficjalna strona formatów DataWeave opisuje trzy read strategies: in-memory, indexed i streaming (Supported Data Formats).
In-memory
Cały dokument ląduje w pamięci. Masz pełny random access — dowolny selektor, dowolna kolejność. Przy dużych payloadach to prosta droga do OOM albo agresywnego GC. Działa dla wszystkich formatów, ale nie jest ścieżką „pod obciążeniem”.
Indexed
Parser buduje indeks i może zrzucać treść na dysk, zachowując random access jak przy in-memory. Docs: indexed readers obsługują pliki do ok. 20 GB; większe → streaming (szybszy, bez limitu rozmiaru inputu w dokumentacji) (Indexed Readers).
Próg przejścia do buforów na dysku to com.mulesoft.dw.max_memory_allocation — domyślnie 1 572 864 bajtów (1.5 MB). Powyżej pojawiają się dw-buffer-input-*.tmp, dw-buffer-output-*.tmp oraz dw-buffer-index-*.tmp (Memory Management).
Streaming
Dane płyną sekwencyjnie; w pamięci jest bieżąca jednostka. Jednostka zależy od formatu: wiersz CSV, element tablicy JSON, element kolekcji XML. Brak random access do całego dokumentu (Streaming in DataWeave). Formaty ze streamingiem: CSV, JSON, XML, Excel (XLSX) — ten artykuł skupia się na CSV/JSON/XML.
Włączanie streamingu na źródle (nie w samym skrypcie)
Przełącznik to MIME type na źródle danych — outputMimeType / mimeType na HTTP Listener, Request, File, Set Payload itd. Sam Transform Message bez flagi na wejściu nie „włącza streamingu magicznie”.
<http:listener doc:name="Listener"
outputMimeType="application/json; streaming=true"
config-ref="HTTP_Listener_config" path="/input"/>
- CSV: jednostka = wiersz; w obrębie rekordu random access jest OK.
- JSON: jednostka = element tablicy. Od Mule 4.3 streaming obejmuje też tablice poza rootem (w 4.2 root musiał być tablicą).
- XML: wymagane oba:
streaming=trueicollectionPath(lokalizacja kolekcji). Brak któregokolwiek = brak streamu.
Przykład z docs (uwaga na casing w przykładzie oficjalnym — collectionpath):
<http:listener
outputMimeType="application/xml; collectionpath=order.order-items; streaming=true"
config-ref="HTTP_Listener_config" path="/input"/>
Źródło: Streaming in DataWeave.
Co psuje streaming (nawet gdy „wygląda OK”)
Streaming = dostęp sekwencyjny do jednostek. Poniższe wzorce wymuszają random access albo „cofanie się” w dokumencie.
Ujemne indeksy i przestawianie kolejności
[payload[-2], payload[-1], payload[3]]
Ten skrypt wymaga dostępu do całego dokumentu w innej kolejności niż napływ — streaming nie zadziała. Wewnątrz pojedynczego rekordu CSV/JSON elementu random access jest dozwolony.
Dwukrotne odwołanie do payload / zła kolejność kluczy JSON
Chcesz jednocześnie payload.family (streamowana tablica) i payload.name albo { a: payload.age, b: payload.family } gdy age jest po family w obiekcie — stream nie wraca wstecz. JSON nie gwarantuje kolejności kluczy: skrypt może działać na tym pliku, a failować na innym.
Odwołanie z zagnieżdżonej lambdy
[1,2,3] map ((item, index) -> payload) — walidator i runtime traktują to jako referencję poza zakresem definicji zmiennej. Kryteria @StreamCapable (poniżej) wprost to wyłapują.
orderBy / groupBy / distinctBy (i podobne redukcje)
Funkcje wymagające całego zbioru przed pierwszym wynikiem wymuszają materializację — streaming sekwencyjny nie przeżywa globalnego sortu/grupowania. Docs: streaming = dostęp sekwencyjny, bez random access do całego dokumentu. Przy GB-owych plikach sort/group → Indexed (świadomie), baza/downstream albo Batch, nie „sprytniejszy” skrypt z deferred=true.
Wzorzec: jednoprzebiegowe map / filter; metadane przed kolekcją w modelu danych; last-element / reorder / global sort-group → indexed albo Batch / dwa przebiegi.
@StreamCapable() — walidator, nie gwarancja runtime
Adnotacja eksperymentalna: sprawdza, czy skrypt może sekwencyjnie czytać zmienną (zwykle payload). Kryteria:
- zmienna referencjonowana raz,
- brak ujemnego indeksu (
[-1]itd.), - brak referencji z zagnieżdżonej lambdy.
Wymaga dyrektywy input z typem MIME, np. input payload application/json.
False fail: skrypt może streamować na konkretnym pliku, a walidator failuje — bo JSON nie gwarantuje kolejności kluczy, a procesor adnotacji nie zakłada stałej kolejności. Traktuj @StreamCapable jako linter pod sekwencyjność, nie jako certyfikat produkcji.
Źródło: Streaming in DataWeave — Validate a Script.
deferred=true: handoff bez dysku — i bez normalnego error handlingu
Writer property deferred=true w dyrektywie output generuje output jako stream i odracza wykonanie skryptu do momentu konsumpcji przez następny procesor:
output application/json deferred=true
Oficjalna NOTE: exceptions aren’t handled gdy deferred=true. W Studio debug: wyjątek loguje się w konsoli, ale flow nie zatrzymuje się na Transform Message — problemy widać dopiero u konsumenta (File Write, HTTP, kolejny komponent).
End-to-end z docs: File listener streaming=true → Transform z deferred=true → File write. To jest sensowny wzorzec, gdy next hop naprawdę konsumuje stream. Gdy potrzebujesz synchronicznego fail-fast na Transform — bez deferred.
deferred jest też property writera formatu binary (Binary Format).
Pamięć Mule vs pamięć DataWeave (dwa buffery)
Repeatable streams (Mule 4)
Domyślnie Mule 4 używa repeatable streams (EE: file-store, start 512 KB in-memory, potem dysk). Większy inMemorySize = mniej I/O dyskowego, ale mniej concurrent requestów. Niezużyty stream trzyma file handles, kursory DB, połączenia HTTP do końca eventu — ryzyko pool exhaustion / OOM. set-payload value="#[payload]" nie konsumuje streamu (Streaming in Mule Apps).
non-repeatable-stream tylko gdy jedno odczytanie i świadomość, że Cache / For Each / niektóre Transform wymagają pełnej konsumpcji.
Bufory DataWeave
Osobna warstwa: dw-buffer-*.tmp w java.io.tmpdir, off-heap pool (com.mulesoft.dw.memory_pool_size, com.mulesoft.dw.directbuffer.disable). Stringi >1.5 MB w JSON/XML są dzielone na chunki tym samym progiem max_memory_allocation — koszt wydajnościowy; wyłączenie com.mulesoft.dw.buffered_char_sequence.enabled tylko przy świadomym zapasie RAM (Indexed Readers, Memory Management).
Antywzorce wydajności (zanim włączysz streaming)
Transform wewnątrz foreach
Oficjalny antywzorzec: iteracja + per-item Transform zamiast jednego map na kolekcji, potem foreach na side-effects (tuning-app-design — DataWeave).
Zbędne hop’y formatów i nadmiarowe pola
Help article How To Improve Dataweave Performance: unikaj niepotrzebnych konwersji JSON↔XML; mapuj tylko potrzebne pola (Help).
indent=false i logowanie
Na dużych outputach indent=false zmniejsza rozmiar i obciążenie klienta. Nie loguj ciężkich wyrażeń DW na każdy request (tuning-app-design).
Parallel For Each
Buforuje wyniki wszystkich tras w listę — przy dużej liczbie elementów ryzyko OOM. Docs wprost: duże payloady → Batch (Parallel For Each). Help powtarza to ostrzeżenie.
Cache scope
Pomaga przy często powtarzanych, rzadko zmieniających się danych. Cache’uje repeatable streams; nie non-repeatable. W prod unikaj default in-memory OS — Object Store + TTL / max entries (Cache Scope, Tuning Caching).
Kiedy Batch wygrywa z DataWeave
Batch Job (EE) = reliable, asynchronous processing larger-than-memory: persistent queues, resume po crash/redeploy, błędy per rekord, steps + aggregator/bulk do SaaS (Batch Processing).
DW streaming = szybka jednoprzebiegowa transformacja dużego dokumentu → kolejny processor (write/HTTP). Bez gwarancji resume i bez natywnego modelu per-record error jak w Batch.
Częsty wzorzec produkcyjny: DW przygotowuje/split format (albo streamuje kształt) → Batch przetwarza rekordy. Help sugeruje Batch dla very large payloads; Parallel For Each docs mówią to samo przy OOM risk.
Nie wrzucaj każdego GB CSV do Batch „na wszelki wypadek” — i nie pchaj ETL z wymaganym resume w czysty DW.
Mini decision tree
| Pytanie | Ścieżka |
|---|---|
| Potrzebujesz random access na payloadzie >1.5 MB? | Indexed (świadomie) albo przeprojektuj skrypt |
| Jednoprzebieg CSV / JSON array / XML collection → write/next hop? | streaming=true na źródle + rozważ deferred=true |
| Resume / DLQ / bulk API / multi-step per record? | Batch Job |
| Ten sam lookup/response wielokrotnie, dane rzadko się zmieniają? | Cache (+ selective map); OS + TTL w prod |
| Niezależne I/O na elementach, ale nie „miliony”? | Parallel For Each z limitem concurrency — nie na GB zbiorach |
Checklista praktyczna
- Źródło: ustaw
streaming=truena listener/request/file — nie zakładaj, że Transform „streamuje sam”. - XML:
collectionPathistreaming=true(sprawdź casing property w Studio vs przykład docs). - Skrypt: jedna referencja do streamowanej zmiennej; bez
payload[-1]/ reorder całego dokumentu; metadane przed kolekcją. - Output:
deferred=truetylko gdy next hop konsumuje stream; testuj failure path u konsumenta, nie zakładaj fail-fast na Transform. - Pomiary: heap,
/tmp(dw-buffer-*.tmp), concurrency po zmianieinMemorySize/max_memory_allocation. - Antywzorce: jeden
mapna kolekcji zamiast Transform wforeach; zero zbędnych JSON↔XML;indent=falsena dużych outputach. - Parallel For Each / Cache: świadomie — Batch przy dużych zbiorach; Cache nie na non-repeatable; OS strategy w prod.
- Durability: jeśli potrzebujesz resume / per-record errors → Batch, nie „sprytniejszy” DW.
FAQ
1. Czy streaming w DataWeave jest włączony domyślnie?
Nie. Musisz ustawić reader property streaming=true na źródle (outputMimeType / mimeType). Bez tego DataWeave może iść ścieżką in-memory lub indexed. Osobno: deferred=true na outputcie odrocza zapis i przekazuje stream dalej.
2. Czym różni się odczyt indexed od streaming?
Indexed parsuje dokument, buduje indeks (często na dysku) i daje random access — limit ok. 20 GB w docs. Streaming czyta sekwencyjnie jednostkami formatu, bez limitu rozmiaru inputu w dokumentacji, ale bez random access do całego dokumentu. Indexed jest kompromisem pamięć↔dysk; streaming jest najszybszy przy jednoprzebiegowych transformacjach.
3. Co robi deferred=true i dlaczego błąd może „zniknąć” z Transform Message?
deferred=true generuje output jako stream i odracza wykonanie do konsumpcji przez następny komponent. Oficjalnie wyjątki nie są obsługiwane jak zwykle — w Studio debug flow nie zatrzymuje się na Transform; błąd widać w logu / u konsumenta. Używaj, gdy next hop czyta stream; bez deferred, gdy potrzebujesz synchronicznego fail-fast.
4. Dlaczego streaming XML wymaga collectionPath i streaming=true?
XML nie ma tablic jak JSON. collectionPath wskazuje lokalizację kolekcji (np. order.order-items); dopiero wtedy jednostką streamu stają się elementy pod tą ścieżką. Docs: brak któregokolwiek z dwóch ustawień = brak streamu.
5. Kiedy skrypt przechodzi runtime, ale @StreamCapable failuje (lub odwrotnie)?
Walidator sprawdza reguły sekwencyjności (jedna referencja, brak ujemnego indeksu, brak nested-lambda). Może failować, choć na danym pliku JSON sekwencja działa — bo kolejność kluczy JSON nie jest gwarantowana. Odwrotnie: skrypt bez adnotacji może „działać” na małym pliku, a pod loadem zejść w indexed/in-memory i zjeść pamięć.
6. Jak payload[-1] / dwukrotne użycie payload psuje streaming?
Ujemny indeks i przestawianie kolejności elementów wymagają random access do całego dokumentu. Dwukrotne payload (np. family + name) wymaga „cofnięcia” streamu. W obu przypadkach tracisz model sekwencyjny — wracasz do indexed/in-memory albo dostajesz błąd walidacji @StreamCapable.
7. Kiedy wybrać Batch Job zamiast DataWeave streaming?
Gdy potrzebujesz niezawodnego async ETL: persistent queues, resume po crash/redeploy, obsługa błędów per rekord, aggregator/bulk do systemów zewnętrznych. DW streaming wygrywa przy jednoprzebiegowej transformacji kształtu dokumentu do kolejnego hopa. Często łączysz oba: DW → Batch.
8. Co oznaczają pliki dw-buffer-*.tmp i parametr com.mulesoft.dw.max_memory_allocation?
To bufory DataWeave na dysku (input/output/index), gdy payload przekracza próg — domyślnie 1.5 MB. Pliki żyją w java.io.tmpdir do zamknięcia streamów / końca eventu. Podnieś próg tylko gdy masz RAM; przy sekwencyjnych GB-plikach preferuj streaming zamiast windowania indexed.
9. Czy Cache scope pomaga przy dużych streamach?
Cache pomaga przy powtarzalnych, rzadko zmieniających się danych (lookup/reference). Cache’uje repeatable streams; nie cache’uje non-repeatable. Duży stream w default in-memory OS w prod może zjeść heap — użyj Object Store ze strategią expiry / max entries. To nie zamiennik streamingu ani Batch.
10. Jak unikać transformacji w foreach na dużych kolekcjach?
Zrób jedną transformację całej kolekcji (map / mapObject), a dopiero potem foreach, jeśli potrzebujesz side-effects (HTTP per item, DB write). Transform per element w pętli generuje zbędne eventy i CPU — to antywzorzec z oficjalnego tuning guide.
Soft CTA
Projektujesz przepływ GB-owych plików albo hot-path transformacji i chcesz ułożyć streaming vs indexed vs Batch bez eksperymentów na produkcji? Solita to nordycki partner MuleSoft z dostawą z Polski (EU-shoring) — pomagamy zespołom dobrać model odczytu i reliability pod realne obciążenie. Bez checklisty marketingowej i bez obietnic „#1”: konkretny przegląd skryptów, buforów i granicy, gdzie DataWeave oddaje pałeczkę Batchowi.
Źródła
Dokumentacja DataWeave / Mule
- Streaming in DataWeave
- Supported Data Formats (read strategies)
- Indexed Readers in DataWeave
- DataWeave Memory Management
- Binary Format (
deferredwriter) - Streaming in Mule Apps
- App Design That Maximizes Process Performance
- Batch Processing
- Parallel For Each Scope
- Cache Scope
- Tuning Caching
Help
Wersja ścieżki z shortlisty (twin)
How-to wideo (zweryfikowane tytuły / oEmbed)
- VirtualMuleys#3: DataWeave 2.0 Language Fundamentals — fundamenty FP/DW (nie deep-dive streamingu; dobry kontekst przed docs)
- codeLive: Processing Large Volumes of Salesforce Data with MuleSoft (Salesforce Developers) — duże wolumeny / Bulk + Mule (Batch adjacency)
Uwaga redakcyjna: dedykowanego, stabilnego how-to „DataWeave streaming=true + deferred” na YouTube nie znaleziono w tej edycji — primary source pozostaje Streaming in DataWeave.