Skip to main content

airbyte-prefect

PyPI

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.

📚 Documentation — API reference for every task, flow and block, generated from the source.

Note: airbyte-prefect is a port of prefect-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 Blocks Catalog.

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

The API reference is published at haybankz.github.io/airbyte-prefect, built from the docstrings in this repository on each release.

If you encounter any 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:

  1. Fork the repository
  2. Clone the forked repository
  3. Install the repository and its dependencies:
pip install -e ".[dev]"
  1. Make desired changes.
  2. Add tests.
  3. Insert an entry to CHANGELOG.md
  4. Install pre-commit to perform quality checks prior to commit:
 pre-commit install
  1. git commit, git push, and create a pull request.

Metadata

Release files for airbyte-prefect 1.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 airbyte-prefect 1.2.0
File Size Uploaded
airbyte_prefect-1.2.0.tar.gz 40.2 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for airbyte-prefect 1.2.0
File Interpreter ABI Platform
airbyte_prefect-1.2.0-py3-none-any.whl Python 3 none any Details

Total release size: 60.4 kB

Release files / airbyte_prefect-1.2.0.tar.gz

Download URL airbyte_prefect-1.2.0.tar.gz
Size 40.2 kB
Tags Source
SHA-256 checksum
How to use checksums
a48d88d299eabd845d3938a3133f1dde4824838af4b235dcb8cf5d989b006a95
BLAKE2b-256 checksum
How to use checksums
a898de6598e44941b2242939658ba73661c5787ecf241989f571651c73d020f1
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Sep 23, 2026.

Transparency log

Release files / airbyte_prefect-1.2.0-py3-none-any.whl

Download URL airbyte_prefect-1.2.0-py3-none-any.whl
Size 20.2 kB
Tags Python 3
SHA-256 checksum
How to use checksums
16ed55fab9c274b88082cab1bc58229ca64f862f4f0611d6fe977e7d127cb9a0
BLAKE2b-256 checksum
How to use checksums
070ff830cd79eb7db5e54259198e98c8bff4e7786052c98b387caca36e1e69d3
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Sep 23, 2026.

Transparency log

Release history Release notifications | RSS feed

This release

1.2.0 This release

2 release files

1.1.0

2 release files

1.0.1

2 release files

1.0.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