Skip to main content

Kestra Python Client

This Python client provides functionality to interact with the Kestra server for sending metrics, outputs, and logs, as well as executing/polling flows.

Installation

pip install kestra

Kestra Class

The Kestra class is responsible for sending metrics, outputs, and logs to the Kestra server.

Methods

  • _send(map_: dict): Sends a message to the Kestra server.
  • format(map_: dict) -> str: Formats a message to be sent to the Kestra server.
  • _metrics(name: str, type_: str, value: int, tags: dict | None = None): Sends a metric to the Kestra server.
  • outputs(map_: dict): Sends outputs to the Kestra server.
  • counter(name: str, value: int, tags: dict | None = None): Sends a counter to the Kestra server.
  • timer(name: str, duration: int | Callable, tags: dict | None = None): Sends a timer to the Kestra server.
  • logger() -> Logger: Retrieves the logger for the Kestra server.

Flow Class

The Flow class is used to execute a Kestra flow and optionally wait for its completion. It can also be used to get the status of an execution and the logs of an execution.

Initialization

flow = Flow(
  wait_for_completion=True, # default is True
  poll_interval=1, # seconds. default is 1  
  labels_from_inputs=False, # default is False
  tenant=None # default is None
)

You can also set the hostname and authentication credentials using environment variables:

export KESTRA_HOSTNAME=http://localhost:8080
export KESTRA_USER=admin
export KESTRA_PASSWORD=admin
export KESTRA_API_TOKEN=my_api_token

It is worth noting that the KESTRA_API_TOKEN or KESTRA_USER and KESTRA_PASSWORD need to be used, you do not need all at once. The possible Authentication patterns are:

  1. KESTRA_API_TOKEN
  2. KESTRA_USER and KESTRA_PASSWORD
  3. No Authentication (not recommended for production environments)

Methods

  • _make_request(method: str, url: str, **kwargs) -> requests.Response: Makes a request to the Kestra server with optional authentication and retries.
  • check_status(execution_id: str) -> requests.Response: Checks the status of an execution.
  • get_logs(execution_id: str) -> requests.Response: Retrieves the logs of an execution.
  • execute(namespace: str, flow: str, inputs: dict = None) -> namedtuple: Executes a Kestra flow and optionally waits for its completion. The namedtuple returned is a namedtuple with the following properties:
    • status: The status of the execution.
    • log: The log of the execution.
    • error: The error of the execution.

Usage Examples

  1. Trigger a flow and wait for its completion:

    from kestra import Flow
    flow = Flow()
    flow.execute('mynamespace', 'myflow', {'param': 'value'})
    
  2. Set labels from inputs:

    from kestra import Flow
    flow = Flow(labels_from_inputs=True)
    flow.execute('mynamespace', 'myflow', {'param': 'value'})
    
  3. Pass a text file to an input of type FILE named 'myfile':

    from kestra import Flow
    flow = Flow()
    with open('example.txt', 'rb') as fh:
        flow.execute('mynamespace', 'myflow', {'files': ('myfile', fh, 'text/plain')})
    
  4. Fire and forget:

    from kestra import Flow
    flow = Flow(wait_for_completion=False)
    flow.execute('mynamespace', 'myflow', {'param': 'value'})
    
  5. Overwrite the username and password:

    from kestra import Flow
    flow = Flow()
    flow.user = 'admin'
    flow.password = 'admin'
    flow.execute('mynamespace', 'myflow')
    
  6. Set the hostname, username, and password using environment variables:

    from kestra import Flow
    import os
    
    os.environ["KESTRA_HOSTNAME"] = "http://localhost:8080"
    os.environ["KESTRA_USER"] = "admin"
    os.environ["KESTRA_PASSWORD"] = "admin"
    flow = Flow()
    flow.execute('mynamespace', 'myflow', {'param': 'value'})
    

Error Handling

The client includes retry logic with exponential backoff for certain HTTP status codes, and raises a FailedExponentialBackoff exception if the request fails after multiple retries.

Kestra Class

Logging

The Kestra class provides a logger that formats logs in JSON format, making it easier to integrate with log management systems.

from kestra import Kestra

Kestra.logger().info("Hello, world!")

Outputs

