Skip to content
This repository was archived by the owner on Sep 9, 2025. It is now read-only.

Repository files navigation

ARCHIVED & NO LONGER ACTIVELY MAINTAINED

JPipe

go report cardgo versiondocumentationGo.Dev referenceMIT license

A user-friendly implementation of the pipeline pattern in Go.

Overview

The pipeline pattern has been described by members of the core Go team several times:

Go provides very powerful concurrency primitives, but implementing the pipeline pattern correctly, with a correct handling of cancellation, requires a very good understanding of those primitives, and some non-negligible amount of boilerplate code. As pipelines become more complex, that boilerplate also starts to weigh heavily on code readability. Enter JPipe.

Features

  • Simple, controlled, per-operator concurrency
  • Ordered FIFO-like concurrency
  • Asynchronous execution and API
  • Safe context-based cancellation
  • Type-safe API
  • Fluent API (as much as allowed by Go generics)

Model

A Pipeline is a directed acyclic graph(DAG), where operators are nodes and Channels are edges:

graph LR;
A["FromSlice\n(operator)"]--Channel-->B["Filter\n(operator)"];
B--Channel-->C["Map\n(operator)"];
C--Channel-->D["ForEach\n(operator)"];
Loading

JPipe has the classic operators Map, Filter, ForEach and many more. Operators may have options. In particular, operators that take a function as parameter(e.g. Map) usually support concurrency, which applies only to that operator, and not the whole pipeline. This allows for fine-grained concurrency control.

Operators are not necessarily linear, so they may have multiple input/output Channels. Merge e.g. takes several inputs and merges them into a single output.

Channels are just a light wrapper over a plain Go channel, and you can reason about them in the same way you do with Go channels. The only exception is that a Channel can only be input to one operator.

Usage

Assume we have an expensive IO operation that takes 1 second to execute:

funcexpensiveIOOperation(idint) {
time.Sleep(time.Second)
}

Imagine this operation must be run for ids 1 through 10. We don't want to wait 10 seconds though, so we decide to do it with a concurrency factor of 5, expecting to get the full operation down to 2 seconds. The full Go code for that would be:

Plain Go version
ids:= []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
channel:=make(chanint)
concurrency:=5varwg sync.WaitGroupfori:=0; i<concurrency; i++ {
wg.Add(1)
gofunc() {
deferwg.Done()
forid:=rangechannel {
expensiveIOOperation(id)
}
}()
}
outer:
for_, id:=rangeids {
select {
// The nested select gives priority to the ctx.Done() signal, so we always exit early if needed// Without it, a single select just has no priority, so a new value could be processed even if the context has been canceledcase<-ctx.Done():
break outer
default:
select {
casechannel<-id:
case<-ctx.Done(): // always check ctx.Done() to avoid leaking the goroutinebreak outer
}
}
}
close(channel)
wg.Wait()

That's a lot of code right there for a simple work pool! We even had to make it collapsable to avoid disrupting the reading flow. Admittedly, most of the complexity comes from cancellation handling, but you don't want to go around leaking your goroutines. Now let's see how the same thing is done with JPipe:

ids:= []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
pipeline:=jpipe.New(ctx)
<-jpipe.FromSlice(pipeline, ids).
ForEach(expensiveIOOperation, jpipe.Concurrent(5))

Complex pipelines

The above is a simple work pool and admittedly the most common use case you'll find for concurrency. But JPipe allows you to build more complex pipelines with its catalog of operators. Imagine the following processing pipeline:

graph LR;
F1["Transaction feed 1"]-->M["Merge"];
F2["Transaction feed 2"]-->M;
M-->F["Filter transactions\nover 50 EUR"]
F-->T["Take first 20\ntransactions"]
T-->E["Enhance transactions\nwith external data\n(IO-bound,\nconcurrency 5,\nFIFO)"]
E-->B["Batch transactions\n(batch size 3,\ntimeout 5 sec)"]
B-->K["Send batches\nto Kafka"]
Loading

The JPipe implementation would be:

pipeline:=jpipe.New(ctx)
feed1:=jpipe.FromGoChannel(pipeline, getTransactionsFromFeed("feed1"))
feed2:=jpipe.FromGoChannel(pipeline, getTransactionsFromFeed("feed2"))
txs:=jpipe.Merge(feed1, feed2).
Filter(func(ftFeedTransaction) bool { returnft.Amount>50 }).
Take(20)
enhancedTxs:=jpipe.Map(txs, enhanceTransaction, jpipe.Concurrent(4), jpipe.Ordered(10))
<-jpipe.Batch(enhancedTxs, 3, 5*time.Second).
ForEach(sendTransactionBatchToKafka)

Documentation

You can find much more details on our official documentation. Some useful links in there:

You can also check the Go.Dev reference.

Similar projects

  • RxGo: Probably the best alternative out there, with lots of operators. It hasn't seen action in two years though, so it hasn't adopted generics: the API still deals with interface{}. A potential drawback for Go developers is its use of the reactive model, which is an abstraction very different from Go channels, and requires understanding of concepts like observer, observable, backpressure strategies, etc. If you'd rather stay within Go's channel abstraction, JPipe may be simpler to use, but RxGo may be your preferred choice if you have a background with the reactive model from other languages.
  • pipeline: An implementation of the pipeline pattern, but with a limited set of operators. It hasn't adopted generics yet either.
  • parapipe: A very simple pipeline implementation, but it has a single Pipe operator and does not support very complex pipelines. It hasn't adopted generics yet either.
  • ordered-concurrently: An implementation of ordered concurrency. It doesn't try to be a complete pipeline pattern library though, and just focuses on that feature.

License

JPipe is open-source software released under the MIT License.

About

Concurrent pipelines for Go

Topics

Resources

Security policy

Stars

21 stars

Watchers

3 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
GitHub - junitechnology/jpipe: Concurrent pipelines for Go · GitHub
Skip to content
This repository was archived by the owner on Sep 9, 2025. It is now read-only.

