Skip to main content

Dbnd Airflow Operator

This plugin was written to provide an explicit way of declaratively passing messages between two airflow operators.

This plugin was inspired by AIP-31. Essentially, this plugin connects between dbnd's implementation of tasks and pipelines to airflow operators.

This implementation uses XCom communication and XCom templates to transfer said messages. This plugin is fully functional, however as soon as AIP-31 is implemented it will support all edge-cases.

Fully tested on airflow 1.10.X.

Code Example

Here is an example of how we achieve our goal:

import logging
from typing import Tuple
from datetime import timedelta, datetime
from airflow import DAG
from airflow.utils.dates import days_ago
from airflow.operators.python_operator import PythonOperator
from dbnd import task

# Define arguments that we will pass to our DAG
default_args = {
    "owner": "airflow",
    "depends_on_past": False,
    "start_date": days_ago(2),
    "retries": 1,
    "retry_delay": timedelta(seconds=10),
}
@task
def my_task(p_int=3, p_str="check", p_int_with_default=0) -> str:
    logging.info("I am running")
    return "success"


@task
def my_multiple_outputs(p_str="some_string") -> Tuple[int, str]:
    return (1, p_str + "_extra_postfix")


def some_python_function(input_path, output_path):
    logging.error("I am running")
    input_value = open(input_path, "r").read()
    with open(output_path, "w") as output_file:
        output_file.write(input_value)
        output_file.write("\n\n")
        output_file.write(str(datetime.now().strftime("%Y-%m-%dT%H:%M:%S")))
    return "success"

# Define DAG context
with DAG(dag_id="dbnd_operators", default_args=default_args) as dag_operators:
    # t1, t2 and t3 are examples of tasks created by instantiating operators
    # All tasks and operators created under this DAG context will be collected as a part of this DAG
    t1 = my_task(2)
    t2, t3 = my_multiple_outputs(t1)
    python_op = PythonOperator(
        task_id="some_python_function",
        python_callable=some_python_function,
        op_kwargs={"input_path": t3, "output_path": "/tmp/output.txt"},
    )
    """
    t3.op describes the operator used to execute my_multiple_outputs
    This call defines the some_python_function task's operator as dependent upon t3's operator
    """
    python_op.set_upstream(t3.op)

As you can see, messages are passed explicitly between all three tasks:

  • t1, the result of the first task is passed to the next task my_multiple_outputs
  • t2 and t3 represent the results of my_multiple_outputs
  • some_python_function is wrapped with an operator
  • The new python operator is defined as dependent upon t3's execution (downstream) - explicitly.

Note: If you run a function marked with the @task decorator without a DAG context, and without using the dbnd library to run it - it will execute absolutely normally!

Using this method to pass arguments between tasks not only improves developer user-experience, but also allows for pipeline execution support for many use-cases. It does not break currently existing DAGs.

Using dbnd_config

Let's look at the example again, but change the default_args defined at the very top:

default_args = {
    "owner": "airflow",
    "depends_on_past": False,
    "start_date": days_ago(2),
    "retries": 1,
    "retry_delay": timedelta(minutes=5),
    'dbnd_config': {
        "my_task.p_int_with_default": 4
    }
}

Added a new key-value pair to the arguments called dbnd_config

dbnd_config is expected to define a dictionary of configuration settings that you can pass to your tasks. For example, the dbnd_config in this code section defines that the int parameter p_int_with_default passed to my_task will be overridden and changed to 4 from the default value 0.

To see further possibilities of changing configuration settings, see our documentation

Release files for dbnd-airflow 1.0.34.1

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

Source distribution (sdist)

Source distribution for dbnd-airflow 1.0.34.1
File Size Uploaded
dbnd_airflow-1.0.34.1.tar.gz 39.3 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for dbnd-airflow 1.0.34.1
File Interpreter ABI Platform
dbnd_airflow-1.0.34.1-py2.py3-none-any.whl Python 3, Python 2 none any Details

Total release size: 86.7 kB

Release files / dbnd_airflow-1.0.34.1.tar.gz

