airbyte-prefect
Welcome!
airbyte-prefect is a collection of prebuilt Prefect tasks and flows that can be used to quickly construct Prefect flows to interact with Airbyte.
Note:
airbyte-prefectis a port ofprefect-airbyte, the original Prefect 2 collection by PrefectHQ, updated to work with Prefect 3. Credit for the original design and implementation goes to its authors and contributors.
Getting Started
Python setup
Requires an installation of Python 3.10+
We recommend using a Python virtual environment manager such as pipenv, conda or virtualenv.
These tasks are designed to work with Prefect 3.2+. For more information about how to use Prefect, please refer to the Prefect documentation.
Airbyte setup
See the airbyte documention on how to get your own instance.
Installation
Install airbyte-prefect
pip install airbyte-prefect
For available blocks and their setup instructions, see the documentation in this repository.
Examples
Create an AirbyteServer block and save it
from airbyte_prefect.server import AirbyteServer
# running airbyte locally at http://localhost:8000 with default auth
local_airbyte_server = AirbyteServer()
# running airbyte remotely at http://<someIP>:<somePort> as user `Marvin`
remote_airbyte_server = AirbyteServer(
username="Marvin",
password="DontPanic42",
server_host="42.42.42.42",
server_port="4242"
)
local_airbyte_server.save("my-local-airbyte-server")
remote_airbyte_server.save("my-remote-airbyte-server")
Trigger a defined connection sync
from prefect import flow
from airbyte_prefect.server import AirbyteServer
from airbyte_prefect.connections import AirbyteConnection
from airbyte_prefect.flows import run_connection_sync
server = AirbyteServer(server_host="localhost", server_port=8000)
connection = AirbyteConnection(
airbyte_server=server,
connection_id="e1b2078f-882a-4f50-9942-cfe34b2d825b",
status_updates=True,
)
@flow
def airbyte_syncs():
# do some setup
sync_result = run_connection_sync(
airbyte_connection=connection,
)
# do some other things, like trigger DBT based on number of records synced
print(f'Number of Records Synced: {sync_result.records_synced}')
❯ python airbyte_syncs.py
03:46:03 | prefect.engine - Created flow run 'thick-seahorse' for flow 'example_trigger_sync_flow'
03:46:03 | Flow run 'thick-seahorse' - Using task runner 'ConcurrentTaskRunner'
03:46:03 | Flow run 'thick-seahorse' - Created task run 'trigger_sync-35f0e9c2-0' for task 'trigger_sync'
03:46:03 | prefect - trigger airbyte connection: e1b2078f-882a-4f50-9942-cfe34b2d825b, poll interval 3 seconds
03:46:03 | prefect - pending
03:46:06 | prefect - running
03:46:09 | prefect - running
03:46:12 | prefect - running
03:46:16 | prefect - running
03:46:19 | prefect - running
03:46:22 | prefect - Job 26 succeeded.
03:46:22 | Task run 'trigger_sync-35f0e9c2-0' - Finished in state Completed(None)
03:46:22 | Flow run 'thick-seahorse' - Finished in state Completed('All states completed.')
Export an Airbyte instance's configuration
NOTE: The API endpoint corresponding to this task is no longer supported by open-source Airbyte versions as of v0.40.7. Check out the Octavia CLI docs for more info.
import gzip
from prefect import flow, task
from airbyte_prefect.configuration import export_configuration
from airbyte_prefect.server import AirbyteServer
@task
def zip_and_write_somewhere(
airbyte_config: bytearray,
somewhere: str,
):
with gzip.open(somewhere, 'wb') as f:
f.write(airbyte_config)
@flow
def example_export_configuration_flow(filepath: str):
# Run other tasks and subflows here
airbyte_config = export_configuration(
airbyte_server=AirbyteServer.load("my-airbyte-server-block")
)
zip_and_write_somewhere(
somewhere=filepath,
airbyte_config=airbyte_config
)
if __name__ == "__main__":
example_export_configuration_flow('*://**/my_destination.gz')
Use with_options to customize options on any existing task or flow
from prefect import flow
from airbyte_prefect.connections import AirbyteConnection
from airbyte_prefect.flows import run_connection_sync
custom_run_connection_sync = run_connection_sync.with_options(
name="Custom Airbyte Sync Flow",
retries=2,
retry_delay_seconds=10,
)
@flow
def some_airbyte_flow():
custom_run_connection_sync(
airbyte_connection=AirbyteConnection.load("my-airbyte-connection-block")
)
some_airbyte_flow()
For more tips on how to use tasks and flows in a Collection, check out Using Collections!
Resources
If you encounter and bugs while using airbyte-prefect, feel free to open an issue in the airbyte-prefect repository.
If you have any questions or issues while using airbyte-prefect, you can find help in either the Prefect Discourse forum or the Prefect Slack community
Feel free to star or watch airbyte-prefect for updates too!
Acknowledgements
This project is derived from prefect-airbyte by PrefectHQ,
released under the Apache 2.0 License.
Contribute
If you'd like to help contribute to fix an issue or add a feature to airbyte-prefect, please propose changes through a pull request from a fork of the repository.
Contribution Steps:
- Fork the repository
- Clone the forked repository
- Install the repository and its dependencies:
pip install -e ".[dev]"
- Make desired changes.
- Add tests.
- Insert an entry to CHANGELOG.md
- Install
pre-committo perform quality checks prior to commit:
pre-commit install
git commit,git push, and create a pull request.
Metadata
Release files for airbyte-prefect 1.1.0
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| airbyte_prefect-1.1.0.tar.gz | 39.1 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| airbyte_prefect-1.1.0-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 58.9 kB
Release files / airbyte_prefect-1.1.0.tar.gz
| Download URL | airbyte_prefect-1.1.0.tar.gz |
|---|---|
| Size | 39.1 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
a7826616c2638abb6bd62a4caae9d0442c9a595c18dcb23b9ec736c984f09917
|
|
BLAKE2b-256 checksum How to use checksums |
e63f1c68798dbb8d754d3c3dd59f777d984507ef1ff797351544afba3842ea7a
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/6.2.0 CPython/3.12.0
|
Release files / airbyte_prefect-1.1.0-py3-none-any.whl
| Download URL | airbyte_prefect-1.1.0-py3-none-any.whl |
|---|---|
| Size | 19.7 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
90c1d6188c0959d0ebb9377d9744aa27c1fefd3cf3029b2efda1112eced9c68c
|
|
BLAKE2b-256 checksum How to use checksums |
b7170f8adcbcd80d19f6dc72786366e9af4d16f0de83e6a1a725ed08792eb12b
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/6.2.0 CPython/3.12.0
|