Repository files navigation

ARCHIVED & NO LONGER ACTIVELY MAINTAINED

JPipe

go report cardgo versiondocumentationGo.Dev referenceMIT license

A user-friendly implementation of the pipeline pattern in Go.

Overview

The pipeline pattern has been described by members of the core Go team several times:

Go provides very powerful concurrency primitives, but implementing the pipeline pattern correctly, with a correct handling of cancellation, requires a very good understanding of those primitives, and some non-negligible amount of boilerplate code. As pipelines become more complex, that boilerplate also starts to weigh heavily on code readability. Enter JPipe.

Features

  • Simple, controlled, per-operator concurrency
  • Ordered FIFO-like concurrency
  • Asynchronous execution and API
  • Safe context-based cancellation
  • Type-safe API
  • Fluent API (as much as allowed by Go generics)

Model

A Pipeline is a directed acyclic graph(DAG), where operators are nodes and Channels are edges:

graph LR;
A["FromSlice\n(operator)"]--Channel-->B["Filter\n(operator)"];
B--Channel-->C["Map\n(operator)"];
C--Channel-->D["ForEach\n(operator)"];
Loading

JPipe has the classic operators Map, Filter, ForEach and many more. Operators may have options. In particular, operators that take a function as parameter(e.g. Map) usually support concurrency, which applies only to that operator, and not the whole pipeline. This allows for fine-grained concurrency control.

Operators are not necessarily linear, so they may have multiple input/output Channels. Merge e.g. takes several inputs and merges them into a single output.

Channels are just a light wrapper over a plain Go channel, and you can reason about them in the same way you do with Go channels. The only exception is that a Channel can only be input to one operator.

Usage

Assume we have an expensive IO operation that takes 1 second to execute:

funcexpensiveIOOperation(idint) {
time.Sleep(time.Second)
}

Imagine this operation must be run for ids 1 through 10. We don't want to wait 10 seconds though, so we decide to do it with a concurrency factor of 5, expecting to get the full operation down to 2 seconds. The full Go code for that would be:

Plain Go version
ids:= []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
channel:=make(chanint)
concurrency:=5varwg sync.WaitGroupfori:=0; i<concurrency; i++ {
wg.Add(1)
gofunc() {
deferwg.Done()
forid:=rangechannel {
expensiveIOOperation(id)
}
}()
}
outer:
for_, id:=rangeids {
select {
// The nested select gives priority to the ctx.Done() signal, so we always exit early if needed// Without it, a single select just has no priority, so a new value could be processed even if the context has been canceledcase<-ctx.Done():
break outer
default:
select {
casechannel<-id:
case<-ctx.Done(): // always check ctx.Done() to avoid leaking the goroutinebreak outer
}
}
}
close(channel)
wg.Wait()

That's a lot of code right there for a simple work pool! We even had to make it collapsable to avoid disrupting the reading flow. Admittedly, most of the complexity comes from cancellation handling, but you don't want to go around leaking your goroutines. Now let's see how the same thing is done with JPipe:

ids:= []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
pipeline:=jpipe.New(ctx)
<-jpipe.FromSlice(pipeline, ids).
ForEach(expensiveIOOperation, jpipe.Concurrent(5))

Complex pipelines

The above is a simple work pool and admittedly the most common use case you'll find for concurrency. But JPipe allows you to build more complex pipelines with its catalog of operators. Imagine the following processing pipeline:

graph LR;
F1["Transaction feed 1"]-->M["Merge"];
F2["Transaction feed 2"]-->M;
M-->F["Filter transactions\nover 50 EUR"]
F-->T["Take first 20\ntransactions"]
T-->E["Enhance transactions\nwith external data\n(IO-bound,\nconcurrency 5,\nFIFO)"]
E-->B["Batch transactions\n(batch size 3,\ntimeout 5 sec)"]
B-->K["Send batches\nto Kafka"]
Loading

The JPipe implementation would be:

pipeline:=jpipe.New(ctx)
feed1:=jpipe.FromGoChannel(pipeline, getTransactionsFromFeed("feed1"))
feed2:=jpipe.FromGoChannel(pipeline, getTransactionsFromFeed("feed2"))
txs:=jpipe.Merge(feed1, feed2).
Filter(func(ftFeedTransaction) bool { returnft.Amount>50 }).
Take(20)
enhancedTxs:=jpipe.Map(txs, enhanceTransaction, jpipe.Concurrent(4), jpipe.Ordered(10))
<-jpipe.Batch(enhancedTxs, 3, 5*time.Second).
ForEach(sendTransactionBatchToKafka)

Documentation

You can find much more details on our official documentation. Some useful links in there:

You can also check the Go.Dev reference.

Similar projects

  • RxGo: Probably the best alternative out there, with lots of operators. It hasn't seen action in two years though, so it hasn't adopted generics: the API still deals with interface{}. A potential drawback for Go developers is its use of the reactive model, which is an abstraction very different from Go channels, and requires understanding of concepts like observer, observable, backpressure strategies, etc. If you'd rather stay within Go's channel abstraction, JPipe may be simpler to use, but RxGo may be your preferred choice if you have a background with the reactive model from other languages.
  • pipeline: An implementation of the pipeline pattern, but with a limited set of operators. It hasn't adopted generics yet either.
  • parapipe: A very simple pipeline implementation, but it has a single Pipe operator and does not support very complex pipelines. It hasn't adopted generics yet either.
  • ordered-concurrently: An implementation of ordered concurrency. It doesn't try to be a complete pipeline pattern library though, and just focuses on that feature.

License

JPipe is open-source software released under the MIT License.

About

Concurrent pipelines for Go

Topics

Resources

Security policy

Stars

21 stars

Watchers

3 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' GitHub - junitechnology/jpipe: Concurrent pipelines for Go · GitHub
Skip to content
This repository was archived by the owner on Sep 9, 2025. It is now read-only.

