Streaming, declarative pipelined data processing library for Python
421
Conduit is a streaming data pipeline framework that lets you build complex data processing workflows using simple YAML configurations instead of brittle scripts.
Replace this:
# Fragile, hard-to-modify script
data = fetch_api_data()
filtered = [item for item in data if item['status'] == 'active']
processed = [transform(item) for item in filtered]
for item in processed:
print(f"Result: {item}")
With this:
- id: conduit.RestApi
url: "https://api.example.com/data"
- id: conduit.Filter
condition: "input.status == 'active'"
- id: conduit.Transform
template: "{{input | process}}"
- id: conduit.Console
format: "Result: {{input}}"
The easiest way to try or deploy Conduit - no installation required!
# Run a simple hello world pipeline
docker run --rm aniondocker/conduit:latest run examples/hello_world.yaml
Pokemon API downloader - Fetches Pokemon data and downloads official artwork images
# Downloads pokemon images to ./pokemon_images (requires volume mount for file access)
docker run --rm -v $(pwd):/data -w /data aniondocker/conduit:latest run examples/pokemon_api.yaml
# Download 20 Pokemon instead of default 10
docker run --rm -v $(pwd):/data -w /data aniondocker/conduit:latest run examples/pokemon_api.yaml --args limit=20
File processing - Analyzes files in your directory and shows their sizes
# Analyzes *.py files in current directory (requires volume mount to access your files)
docker run --rm -v $(pwd):/data -w /data aniondocker/conduit:latest run examples/file_sizes.yaml
Pokemon evolution chains - Complex workflow showing parallel data fetching and processing
# Fetches Pokemon data and maps out evolution relationships
docker run --rm -v $(pwd):/data -w /data aniondocker/conduit:latest run examples/pokemon_evolution.yaml
Data grouping - Demonstrates grouping and aggregation of data sets
# Groups sample data by category and shows aggregated results
docker run --rm -v $(pwd):/data -w /data aniondocker/conduit:latest run examples/groupby_example.yaml
Data sorting - Shows how to sort data using custom keys
# Sorts sample data by different criteria
docker run --rm -v $(pwd):/data -w /data aniondocker/conduit:latest run examples/sort_example.yaml
Pokemon filtering - Filters Pokemon by specific criteria (type, stats, etc.)
# Filters Pokemon list based on criteria like type or generation
docker run --rm -v $(pwd):/data -w /data aniondocker/conduit:latest run examples/pokemon_filter.yaml
SFTP file operations - List and download files from SFTP servers (uses public test server)
# Lists files from test SFTP server, filters by size, and downloads
docker run --rm -v $(pwd):/data -w /data aniondocker/conduit:latest run examples/sftp_rebex_example.yaml
CSV processing - Download and process CSV files with grouping and aggregation
# Downloads CSV from Google Drive, groups by country, and displays results
docker run --rm aniondocker/conduit:latest run examples/csv_example.yaml
Custom elements - Demonstrates extending Conduit with user-defined processing elements
# Shows how custom elements work alongside built-in ones
docker run --rm -v $(pwd):/data -w /data aniondocker/conduit:latest run examples/custom_element_example.yaml
# Run your own pipeline files
docker run --rm -v $(pwd):/data -w /data aniondocker/conduit:latest run your-pipeline.yaml
# Start the server (accessible on http://localhost:8000)
docker run --rm -p 8000:8000 aniondocker/conduit:latest serve --host 0.0.0.0
# Test the server
curl -X POST http://localhost:8000/run \
-H "Content-Type: application/json" \
-d '{"pipeline": [
{"id": "conduit.Input", "data": [{"name": "Docker"}]},
{"id": "conduit.Console", "format": "Hello from {{input.name}}!"}
]}'
Mount your custom elements directory to extend Conduit:
# Create your custom element
mkdir my-elements
cat > my-elements/MyCustom.py << 'EOF'
from conduit.pipelineElement import PipelineElement
from typing import Iterator
from dataclasses import dataclass
@dataclass
class MyCustomInput:
message: str = "Hello from my custom element!"
class MyCustom(PipelineElement):
def process(self, input: Iterator[MyCustomInput]) -> Iterator[dict]:
for item in input:
yield {"custom_output": f"🚀 {item.message}"}
EOF
# Create package init file
echo "from .MyCustom import MyCustom" > my-elements/__init__.py
# Use your custom element
docker run --rm -v $(pwd)/my-elements:/elements aniondocker/conduit:latest run - << 'EOF'
- id: conduit.Input
data: [{"message": "Docker + Custom Elements!"}]
- id: elements.MyCustom
- id: conduit.Console
format: "{{input.custom_output}}"
EOF
If you prefer to install locally:
pip install git+https://github.com/aniongithub/conduit.git
Create hello.yaml:
- id: conduit.Input
data: [{message: "Hello, Conduit!"}]
- id: conduit.Console
format: "{{input.message}}"
Run it:
conduit-cli run hello.yaml
# Output: Hello, Conduit!
Start the API server:
conduit-cli serve --host 0.0.0.0 --port 8000
Execute pipelines via REST:
curl -X POST http://localhost:8000/run \
-H "Content-Type: application/json" \
-d '{"pipeline": [
{"id": "conduit.Input", "data": [{"name": "World"}]},
{"id": "conduit.Console", "format": "Hello, {{input.name}}!"}
]}'
Note: This approach is not recommended, use docker or devcontainers for the most supported and easiest path
# Find and analyze Python files
- id: conduit.Glob
pattern: "**/*.py"
- id: conduit.FileInfo
- id: conduit.Console
format: "{{input.name}}: {{input.size}} bytes"
# Get Pokemon data with custom limit
- id: conduit.RestApi
url: "https://pokeapi.co/api/v2/pokemon?limit=${limit:-5}"
- id: conduit.JsonQuery
query: ".results[].name"
- id: conduit.Console
format: "Pokemon: {{input}}"
Run with arguments:
conduit-cli pokemon.yaml --args limit=10
# Parallel processing with Fork
- id: conduit.RestApi
url: "https://api.example.com/data"
- id: conduit.Fork
paths:
summary:
- id: conduit.JsonQuery
query: ".metadata.title"
details:
- id: conduit.JsonQuery
query: ".content"
- id: conduit.Filter
condition: "len(input) > 100"
- id: conduit.Console
format: "{{input.summary}}: {{input.details[:50]}}..."
# Convert YAML to API request with yq
yq eval '{"pipeline": ., "args": {"limit": "10"}}' examples/pokemon_evolution.yaml -o=json | \
curl -X POST http://localhost:8000/run -H "Content-Type: application/json" -d @-
| Category | Elements | Purpose |
|---|---|---|
| Input | Input, RestApi, Random, Glob | Data sources and generation |
| Transform | Filter, JsonQuery, Extract, Format | Data processing and extraction |
| Data | CsvReader, GroupBy, Sort | Data parsing and organization |
| Flow | Fork, Iterate, Identity, Empty | Control flow and parallelization |
| Output | Console, DownloadFile | Results and file operations |
| Network | SftpList, SftpDownload | SFTP file listing and transfer |
| System | Cli, FileInfo, Find, Path | System integration |
Need more? Check the full element reference or create custom elements
This repository is meant for development with a VS Code dev container:
from dataclasses import dataclass
from typing import Generator, Iterator
from conduit import PipelineElement
@dataclass
class MyInput:
value: str
class MyElement(PipelineElement):
def process(self, input: Iterator[MyInput]) -> Generator[str, None, None]:
for item in input:
yield f"Processed: {item.value}"
{
"success": true,
"results": ["item1", "item2", "item3"],
"stdout": ["Console output line 1", "Console output line 2"],
"stderr": [],
"stats": {
"duration": 1.23,
"total_items_processed": 3,
"throughput": 2.44,
"element_metrics": [...]
}
}
MIT License - see LICENSE for details.
Content type
Image
Digest
sha256:704723a9f…
Size
469.3 MB
Last updated
11 months ago
docker pull aniondocker/conduit