# ETL Pipelines with Node-RED?

**URL:** <https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347>\
**Category:** Share Your Projects\
**Created:** [6 January 2024 19:09 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347 "2024-01-06T19:09:33Z")\
**Posts on this page:** 17\
**Page:** 4

<div class="post-metadata">

**Author:** ![gregorius](https://sea2.discourse-cdn.com/flex026/user_avatar/discourse.nodered.org/gregorius/32/73816_2.png) [@gregorius](https://discourse.nodered.org/u/gregorius)\
**Post date:** [12 January 2024 15:15 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347/61 "2024-01-12T15:15:12Z")

</div>

> [@TotallyInformation](#):
>
> I was wondering whether it would be better to move the actual process from PipeEnd to a config node. Then every stream capable node would have a reference to that config node and, by definition, anything that didn't have such a reference would not be a streaming node?

I not quite sure how that would work since the PipeEnd node also creates the pipeline, so it will always be needed. Are you saying that for the PipeEnd to know something is a streaming node, it would check whether a node has that config node as reference? I.e. the config node would act as a type of marker?

I would definitely not move the implementation of the stream to that config node. The code for creating the stream should be on the node itself, e.g., [CsvStream](https://github.com/gorenje/node-red-streaming/blob/45f037735cdaf0124c6c2891fafeb86792359d21/nodes/csvstream.js#L10-L23C1).

This binding of stream creation to the node allows others to create stream nodes without having to touch the original streaming package.

If Joe 'Bob' Bloggs were to implement a StreamBanana node then she would only need to a) add the node id to the [\_streamPipeline array](https://github.com/gorenje/node-red-streaming/blob/45f037735cdaf0124c6c2891fafeb86792359d21/nodes/csvstream.js#L34-L36) on 'input' and b) add a [`createStream`](https://github.com/gorenje/node-red-streaming/blob/45f037735cdaf0124c6c2891fafeb86792359d21/nodes/csvstream.js#L10) function to the nodes JS file.

PipeEnd goes through the [\_streamPipeline array](https://github.com/gorenje/node-red-streaming/blob/45f037735cdaf0124c6c2891fafeb86792359d21/nodes/pipeend.js#L20-L29) and sequentially calls createStream, finally creating a [pipeline](https://github.com/gorenje/node-red-streaming/blob/45f037735cdaf0124c6c2891fafeb86792359d21/nodes/pipeend.js#L31-L33) from those streams.

There are some subtleties around how to handle evaluateNodeProperty, i.e. when that should be called. My nodes call it when the [first message comes through](https://github.com/gorenje/node-red-streaming/blob/45f037735cdaf0124c6c2891fafeb86792359d21/nodes/httprequeststream.js#L234-L243) but a [second msg is passed to the createStream](https://github.com/gorenje/node-red-streaming/blob/45f037735cdaf0124c6c2891fafeb86792359d21/nodes/httprequeststream.js#L35) function. That msg object is the one that the PipeEnd node received. So in theory, the node _could_ call evaluateNodeProperty in createStream ...

---

<div class="post-metadata">

**Author:** ![shrickus](https://sea2.discourse-cdn.com/flex026/user_avatar/discourse.nodered.org/shrickus/32/517_2.png) [@shrickus](https://discourse.nodered.org/u/shrickus)\
**Post date:** [12 January 2024 17:23 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347/62 "2024-01-12T17:23:54Z")

</div>

Fascinating discussion! I have to admit I've not thought through all the ramifications, but my first thought after looking at your pipe flows above was that the `pipestart` and `pipeend` nodes are "just" bounding the stream processing (i know there is way more than that! ;\*).

To simplify the node-red diagrams, what about adding streaming support to either "groups" or "subflows"? Yes it would mean extending the core with a different internal processing implementation -- but then you could build whatever flows you want inside that streaming container. The only caveat would be that all those included nodes would have to support streaming mode.

Initially it may be just these new nodes that you have working -- but over time the core nodes could be updated to include streaming support (or verify that they already work with streams?). For instance, the `debug` node `onMessage(...)` function takes the incoming `msg` object and writes it to the sidebar (console, status, whatever) which seems to "inherently" resemble streaming. So mark that node with stream support, and let the core processing listen on the stream and feed each incoming message to it like usual -- no real code changes, I would think?

Perhaps this is greatly over-simplified, but what a huge leap forward for node-red if it supported both paths. Feels similar to how Promise chaining in JS simplified our handling of ["callback hell"](http://callbackhell.com/).

---

<div class="post-metadata">

**Author:** ![TotallyInformation](https://sea2.discourse-cdn.com/flex026/user_avatar/discourse.nodered.org/totallyinformation/32/31_2.png) [@TotallyInformation](https://discourse.nodered.org/u/TotallyInformation)\
**Post date:** [12 January 2024 17:37 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347/63 "2024-01-12T17:37:10Z")

</div>

> [@gregorius](#):
>
> I not quite sure how that would work since the PipeEnd node also creates the pipeline, so it will always be needed. Are you saying that for the PipeEnd to know something is a streaming node, it would check whether a node has that config node as reference? I.e. the config node would act as a type of marker?
> 
> I would definitely not move the implementation of the stream to that config node. The code for creating the stream should be on the node itself, e.g., [CsvStream](https://github.com/gorenje/node-red-streaming/blob/45f037735cdaf0124c6c2891fafeb86792359d21/nodes/csvstream.js#L10-L23C1).

OK, it was just a thought. I was thinking that the stream management would all be in the config node. Not to worry. No point in having the complexity of a config node if it is only being used for a flag.

---

<div class="post-metadata">

**Author:** ![gregorius](https://sea2.discourse-cdn.com/flex026/user_avatar/discourse.nodered.org/gregorius/32/73816_2.png) [@gregorius](https://discourse.nodered.org/u/gregorius)\
**Post date:** [12 January 2024 17:50 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347/64 "2024-01-12T17:50:42Z")

</div>

> [@shrickus](#):
>
> (i'm sure there is way more than that!).

PipeStart just adds an array to the message and PipeEnd executes the streams - so not too much more 🙂

> [@shrickus](#):
>
> either stream "groups" or "subflows"

In the ETL flow, I now use subflows to simplify the flow --\> a somewhat too colourful [comparison](https://flowhub.org/f/c520d9da20ad7f1d?v2=96ea2306471408d0d760254f9c77a763bea633bf&v1=68222eff66a1ff2e8294ca2553cbd93222deda99) of the two flows.

The point is that I have actually created several subflows to encapsulate certain common behaviour --\> [the Get2Disk node](https://cdn.flowhub.org/?t=0&fhid=c520d9da20ad7f1d#node/8dd8e84ce59f090a/edit) is a subflow that contains the PipeStart and PipeEnd (_edit subflow template_ works) nodes. So it's possible to use subflows and not really necessary to create a new concept of wrapping streaming nodes.

Also in creating subflows, I begin to create "ETL" behaviour (i.e. retrieve and store data) instead of just streaming A to B. There is also a [Path2JsonL](https://cdn.flowhub.org/?fhid=c520d9da20ad7f1d#flow/2c7ddaaab6869956) subflow which assumes there is a path attribute with a .jsonl file that is then streamed into a JsonLStream node --\> that's ETL semantics, not streaming.

> [@TotallyInformation](#):
>
> OK, it was just a thought.

👍 Brainstorming ideas is great!

---

<div class="post-metadata">

**Author:** ![gregorius](https://sea2.discourse-cdn.com/flex026/user_avatar/discourse.nodered.org/gregorius/32/73816_2.png) [@gregorius](https://discourse.nodered.org/u/gregorius)\
**Post date:** [12 January 2024 18:02 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347/65 "2024-01-12T18:02:26Z")

</div>

> [@shrickus](#):
>
> but over time the core nodes could be updated to include streaming support

Here we begin to mix terminology, I should say that the "streaming" I'm talking about here is the [streaming API](https://nodejs.org/dist/latest-v18.x/docs/api/stream.html) defined by NodeJS.

Node-RED is a message streaming tool. So it certainly does stream already but it does not support the streaming API from NodeJS.

The streaming API is more about having byte streams interconnected between javascript components with each component doing something with the data on-the-fly. I.e. the entire file is never completely in memory, only chunks of it.

What the intention of my work here is to "data to message streamification": large CSV/JsonL files are converted into a stream of Node-RED msg objects _without_ generating large memory footprint or large msg objects flowing through Node-RED. CSV/JsonL files are line based, meaning they consist of lines of independent data, each line can be represented by a single msg object.

---

<div class="post-metadata">

**Author:** ![shrickus](https://sea2.discourse-cdn.com/flex026/user_avatar/discourse.nodered.org/shrickus/32/517_2.png) [@shrickus](https://discourse.nodered.org/u/shrickus)\
**Post date:** [12 January 2024 20:24 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347/66 "2024-01-12T20:24:30Z")

</div>

Thanks for that clarification... so even though the `csv` node already has an option to "output each line as individual messages" it still reads the entire .csv file into memory before the first line is sent out?

---

<div class="post-metadata">

**Author:** ![TotallyInformation](https://sea2.discourse-cdn.com/flex026/user_avatar/discourse.nodered.org/totallyinformation/32/31_2.png) [@TotallyInformation](https://discourse.nodered.org/u/TotallyInformation)\
**Post date:** [12 January 2024 20:34 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347/67 "2024-01-12T20:34:20Z")

</div>

> [@shrickus](#):
>
> it still reads the entire .csv file into memory before the first line is sent out?

Yes, it is a severe limitation of most of the nodes. Someone created the "bigxxxxxx" nodes a long time back where I think he tried to work around these issues but this is a better solution for sure.

---

<div class="post-metadata">

**Author:** ![gregorius](https://sea2.discourse-cdn.com/flex026/user_avatar/discourse.nodered.org/gregorius/32/73816_2.png) [@gregorius](https://discourse.nodered.org/u/gregorius)\
**Post date:** [12 January 2024 20:46 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347/68 "2024-01-12T20:46:59Z")

</div>

> [@gregorius](#):
>
> ... also a [Path2JsonL](https://cdn.flowhub.org/?fhid=c520d9da20ad7f1d#flow/2c7ddaaab6869956) subflow ...

> [@gregorius](#):
>
> behaviour --\> [the Get2Disk node](https://cdn.flowhub.org/?t=0&fhid=c520d9da20ad7f1d#node/8dd8e84ce59f090a/edit) is a s

I have just updated the serverless Node-RED installation to include the pipestream nodes, so the above two links can be used to explore the pipestream nodes in the the ETL pipeline. i.e., their are no longer red marked as being missing.

---

<div class="post-metadata">

**Author:** ![dceejay](https://sea2.discourse-cdn.com/flex026/user_avatar/discourse.nodered.org/dceejay/32/38_2.png) [@dceejay](https://discourse.nodered.org/u/dceejay)\
**Post date:** [12 January 2024 22:35 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347/69 "2024-01-12T22:35:11Z")

</div>

Though the file node can read in a line at a time to then send to the csv node, which knows that the parts come from a single file, so can do the right thing with headers/columns, but yes - other nodes , not so much.

---

<div class="post-metadata">

**Author:** ![dceejay](https://sea2.discourse-cdn.com/flex026/user_avatar/discourse.nodered.org/dceejay/32/38_2.png) [@dceejay](https://discourse.nodered.org/u/dceejay)\
**Post date:** [12 January 2024 22:41 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347/70 "2024-01-12T22:41:03Z")

</div>

This is all great work, and is something I have wondered about for a long time, it’s great to see this progress, but my head scratcher is how to “cross the streams” so to speak, for example how can I search in a stream for a sequence of header bytes and then send a chunk to say mqtt ?

IE how do I jump from the stream world out to flow world and indeed maybe back in.

---

<div class="post-metadata">

**Author:** ![gregorius](https://sea2.discourse-cdn.com/flex026/user_avatar/discourse.nodered.org/gregorius/32/73816_2.png) [@gregorius](https://discourse.nodered.org/u/gregorius)\
**Post date:** [12 January 2024 23:37 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347/71 "2024-01-12T23:37:40Z")

</div>

> [@dceejay](#):
>
> IE how do I jump from the stream world out to flow world and indeed maybe back in.

For that, check out the code for the [LineStream](https://github.com/gorenje/node-red-streaming/blob/a849576c86fec5a97ac18b8900eb2a167fc08ca2/nodes/linestream.js#L44-L54) node. What happens is that there is a [ByLine](https://github.com/gorenje/node-red-streaming/blob/a849576c86fec5a97ac18b8900eb2a167fc08ca2/nodes/linestream.js#L8-L34) streamer that does minimal buffering before passing on a line to the next stream (in this case the inline Transformer that sends the message).

This node is inserted in a pipeline after something like a FileStream which is a Readable stream.

There are three types of streams: [Writable](https://nodejs.org/dist/latest-v18.x/docs/api/stream.html#class-streamwritable), [Readable](https://nodejs.org/dist/latest-v18.x/docs/api/stream.html#class-streamreadable) and [Transformer](https://nodejs.org/dist/latest-v18.x/docs/api/stream.html#class-streamtransform) - all pipelines do something like Readable =\> Transformer =\> Writable whereby there can be as many Transformers as required. Pipelines are just a collection of streams. The LineStream node adds two Transform nodes to the pipeline: one to create lines and the other to send a msg object into Node-RED for each line that comes along.

This is all explained with my very high-level understanding of Streaming in JS. I suspect, for example, that the LineStream node is actually wrong because it does not pass on its content - that should could done better. I don't quite understand how to complete a pipeline or how to split pipelines along different paths.

---

<div class="post-metadata">

**Author:** ![dceejay](https://sea2.discourse-cdn.com/flex026/user_avatar/discourse.nodered.org/dceejay/32/38_2.png) [@dceejay](https://discourse.nodered.org/u/dceejay)\
**Post date:** [13 January 2024 08:56 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347/72 "2024-01-13T08:56:51Z")

</div>

I understand when you say "a line" - but what defines "a line" in non-text data ? eg video frames ? Is that defined in the Readable part ?

---

<div class="post-metadata">

**Author:** ![kuema](https://sea2.discourse-cdn.com/flex026/user_avatar/discourse.nodered.org/kuema/32/6542_2.png) [@kuema](https://discourse.nodered.org/u/kuema)\
**Post date:** [13 January 2024 09:12 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347/73 "2024-01-13T09:12:17Z")

</div>

That example transformer is strictly for line separated textual data.

You'd need to implement one for your specific data format, if you need to split it into meaningful chunks.

I use that for example to decode and create messages from binary data streams, where the data length is encoded in a header for each message.

---

<div class="post-metadata">

**Author:** ![gregorius](https://sea2.discourse-cdn.com/flex026/user_avatar/discourse.nodered.org/gregorius/32/73816_2.png) [@gregorius](https://discourse.nodered.org/u/gregorius)\
**Post date:** [13 January 2024 09:48 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347/74 "2024-01-13T09:48:24Z")

</div>

> [@dceejay](#):
>
> Is that defined in the Readable part

It's not the Readable that defines (at least in my experience) the meaning of the data, all the Readable does is generate chunks of data until there is no more data.

It's the Transformer that assigns "meaning" to the data, i.e., frame or line or packet or whatever.

In the context of ETL pipelines, I could imagine that a line is defined, then transformed to CSV object and then the CSV object is modified according to some business logic and finally the csv object is streamed into Node-RED. Or perhaps the CSV object is streamed directly into some data sink (i.e. database, data warehouse or some message bus). So I could imagine that the ETL pipeline is created, maintained and modified in Node-RED but, when executed, the data is streamed completely passed NR and directly from web/data source to database/data store without touching NR.

---

<div class="post-metadata">

**Author:** ![gregorius](https://sea2.discourse-cdn.com/flex026/user_avatar/discourse.nodered.org/gregorius/32/73816_2.png) [@gregorius](https://discourse.nodered.org/u/gregorius)\
**Post date:** [13 January 2024 15:55 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347/75 "2024-01-13T15:55:32Z")

</div>

Just a heads up, the package is now called [node-red-streaming](https://flows.nodered.org/node/@gregoriusrippenstein/node-red-streaming) and has gotten a more comprehensive readme describing much of what we have discussed here.

---

<div class="post-metadata">

**Author:** ![TotallyInformation](https://sea2.discourse-cdn.com/flex026/user_avatar/discourse.nodered.org/totallyinformation/32/31_2.png) [@TotallyInformation](https://discourse.nodered.org/u/TotallyInformation)\
**Post date:** [13 January 2024 16:01 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347/76 "2024-01-13T16:01:31Z")

</div>

> [@gregorius](#):
>
> In the context of ETL pipelines, I could imagine that a line is defined, then transformed to CSV object and then the CSV object is modified according to some business logic and finally the csv object is streamed into Node-RED. Or perhaps the CSV object is streamed directly into some data sink (i.e. database, data warehouse or some message bus). So I could imagine that the ETL pipeline is created, maintained and modified in Node-RED but, when executed, the data is streamed completely passed NR and directly from web/data source to database/data store without touching NR.

Most common requirements in my world is that the source is already CSV, or sometimes multiple CSV's and the pipeline would convert a line to an object and then a "business" transformation applied - most likely filter and/or grouping. Final stage would either be output to a new file or direct to HTML display. HTML output most likely because otherwise I'd probably be using a different tool to be honest.

For me, I can't think of a reason I'd be streaming anything other than some structured data. There are better tools for handling streaming media I think and I'd only use Node-RED to manage a menu or some other metadata display.

---

<div class="post-metadata">

**Author:** ![system](https://us1.discourse-cdn.com/flex026/uploads/nodered/original/1X/d073cd938eafa2e558d7c2cd59003b3ef4963033.png) [@system](https://discourse.nodered.org/u/system)\
**Post date:** [10 December 2024 08:23 UTC](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347/77 "2024-12-10T08:23:46Z")

</div>

This topic was automatically closed 14 days after the last reply. New replies are no longer allowed.

[Previous page](https://discourse.nodered.org/t/etl-pipelines-with-node-red/84347.md?page=3)