Repository files navigation

ARCHIVED & NO LONGER ACTIVELY MAINTAINED

JPipe

go report cardgo versiondocumentationGo.Dev referenceMIT license

A user-friendly implementation of the pipeline pattern in Go.

Overview

The pipeline pattern has been described by members of the core Go team several times:

Go provides very powerful concurrency primitives, but implementing the pipeline pattern correctly, with a correct handling of cancellation, requires a very good understanding of those primitives, and some non-negligible amount of boilerplate code. As pipelines become more complex, that boilerplate also starts to weigh heavily on code readability. Enter JPipe.

Features

  • Simple, controlled, per-operator concurrency
  • Ordered FIFO-like concurrency
  • Asynchronous execution and API
  • Safe context-based cancellation
  • Type-safe API
  • Fluent API (as much as allowed by Go generics)

Model

A Pipeline is a directed acyclic graph(DAG), where operators are nodes and Channels are edges:

graph LR;
A["FromSlice\n(operator)"]--Channel-->B["Filter\n(operator)"];
B--Channel-->C["Map\n(operator)"];
C--Channel-->D["ForEach\n(operator)"];
Loading

JPipe has the classic operators Map, Filter, ForEach and many more. Operators may have options. In particular, operators that take a function as parameter(e.g. Map) usually support concurrency, which applies only to that operator, and not the whole pipeline. This allows for fine-grained concurrency control.

Operators are not necessarily linear, so they may have multiple input/output Channels. Merge e.g. takes several inputs and merges them into a single output.

Channels are just a light wrapper over a plain Go channel, and you can reason about them in the same way you do with Go channels. The only exception is that a Channel can only be input to one operator.

Usage

Assume we have an expensive IO operation that takes 1 second to execute:

funcexpensiveIOOperation(idint) {
time.Sleep(time.Second)
}

Imagine this operation must be run for ids 1 through 10. We don't want to wait 10 seconds though, so we decide to do it with a concurrency factor of 5, expecting to get the full operation down to 2 seconds. The full Go code for that would be:

Plain Go version
ids:= []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
channel:=make(chanint)
concurrency:=5varwg sync.WaitGroupfori:=0; i<concurrency; i++ {
wg.Add(1)
gofunc() {
deferwg.Done()
forid:=rangechannel {
expensiveIOOperation(id)
}
}()
}
outer:
for_, id:=rangeids {
select {
// The nested select gives priority to the ctx.Done() signal, so we always exit early if needed// Without it, a single select just has no priority, so a new value could be processed even if the context has been canceledcase<-ctx.Done():
break outer
default:
select {
casechannel<-id:
case<-ctx.Done(): // always check ctx.Done() to avoid leaking the goroutinebreak outer
}
}
}
close(channel)
wg.Wait()

That's a lot of code right there for a simple work pool! We even had to make it collapsable to avoid disrupting the reading flow. Admittedly, most of the complexity comes from cancellation handling, but you don't want to go around leaking your goroutines. Now let's see how the same thing is done with JPipe:

ids:= []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
pipeline:=jpipe.New(ctx)
<-jpipe.FromSlice(pipeline, ids).
ForEach(expensiveIOOperation, jpipe.Concurrent(5))

Complex pipelines

The above is a simple work pool and admittedly the most common use case you'll find for concurrency. But JPipe allows you to build more complex pipelines with its catalog of operators. Imagine the following processing pipeline:

graph LR;
F1["Transaction feed 1"]-->M["Merge"];
F2["Transaction feed 2"]-->M;
M-->F["Filter transactions\nover 50 EUR"]
F-->T["Take first 20\ntransactions"]
T-->E["Enhance transactions\nwith external data\n(IO-bound,\nconcurrency 5,\nFIFO)"]
E-->B["Batch transactions\n(batch size 3,\ntimeout 5 sec)"]
B-->K["Send batches\nto Kafka"]
Loading

The JPipe implementation would be:

pipeline:=jpipe.New(ctx)
feed1:=jpipe.FromGoChannel(pipeline, getTransactionsFromFeed("feed1"))
feed2:=jpipe.FromGoChannel(pipeline, getTransactionsFromFeed("feed2"))
txs:=jpipe.Merge(feed1, feed2).
Filter(func(ftFeedTransaction) bool { returnft.Amount>50 }).
Take(20)
enhancedTxs:=jpipe.Map(txs, enhanceTransaction, jpipe.Concurrent(4), jpipe.Ordered(10))
<-jpipe.Batch(enhancedTxs, 3, 5*time.Second).
ForEach(sendTransactionBatchToKafka)

Documentation

You can find much more details on our official documentation. Some useful links in there:

You can also check the Go.Dev reference.

Similar projects

  • RxGo: Probably the best alternative out there, with lots of operators. It hasn't seen action in two years though, so it hasn't adopted generics: the API still deals with interface{}. A potential drawback for Go developers is its use of the reactive model, which is an abstraction very different from Go channels, and requires understanding of concepts like observer, observable, backpressure strategies, etc. If you'd rather stay within Go's channel abstraction, JPipe may be simpler to use, but RxGo may be your preferred choice if you have a background with the reactive model from other languages.
  • pipeline: An implementation of the pipeline pattern, but with a limited set of operators. It hasn't adopted generics yet either.
  • parapipe: A very simple pipeline implementation, but it has a single Pipe operator and does not support very complex pipelines. It hasn't adopted generics yet either.
  • ordered-concurrently: An implementation of ordered concurrency. It doesn't try to be a complete pipeline pattern library though, and just focuses on that feature.

License

JPipe is open-source software released under the MIT License.

About

Concurrent pipelines for Go

Topics

Resources

Security policy

Stars

21 stars

Watchers

