A minimal functional pipeline for Python.
Project description
miniplumber
A minimal functional pipeline for Python.
from miniplumber import pipe, sort, field
scores > pipe // field("score") / sort(reverse=True)
words = sentences > pipe // str.split / flatten // str.lower @ str.isalpha / unique / sort()
Quick reference
value > pipe fire pipeline, return raw value
/ func pass whole value to func
/ pipeline compose two pipelines sequentially
// func map over each element
@ func filter — keep where func(x) is truthy
+ pipeline fork — always inside parens: / (a + b) /
All of / // @ + share the same precedence — left-to-right, no exceptions.
> is lower precedence — pipeline always builds fully before firing.
Philosophy
Every data pipeline is a combination of three operations:
// map — transform each element (same number of inputs in and out)
@ filter — select elements (less or equal)
/ pass — send the whole value to a function as-is
/ does not imply a single output. It passes whatever you have — a list, a dict, a string — to the next function whole. The function decides what comes out:
pipe / sorted # list → list (reordered)
pipe / len # list → int (counted)
pipe / " ".join # list → str (joined)
miniplumber gives these three operations clean operator syntax and one rule: all operators share the same precedence, always left-to-right. No brackets to manage precedence. No surprises.
Everything else — flatten, sort, group, twist — is a plain function you pass with / or //. The library is the glue, not the logic. Any function that takes one value and returns one value is a valid step. No wrappers, no base classes, no registration.
The pipeline is lazy. Nothing executes until > fires it. This means pipelines are values — name them, compose them, pass them around, reuse them anywhere.
Operators
/ — pass
Pass the whole value to a function:
["hello", "world"] > pipe / len # → 2
["hello", "world"] > pipe / " ".join # → "hello world"
["hello", "world"] > pipe / sorted # → ["hello", "world"]
Compose two named pipelines sequentially:
clean = pipe // str.strip // str.lower
shout = pipe // str.upper
process = clean / shout
// — map
Apply a function to each element. Always same length in, same length out:
[1, 2, 3] > pipe // double # → [2, 4, 6]
["hello", "world"] > pipe // len # → [5, 5]
["hello", "world"] > pipe // str.upper # → ["HELLO", "WORLD"]
// is polymorphic — it works on dicts too, mapping over values and preserving keys:
{"a": 1, "b": 2, "c": 3} > pipe // double
# → {"a": 2, "b": 4, "c": 6}
@ — filter
Keep elements where func(x) is truthy:
[3, 0, 1, -1] > pipe @ bool # → [3, 1]
words > pipe @ having(pos="noun")# → [only nouns]
@ also works on scalars — returns the value if truthy, None if not:
3 > pipe @ bool # → 3
0 > pipe @ bool # → None
+ — fork
Split one input into parallel pipelines. Always wrap in parentheses:
import statistics
data > pipe / (
pipe / statistics.mean +
pipe / statistics.median +
pipe / statistics.stdev
)
# → [mean, median, stdev]
Fork then merge — the step after ) receives the list of branch results:
"photo.jpg" > load / preproc / (edges + blurred) / np.hstack / save("compare.jpg")
> — fire
Fires the pipeline and returns the raw value:
result = data > pipe // str.upper / " ".join
Named pipelines
Pipelines are values. Name them, reuse them, compose them with /:
tokenize = pipe // str.split / flatten
clean = pipe // str.strip // str.lower
join = pipe / " ".join
process = tokenize / clean / join
[" Hello World "] > process # → "hello world"
[" FOO BAR BAZ "] > process # → "foo bar baz"
Test each piece independently. Compose freely.
Writing steps
Any def function works as a pipeline step:
def remove_stopwords(words):
stopwords = {"the", "a", "an"}
return [w for w in words if w not in stopwords]
sentences > pipe // str.split / flatten / remove_stopwords
When a step needs configuration, use a closure. The outer call captures configuration at build time; the inner function receives the value at fire time:
def blur(sigma):
def _blur(img):
k = max(3, int(6 * sigma) | 1)
return cv2.GaussianBlur(img, (k, k), sigma)
return _blur
def above(threshold):
return lambda x: x > threshold
pipe / blur(5.0) # pre-configured step
pipe @ above(100) # pre-configured predicate
Utils
Sequence
pipe / flatten # one level: [[1,2],[3,4]] → [1,2,3,4]
pipe / flatten_deep # any depth: [1,[2,[3]]] → [1,2,3]
pipe / sort() # alphabetical
pipe / sort(key=len, reverse=True) # by length descending
pipe / unique # deduplicate preserving order
pipe / take(3) # [:3] first 3 elements
pipe / take(3, None) # [3:] skip first 3
pipe / take(1, 5) # [1:5] elements 1 to 4
pipe / take(None, None, -1) # [::-1] reverse
pipe / chunk(2) # [[1,2],[3,4],[5]]
pipe / window(2) # [(1,2),(2,3),(3,4)]
pipe / group(key) # → dict grouped by key function
Dict and object access
users > pipe // field("name") # extract dict key
users > pipe // field("age", default=0) # with fallback
objects > pipe // attr("created_at") # extract object attribute
Fork utilities
pipe / twist(2) # branch-major → item-major
pipe / named(["micro", "meso", "macro"]) # flat list → named dict
Predicates for @
pipe @ instance(str) # isinstance check
pipe @ matching("^[A-Z]") # substring or regex match
pipe @ having(status="active") # dict key/value match
Debug
pipe / tap(print) # print and pass through
pipe / tap(lambda x: print("after:", x)) # with label
pipe / tap(log_to_file) # any side effect
Patterns
Growing state with dicts
For pipelines where each step enriches a shared state, pass a dict and grow it at each step. The convention: every step returns {**state, "new_key": value}.
def load(state):
return {**state, "img": cv2.imread(state["path"])}
def preprocess(state):
gray = cv2.cvtColor(state["img"], cv2.COLOR_BGR2GRAY)
return {**state, "gray": gray}
def segment(state):
return {**state, "letters": find_letters(state["gray"])}
{"path": "image.jpg"} > pipe / load / preprocess / segment
# → {"path": ..., "img": ..., "gray": ..., "letters": [...]}
Each step reads what it needs, adds what it produces, passes everything forward. Nothing is lost. For type safety and autocomplete, use a dataclass with dataclasses.replace instead — same pattern, typed.
Fork → twist → named
A fork produces results in branch-major order: all results of branch A, then all of branch B. twist reorders them to item-major: all branch results for item 0, then item 1, and so on. named then restores meaning to the flat list.
energy_before = pipe // sobel_energy
energy_after = pipe // (probe_blur / sobel_energy)
bands > pipe / dict.values / list
/ (energy_before + energy_after)
/ twist(2) # [[e0,pe0], [e1,pe1], [e2,pe2]]
// divide # [ratio0, ratio1, ratio2]
/ named(["micro", "meso", "macro"]) # {"micro": r0, "meso": r1, "macro": r2}
The physical model lives in the operator structure — two energy measurements per band, one ratio each.
Preserving a value across a transformation
Use pipe as the identity branch in a fork to carry a value forward untouched:
result = data > pipe / step1 / (transform + pipe) / merge
# ↑ value after step1, unchanged
Error handling
Handle errors where they belong — inside the function that knows what failure means, or around the whole pipeline:
# inside the function — owns its own error contract
def parse(x):
try:
return int(x)
except ValueError:
return 0
# around the pipeline — one place for the whole flow
try:
result = data > pipe / step1 / step2 / step3
except ValueError as e:
result = fallback
Installation
pip install miniplumber
from miniplumber import pipe # minimum
from miniplumber import pipe, flatten, sort, field, tap # with utilities
from miniplumber import * # everything
Package structure
miniplumber/
__init__.py # re-exports everything
core.py # Pipeline class and pipe sentinel — zero dependencies
utils.py # flatten, sort, twist, named, field, and friends
core.py is self-contained. Copy just that file if you want the pipeline with no utilities.
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file miniplumber-0.2.0.tar.gz.
File metadata
- Download URL: miniplumber-0.2.0.tar.gz
- Upload date:
- Size: 14.0 kB
- Tags: Source
- Uploaded using Trusted Publishing? Yes
- Uploaded via: twine/6.1.0 CPython/3.13.7
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
db76ab1d1d57c884f040bc49bb963e155b742afcc1f4b50ef004d99a8f935de5
|
|
| MD5 |
9aa4b9fea22c434d51b3cfbfbffe5df2
|
|
| BLAKE2b-256 |
3b3447e9f15710c4b91d1ab42b2efa1f9478de9345acc29f6a09e218e09a03c0
|
Provenance
The following attestation bundles were made for miniplumber-0.2.0.tar.gz:
Publisher:
python-publish.yml on aleixfort/miniplumber
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
miniplumber-0.2.0.tar.gz -
Subject digest:
db76ab1d1d57c884f040bc49bb963e155b742afcc1f4b50ef004d99a8f935de5 - Sigstore transparency entry: 1186815989
- Sigstore integration time:
-
Permalink:
aleixfort/miniplumber@0178d6009b4c3f546906be4bc85aec1ed18ec490 -
Branch / Tag:
refs/tags/v0.2.0 - Owner: https://github.com/aleixfort
-
Access:
public
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
python-publish.yml@0178d6009b4c3f546906be4bc85aec1ed18ec490 -
Trigger Event:
release
-
Statement type:
File details
Details for the file miniplumber-0.2.0-py3-none-any.whl.
File metadata
- Download URL: miniplumber-0.2.0-py3-none-any.whl
- Upload date:
- Size: 10.5 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? Yes
- Uploaded via: twine/6.1.0 CPython/3.13.7
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
de65e42f087c5c4d082506c03dcf061186c2befc5bb633ea120b8b624e71c709
|
|
| MD5 |
3cfa6d8e807def28fb0142bc8ac99d4f
|
|
| BLAKE2b-256 |
37fdb494d464b91e7fd3d8643a4112b7b175255bdbb06d5e25fe773053a44e40
|
Provenance
The following attestation bundles were made for miniplumber-0.2.0-py3-none-any.whl:
Publisher:
python-publish.yml on aleixfort/miniplumber
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
miniplumber-0.2.0-py3-none-any.whl -
Subject digest:
de65e42f087c5c4d082506c03dcf061186c2befc5bb633ea120b8b624e71c709 - Sigstore transparency entry: 1186815992
- Sigstore integration time:
-
Permalink:
aleixfort/miniplumber@0178d6009b4c3f546906be4bc85aec1ed18ec490 -
Branch / Tag:
refs/tags/v0.2.0 - Owner: https://github.com/aleixfort
-
Access:
public
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
python-publish.yml@0178d6009b4c3f546906be4bc85aec1ed18ec490 -
Trigger Event:
release
-
Statement type: