Wróć do bloga

DataWeave pod obciążeniem: streaming, deferred i kiedy lepiej użyć Batch

2026-09-17

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=true i collectionPath (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:

  1. zmienna referencjonowana raz,
  2. brak ujemnego indeksu ([-1] itd.),
  3. 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

  1. Źródło: ustaw streaming=true na listener/request/file — nie zakładaj, że Transform „streamuje sam”.
  2. XML: collectionPath i streaming=true (sprawdź casing property w Studio vs przykład docs).
  3. Skrypt: jedna referencja do streamowanej zmiennej; bez payload[-1] / reorder całego dokumentu; metadane przed kolekcją.
  4. Output: deferred=true tylko gdy next hop konsumuje stream; testuj failure path u konsumenta, nie zakładaj fail-fast na Transform.
  5. Pomiary: heap, /tmp (dw-buffer-*.tmp), concurrency po zmianie inMemorySize / max_memory_allocation.
  6. Antywzorce: jeden map na kolekcji zamiast Transform w foreach; zero zbędnych JSON↔XML; indent=false na dużych outputach.
  7. Parallel For Each / Cache: świadomie — Batch przy dużych zbiorach; Cache nie na non-repeatable; OS strategy w prod.
  8. 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

Help

Wersja ścieżki z shortlisty (twin)

How-to wideo (zweryfikowane tytuły / oEmbed)

Uwaga redakcyjna: dedykowanego, stabilnego how-to „DataWeave streaming=true + deferred” na YouTube nie znaleziono w tej edycji — primary source pozostaje Streaming in DataWeave.