3 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' GitHub - junitechnology/jpipe: Concurrent pipelines for Go · GitHub
Skip to content
This repository was archived by the owner on Sep 9, 2025. It is now read-only.

Repository files navigation

ARCHIVED & NO LONGER ACTIVELY MAINTAINED

JPipe

go report cardgo versiondocumentationGo.Dev referenceMIT license

A user-friendly implementation of the pipeline pattern in Go.

Overview

The pipeline pattern has been described by members of the core Go team several times:

Go provides very powerful concurrency primitives, but implementing the pipeline pattern correctly, with a correct handling of cancellation, requires a very good understanding of those primitives, and some non-negligible amount of boilerplate code. As pipelines become more complex, that boilerplate also starts to weigh heavily on code readability. Enter JPipe.

Features

  • Simple, controlled, per-operator concurrency
  • Ordered FIFO-like concurrency
  • Asynchronous execution and API
  • Safe context-based cancellation
  • Type-safe API
  • Fluent API (as much as allowed by Go generics)

Model

A Pipeline is a directed acyclic graph(DAG), where operators are nodes and Channels are edges:

graph LR;
A["FromSlice\n(operator)"]--Channel-->B["Filter\n(operator)"];
B--Channel-->C["Map\n(operator)"];
C--Channel-->D["ForEach\n(operator)"];
Loading

JPipe has the classic operators Map, Filter, ForEach and many more. Operators may have options. In particular, operators that take a function as parameter(e.g. Map) usually support concurrency, which applies only to that operator, and not the whole pipeline. This allows for fine-grained concurrency control.

Operators are not necessarily linear, so they may have multiple input/output Channels. Merge e.g. takes several inputs and merges them into a single output.

Channels are just a light wrapper over a plain Go channel, and you can reason about them in the same way you do with Go channels. The only exception is that a Channel can only be input to one operator.

Usage

Assume we have an expensive IO operation that takes 1 second to execute:

funcexpensiveIOOperation(idint) {
time.Sleep(time.Second)
}

Imagine this operation must be run for ids 1 through 10. We don't want to wait 10 seconds though, so we decide to do it with a concurrency factor of 5, expecting to get the full operation down to 2 seconds. The full Go code for that would be:

Plain Go version
ids:= []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
channel:=make(chanint)
concurrency:=5varwg sync.WaitGroupfori:=0; i<concurrency; i++ {
wg.Add(1)
gofunc() {
deferwg.Done()
forid:=rangechannel {
expensiveIOOperation(id)
}
}()
}
outer:
for_, id:=rangeids {
select {
// The nested select gives priority to the ctx.Done() signal, so we always exit early if needed// Without it, a single select just has no priority, so a new value could be processed even if the context has been canceledcase<-ctx.Done():
break outer
default:
select {
casechannel<-id:
case<-ctx.Done(): // always check ctx.Done() to avoid leaking the goroutinebreak outer
}
}
}
close(channel)
wg.Wait()

That's a lot of code right there for a simple work pool! We even had to make it collapsable to avoid disrupting the reading flow. Admittedly, most of the complexity comes from cancellation handling, but you don't want to go around leaking your goroutines. Now let's see how the same thing is done with JPipe:

ids:= []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
pipeline:=jpipe.New(ctx)
<-jpipe.FromSlice(pipeline, ids).
ForEach(expensiveIOOperation, jpipe.Concurrent(5))

Complex pipelines

The above is a simple work pool and admittedly the most common use case you'll find for concurrency. But JPipe allows you to build more complex pipelines with its catalog of operators. Imagine the following processing pipeline:

graph LR;
F1["Transaction feed 1"]-->M["Merge"];
F2["Transaction feed 2"]-->M;
M-->F["Filter transactions\nover 50 EUR"]
F-->T["Take first 20\ntransactions"]
T-->E["Enhance transactions\nwith external data\n(IO-bound,\nconcurrency 5,\nFIFO)"]
E-->B["Batch transactions\n(batch size 3,\ntimeout 5 sec)"]
B-->K["Send batches\nto Kafka"]
Loading

The JPipe implementation would be:

pipeline:=jpipe.New(ctx)
feed1:=jpipe.FromGoChannel(pipeline, getTransactionsFromFeed("feed1"))
feed2:=jpipe.FromGoChannel(pipeline, getTransactionsFromFeed("feed2"))
txs:=jpipe.Merge(feed1, feed2).
Filter(func(ftFeedTransaction) bool { returnft.Amount>50 }).
Take(20)
enhancedTxs:=jpipe.Map(txs, enhanceTransaction, jpipe.Concurrent(4), jpipe.Ordered(10))
<-jpipe.Batch(enhancedTxs, 3, 5*time.Second).
ForEach(sendTransactionBatchToKafka)

Documentation

You can find much more details on our official documentation. Some useful links in there:

You can also check the Go.Dev reference.

Similar projects

  • RxGo: Probably the best alternative out there, with lots of operators. It hasn't seen action in two years though, so it hasn't adopted generics: the API still deals with interface{}. A potential drawback for Go developers is its use of the reactive model, which is an abstraction very different from Go channels, and requires understanding of concepts like observer, observable, backpressure strategies, etc. If you'd rather stay within Go's channel abstraction, JPipe may be simpler to use, but RxGo may be your preferred choice if you have a background with the reactive model from other languages.
  • pipeline: An implementation of the pipeline pattern, but with a limited set of operators. It hasn't adopted generics yet either.
  • parapipe: A very simple pipeline implementation, but it has a single Pipe operator and does not support very complex pipelines. It hasn't adopted generics yet either.
  • ordered-concurrently: An implementation of ordered concurrency. It doesn't try to be a complete pipeline pattern library though, and just focuses on that feature.

License

JPipe is open-source software released under the MIT License.

About

Concurrent pipelines for Go

Topics

Resources

Security policy

Stars

21 stars

Watchers