The Kestra class provides a method to send key-value-based outputs to the Kestra server. If you want to output large objects, write them to a file and specify them within the outputFiles property of the Python script task.

Kestra.outputs({"my_output": "my_value"})

Counters

The Kestra class provides a method to send counter metrics to the Kestra server.

Kestra.counter("my_counter", 1)

Timers

The Kestra class provides a method to send timer metrics to the Kestra server.

Kestra.timer("my_timer", 1)

Gauges

The Kestra class provides a method to send gauge metrics to the Kestra server.

Kestra.gauge("my_gauge", 42.5)

Execution Context

The context object exposes the metadata of the current execution, so scripts stay plain Python instead of embedding Pebble expressions:

import pandas as pd
from kestra import context

my_labels = context.labels
my_inputs = context.inputs
my_vars = context.vars
start_date = context.trigger.startDate
csv_path = context.outputs.prev_task.uri  # outputs from previous tasks

df = pd.read_csv(csv_path)

The context is read from the .kestra-execution-context.json file that Kestra injects into the task working directory. Use KESTRA_EXECUTION_CONTEXT_FILE to read it from another path, or set the base64-encoded KESTRA_CONTEXT environment variable when no file is available.

It is loaded lazily on first access and raises a FileNotFoundError when the script does not run inside a Kestra task.

Nested objects are wrapped, including those inside lists, so any depth can be traversed with attributes: context.inputs.my_list[0].my_key.

context is a read-only Mapping, so len(), in, iteration, keys(), values(), items(), get() and dict(context) all work as expected.

Methods

  • context[name]: Same as attribute access. Required for keys that are not valid Python identifiers, or that collide with a method name such as context["items"].
  • context.get(name, default=None): Returns default instead of raising when the key is missing.
  • context.to_dict(): Returns the raw context as a plain dict, ready for json.dumps.
  • load_execution_context(path=None) -> ExecutionContext: Loads the context explicitly, from path when given.

Kestra Ion

The Kestra ION extra provides a method to read files and convert them to a list of dictionaries.

Installation

pip install kestra[ion]

Methods

  • read(path_: str) -> list[dict[str, Any]]: Reads an Ion file and converts it to a list of dictionaries.

Usage Example

import pandas as pd
import requests
from kestra import Kestra

file_path = "employees.ion"
url = "https://huggingface.co/datasets/kestra/datasets/resolve/main/ion/employees.ion"
response = requests.get(url)
if response.status_code == 200:
    with open(file_path, "wb") as file:
        file.write(response.content)
else:
    print(f"Failed to download the file. Status code: {response.status_code}")


data = Kestra.read(file_path)
df = pd.DataFrame(data)
print(df.info())

Release files for kestra 2.0.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 kestra 2.0.0
File Size Uploaded
kestra-2.0.0.tar.gz 14.3 kB Details

Built distribution (wheel)

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

Total release size: 25.3 kB

Release files / kestra-2.0.0.tar.gz

Download URL kestra-2.0.0.tar.gz
Size 14.3 kB
Tags Source
SHA-256 checksum
How to use checksums
5319034874395ea15bc2ee798c1adba2a90d546e704d43039b25e3c9429d48a1
BLAKE2b-256 checksum
How to use checksums
df2db3d91d447b7001e9fa5908d194e101417093a43e303545488d802cab119f
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.13.14

Release files / kestra-2.0.0-py3-none-any.whl

Download URL kestra-2.0.0-py3-none-any.whl
Size 11.0 kB
Tags Python 3
SHA-256 checksum
How to use checksums
7affc73944f17e0c663f927aadb0a01d5828e1ccbd0e1c549bdf97313e441e42
BLAKE2b-256 checksum
How to use checksums
899ffa5acc866229b418da7bd443939d0038a0b0937fc2b66c31e020a078668c
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.13.14

Release history Release notifications | RSS feed

This release

2.0.0 This release

2 release files

1.3.0

2 release files

1.2.0

2 release files

1.1.0

2 release files

1.0.0

2 release files

0.23.0

2 release files

0.21.6

2 release files

0.21.0

2 release files

0.18.2

2 release files

0.18.1

2 release files

0.18.0

2 release files

0.10.1

2 release files

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