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:
- KESTRA_API_TOKEN
- KESTRA_USER and KESTRA_PASSWORD
- 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
-
Trigger a flow and wait for its completion:
from kestra import Flow flow = Flow() flow.execute('mynamespace', 'myflow', {'param': 'value'})
-
Set labels from inputs:
from kestra import Flow flow = Flow(labels_from_inputs=True) flow.execute('mynamespace', 'myflow', {'param': 'value'})
-
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')})
-
Fire and forget:
from kestra import Flow flow = Flow(wait_for_completion=False) flow.execute('mynamespace', 'myflow', {'param': 'value'})
-
Overwrite the username and password:
from kestra import Flow flow = Flow() flow.user = 'admin' flow.password = 'admin' flow.execute('mynamespace', 'myflow')
-
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
defaultinstead of raising when the key is missing. - context.to_dict(): Returns the raw context as a plain
dict, ready forjson.dumps. - load_execution_context(path=None) -> ExecutionContext: Loads the context explicitly, from
pathwhen 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)
| File | Size | Uploaded | |
|---|---|---|---|
| kestra-2.0.0.tar.gz | 14.3 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| 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
|