3 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + ' GitHub - junitechnology/jpipe: Concurrent pipelines for Go · GitHub
Skip to content
This repository was archived by the owner on Sep 9, 2025. It is now read-only.

Repository files navigation

ARCHIVED & NO LONGER ACTIVELY MAINTAINED

JPipe

go report cardgo versiondocumentationGo.Dev referenceMIT license

A user-friendly implementation of the pipeline pattern in Go.

Overview

The pipeline pattern has been described by members of the core Go team several times:

Go provides very powerful concurrency primitives, but implementing the pipeline pattern correctly, with a correct handling of cancellation, requires a very good understanding of those primitives, and some non-negligible amount of boilerplate code. As pipelines become more complex, that boilerplate also starts to weigh heavily on code readability. Enter JPipe.

Features

  • Simple, controlled, per-operator concurrency
  • Ordered FIFO-like concurrency
  • Asynchronous execution and API
  • Safe context-based cancellation
  • Type-safe API
  • Fluent API (as much as allowed by Go generics)

Model

A Pipeline is a directed acyclic graph(DAG), where operators are nodes and Channels are edges:

graph LR;
A["FromSlice\n(operator)"]--Channel-->B["Filter\n(operator)"];
B--Channel-->C["Map\n(operator)"];
C--Channel-->D["ForEach\n(operator)"];
Loading

JPipe has the classic operators Map, Filter, ForEach and many more. Operators may have options. In particular, operators that take a function as parameter(e.g. Map) usually support concurrency, which applies only to that operator, and not the whole pipeline. This allows for fine-grained concurrency control.

Operators are not necessarily linear, so they may have multiple input/output Channels. Merge e.g. takes several inputs and merges them into a single output.

Channels are just a light wrapper over a plain Go channel, and you can reason about them in the same way you do with Go channels. The only exception is that a Channel can only be input to one operator.

Usage

Assume we have an expensive IO operation that takes 1 second to execute:

funcexpensiveIOOperation(idint) {
time.Sleep(time.Second)
}

Imagine this operation must be run for ids 1 through 10. We don't want to wait 10 seconds though, so we decide to do it with a concurrency factor of 5, expecting to get the full operation down to 2 seconds. The full Go code for that would be:

Plain Go version
ids:= []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
channel:=make(chanint)
concurrency:=5varwg sync.WaitGroupfori:=0; i<concurrency; i++ {
wg.Add(1)
gofunc() {
deferwg.Done()
forid:=rangechannel {
expensiveIOOperation(id)
}
}()
}
outer:
for_, id:=rangeids {
select {
// The nested select gives priority to the ctx.Done() signal, so we always exit early if needed// Without it, a single select just has no priority, so a new value could be processed even if the context has been canceledcase<-ctx.Done():
break outer
default:
select {
casechannel<-id:
case<-ctx.Done(): // always check ctx.Done() to avoid leaking the goroutinebreak outer
}
}
}
close(channel)
wg.Wait()

That's a lot of code right there for a simple work pool! We even had to make it collapsable to avoid disrupting the reading flow. Admittedly, most of the complexity comes from cancellation handling, but you don't want to go around leaking your goroutines. Now let's see how the same thing is done with JPipe:

ids:= []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
pipeline:=jpipe.New(ctx)
<-jpipe.FromSlice(pipeline, ids).
ForEach(expensiveIOOperation, jpipe.Concurrent(5))

Complex pipelines

The above is a simple work pool and admittedly the most common use case you'll find for concurrency. But JPipe allows you to build more complex pipelines with its catalog of operators. Imagine the following processing pipeline:

graph LR;
F1["Transaction feed 1"]-->M["Merge"];
F2["Transaction feed 2"]-->M;
M-->F["Filter transactions\nover 50 EUR"]
F-->T["Take first 20\ntransactions"]
T-->E["Enhance transactions\nwith external data\n(IO-bound,\nconcurrency 5,\nFIFO)"]
E-->B["Batch transactions\n(batch size 3,\ntimeout 5 sec)"]
B-->K["Send batches\nto Kafka"]
Loading

The JPipe implementation would be:

pipeline:=jpipe.New(ctx)
feed1:=jpipe.FromGoChannel(pipeline, getTransactionsFromFeed("feed1"))
feed2:=jpipe.FromGoChannel(pipeline, getTransactionsFromFeed("feed2"))
txs:=jpipe.Merge(feed1, feed2).
Filter(func(ftFeedTransaction) bool { returnft.Amount>50 }).
Take(20)
enhancedTxs:=jpipe.Map(txs, enhanceTransaction, jpipe.Concurrent(4), jpipe.Ordered(10))
<-jpipe.Batch(enhancedTxs, 3, 5*time.Second).
ForEach(sendTransactionBatchToKafka)

Documentation

You can find much more details on our official documentation. Some useful links in there:

You can also check the Go.Dev reference.

Similar projects

  • RxGo: Probably the best alternative out there, with lots of operators. It hasn't seen action in two years though, so it hasn't adopted generics: the API still deals with interface{}. A potential drawback for Go developers is its use of the reactive model, which is an abstraction very different from Go channels, and requires understanding of concepts like observer, observable, backpressure strategies, etc. If you'd rather stay within Go's channel abstraction, JPipe may be simpler to use, but RxGo may be your preferred choice if you have a background with the reactive model from other languages.
  • pipeline: An implementation of the pipeline pattern, but with a limited set of operators. It hasn't adopted generics yet either.
  • parapipe: A very simple pipeline implementation, but it has a single Pipe operator and does not support very complex pipelines. It hasn't adopted generics yet either.
  • ordered-concurrently: An implementation of ordered concurrency. It doesn't try to be a complete pipeline pattern library though, and just focuses on that feature.

License

JPipe is open-source software released under the MIT License.

About

Concurrent pipelines for Go

Topics

Resources

Security policy

Stars

21 stars

Watchers

