Skip to main content

dbpipe

dbpipe is a lightweight and simple way to manage data pipelines.

graph LR
A(Endpoints)-->B(Pipes)
B-->C(Jobs)
D(Schedules)-->C
C-.->E(Clusters)

Creating Endpoints

from dbpipe import EndPoint


facebook = EndPoint('Facebook','API','https://facebook.com/Posts')
facebook
{'name': 'Facebook', 'type': 'API', 'location': 'https://facebook.com/Posts'}
facebook.save()
posttable = EndPoint('DW.Facebook.Posts','Database','ServerName')
posttable
{'name': 'DW.Facebook.Posts', 'type': 'Database', 'location': 'ServerName'}
posttable.save()

Creating a Pipe

from dbpipe import Pipe


pipe = Pipe(
        name='DW',
        sources=[facebook],
        destination=posttable,
        processfile="Test.py"
    )

pipe
{'name': 'DW', 'sources': [{'name': 'Facebook', 'type': 'API', 'location': 'https://facebook.com/Posts'}], 'destination': {'name': 'DW.Facebook.Posts', 'type': 'Database', 'location': 'ServerName'}, 'logfile': None, 'processfile': 'Test.py'}
pipe.to_dict()
{'name': 'DW',
 'sources': [{'name': 'Facebook',
   'type': 'API',
   'location': 'https://facebook.com/Posts'}],
 'destination': {'name': 'DW.Facebook.Posts',
  'type': 'Database',
  'location': 'ServerName'},
 'logfile': None,
 'processfile': 'Test.py'}
pipe.save()

Creating a Schedule

from dbpipe import Schedule

schedule = Schedule(frequency="Daily", start_time="8:00AM")

schedule
{'frequency': 'Daily', 'start_time': '8:00AM', 'end_time': None, 'time_zone': 'UTC'}
schedule.to_dict()
{'frequency': 'Daily',
 'start_time': '8:00AM',
 'end_time': None,
 'time_zone': 'UTC'}

Creating a Pipe Cluster

from dbpipe.core.pipes import Cluster


clstr = Cluster([pipe,pipe])
clstr
[{'name': 'DW', 'sources': ['AdSpend', 'SocialStats'], 'destination': 'DW', 'logfile': None, 'processfile': 'Test.py'}, {'name': 'DW', 'sources': ['AdSpend', 'SocialStats'], 'destination': 'DW', 'logfile': None, 'processfile': 'Test.py'}]

Creating a Job

from dbpipe import Job


job = Job('My Job',schedule=schedule,jobs=clstr)
job
{'name': 'My Job', 'schedule': {'frequency': 'Daily', 'start_time': '8:00AM', 'end_time': None, 'time_zone': 'UTC'}, 'jobs': [{'name': 'DW', 'sources': ['AdSpend', 'SocialStats'], 'destination': 'DW', 'logfile': None, 'processfile': 'Test.py'}, {'name': 'DW', 'sources': ['AdSpend', 'SocialStats'], 'destination': 'DW', 'logfile': None, 'processfile': 'Test.py'}]}
job.save()

Reading a Pipe

from dbpipe import read_pipe


pipe = read_pipe('pipes/DW.json')
pipe
{'name': 'DW', 'sources': ['AdSpend', 'SocialStats'], 'destination': 'DW', 'logfile': None, 'processfile': 'Test.py'}
pipe.to_dict()
{'name': 'DW',
 'sources': ['AdSpend', 'SocialStats'],
 'destination': 'DW',
 'logfile': None,
 'processfile': 'Test.py'}

Reading a Job

from dbpipe import read_job

job = read_job('jobs/My Job.json')
job
{'name': 'My Job', 'schedule': {'frequency': 'Daily', 'start_time': '8:00AM', 'end_time': None, 'time_zone': 'UTC'}, 'jobs': [{'name': 'DW', 'sources': ['AdSpend', 'SocialStats'], 'destination': 'DW', 'logfile': None, 'processfile': 'Test.py'}, {'name': 'DW', 'sources': ['AdSpend', 'SocialStats'], 'destination': 'DW', 'logfile': None, 'processfile': 'Test.py'}]}
job.to_dict()
{'name': 'My Job',
 'schedule': {'frequency': 'Daily',
  'start_time': '8:00AM',
  'end_time': None,
  'time_zone': 'UTC'},
 'jobs': [{'name': 'DW',
   'sources': ['AdSpend', 'SocialStats'],
   'destination': 'DW',
   'logfile': None,
   'processfile': 'Test.py'},
  {'name': 'DW',
   'sources': ['AdSpend', 'SocialStats'],
   'destination': 'DW',
   'logfile': None,
   'processfile': 'Test.py'}]}

Lineage

from dbpipe.lineage.mermaid import generate_mermaid_markdown_file


generate_mermaid_markdown_file('pipes','test.md')

Metadata

Release files for dbpipe 0.2.0

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for dbpipe 0.2.0
File Size Uploaded
dbpipe-0.2.0.tar.gz 6.0 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for dbpipe 0.2.0
File Interpreter ABI Platform
dbpipe-0.2.0-py3-none-any.whl Python 3 none any Details

Total release size: 12.1 kB

Release files / dbpipe-0.2.0.tar.gz

Download URL dbpipe-0.2.0.tar.gz
Size 6.0 kB
Tags Source
SHA-256 checksum
How to use checksums
159a801a8ae20e35d433657bd4a66bb8607eb53688fc098704984cdf2aef21b5
BLAKE2b-256 checksum
How to use checksums
415cee55a4d5d90232cbed1cd3486a28ace842b9a302219788a69bd44c1c7ce3
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/5.0.0 CPython/3.11.4

Release files / dbpipe-0.2.0-py3-none-any.whl

Download URL dbpipe-0.2.0-py3-none-any.whl
Size 6.1 kB
Tags Python 3
SHA-256 checksum
How to use checksums
c679398b03330d7dc8f4d2d4e1d371d38db27ffddcbc9e12b3e2b40b8ffe39e0
BLAKE2b-256 checksum
How to use checksums
b1a35f348e275f9870c451b41dc5e8227b9944e1d596a65b5e00d2f1ac25ee7a
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/5.0.0 CPython/3.11.4

Release history Release notifications | RSS feed

This release

0.2.0 This release

2 release files

0.1.7

2 release files

0.1.6

2 release files

0.1.5

2 release files

0.1.4

2 release files

0.1.3

2 release files

0.1.2

2 release files

0.1.1

2 release files

0.1.0

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page