Download URL dbnd_airflow-1.0.34.1.tar.gz
Size 39.3 kB
Tags Source
SHA-256 checksum
How to use checksums
2b8f06595afccf3c53abc1bd5a6c9dcf3bba24e83402ded0a837206497b87fb9
BLAKE2b-256 checksum
How to use checksums
60b779b2044fa2afb0945373ef3f61a231e260c0c9b1e157998794f8d66fb86b
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.2.0 CPython/3.10.9

Release files / dbnd_airflow-1.0.34.1-py2.py3-none-any.whl

Download URL dbnd_airflow-1.0.34.1-py2.py3-none-any.whl
Size 47.3 kB
Tags Python 2 Python 3
SHA-256 checksum
How to use checksums
90e524713a1ca7824ab81353ccf73df54f7c43f8d4394f3757acd0c4d52f340b
BLAKE2b-256 checksum
How to use checksums
a464846d2569bddf824e6be998d8abbb27eacc3a68eacb8ebe9473c15c8f2ba4
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.2.0 CPython/3.10.9

Release history Release notifications | RSS feed

This release

1.0.34.1 This release

2 release files

0.89.6

2 release files

0.89.5

2 release files

0.89.4

2 release files

0.89.3

2 release files

0.89.2

2 release files

0.88.5

2 release files

0.88.4

2 release files

0.88.3

2 release files

0.88.2

2 release files

0.88.1

2 release files

0.87.4

2 release files

0.87.3

2 release files

0.87.2

2 release files

0.87.1

2 release files

0.86.3

2 release files

0.85.9

2 release files

0.85.8

2 release files

0.83.4

2 release files

0.83.3

2 release files

0.83.2

2 release files

0.83.1

2 release files

0.82.9

1 release file

0.82.8

2 release files

0.82.5

2 release files

0.82.4

2 release files

0.82.3

2 release files

0.82.2

2 release files

0.82.1

2 release files

0.81.4

2 release files

0.81.2

2 release files

0.81.1

2 release files

0.80.9

2 release files

0.79.9

2 release files

0.79.8

2 release files

0.79.7

2 release files

0.79.6

2 release files

0.79.4

2 release files

0.79.3

2 release files

0.79.2

2 release files

0.79.1

2 release files

0.78.9

2 release files

0.78.8

2 release files

0.78.7

2 release files

0.78.6

2 release files

0.78.5

2 release files

0.78.4

2 release files

0.78.3

2 release files

0.78.2

2 release files

0.78.1

2 release files

0.76.9

2 release files

0.76.7

2 release files

0.76.6

2 release files

0.76.5

2 release files

0.76.4

2 release files

0.75.4

2 release files

0.75.3

2 release files

0.75.2

2 release files

0.74.9

2 release files

0.74.8

2 release files

0.74.7

2 release files

0.74.6

2 release files

0.74.5

2 release files

0.74.4

2 release files

0.74.3

2 release files

0.74.2

2 release files

0.74.1

2 release files

0.73.9

2 release files

0.73.8

2 release files

0.73.7

2 release files

0.73.6

2 release files

0.73.5

2 release files

0.73.4

2 release files

0.71.7

2 release files

0.71.6

2 release files

0.71.5

2 release files

0.71.4

2 release files

0.71.3

2 release files

0.71.2

2 release files

0.71.1

2 release files

0.70.9

2 release files

0.70.8

2 release files

0.70.7

2 release files

0.70.6

2 release files

0.70.5

2 release files

0.70.4

2 release files

0.70.3

2 release files

0.70.2

2 release files

0.70.1

2 release files

0.69.6

2 release files

0.68.1

2 release files

0.67.4

2 release files

0.67.3

2 release files

0.67.2

2 release files

0.67.1

2 release files

0.66.9

2 release files

0.66.8

2 release files

0.66.7

2 release files

0.66.6

2 release files

0.66.5

2 release files

0.66.4

2 release files

0.66.3

2 release files

0.65.8

2 release files

0.65.7

2 release files

0.64.1

2 release files

0.63.5

2 release files

0.63.4

2 release files

0.63.3

2 release files

0.63.2

2 release files

0.63.1

2 release files

0.61.3

2 release files

0.61.2

2 release files

0.61.1

2 release files

0.60.3

2 release files

0.60.2

2 release files