3 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' GitHub - junitechnology/jpipe: Concurrent pipelines for Go · GitHub
Skip to content
This repository was archived by the owner on Sep 9, 2025. It is now read-only.

Repository files navigation

ARCHIVED & NO LONGER ACTIVELY MAINTAINED

JPipe

go report cardgo versiondocumentationGo.Dev referenceMIT license

A user-friendly implementation of the pipeline pattern in Go.

Overview

The pipeline pattern has been described by members of the core Go team several times:

Go provides very powerful concurrency primitives, but implementing the pipeline pattern correctly, with a correct handling of cancellation, requires a very good understanding of those primitives, and some non-negligible amount of boilerplate code. As pipelines become more complex, that boilerplate also starts to weigh heavily on code readability. Enter JPipe.

Features

  • Simple, controlled, per-operator concurrency
  • Ordered FIFO-like concurrency
  • Asynchronous execution and API
  • Safe context-based cancellation
  • Type-safe API
  • Fluent API (as much as allowed by Go generics)

Model

A Pipeline is a directed acyclic graph(DAG), where operators are nodes and Channels are edges:

graph LR;
A["FromSlice\n(operator)"]--Channel-->B["Filter\n(operator)"];
B--Channel-->C["Map\n(operator)"];
C--Channel-->D["ForEach\n(operator)"];
Loading

JPipe has the classic operators Map, Filter, ForEach and many more. Operators may have options. In particular, operators that take a function as parameter(e.g. Map) usually support concurrency, which applies only to that operator, and not the whole pipeline. This allows for fine-grained concurrency control.

Operators are not necessarily linear, so they may have multiple input/output Channels. Merge e.g. takes several inputs and merges them into a single output.

Channels are just a light wrapper over a plain Go channel, and you can reason about them in the same way you do with Go channels. The only exception is that a Channel can only be input to one operator.

Usage

Assume we have an expensive IO operation that takes 1 second to execute:

funcexpensiveIOOperation(idint) {
time.Sleep(time.Second)
}

Imagine this operation must be run for ids 1 through 10. We don't want to wait 10 seconds though, so we decide to do it with a concurrency factor of 5, expecting to get the full operation down to 2 seconds. The full Go code for that would be:

Plain Go version
ids:= []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
channel:=make(chanint)
concurrency:=5varwg sync.WaitGroupfori:=0; i<concurrency; i++ {
wg.Add(1)
gofunc() {
deferwg.Done()
forid:=rangechannel {
expensiveIOOperation(id)
}
}()
}
outer:
for_, id:=rangeids {
select {
// The nested select gives priority to the ctx.Done() signal, so we always exit early if needed// Without it, a single select just has no priority, so a new value could be processed even if the context has been canceledcase<-ctx.Done():
break outer
default:
select {
casechannel<-id:
case<-ctx.Done(): // always check ctx.Done() to avoid leaking the goroutinebreak outer
}
}
}
close(channel)
wg.Wait()

That's a lot of code right there for a simple work pool! We even had to make it collapsable to avoid disrupting the reading flow. Admittedly, most of the complexity comes from cancellation handling, but you don't want to go around leaking your goroutines. Now let's see how the same thing is done with JPipe:

ids:= []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
pipeline:=jpipe.New(ctx)
<-jpipe.FromSlice(pipeline, ids).
ForEach(expensiveIOOperation, jpipe.Concurrent(5))

Complex pipelines

The above is a simple work pool and admittedly the most common use case you'll find for concurrency. But JPipe allows you to build more complex pipelines with its catalog of operators. Imagine the following processing pipeline:

graph LR;
F1["Transaction feed 1"]-->M["Merge"];
F2["Transaction feed 2"]-->M;
M-->F["Filter transactions\nover 50 EUR"]
F-->T["Take first 20\ntransactions"]
T-->E["Enhance transactions\nwith external data\n(IO-bound,\nconcurrency 5,\nFIFO)"]
E-->B["Batch transactions\n(batch size 3,\ntimeout 5 sec)"]
B-->K["Send batches\nto Kafka"]
Loading

The JPipe implementation would be:

pipeline:=jpipe.New(ctx)
feed1:=jpipe.FromGoChannel(pipeline, getTransactionsFromFeed("feed1"))
feed2:=jpipe.FromGoChannel(pipeline, getTransactionsFromFeed("feed2"))
txs:=jpipe.Merge(feed1, feed2).
Filter(func(ftFeedTransaction) bool { returnft.Amount>50 }).
Take(20)
enhancedTxs:=jpipe.Map(txs, enhanceTransaction, jpipe.Concurrent(4), jpipe.Ordered(10))
<-jpipe.Batch(enhancedTxs, 3, 5*time.Second).
ForEach(sendTransactionBatchToKafka)

Documentation

You can find much more details on our official documentation. Some useful links in there:

You can also check the Go.Dev reference.

Similar projects

  • RxGo: Probably the best alternative out there, with lots of operators. It hasn't seen action in two years though, so it hasn't adopted generics: the API still deals with interface{}. A potential drawback for Go developers is its use of the reactive model, which is an abstraction very different from Go channels, and requires understanding of concepts like observer, observable, backpressure strategies, etc. If you'd rather stay within Go's channel abstraction, JPipe may be simpler to use, but RxGo may be your preferred choice if you have a background with the reactive model from other languages.
  • pipeline: An implementation of the pipeline pattern, but with a limited set of operators. It hasn't adopted generics yet either.
  • parapipe: A very simple pipeline implementation, but it has a single Pipe operator and does not support very complex pipelines. It hasn't adopted generics yet either.
  • ordered-concurrently: An implementation of ordered concurrency. It doesn't try to be a complete pipeline pattern library though, and just focuses on that feature.

License

JPipe is open-source software released under the MIT License.

About

Concurrent pipelines for Go

Topics

Resources

Security policy

Stars

21 stars

Watchers

