This is an automated email from the ASF dual-hosted git repository.
eladkal pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 985e5007dc1 Add environments and trigger named parameters to
Databricks operators (#72505)
985e5007dc1 is described below
commit 985e5007dc1b9e5d1b5219e6ebb95ef0dad9a099
Author: Noritaka Sekiyama <[email protected]>
AuthorDate: Sat Sep 19 16:09:26 2026 +0900
Add environments and trigger named parameters to Databricks operators
(#72505)
Serverless task dependencies and event-driven triggers are how Databricks
jobs are
written today, but reaching either from Airflow meant hand-building the raw
payload
in ``json`` — the two fields were the only mainstream top-level ones
without a named
parameter, so the operator signature no longer showed what a job can be
given.
---
.../databricks/docs/operators/jobs_create.rst | 2 +
providers/databricks/docs/operators/submit_run.rst | 1 +
.../providers/databricks/operators/databricks.py | 23 ++++
.../unit/databricks/operators/test_databricks.py | 117 +++++++++++++++++++++
4 files changed, 143 insertions(+)
diff --git a/providers/databricks/docs/operators/jobs_create.rst
b/providers/databricks/docs/operators/jobs_create.rst
index 9c2a56f170b..b9d1670d857 100644
--- a/providers/databricks/docs/operators/jobs_create.rst
+++ b/providers/databricks/docs/operators/jobs_create.rst
@@ -48,11 +48,13 @@ Currently the named parameters that
``DatabricksCreateJobsOperator`` supports ar
- ``tags``
- ``tasks``
- ``job_clusters``
+ - ``environments``
- ``email_notifications``
- ``webhook_notifications``
- ``notification_settings``
- ``timeout_seconds``
- ``schedule``
+ - ``trigger``
- ``max_concurrent_runs``
- ``git_source``
- ``access_control_list``
diff --git a/providers/databricks/docs/operators/submit_run.rst
b/providers/databricks/docs/operators/submit_run.rst
index b5e59f50cce..ba7f07881d4 100644
--- a/providers/databricks/docs/operators/submit_run.rst
+++ b/providers/databricks/docs/operators/submit_run.rst
@@ -78,6 +78,7 @@ Currently the named parameters that
``DatabricksSubmitRunOperator`` supports are
- ``new_cluster``
- ``existing_cluster_id``
- ``libraries``
+ - ``environments``
- ``run_name``
- ``timeout_seconds``
- ``performance_target``
diff --git
a/providers/databricks/src/airflow/providers/databricks/operators/databricks.py
b/providers/databricks/src/airflow/providers/databricks/operators/databricks.py
index 540ae06cb6b..f65c414c3e7 100644
---
a/providers/databricks/src/airflow/providers/databricks/operators/databricks.py
+++
b/providers/databricks/src/airflow/providers/databricks/operators/databricks.py
@@ -451,11 +451,16 @@ class DatabricksCreateJobsOperator(BaseOperator):
Array of objects (JobTaskSettings).
:param job_clusters: A list of job cluster specifications that can be
shared and reused by
tasks of this job. Array of objects (JobCluster).
+ :param environments: A list of task execution environments that the
serverless tasks of this
+ job reference through their ``environment_key``. Array of objects
(JobEnvironment).
:param email_notifications: Object (JobEmailNotifications).
:param webhook_notifications: Object (WebhookNotifications).
:param notification_settings: Optional notification settings.
:param timeout_seconds: An optional timeout applied to each run of this
job.
:param schedule: Object (CronSchedule).
+ :param trigger: A configuration to trigger a run when certain conditions
are met, such as a
+ file arriving in an external location, a monitored table being
updated, or a periodic
+ interval elapsing. Object (TriggerSettings).
:param max_concurrent_runs: An optional maximum allowed number of
concurrent runs of the job.
:param git_source: An optional specification for a remote repository
containing the notebooks
used by this job's notebook tasks. Object (GitSource).
@@ -503,11 +508,13 @@ class DatabricksCreateJobsOperator(BaseOperator):
"tags",
"tasks",
"job_clusters",
+ "environments",
"email_notifications",
"webhook_notifications",
"notification_settings",
"timeout_seconds",
"schedule",
+ "trigger",
"max_concurrent_runs",
"git_source",
"access_control_list",
@@ -527,11 +534,13 @@ class DatabricksCreateJobsOperator(BaseOperator):
tags: dict[str, str] | None = None,
tasks: list[dict] | None = None,
job_clusters: list[dict] | None = None,
+ environments: list[dict] | None = None,
email_notifications: dict | None = None,
webhook_notifications: dict | None = None,
notification_settings: dict | None = None,
timeout_seconds: int | None = None,
schedule: dict | None = None,
+ trigger: dict | None = None,
max_concurrent_runs: int | None = None,
git_source: dict | None = None,
access_control_list: list[dict] | None = None,
@@ -551,11 +560,13 @@ class DatabricksCreateJobsOperator(BaseOperator):
self.tags = tags
self.tasks = tasks
self.job_clusters = job_clusters
+ self.environments = environments
self.email_notifications = email_notifications
self.webhook_notifications = webhook_notifications
self.notification_settings = notification_settings
self.timeout_seconds = timeout_seconds
self.schedule = schedule
+ self.trigger = trigger
self.max_concurrent_runs = max_concurrent_runs
self.git_source = git_source
self.access_control_list = access_control_list
@@ -573,11 +584,13 @@ class DatabricksCreateJobsOperator(BaseOperator):
"tags": self.tags,
"tasks": self.tasks,
"job_clusters": self.job_clusters,
+ "environments": self.environments,
"email_notifications": self.email_notifications,
"webhook_notifications": self.webhook_notifications,
"notification_settings": self.notification_settings,
"timeout_seconds": self.timeout_seconds,
"schedule": self.schedule,
+ "trigger": self.trigger,
"max_concurrent_runs": self.max_concurrent_runs,
"git_source": self.git_source,
"access_control_list": self.access_control_list,
@@ -700,6 +713,12 @@ class DatabricksSubmitRunOperator(ResumableJobMixin,
BaseOperator):
.. seealso::
https://docs.databricks.com/dev-tools/api/2.0/jobs.html#managedlibrarieslibrary
+ :param environments: A list of task execution environments that the
serverless tasks of this
+ run reference through their ``environment_key``. Array of objects
(JobEnvironment).
+ This field will be templated.
+
+ .. seealso::
+ https://docs.databricks.com/api/workspace/jobs/submit#environments
:param run_name: The run name used for this task.
By default this will be set to the Airflow ``task_id``. This
``task_id`` is a
required parameter of the superclass ``BaseOperator``.
@@ -781,6 +800,7 @@ class DatabricksSubmitRunOperator(ResumableJobMixin,
BaseOperator):
"new_cluster",
"existing_cluster_id",
"libraries",
+ "environments",
"run_name",
"timeout_seconds",
"idempotency_token",
@@ -809,6 +829,7 @@ class DatabricksSubmitRunOperator(ResumableJobMixin,
BaseOperator):
new_cluster: dict[str, object] | None = None,
existing_cluster_id: str | None = None,
libraries: list[dict[str, Any]] | None = None,
+ environments: list[dict[str, Any]] | None = None,
run_name: str | None = None,
timeout_seconds: int | None = None,
databricks_conn_id: str = "databricks_default",
@@ -849,6 +870,7 @@ class DatabricksSubmitRunOperator(ResumableJobMixin,
BaseOperator):
self.new_cluster = new_cluster
self.existing_cluster_id = existing_cluster_id
self.libraries = libraries
+ self.environments = environments
self.run_name = run_name
self.timeout_seconds = timeout_seconds
self.idempotency_token = idempotency_token
@@ -881,6 +903,7 @@ class DatabricksSubmitRunOperator(ResumableJobMixin,
BaseOperator):
"new_cluster": self.new_cluster,
"existing_cluster_id": self.existing_cluster_id,
"libraries": self.libraries,
+ "environments": self.environments,
"run_name": self.run_name,
"timeout_seconds": self.timeout_seconds,
"idempotency_token": self.idempotency_token,
diff --git
a/providers/databricks/tests/unit/databricks/operators/test_databricks.py
b/providers/databricks/tests/unit/databricks/operators/test_databricks.py
index d05be38716d..3aad7c9b933 100644
--- a/providers/databricks/tests/unit/databricks/operators/test_databricks.py
+++ b/providers/databricks/tests/unit/databricks/operators/test_databricks.py
@@ -293,6 +293,36 @@ ACCESS_CONTROL_LIST = [
"permission_level": "CAN_MANAGE",
}
]
+ENVIRONMENTS = [
+ {
+ "environment_key": "serverless_default",
+ "spec": {"environment_version": "3", "dependencies":
["simplejson==3.19.3"]},
+ }
+]
+TRIGGER = {
+ "pause_status": "UNPAUSED",
+ "file_arrival": {"url": "/Volumes/main/default/landing/",
"min_time_between_triggers_seconds": 60},
+}
+TEMPLATED_ENVIRONMENTS = [
+ {
+ "environment_key": "serverless_default",
+ "spec": {"environment_version": "3", "dependencies": ["simplejson=={{
ds }}"]},
+ }
+]
+RENDERED_TEMPLATED_ENVIRONMENTS = [
+ {
+ "environment_key": "serverless_default",
+ "spec": {"environment_version": "3", "dependencies":
[f"simplejson=={DATE}"]},
+ }
+]
+TEMPLATED_TRIGGER = {
+ "pause_status": "UNPAUSED",
+ "file_arrival": {"url": "/Volumes/main/default/landing/{{ ds }}/"},
+}
+RENDERED_TEMPLATED_TRIGGER = {
+ "pause_status": "UNPAUSED",
+ "file_arrival": {"url": f"/Volumes/main/default/landing/{DATE}/"},
+}
JOB_PARAMS = [{"name": "param1", "default": "value1"}]
@@ -477,6 +507,56 @@ class TestDatabricksCreateJobsOperator:
assert expected == utils.normalise_json_content(op._get_merged_json())
+ @pytest.mark.parametrize(
+ "json",
+ [
+ pytest.param(None, id="named-parameters-only"),
+ pytest.param({"environments": [], "trigger": {}},
id="named-parameters-override-json"),
+ ],
+ )
+ def test_init_with_environments_and_trigger_named_parameters(self, json):
+ """
+ Test the initializer merges ``environments`` and ``trigger`` into the
create payload.
+ """
+ op = DatabricksCreateJobsOperator(
+ task_id=TASK_ID,
+ json=json,
+ name=JOB_NAME,
+ tasks=TASKS,
+ environments=ENVIRONMENTS,
+ trigger=TRIGGER,
+ )
+ expected = utils.normalise_json_content(
+ {
+ "name": JOB_NAME,
+ "tasks": TASKS,
+ "environments": ENVIRONMENTS,
+ "trigger": TRIGGER,
+ }
+ )
+
+ assert expected == utils.normalise_json_content(op._get_merged_json())
+
+ def test_environments_and_trigger_are_templated(self):
+ dag = DAG("test", schedule=None, start_date=datetime.now())
+ op = DatabricksCreateJobsOperator(
+ dag=dag,
+ task_id=TASK_ID,
+ name=JOB_NAME,
+ environments=TEMPLATED_ENVIRONMENTS,
+ trigger=TEMPLATED_TRIGGER,
+ )
+ op.render_template_fields(context={"ds": DATE})
+ expected = utils.normalise_json_content(
+ {
+ "name": JOB_NAME,
+ "environments": RENDERED_TEMPLATED_ENVIRONMENTS,
+ "trigger": RENDERED_TEMPLATED_TRIGGER,
+ }
+ )
+
+ assert expected == utils.normalise_json_content(op._get_merged_json())
+
def test_init_with_templating(self):
json = {"name": "test-{{ ds }}"}
@@ -789,6 +869,43 @@ class TestDatabricksSubmitRunOperator:
assert expected == utils.normalise_json_content(op._get_merged_json())
+ @pytest.mark.parametrize(
+ "json",
+ [
+ pytest.param({"notebook_task": NOTEBOOK_TASK},
id="named-parameters-only"),
+ pytest.param(
+ {"notebook_task": NOTEBOOK_TASK, "environments": []},
+ id="named-parameters-override-json",
+ ),
+ ],
+ )
+ def test_init_with_environments_named_parameter(self, json):
+ """
+ Test the initializer merges ``environments`` into the submit payload.
+ """
+ op = DatabricksSubmitRunOperator(task_id=TASK_ID, json=json,
environments=ENVIRONMENTS)
+ expected = utils.normalise_json_content(
+ {"notebook_task": NOTEBOOK_TASK, "environments": ENVIRONMENTS,
"run_name": TASK_ID}
+ )
+
+ assert expected == utils.normalise_json_content(op._get_merged_json())
+
+ def test_environments_is_templated(self):
+ dag = DAG("test", schedule=None, start_date=datetime.now())
+ op = DatabricksSubmitRunOperator(
+ dag=dag, task_id=TASK_ID, notebook_task=NOTEBOOK_TASK,
environments=TEMPLATED_ENVIRONMENTS
+ )
+ op.render_template_fields(context={"ds": DATE})
+ expected = utils.normalise_json_content(
+ {
+ "notebook_task": NOTEBOOK_TASK,
+ "environments": RENDERED_TEMPLATED_ENVIRONMENTS,
+ "run_name": TASK_ID,
+ }
+ )
+
+ assert expected == utils.normalise_json_content(op._get_merged_json())
+
def test_init_with_spark_python_task_named_parameters(self):
"""
Test the initializer with the named parameters.