0.60.1

2 release files

0.59.3

2 release files

0.58.6

2 release files

0.58.5

2 release files

0.58.4

2 release files

0.58.3

2 release files

0.57.3

2 release files

0.57.2

2 release files

0.57.1

2 release files

0.56.7

2 release files

0.56.6

2 release files

0.56.5

2 release files

0.56.4

2 release files

0.56.3

2 release files

0.56.2

2 release files

0.54.2

2 release files

0.54.1

2 release files

0.53.8

2 release files

0.53.7

2 release files

0.53.6

2 release files

0.53.5

2 release files

0.53.4

2 release files

0.53.3

2 release files

0.52.1

2 release files

0.51.4

2 release files

0.51.3

2 release files

0.51.2

2 release files

0.51.1

2 release files

0.51.0

2 release files

0.50.6

2 release files

0.50.5

2 release files

0.50.4

2 release files

0.49.6

2 release files

0.49.4

2 release files

0.49.3

2 release files

0.49.2

2 release files

0.49.0

2 release files

0.48.7

2 release files

0.48.6

2 release files

0.48.0

2 release files

0.47.4

2 release files

0.47.3

2 release files

0.47.2

2 release files

0.47.1

2 release files

0.47.0

2 release files

0.46.5

2 release files

0.46.4

2 release files

0.45.9

2 release files

0.45.8

2 release files

0.45.7

2 release files

0.45.6

2 release files

0.45.5

2 release files

0.45.4

2 release files

0.45.3

2 release files

0.45.2

2 release files

0.45.1

2 release files

0.45.0

2 release files

0.44.7

2 release files

0.44.6

2 release files

0.44.5

2 release files

0.44.4

2 release files

0.44.3

2 release files

0.44.2

2 release files

0.43.4

2 release files

0.43.3

2 release files

0.43.2

2 release files

0.43.1

2 release files

0.42.5

2 release files

0.42.4

2 release files

0.42.3

2 release files

0.42.2

2 release files

0.42.1

2 release files

0.41.3

2 release files

0.41.2

2 release files

0.41.1

2 release files

0.41.0

2 release files

0.40.3

2 release files

0.40.2

2 release files

0.40.1

2 release files

0.40.0

2 release files

0.39.1

2 release files

0.39.0

2 release files

0.38.7

2 release files

0.38.6

2 release files

0.38.5

2 release files

0.38.4

2 release files

0.38.3

2 release files

0.38.2

2 release files

0.38.1

2 release files

0.38.0

2 release files

0.37.5

2 release files

0.37.4

2 release files

0.37.1

2 release files

0.37.0

2 release files

0.36.5

2 release files

0.36.4

2 release files

0.36.3

2 release files

0.36.2

2 release files

0.36.1

2 release files

0.36.0

2 release files

0.35.2

2 release files

0.34.7

2 release files

0.34.6

2 release files

0.34.5

2 release files

0.34.4

2 release files

0.34.3

2 release files

0.34.2

2 release files

0.33.1

2 release files

0.33.0

2 release files

0.32.6

2 release files

0.32.5

2 release files

0.32.4

2 release files

0.32.3

2 release files

0.32.2

2 release files

0.32.1

2 release files

0.32.0

2 release files

0.31.2

2 release files

0.31.1

2 release files

0.30.8

2 release files

0.30.6

2 release files

0.30.5

2 release files

0.30.4

2 release files

0.30.3

2 release files

0.30.2

2 release files

0.30.1

2 release files

0.30.0

2 release files

0.29.6

2 release files

0.29.5

2 release files

0.29.3

2 release files

0.29.2

2 release files

0.29.1

2 release files

0.29.0

2 release files

0.28.5

2 release files

0.28.3

2 release files

0.28.0

2 release files

0.27.1

2 release files

0.27.0

2 release files

0.26.4

2 release files

0.26.3

2 release files

0.26.2

2 release files

0.26.1

2 release files

0.26.0

2 release files

0.25.8

2 release files

0.25.7

2 release files

0.25.6

2 release files

0.25.5

2 release files

0.25.4

2 release files

0.25.3

2 release files

0.25.2

2 release files

0.25.1

2 release files

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