3 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' GitHub - junitechnology/jpipe: Concurrent pipelines for Go · GitHub
Skip to content
This repository was archived by the owner on Sep 9, 2025. It is now read-only.

Repository files navigation

ARCHIVED & NO LONGER ACTIVELY MAINTAINED

JPipe

go report cardgo versiondocumentationGo.Dev referenceMIT license

A user-friendly implementation of the pipeline pattern in Go.

Overview

The pipeline pattern has been described by members of the core Go team several times:

Go provides very powerful concurrency primitives, but implementing the pipeline pattern correctly, with a correct handling of cancellation, requires a very good understanding of those primitives, and some non-negligible amount of boilerplate code. As pipelines become more complex, that boilerplate also starts to weigh heavily on code readability. Enter JPipe.

Features

  • Simple, controlled, per-operator concurrency
  • Ordered FIFO-like concurrency
  • Asynchronous execution and API
  • Safe context-based cancellation
  • Type-safe API
  • Fluent API (as much as allowed by Go generics)

Model

A Pipeline is a directed acyclic graph(DAG), where operators are nodes and Channels are edges:

graph LR;
A["FromSlice\n(operator)"]--Channel-->B["Filter\n(operator)"];
B--Channel-->C["Map\n(operator)"];
C--Channel-->D["ForEach\n(operator)"];
Loading

JPipe has the classic operators Map, Filter, ForEach and many more. Operators may have options. In particular, operators that take a function as parameter(e.g. Map) usually support concurrency, which applies only to that operator, and not the whole pipeline. This allows for fine-grained concurrency control.

Operators are not necessarily linear, so they may have multiple input/output Channels. Merge e.g. takes several inputs and merges them into a single output.

Channels are just a light wrapper over a plain Go channel, and you can reason about them in the same way you do with Go channels. The only exception is that a Channel can only be input to one operator.

Usage

Assume we have an expensive IO operation that takes 1 second to execute:

funcexpensiveIOOperation(idint) {
time.Sleep(time.Second)
}

Imagine this operation must be run for ids 1 through 10. We don't want to wait 10 seconds though, so we decide to do it with a concurrency factor of 5, expecting to get the full operation down to 2 seconds. The full Go code for that would be:

Plain Go version
ids:= []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
channel:=make(chanint)
concurrency:=5varwg sync.WaitGroupfori:=0; i<concurrency; i++ {
wg.Add(1)
gofunc() {
deferwg.Done()
forid:=rangechannel {
expensiveIOOperation(id)
}
}()
}
outer:
for_, id:=rangeids {
select {
// The nested select gives priority to the ctx.Done() signal, so we always exit early if needed// Without it, a single select just has no priority, so a new value could be processed even if the context has been canceledcase<-ctx.Done():
break outer
default:
select {
casechannel<-id:
case<-ctx.Done(): // always check ctx.Done() to avoid leaking the goroutinebreak outer
}
}
}
close(channel)
wg.Wait()

That's a lot of code right there for a simple work pool! We even had to make it collapsable to avoid disrupting the reading flow. Admittedly, most of the complexity comes from cancellation handling, but you don't want to go around leaking your goroutines. Now let's see how the same thing is done with JPipe:

ids:= []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
pipeline:=jpipe.New(ctx)
<-jpipe.FromSlice(pipeline, ids).
ForEach(expensiveIOOperation, jpipe.Concurrent(5))

Complex pipelines

The above is a simple work pool and admittedly the most common use case you'll find for concurrency. But JPipe allows you to build more complex pipelines with its catalog of operators. Imagine the following processing pipeline:

graph LR;
F1["Transaction feed 1"]-->M["Merge"];
F2["Transaction feed 2"]-->M;
M-->F["Filter transactions\nover 50 EUR"]
F-->T["Take first 20\ntransactions"]
T-->E["Enhance transactions\nwith external data\n(IO-bound,\nconcurrency 5,\nFIFO)"]
E-->B["Batch transactions\n(batch size 3,\ntimeout 5 sec)"]
B-->K["Send batches\nto Kafka"]
Loading

The JPipe implementation would be:

pipeline:=jpipe.New(ctx)
feed1:=jpipe.FromGoChannel(pipeline, getTransactionsFromFeed("feed1"))
feed2:=jpipe.FromGoChannel(pipeline, getTransactionsFromFeed("feed2"))
txs:=jpipe.Merge(feed1, feed2).
Filter(func(ftFeedTransaction) bool { returnft.Amount>50 }).
Take(20)
enhancedTxs:=jpipe.Map(txs, enhanceTransaction, jpipe.Concurrent(4), jpipe.Ordered(10))
<-jpipe.Batch(enhancedTxs, 3, 5*time.Second).
ForEach(sendTransactionBatchToKafka)

Documentation

You can find much more details on our official documentation. Some useful links in there:

You can also check the Go.Dev reference.

Similar projects

  • RxGo: Probably the best alternative out there, with lots of operators. It hasn't seen action in two years though, so it hasn't adopted generics: the API still deals with interface{}. A potential drawback for Go developers is its use of the reactive model, which is an abstraction very different from Go channels, and requires understanding of concepts like observer, observable, backpressure strategies, etc. If you'd rather stay within Go's channel abstraction, JPipe may be simpler to use, but RxGo may be your preferred choice if you have a background with the reactive model from other languages.
  • pipeline: An implementation of the pipeline pattern, but with a limited set of operators. It hasn't adopted generics yet either.
  • parapipe: A very simple pipeline implementation, but it has a single Pipe operator and does not support very complex pipelines. It hasn't adopted generics yet either.
  • ordered-concurrently: An implementation of ordered concurrency. It doesn't try to be a complete pipeline pattern library though, and just focuses on that feature.

License

JPipe is open-source software released under the MIT License.

About

Concurrent pipelines for Go

Topics

Resources

Security policy

Stars

21 stars

Watchers

