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.

Reply via email to