3 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })(); GitHub - junitechnology/jpipe: Concurrent pipelines for Go · GitHub
Skip to content
This repository was archived by the owner on Sep 9, 2025. It is now read-only.

Repository files navigation

ARCHIVED & NO LONGER ACTIVELY MAINTAINED

JPipe

go report cardgo versiondocumentationGo.Dev referenceMIT license

A user-friendly implementation of the pipeline pattern in Go.

Overview

The pipeline pattern has been described by members of the core Go team several times:

Go provides very powerful concurrency primitives, but implementing the pipeline pattern correctly, with a correct handling of cancellation, requires a very good understanding of those primitives, and some non-negligible amount of boilerplate code. As pipelines become more complex, that boilerplate also starts to weigh heavily on code readability. Enter JPipe.

Features

  • Simple, controlled, per-operator concurrency
  • Ordered FIFO-like concurrency
  • Asynchronous execution and API
  • Safe context-based cancellation
  • Type-safe API
  • Fluent API (as much as allowed by Go generics)

Model

A Pipeline is a directed acyclic graph(DAG), where operators are nodes and Channels are edges:

graph LR;
A["FromSlice\n(operator)"]--Channel-->B["Filter\n(operator)"];
B--Channel-->C["Map\n(operator)"];
C--Channel-->D["ForEach\n(operator)"];
Loading

JPipe has the classic operators Map, Filter, ForEach and many more. Operators may have options. In particular, operators that take a function as parameter(e.g. Map) usually support concurrency, which applies only to that operator, and not the whole pipeline. This allows for fine-grained concurrency control.

Operators are not necessarily linear, so they may have multiple input/output Channels. Merge e.g. takes several inputs and merges them into a single output.

Channels are just a light wrapper over a plain Go channel, and you can reason about them in the same way you do with Go channels. The only exception is that a Channel can only be input to one operator.

Usage

Assume we have an expensive IO operation that takes 1 second to execute:

funcexpensiveIOOperation(idint) {
time.Sleep(time.Second)
}

Imagine this operation must be run for ids 1 through 10. We don't want to wait 10 seconds though, so we decide to do it with a concurrency factor of 5, expecting to get the full operation down to 2 seconds. The full Go code for that would be:

Plain Go version
ids:= []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
channel:=make(chanint)
concurrency:=5varwg sync.WaitGroupfori:=0; i<concurrency; i++ {
wg.Add(1)
gofunc() {
deferwg.Done()
forid:=rangechannel {
expensiveIOOperation(id)
}
}()
}
outer:
for_, id:=rangeids {
select {
// The nested select gives priority to the ctx.Done() signal, so we always exit early if needed// Without it, a single select just has no priority, so a new value could be processed even if the context has been canceledcase<-ctx.Done():
break outer
default:
select {
casechannel<-id:
case<-ctx.Done(): // always check ctx.Done() to avoid leaking the goroutinebreak outer
}
}
}
close(channel)
wg.Wait()

That's a lot of code right there for a simple work pool! We even had to make it collapsable to avoid disrupting the reading flow. Admittedly, most of the complexity comes from cancellation handling, but you don't want to go around leaking your goroutines. Now let's see how the same thing is done with JPipe:

ids:= []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
pipeline:=jpipe.New(ctx)
<-jpipe.FromSlice(pipeline, ids).
ForEach(expensiveIOOperation, jpipe.Concurrent(5))

Complex pipelines

The above is a simple work pool and admittedly the most common use case you'll find for concurrency. But JPipe allows you to build more complex pipelines with its catalog of operators. Imagine the following processing pipeline:

graph LR;
F1["Transaction feed 1"]-->M["Merge"];
F2["Transaction feed 2"]-->M;
M-->F["Filter transactions\nover 50 EUR"]
F-->T["Take first 20\ntransactions"]
T-->E["Enhance transactions\nwith external data\n(IO-bound,\nconcurrency 5,\nFIFO)"]
E-->B["Batch transactions\n(batch size 3,\ntimeout 5 sec)"]
B-->K["Send batches\nto Kafka"]
Loading

The JPipe implementation would be:

pipeline:=jpipe.New(ctx)
feed1:=jpipe.FromGoChannel(pipeline, getTransactionsFromFeed("feed1"))
feed2:=jpipe.FromGoChannel(pipeline, getTransactionsFromFeed("feed2"))
txs:=jpipe.Merge(feed1, feed2).
Filter(func(ftFeedTransaction) bool { returnft.Amount>50 }).
Take(20)
enhancedTxs:=jpipe.Map(txs, enhanceTransaction, jpipe.Concurrent(4), jpipe.Ordered(10))
<-jpipe.Batch(enhancedTxs, 3, 5*time.Second).
ForEach(sendTransactionBatchToKafka)

Documentation

You can find much more details on our official documentation. Some useful links in there:

You can also check the Go.Dev reference.

Similar projects

  • RxGo: Probably the best alternative out there, with lots of operators. It hasn't seen action in two years though, so it hasn't adopted generics: the API still deals with interface{}. A potential drawback for Go developers is its use of the reactive model, which is an abstraction very different from Go channels, and requires understanding of concepts like observer, observable, backpressure strategies, etc. If you'd rather stay within Go's channel abstraction, JPipe may be simpler to use, but RxGo may be your preferred choice if you have a background with the reactive model from other languages.
  • pipeline: An implementation of the pipeline pattern, but with a limited set of operators. It hasn't adopted generics yet either.
  • parapipe: A very simple pipeline implementation, but it has a single Pipe operator and does not support very complex pipelines. It hasn't adopted generics yet either.
  • ordered-concurrently: An implementation of ordered concurrency. It doesn't try to be a complete pipeline pattern library though, and just focuses on that feature.

License

JPipe is open-source software released under the MIT License.

About

Concurrent pipelines for Go

Topics

Resources

Security policy

Stars

21 stars

Watchers

3 watching

Forks

Releases

Packages

Used by

Contributors

Languages