This is an automated email from the ASF dual-hosted git repository.

pankajastro 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 5487c080bea Add Azure Blob Storage support to the DataFusion object 
storage layer (#73374)
5487c080bea is described below

commit 5487c080bea96f229b13f24eba49ce56bc560caa
Author: Pankaj Singh <[email protected]>
AuthorDate: Mon Oct 5 13:10:26 2026 +0530

    Add Azure Blob Storage support to the DataFusion object storage layer 
(#73374)
    
    * Add Azure Blob Storage support to the DataFusion object storage layer
    
    DataFusion credential resolution only handled AWS and GCS, so a
    wasb-backed DataSourceConfig raised Unknown connection type instead of
    working, even though DataFusion's storage layer supports Azure. Add a
    StorageType.AZURE, an AzureObjectStorageProvider using DataFusion's
    MicrosoftAzure, and a wasb credential branch resolving the connection's
    shared key, SAS token, or Azure AD service principal, falling back to
    ambient auth when none are set.
    
    The service principal fields (client_id, client_secret, tenant_id)
    must all be set together or not at all -- DataFusion's binding panics
    on a partial combination -- so they're only forwarded when the
    connection has all three (tenant_id extra, login, and password).
    
    Co-Authored-By: Claude Sonnet 5 <[email protected]>
    
    * Remove unused WasbHook install guard from DataFusion Azure credentials
    
    The wasb credential branch never called WasbHook -- it only imported it
    to detect whether apache-airflow-providers-microsoft-azure was installed,
    unlike the aws branch which actually calls AwsGenericHook. Since Azure
    credentials come entirely from the generic Connection object, this guard
    and its extra dependency are unnecessary, matching how the GCS branch
    already works without a provider-specific import.
    
    * Reject unsupported Azure identity fields in DataFusion credential 
resolution
    
    DataFusion's Azure binding has no equivalent for a raw connection
    string, a specific managed/workload identity, or a URL-form SAS
    token, so these connection extra fields were silently dropped and
    the code fell back to whatever ambient credentials were available.
    That can authenticate as a different identity than the one the
    connection was configured for. Raise a clear error instead, matching
    the same guard already used for unsupported GCS identity fields.
    
    * Fix ambiguous object-storage list in LLMSchemaCompareOperator docs
    
    Missing comma made "Azure Blob Storage Parquet" read as one item
    instead of two separate examples in the supported source list.
    
    * Document Azure Blob Storage support in the Analytics operator docs
    
    The Analytics operator docs listed Azure as supported in the intro
    line but not in the Supported Storage Systems list, still flagged it
    under "not yet supported", and had no usage example or credential
    docs, unlike S3 and GCS.
    
    * Fix DataFusion Azure account resolution ignoring connection host
    
    The wasb credential branch took the storage account name unconditionally
    from the connection's login field. For service-principal auth, login holds
    the client_id, not the account name, so a connection whose host correctly
    named the storage account still had the object store target the client_id
    GUID instead. A connection with no login at all (key/SAS auth) silently
    fell back to the AZURE_STORAGE_ACCOUNT_NAME environment variable.
    
    Resolve the account from host first, matching WasbHook's own behavior, and
    fall back to login only when host is empty. Also format the storage type
    name consistently with .value (avoids drifting str-Enum formatting across
    Python versions) and move the sas_token parsing import to the top of the
    file.
    
    * Document that DataFusion's Azure storage only recognizes az:// URIs
    
    abfs:// and abfss:// embed the storage account in the URI itself, a
    different shape from az:// that engine.py does not parse; using them
    currently fails with an unhelpful "Unsupported storage type" error with
    no hint that Azure support exists under a different scheme.
    
    * Fix silent identity/endpoint mismatches in DataFusion Azure credentials
    
    Four issues surfaced by review, each verified against the real datafusion
    binding and object_store source:
    
    - An empty connection (no host or login -- the shape of the wasb_default
      connection airflow db creates) sent the literal string "None" as the
      account instead of omitting it, breaking the fallback to
      AZURE_STORAGE_ACCOUNT_NAME that worked before the host-priority fix.
    
    - A partial service-principal config (tenant_id set but login or password
      missing) silently fell through to ambient auth, or sent the client
      secret as a shared key, instead of raising -- authenticating as a
      different identity than the one configured.
    
    - The account name kept only the first DNS label, so a sovereign-cloud
      host (*.chinacloudapi.cn, *.usgovcloudapi.net) or the Azurite emulator's
      loopback URL was silently redirected to the public
      *.blob.core.windows.net account of the same name; the binding has no
      endpoint override to correct this after the fact.
    
    - DataFusion's MicrosoftAzure always reads AZURE_* worker environment
      variables before overlaying the connection's credential, and checks an
      environment-derived access key or workload-identity token ahead of the
      connection's SAS token or client secret. In the common AKS
      workload-identity deployment, this meant the connection's credential was
      silently ignored in favor of the pod's identity. The binding has no way
      to skip its own environment read, so this can only be caught, not fixed.
    
    * Read the legacy extra__<conn_type>__ prefix for wasb and GCS credentials
    
    WasbHook and GoogleBaseHook still fall back to fields written as
    extra__wasb__<name> / extra__google_cloud_platform__<name> -- the spelling
    older Airflow connection UIs used for custom extra fields. This code only
    looked at the bare key, so a connection created that way had its sas_token
    or connection_string silently invisible: the unsupported-field guard never
    saw a prefixed connection_string, and a prefixed sas_token fell through to
    ambient auth instead of being used.
    
    * Trim two overlong inline comments in Azure credential resolution
    
    Same facts, fewer words -- the mechanism and consequence were spelled out
    twice in slightly different phrasing.
    
    * Simplify the Azure env-precedence docs paragraph
    
    Same facts, plainer words -- "outrank", "ahead of", and "consulted last"
    weren't adding anything a simpler verb wouldn't.
    
    * Attribute AZURE_USE_AZURE_CLI to object_store, not Airflow
    
    Names the actual library so it's clear this is DataFusion's underlying
    binding's behavior, not an Airflow-defined setting, and contrasts it
    directly with WasbHook's automatic az login to explain why it's worth
    calling out at all.
    
    * Fix docs spell-check failure on British spelling "honours"
    
    The docs spellcheck job uses an American-English dictionary; "honors" is
    already the spelling used elsewhere in this same doc tree.
    
    * Simplify wasb credential resolution per review
    
    Collapse the legacy-field lookup into one dict.get() call, and use the
    walrus operator so tenant_id/sas_token are only computed where they're
    actually checked, instead of unconditionally ahead of the branch. Also
    consolidate the three shared-key test variants (password, shared_access_key
    extra, account_key extra) into one parametrized test, picking up account_key
    coverage that was previously untested.
    
    * Fix Azure Blob Storage credential precedence gaps in DataFusion
    
    The env-var precedence guard was missing two real object_store aliases
    (AZURE_STORAGE_MASTER_KEY and the AZURE_CLIENT_ID/SECRET/TENANT_ID
    triple, including the AKS workload-identity webhook's own env vars)
    that silently outrank a connection's SAS token, while AZURE_STORAGE_SAS_KEY
    never can and was raising for nothing.
    
    The account-resolution check also rejected a non-public-cloud host even
    when the worker had already set AZURE_STORAGE_ENDPOINT to fix exactly
    that, so the documented workaround never actually worked.
    
    Finally, the SAS-token branch was checked ahead of the shared-key extra,
    the opposite of WasbHook.get_conn's order -- a connection with a current
    shared_access_key but an expired SAS token would work through the hook
    and fail here.
    
    Co-Authored-By: Claude <[email protected]>
    
    * Move Azure env-var precedence constants next to their only caller
    
    They were declared next to the unrelated _AZURE_PUBLIC_SUFFIX constant,
    ~180 lines from _get_credentials, the only place that reads them --
    inconsistent with the file's existing convention of keeping a constant
    next to its sole user.
    
    Co-Authored-By: Claude <[email protected]>
    
    * Make the Azure credential env-var guard aware of which tier it protects
    
    The flat guard list raised for env vars that could never actually
    outrank a connection's own credential in object_store's precedence
    chain -- most notably AZURE_FEDERATED_TOKEN_FILE against a plain
    shared-key connection, which blocks the single most common Azure
    Kubernetes deployment shape (the workload-identity webhook injects it
    into every labelled pod) even though the connection's own key always
    wins there.
    
    Also closes a gap in the sovereign-cloud/endpoint-override check: an
    Azurite-style host (bare hostname or host:port, with no login) fell
    through an unrelated Active-Directory-ID fallback and silently
    resolved to the literal string "None" instead of raising.
    
    Co-Authored-By: Claude <[email protected]>
    
    * Consolidate Azure env-precedence tests into one parametrized case per tier
    
    Each credential tier (access_key, service_principal, SAS) had a
    "rejects when outranked" test and a separate "ignores when it can't be
    outranked" test that only differed in which env vars were set and
    whether the call was expected to raise. Folding both into a single
    should_raise-parametrized test per tier keeps the same coverage with
    far less repetition.
    
    Co-Authored-By: Claude <[email protected]>
    
    * De-nest the endpoint-override check in _resolve_wasb_account
    
    Left over from incrementally adding the emulator-host colon check;
    the suffix check and the endpoint-override check were always meant
    to be one combined condition, not a nested if.
    
    Co-Authored-By: Claude <[email protected]>
    
    * Fix inaccurate "dead" claim about AZURE_STORAGE_SAS_KEY in test docstring
    
    object_store's resolve_sas_token() does read this variable into
    self.sas_key; it's just checked after sas_query_pairs within that same
    step, so it never wins once the connection sets sas_query_pairs. "Dead"
    overstated that as unused entirely, which could mislead a future reader
    into thinking the safety argument doesn't depend on precedence order.
    
    Co-Authored-By: Claude <[email protected]>
    
    ---------
    
    Co-authored-by: Claude Sonnet 5 <[email protected]>
---
 .../ai/docs/operators/llm_schema_compare.rst       |   2 +-
 providers/common/ai/docs/operators/llm_sql.rst     |   2 +-
 providers/common/ai/docs/toolsets/datafusion.rst   |   2 +-
 providers/common/sql/docs/operators.rst            |  60 ++-
 .../sql/src/airflow/providers/common/sql/config.py |   3 +
 .../providers/common/sql/datafusion/engine.py      | 185 +++++++++
 .../sql/datafusion/object_storage_provider.py      |  36 +-
 .../common/sql/example_dags/example_analytics.py   |  13 +
 .../unit/common/sql/datafusion/test_engine.py      | 447 +++++++++++++++++++++
 .../sql/datafusion/test_object_storage_provider.py |  48 +++
 .../sql/tests/unit/common/sql/test_config.py       |   1 +
 11 files changed, 793 insertions(+), 6 deletions(-)

diff --git a/providers/common/ai/docs/operators/llm_schema_compare.rst 
b/providers/common/ai/docs/operators/llm_schema_compare.rst
index e825c873181..cafde3900de 100644
--- a/providers/common/ai/docs/operators/llm_schema_compare.rst
+++ b/providers/common/ai/docs/operators/llm_schema_compare.rst
@@ -70,7 +70,7 @@ With Object Storage or a Database Table
 
 Use ``data_sources`` with
 :class:`~airflow.providers.common.sql.config.DataSourceConfig` to include
-object-storage sources (S3 or GCS Parquet, CSV, Iceberg, etc.) in the 
comparison.
+object-storage sources (S3, GCS, Azure Blob Storage, Parquet, CSV, Iceberg, 
etc.) in the comparison.
 These can be freely combined with ``db_conn_ids``. Whether an entry is
 introspected via ``DbApiHook`` or DataFusion depends on what its ``conn_id``
 resolves to, not on its ``uri``/``format`` fields: a ``DataSourceConfig``
diff --git a/providers/common/ai/docs/operators/llm_sql.rst 
b/providers/common/ai/docs/operators/llm_sql.rst
index 6f3b74b5aab..5d06ec97e39 100644
--- a/providers/common/ai/docs/operators/llm_sql.rst
+++ b/providers/common/ai/docs/operators/llm_sql.rst
@@ -65,7 +65,7 @@ With Object Storage
 -------------------
 
 Use ``datasource_config`` to generate queries for data stored in object storage
-(e.g., S3, GCS, local filesystem) via `DataFusion 
<https://datafusion.apache.org/>`_.
+(e.g., S3, GCS, Azure Blob Storage, local filesystem) via `DataFusion 
<https://datafusion.apache.org/>`_.
 The operator uses 
:class:`~airflow.providers.common.sql.config.DataSourceConfig`
 to register the object storage source as a table so the LLM can include it in
 the schema context.
diff --git a/providers/common/ai/docs/toolsets/datafusion.rst 
b/providers/common/ai/docs/toolsets/datafusion.rst
index dcf5831891a..5c494caae74 100644
--- a/providers/common/ai/docs/toolsets/datafusion.rst
+++ b/providers/common/ai/docs/toolsets/datafusion.rst
@@ -26,7 +26,7 @@ Files with DataFusion: ``DataFusionToolset``
 Curated toolset wrapping
 :class:`~airflow.providers.common.sql.datafusion.engine.DataFusionEngine`
 with three tools (``list_tables``, ``get_schema``, and ``query``) for
-querying files on object stores (S3, GCS, local filesystem, Iceberg) via 
Apache DataFusion.
+querying files on object stores (S3, GCS, Azure Blob Storage, local 
filesystem, Iceberg) via Apache DataFusion.
 
 .. list-table::
    :header-rows: 1
diff --git a/providers/common/sql/docs/operators.rst 
b/providers/common/sql/docs/operators.rst
index 85c501e3679..f4d2c251db9 100644
--- a/providers/common/sql/docs/operators.rst
+++ b/providers/common/sql/docs/operators.rst
@@ -297,10 +297,11 @@ Supported Storage Systems
 -------------------------
 - S3
 - GCS
+- Azure Blob Storage
 - Local File System
 
 .. note::
-   Azure, HTTP, Delta are not yet supported but will be added in the future.
+   HTTP, Delta are not yet supported but will be added in the future.
 
 
 
@@ -362,6 +363,63 @@ resolved in this order:
     :start-after: [START howto_analytics_operator_with_gcs]
     :end-before: [END howto_analytics_operator_with_gcs]
 
+Azure Storage
+-------------
+Use an ``az://`` URI with a ``conn_id`` pointing to a ``wasb`` connection.
+``abfs://`` and ``abfss://`` URIs are not recognized yet. The account name
+comes from ``host`` (its first DNS label) when set, falling back to
+``login`` only when ``host`` is empty; only the public
+``*.blob.core.windows.net`` cloud is supported, since DataFusion's binding
+otherwise has no endpoint override -- unless the worker sets
+``AZURE_STORAGE_ENDPOINT``/``AZURE_ENDPOINT`` for a real sovereign-cloud
+hostname. A ``host:port`` address (the Azurite emulator's shape, e.g.
+``azurite:10000``) is never a real DNS hostname and always raises instead,
+even with those variables set: put the account name in ``login`` with
+``host`` left empty for Azurite, alongside those same variables and
+``AZURE_ALLOW_HTTP=true``. ``client_secret_auth_config`` (the authority
+override ``WasbHook`` honors) is not read here.
+
+The connection supplies one of the following credentials, checked in this
+order (matching ``WasbHook.get_conn``):
+
+1. Azure AD service principal -- ``tenant_id`` extra, with ``login`` as the
+   client ID and ``password`` as the client secret (both required together)
+2. Shared key -- the ``shared_access_key`` extra
+3. SAS token -- ``sas_token`` extra, as a query string
+4. Shared key -- ``password``, or the ``account_key`` extra
+5. None of the above -- ambient auth (see below)
+
+**A worker environment variable can override the connection, but only from a
+higher-priority tier.** DataFusion reads ``AZURE_*`` environment variables
+first, and resolves credentials in this priority order:
+
+1. Bearer token
+2. Access key -- the shared-key connection's tier
+3. Workload identity
+4. Client secret -- the service-principal connection's tier
+5. SAS -- the SAS connection's tier
+
+Only a complete credential from an earlier tier can override the
+connection -- a single bearer or access-key variable, or every variable a
+multi-field tier needs. This raises when found, naming the variables,
+instead of silently using the wrong identity.
+
+With no explicit credential, authentication falls back to ``AZURE_*``
+environment variables, managed identity, or workload identity. Unlike
+``WasbHook`` (which tries ``az login`` automatically), DataFusion's
+underlying ``object_store`` binding only tries the Azure CLI if
+``AZURE_USE_AZURE_CLI=true`` is set; otherwise it defaults straight to
+IMDS managed identity.
+
+``connection_string``, ``managed_identity_client_id``, 
``workload_identity_tenant_id``,
+and a URL-form ``sas_token`` are not supported.
+
+.. exampleinclude:: 
/../../sql/src/airflow/providers/common/sql/example_dags/example_analytics.py
+    :language: python
+    :dedent: 4
+    :start-after: [START howto_analytics_operator_with_azure]
+    :end-before: [END howto_analytics_operator_with_azure]
+
 Local File System Storage
 -------------------------
 .. exampleinclude:: 
/../../sql/src/airflow/providers/common/sql/example_dags/example_analytics.py
diff --git a/providers/common/sql/src/airflow/providers/common/sql/config.py 
b/providers/common/sql/src/airflow/providers/common/sql/config.py
index a1a856a17f3..ad1d2d82eb5 100644
--- a/providers/common/sql/src/airflow/providers/common/sql/config.py
+++ b/providers/common/sql/src/airflow/providers/common/sql/config.py
@@ -48,6 +48,7 @@ class StorageType(str, Enum):
 
     S3 = "s3"
     GCS = "gcs"
+    AZURE = "azure"
     LOCAL = "local"
 
 
@@ -117,6 +118,8 @@ class DataSourceConfig:
             return StorageType.S3
         if self.uri.startswith("gs://"):
             return StorageType.GCS
+        if self.uri.startswith("az://"):
+            return StorageType.AZURE
         if self.uri.startswith("file://"):
             return StorageType.LOCAL
         raise ValueError(f"Unsupported storage type for URI: {self.uri}")
diff --git 
a/providers/common/sql/src/airflow/providers/common/sql/datafusion/engine.py 
b/providers/common/sql/src/airflow/providers/common/sql/datafusion/engine.py
index 1b8a3998ef6..3eec4d0f9c7 100644
--- a/providers/common/sql/src/airflow/providers/common/sql/datafusion/engine.py
+++ b/providers/common/sql/src/airflow/providers/common/sql/datafusion/engine.py
@@ -18,6 +18,7 @@ from __future__ import annotations
 
 import os
 from typing import TYPE_CHECKING, Any
+from urllib.parse import parse_qsl, urlsplit
 
 from datafusion import SessionContext
 
@@ -161,6 +162,12 @@ class DataFusionEngine(LoggingMixin):
                 return extra_dejson[field_name]
             return 
extra_dejson.get(f"extra__google_cloud_platform__{field_name}")
 
+        def _get_wasb_extra_field(extra_dejson: dict[str, Any], field_name: 
str) -> Any:
+            # Older Airflow connection UIs wrote custom extra fields as
+            # extra__wasb__<field_name> instead of the bare key; WasbHook 
still reads that
+            # legacy spelling as a fallback, so this must too.
+            return extra_dejson.get(field_name, 
extra_dejson.get(f"extra__wasb__{field_name}"))
+
         match conn.conn_type:
             case "aws":
                 try:
@@ -204,6 +211,66 @@ class DataFusionEngine(LoggingMixin):
                     key_path = os.environ.get("GOOGLE_APPLICATION_CREDENTIALS")
                 credentials = self._remove_none_values({"key_path": key_path, 
"keyfile_dict": keyfile_dict})
 
+            case "wasb":
+                extra_dejson = conn.extra_dejson
+                for unsupported_field in (
+                    "connection_string",
+                    "managed_identity_client_id",
+                    "workload_identity_tenant_id",
+                ):
+                    if _get_wasb_extra_field(extra_dejson, unsupported_field):
+                        raise ValueError(
+                            f"Connection field {unsupported_field!r} is not 
supported for DataFusion "
+                            "Azure Blob Storage access; only 
tenant_id+login+password (service "
+                            "principal), sas_token, 
shared_access_key/account_key/password, or ambient "
+                            "credentials (AZURE_* environment variables, 
managed identity, workload "
+                            "identity, or az login) are used."
+                        )
+                credentials = {"account": 
self._resolve_wasb_account(conn.host, conn.login)}
+                credential_tier: str | None = None
+                if tenant_id := _get_wasb_extra_field(extra_dejson, 
"tenant_id"):
+                    if not conn.login or not conn.password:
+                        # Falling through here would silently switch identity 
(ambient auth, or
+                        # the client secret sent as a shared key) instead of 
failing clearly.
+                        missing = "login (client_id)" if not conn.login else 
"password (client_secret)"
+                        raise ValueError(
+                            f"Connection extra 'tenant_id' is set for 
DataFusion Azure Blob Storage "
+                            f"service-principal auth, but {missing} is not."
+                        )
+                    credentials.update(
+                        {"client_id": conn.login, "client_secret": 
conn.password, "tenant_id": tenant_id}
+                    )
+                    credential_tier = "client_secret"
+                elif shared_access_key := _get_wasb_extra_field(extra_dejson, 
"shared_access_key"):
+                    # Checked ahead of sas_token to match WasbHook.get_conn's 
precedence.
+                    credentials["access_key"] = shared_access_key
+                    credential_tier = "access_key"
+                elif sas_token := _get_wasb_extra_field(extra_dejson, 
"sas_token"):
+                    if sas_token.startswith("http"):
+                        raise ValueError(
+                            "A URL-form `sas_token` is not supported for 
DataFusion Azure Blob Storage "
+                            "access; provide the SAS token as a query string 
instead."
+                        )
+                    credentials["sas_query_pairs"] = 
parse_qsl(sas_token.lstrip("?"))
+                    credential_tier = "sas"
+                else:
+                    access_key = conn.password or 
_get_wasb_extra_field(extra_dejson, "account_key")
+                    if access_key:
+                        credentials["access_key"] = access_key
+                        credential_tier = "access_key"
+
+                if credential_tier is not None:
+                    conflicting_env_vars = 
self._find_conflicting_azure_env_vars(credential_tier)
+                    if conflicting_env_vars:
+                        raise ValueError(
+                            f"Worker environment variable(s) {', 
'.join(conflicting_env_vars)} would "
+                            "silently take precedence over this connection's 
explicit credential in "
+                            "DataFusion's Azure Blob Storage binding. Unset 
them on the worker, or "
+                            "remove the explicit credential from this 
connection to rely on the "
+                            "environment instead."
+                        )
+                credentials = self._remove_none_values(credentials)
+
             case _:
                 raise ValueError(f"Unknown connection type {conn.conn_type}")
         return credentials, extra_config
@@ -213,6 +280,124 @@ class DataFusionEngine(LoggingMixin):
         """Filter out None values from the dictionary."""
         return {k: v for k, v in params.items() if v is not None}
 
+    _AZURE_PUBLIC_SUFFIX = ".blob.core.windows.net"
+
+    # object_store's build() precedence, high to low: bearer token > access 
key > workload
+    # identity (client_id+tenant_id+federated_token_file) > client secret
+    # (client_id+client_secret+tenant_id) > SAS. Used by 
_find_conflicting_azure_env_vars.
+    _AZURE_ENV_BEARER_VARS = ("AZURE_STORAGE_TOKEN",)
+    _AZURE_ENV_ACCESS_KEY_VARS = (
+        "AZURE_STORAGE_ACCOUNT_KEY",
+        "AZURE_STORAGE_ACCESS_KEY",
+        "AZURE_STORAGE_MASTER_KEY",
+    )
+    _AZURE_ENV_CLIENT_ID_VARS = ("AZURE_STORAGE_CLIENT_ID", "AZURE_CLIENT_ID")
+    _AZURE_ENV_CLIENT_SECRET_VARS = ("AZURE_STORAGE_CLIENT_SECRET", 
"AZURE_CLIENT_SECRET")
+    _AZURE_ENV_TENANT_ID_VARS = (
+        "AZURE_STORAGE_TENANT_ID",
+        "AZURE_STORAGE_AUTHORITY_ID",
+        "AZURE_TENANT_ID",
+        "AZURE_AUTHORITY_ID",
+    )
+    _AZURE_ENV_FEDERATED_TOKEN_FILE_VAR = "AZURE_FEDERATED_TOKEN_FILE"
+
+    @classmethod
+    def _find_conflicting_azure_env_vars(cls, credential_tier: str) -> 
list[str]:
+        """
+        Return worker env vars that would silently outrank the connection's 
own credential.
+
+        The binding always calls `from_env()` with no way to skip it, so this 
can only be
+        caught here, not avoided. The connection's fields overwrite the same 
fields from
+        `from_env()`, so only a tier above the connection's own, or an 
unoccupied tier
+        (workload identity or client secret, for SAS) with *all* its fields 
present, can
+        actually take over.
+        """
+
+        def env_set(*var_groups: tuple[str, ...]) -> list[str]:
+            return [var for group in var_groups for var in group if 
os.environ.get(var)]
+
+        conflicting = env_set(cls._AZURE_ENV_BEARER_VARS)
+        if credential_tier == "access_key":
+            return conflicting
+        conflicting += env_set(cls._AZURE_ENV_ACCESS_KEY_VARS)
+        if credential_tier == "client_secret":
+            # client_id and tenant_id are already the connection's own; only 
the federated
+            # token file is left for env to complete the workload-identity 
triple with.
+            if os.environ.get(cls._AZURE_ENV_FEDERATED_TOKEN_FILE_VAR):
+                conflicting.append(cls._AZURE_ENV_FEDERATED_TOKEN_FILE_VAR)
+            return conflicting
+        # SAS occupies none of these fields, so each whole triple must come 
from env.
+        if (
+            any(os.environ.get(var) for var in cls._AZURE_ENV_CLIENT_ID_VARS)
+            and any(os.environ.get(var) for var in 
cls._AZURE_ENV_TENANT_ID_VARS)
+            and os.environ.get(cls._AZURE_ENV_FEDERATED_TOKEN_FILE_VAR)
+        ):
+            conflicting += env_set(cls._AZURE_ENV_CLIENT_ID_VARS, 
cls._AZURE_ENV_TENANT_ID_VARS)
+            conflicting.append(cls._AZURE_ENV_FEDERATED_TOKEN_FILE_VAR)
+        if (
+            any(os.environ.get(var) for var in cls._AZURE_ENV_CLIENT_ID_VARS)
+            and any(os.environ.get(var) for var in 
cls._AZURE_ENV_CLIENT_SECRET_VARS)
+            and any(os.environ.get(var) for var in 
cls._AZURE_ENV_TENANT_ID_VARS)
+        ):
+            conflicting += env_set(
+                cls._AZURE_ENV_CLIENT_ID_VARS,
+                cls._AZURE_ENV_CLIENT_SECRET_VARS,
+                cls._AZURE_ENV_TENANT_ID_VARS,
+            )
+        return list(dict.fromkeys(conflicting))
+
+    @classmethod
+    def _resolve_wasb_account(cls, host: str | None, login: str | None) -> str 
| None:
+        """
+        Return the storage account name the way WasbHook resolves it.
+
+        From ``host`` (its netloc's first label) when set, falling back to 
``login`` only when
+        ``host`` is empty -- login holds the service-principal client_id 
otherwise, not the
+        account name. Returns ``None`` when neither is set, so the binding 
falls back to
+        ``AZURE_STORAGE_ACCOUNT_NAME`` rather than the literal string 
``"None"``. Reimplemented
+        locally instead of importing
+        ``airflow.providers.microsoft.azure.utils.parse_blob_account_url``, to 
avoid pulling in
+        the microsoft-azure provider's Azure SDK dependency for one stdlib 
string operation.
+
+        Only the public ``*.blob.core.windows.net`` cloud is supported, unless
+        ``AZURE_STORAGE_ENDPOINT``/``AZURE_ENDPOINT`` is set for a real 
sovereign-cloud hostname
+        (DataFusion's binding otherwise has no endpoint override). A 
``host:port`` netloc (the
+        Azurite emulator's shape) is never a real hostname, so it always 
raises instead --
+        checked before the dotless-host fallback below, which would otherwise 
mask it.
+        """
+        if not host and not login:
+            return None
+        netloc = urlsplit(host if host else 
f"https://{login}.blob.core.windows.net/";).netloc
+        if not netloc:
+            # No scheme was given (e.g. a bare DNS name); urlsplit put it all 
in the path instead.
+            netloc = urlsplit(f"https://{host}";).netloc
+        if host and ":" in netloc:
+            raise ValueError(
+                f"Connection host {host!r} looks like an emulator address 
(host:port), which "
+                "DataFusion's Azure Blob Storage binding cannot resolve an 
account name from -- "
+                "even with AZURE_STORAGE_ENDPOINT set. Put the account name in 
`login` with "
+                "`host` empty instead, alongside AZURE_STORAGE_ENDPOINT and 
AZURE_ALLOW_HTTP=true."
+            )
+        if "." not in netloc:
+            if not login:
+                raise ValueError(
+                    f"Connection host {host!r} is not a full URL or DNS name, 
and no `login` was "
+                    "given to resolve it as an Active Directory ID instead."
+                )
+            # Only an Active Directory ID was given, not a full URL or DNS 
name.
+            netloc = f"{login}.blob.core.windows.net"
+        if not netloc.endswith(cls._AZURE_PUBLIC_SUFFIX) and not (
+            os.environ.get("AZURE_STORAGE_ENDPOINT") or 
os.environ.get("AZURE_ENDPOINT")
+        ):
+            raise ValueError(
+                f"Connection host {host!r} does not resolve to the public 
{cls._AZURE_PUBLIC_SUFFIX} "
+                "cloud, which is the only one DataFusion's Azure Blob Storage 
binding can target (it "
+                "has no endpoint override). A sovereign cloud is supported 
once the "
+                "AZURE_STORAGE_ENDPOINT environment variable is set."
+            )
+        # Azure storage account names are capped at 24 characters.
+        return netloc.split(".", 1)[0][:24]
+
     def get_schema(self, table_name: str):
         """Get the schema of a table."""
         schema = str(self.session_context.table(table_name).schema())
diff --git 
a/providers/common/sql/src/airflow/providers/common/sql/datafusion/object_storage_provider.py
 
b/providers/common/sql/src/airflow/providers/common/sql/datafusion/object_storage_provider.py
index 73b70079a84..9799cbe00e2 100644
--- 
a/providers/common/sql/src/airflow/providers/common/sql/datafusion/object_storage_provider.py
+++ 
b/providers/common/sql/src/airflow/providers/common/sql/datafusion/object_storage_provider.py
@@ -20,7 +20,7 @@ import json
 import tempfile
 from pathlib import Path
 
-from datafusion.object_store import AmazonS3, GoogleCloud, LocalFileSystem
+from datafusion.object_store import AmazonS3, GoogleCloud, LocalFileSystem, 
MicrosoftAzure
 
 from airflow.providers.common.sql.config import ConnectionConfig, StorageType
 from airflow.providers.common.sql.datafusion.base import ObjectStorageProvider
@@ -107,6 +107,37 @@ class GCSObjectStorageProvider(ObjectStorageProvider):
         return "gs://"
 
 
+class AzureObjectStorageProvider(ObjectStorageProvider):
+    """Azure Object Storage Provider using DataFusion's MicrosoftAzure."""
+
+    @property
+    def get_storage_type(self) -> StorageType:
+        """Return the storage type."""
+        return StorageType.AZURE
+
+    def create_object_store(self, path: str, connection_config: 
ConnectionConfig | None = None):
+        """Create an Azure object store using DataFusion's MicrosoftAzure."""
+        if connection_config is None:
+            raise ValueError(f"connection_config must be provided for 
{self.get_storage_type.value}")
+
+        try:
+            credentials = connection_config.credentials
+            container = self.get_bucket(path)
+
+            azure_store = MicrosoftAzure(container_name=container, 
**credentials)
+            self.log.info("Created Azure object store for container %s", 
container)
+
+            return azure_store
+
+        except BaseException as e:
+            # A bad credential combination panics as 
pyo3_runtime.PanicException, not a plain Exception.
+            raise ObjectStoreCreationException(f"Failed to create Azure object 
store: {e}")
+
+    def get_scheme(self) -> str:
+        """Return the scheme for Azure."""
+        return "az://"
+
+
 class LocalObjectStorageProvider(ObjectStorageProvider):
     """Local Object Storage Provider using DataFusion's LocalFileSystem."""
 
@@ -126,10 +157,11 @@ class LocalObjectStorageProvider(ObjectStorageProvider):
 
 def get_object_storage_provider(storage_type: StorageType) -> 
ObjectStorageProvider:
     """Get an object storage provider based on the storage type."""
-    # TODO: Add support for Azure, HTTP: 
https://datafusion.apache.org/python/autoapi/datafusion/object_store/index.html
+    # TODO: Add support for HTTP: 
https://datafusion.apache.org/python/autoapi/datafusion/object_store/index.html
     providers: dict[StorageType, type] = {
         StorageType.S3: S3ObjectStorageProvider,
         StorageType.GCS: GCSObjectStorageProvider,
+        StorageType.AZURE: AzureObjectStorageProvider,
         StorageType.LOCAL: LocalObjectStorageProvider,
     }
 
diff --git 
a/providers/common/sql/src/airflow/providers/common/sql/example_dags/example_analytics.py
 
b/providers/common/sql/src/airflow/providers/common/sql/example_dags/example_analytics.py
index 0f80fc6100e..b0d04733939 100644
--- 
a/providers/common/sql/src/airflow/providers/common/sql/example_dags/example_analytics.py
+++ 
b/providers/common/sql/src/airflow/providers/common/sql/example_dags/example_analytics.py
@@ -34,6 +34,10 @@ datasource_config_gcs = DataSourceConfig(
     conn_id="google_cloud_default", table_name="users_data", 
uri="gs://bucket/path/", format="parquet"
 )
 
+datasource_config_azure = DataSourceConfig(
+    conn_id="wasb_default", table_name="users_data", 
uri="az://container/path/", format="parquet"
+)
+
 datasource_config_iceberg = DataSourceConfig(
     conn_id="iceberg_default",
     table_name="users_data",
@@ -85,6 +89,15 @@ with DAG(
     analytics_with_s3 >> analytics_with_gcs
     # [END howto_analytics_operator_with_gcs]
 
+    # [START howto_analytics_operator_with_azure]
+    analytics_with_azure = AnalyticsOperator(
+        task_id="analytics_with_azure",
+        datasource_configs=[datasource_config_azure],
+        queries=["SELECT * FROM users_data", "SELECT count(*) FROM 
users_data"],
+    )
+    analytics_with_s3 >> analytics_with_azure
+    # [END howto_analytics_operator_with_azure]
+
     # [START howto_analytics_operator_with_local]
     analytics_with_local = AnalyticsOperator(
         task_id="analytics_with_local",
diff --git 
a/providers/common/sql/tests/unit/common/sql/datafusion/test_engine.py 
b/providers/common/sql/tests/unit/common/sql/datafusion/test_engine.py
index 10f01b43cb9..c80434ef578 100644
--- a/providers/common/sql/tests/unit/common/sql/datafusion/test_engine.py
+++ b/providers/common/sql/tests/unit/common/sql/datafusion/test_engine.py
@@ -85,6 +85,7 @@ class TestDataFusionEngine:
             ("s3", "csv", "s3"),
             ("s3", "avro", "s3"),
             ("gcs", "parquet", "gs"),
+            ("azure", "parquet", "az"),
         ],
     )
     
@patch("airflow.providers.common.sql.datafusion.engine.get_object_storage_provider",
 autospec=True)
@@ -394,6 +395,452 @@ class TestDataFusionEngine:
         with pytest.raises(ValueError, match="'impersonation_chain' is not 
supported"):
             engine._get_credentials(mock_conn)
 
+    @pytest.mark.parametrize(
+        ("password", "extra_dejson", "expected_access_key"),
+        [
+            ("mykey", {}, "mykey"),
+            (None, {"shared_access_key": "extra-key"}, "extra-key"),
+            (None, {"account_key": "extra-key"}, "extra-key"),
+        ],
+        ids=["password", "shared_access_key_extra", "account_key_extra"],
+    )
+    def test_get_credentials_azure_with_shared_key(self, password, 
extra_dejson, expected_access_key):
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = "myaccount"
+        mock_conn.password = password
+        mock_conn.extra_dejson = extra_dejson
+        engine = DataFusionEngine()
+
+        credentials, extra_config = engine._get_credentials(mock_conn)
+
+        assert credentials == {"account": "myaccount", "access_key": 
expected_access_key}
+        assert extra_config == {}
+
+    def test_get_credentials_azure_with_service_principal(self):
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = "client-id"
+        mock_conn.password = "client-secret"
+        mock_conn.extra_dejson = {"tenant_id": "tenant-id"}
+        engine = DataFusionEngine()
+
+        credentials, extra_config = engine._get_credentials(mock_conn)
+
+        assert credentials == {
+            "account": "client-id",
+            "client_id": "client-id",
+            "client_secret": "client-secret",
+            "tenant_id": "tenant-id",
+        }
+        assert extra_config == {}
+
+    def 
test_get_credentials_azure_with_service_principal_and_host_prefers_host_account(self):
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = "realaccount.blob.core.windows.net"
+        mock_conn.login = "11111111-2222-3333-4444-555555555555"
+        mock_conn.password = "client-secret"
+        mock_conn.extra_dejson = {"tenant_id": "tenant-id"}
+        engine = DataFusionEngine()
+
+        credentials, extra_config = engine._get_credentials(mock_conn)
+
+        assert credentials == {
+            "account": "realaccount",
+            "client_id": "11111111-2222-3333-4444-555555555555",
+            "client_secret": "client-secret",
+            "tenant_id": "tenant-id",
+        }
+        assert extra_config == {}
+
+    @pytest.mark.parametrize(
+        ("login", "password", "missing"),
+        [
+            (None, "client-secret", "login"),
+            ("client-id", None, "password"),
+            (None, None, "login"),
+        ],
+    )
+    def test_get_credentials_azure_partial_service_principal_raises(self, 
login, password, missing):
+        """A partial service-principal config must raise, not silently 
authenticate with a
+        different identity (ambient auth, or the client secret sent as a 
shared key) --
+        DataFusion's binding also panics on a partial 
client_id/client_secret/tenant_id
+        combination."""
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = login
+        mock_conn.password = password
+        mock_conn.extra_dejson = {"tenant_id": "tenant-id"}
+        engine = DataFusionEngine()
+
+        with pytest.raises(ValueError, match=f"{missing}.*is not"):
+            engine._get_credentials(mock_conn)
+
+    def test_get_credentials_azure_fully_empty_connection_omits_account(self):
+        """Neither host nor login set (the shape of the ``wasb_default`` 
connection ``airflow
+        db`` creates) must drop `account` entirely, not send the literal 
string 'None' --
+        the binding then falls back to AZURE_STORAGE_ACCOUNT_NAME."""
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = None
+        mock_conn.password = None
+        mock_conn.extra_dejson = {}
+        engine = DataFusionEngine()
+
+        credentials, extra_config = engine._get_credentials(mock_conn)
+
+        assert credentials == {}
+        assert extra_config == {}
+
+    def 
test_get_credentials_azure_shared_access_key_takes_priority_over_sas_token(self):
+        """Matches WasbHook.get_conn, which checks the `shared_access_key` 
extra before
+        `sas_token`."""
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = "myaccount"
+        mock_conn.password = None
+        mock_conn.extra_dejson = {
+            "sas_token": "?sv=2020-08-04&sp=rl&sig=abc",
+            "shared_access_key": "extra-key",
+        }
+        engine = DataFusionEngine()
+
+        credentials, extra_config = engine._get_credentials(mock_conn)
+
+        assert credentials == {"account": "myaccount", "access_key": 
"extra-key"}
+        assert extra_config == {}
+
+    def test_get_credentials_azure_with_sas_token(self):
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = "myaccount"
+        mock_conn.password = None
+        mock_conn.extra_dejson = {"sas_token": "?sv=2020-08-04&sp=rl&sig=abc"}
+        engine = DataFusionEngine()
+
+        credentials, extra_config = engine._get_credentials(mock_conn)
+
+        assert credentials == {
+            "account": "myaccount",
+            "sas_query_pairs": [("sv", "2020-08-04"), ("sp", "rl"), ("sig", 
"abc")],
+        }
+        assert extra_config == {}
+
+    def test_get_credentials_azure_without_credentials_uses_ambient_auth(self):
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = "myaccount"
+        mock_conn.password = None
+        mock_conn.extra_dejson = {}
+        engine = DataFusionEngine()
+
+        credentials, extra_config = engine._get_credentials(mock_conn)
+
+        assert credentials == {"account": "myaccount"}
+        assert extra_config == {}
+
+    @pytest.mark.parametrize(
+        "unsupported_field",
+        ["connection_string", "managed_identity_client_id", 
"workload_identity_tenant_id"],
+    )
+    def test_get_credentials_azure_rejects_unsupported_identity_fields(self, 
unsupported_field):
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.extra_dejson = {unsupported_field: "some-value"}
+        engine = DataFusionEngine()
+
+        with pytest.raises(ValueError, match=f"{unsupported_field!r} is not 
supported"):
+            engine._get_credentials(mock_conn)
+
+    def test_get_credentials_azure_reads_legacy_extra_prefixed_sas_token(self):
+        """Older Airflow connection UIs wrote custom extra fields as
+        extra__wasb__<field>; WasbHook still reads that spelling as a 
fallback."""
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = "myaccount"
+        mock_conn.password = None
+        mock_conn.extra_dejson = {"extra__wasb__sas_token": 
"?sv=2020-08-04&sp=rl&sig=abc"}
+        engine = DataFusionEngine()
+
+        credentials, extra_config = engine._get_credentials(mock_conn)
+
+        assert credentials == {
+            "account": "myaccount",
+            "sas_query_pairs": [("sv", "2020-08-04"), ("sp", "rl"), ("sig", 
"abc")],
+        }
+        assert extra_config == {}
+
+    def 
test_get_credentials_azure_rejects_legacy_extra_prefixed_unsupported_field(self):
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = "myaccount"
+        mock_conn.password = None
+        mock_conn.extra_dejson = {"extra__wasb__connection_string": 
"some-conn-string"}
+        engine = DataFusionEngine()
+
+        with pytest.raises(ValueError, match="'connection_string' is not 
supported"):
+            engine._get_credentials(mock_conn)
+
+    def test_get_credentials_azure_rejects_url_form_sas_token(self):
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = "myaccount"
+        mock_conn.password = None
+        mock_conn.extra_dejson = {"sas_token": 
"https://myaccount.blob.core.windows.net/?sv=2020-08-04"}
+        engine = DataFusionEngine()
+
+        with pytest.raises(ValueError, match="URL-form `sas_token` is not 
supported"):
+            engine._get_credentials(mock_conn)
+
+    @pytest.mark.parametrize(
+        ("env_vars", "should_raise"),
+        [
+            (["AZURE_STORAGE_ACCOUNT_KEY"], True),
+            (["AZURE_STORAGE_ACCESS_KEY"], True),
+            (["AZURE_STORAGE_MASTER_KEY"], True),
+            (["AZURE_STORAGE_TOKEN"], True),
+            (["AZURE_STORAGE_SAS_KEY"], False),
+            (["AZURE_CLIENT_ID", "AZURE_TENANT_ID", 
"AZURE_FEDERATED_TOKEN_FILE"], True),
+            (["AZURE_FEDERATED_TOKEN_FILE"], False),
+            (["AZURE_CLIENT_ID", "AZURE_CLIENT_SECRET", "AZURE_TENANT_ID"], 
True),
+            (["AZURE_STORAGE_CLIENT_ID", "AZURE_STORAGE_CLIENT_SECRET", 
"AZURE_STORAGE_TENANT_ID"], True),
+            (["AZURE_STORAGE_CLIENT_ID", "AZURE_STORAGE_CLIENT_SECRET", 
"AZURE_STORAGE_AUTHORITY_ID"], True),
+            (["AZURE_CLIENT_ID", "AZURE_CLIENT_SECRET", "AZURE_AUTHORITY_ID"], 
True),
+            (["AZURE_CLIENT_ID"], False),
+        ],
+        ids=[
+            "account-key",
+            "access-key",
+            "master-key",
+            "bearer",
+            "sas-key-checked-after-query-pairs",
+            "workload-identity-triple",
+            "workload-identity-partial",
+            "client-secret-triple-bare",
+            "client-secret-triple-storage-prefixed",
+            "client-secret-triple-storage-authority",
+            "client-secret-triple-bare-authority",
+            "client-secret-partial",
+        ],
+    )
+    def test_get_credentials_azure_sas_env_precedence(self, env_vars, 
should_raise, monkeypatch):
+        """A SAS connection sits at the bottom of object_store's precedence 
order and occupies
+        none of the higher tiers' fields, so a bearer token or any access-key 
spelling always
+        outranks it, and a complete workload-identity or client-secret triple 
(any recognized
+        spelling) does too -- but a partial triple, or AZURE_STORAGE_SAS_KEY 
(checked only after
+        the connection's own sas_query_pairs within SAS resolution itself), 
must not raise."""
+        for var in env_vars:
+            monkeypatch.setenv(var, "some-value")
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = "myaccount"
+        mock_conn.password = None
+        mock_conn.extra_dejson = {"sas_token": "?sv=2020-08-04&sp=rl&sig=abc"}
+        engine = DataFusionEngine()
+
+        if should_raise:
+            with pytest.raises(ValueError, match=", ".join(env_vars)):
+                engine._get_credentials(mock_conn)
+        else:
+            credentials, extra_config = engine._get_credentials(mock_conn)
+            assert credentials == {
+                "account": "myaccount",
+                "sas_query_pairs": [("sv", "2020-08-04"), ("sp", "rl"), 
("sig", "abc")],
+            }
+            assert extra_config == {}
+
+    @pytest.mark.parametrize(
+        ("env_vars", "should_raise"),
+        [
+            (["AZURE_STORAGE_TOKEN"], True),
+            (["AZURE_STORAGE_ACCOUNT_KEY"], False),
+            (["AZURE_STORAGE_ACCESS_KEY"], False),
+            (["AZURE_STORAGE_MASTER_KEY"], False),
+            (["AZURE_FEDERATED_TOKEN_FILE"], False),
+            (["AZURE_CLIENT_ID", "AZURE_CLIENT_SECRET", "AZURE_TENANT_ID"], 
False),
+        ],
+        ids=[
+            "bearer",
+            "account-key",
+            "access-key",
+            "master-key",
+            "federated-token-file",
+            "client-secret-triple",
+        ],
+    )
+    def test_get_credentials_azure_access_key_env_precedence(self, env_vars, 
should_raise, monkeypatch):
+        """A shared-key connection overwrites the same field object_store's 
build() would
+        otherwise read from env, and build() picks access_key ahead of every 
lower tier -- so
+        only a bearer token, the one tier above access_key, can actually 
outrank it, even on an
+        AKS pod where the workload-identity webhook injects 
AZURE_FEDERATED_TOKEN_FILE into
+        every labelled pod."""
+        for var in env_vars:
+            monkeypatch.setenv(var, "some-value")
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = "myaccount"
+        mock_conn.password = "mykey"
+        mock_conn.extra_dejson = {}
+        engine = DataFusionEngine()
+
+        if should_raise:
+            with pytest.raises(ValueError, match=", ".join(env_vars)):
+                engine._get_credentials(mock_conn)
+        else:
+            credentials, extra_config = engine._get_credentials(mock_conn)
+            assert credentials == {"account": "myaccount", "access_key": 
"mykey"}
+            assert extra_config == {}
+
+    @pytest.mark.parametrize(
+        ("env_vars", "should_raise"),
+        [
+            (["AZURE_STORAGE_TOKEN"], True),
+            (["AZURE_STORAGE_ACCOUNT_KEY"], True),
+            (["AZURE_STORAGE_ACCESS_KEY"], True),
+            (["AZURE_STORAGE_MASTER_KEY"], True),
+            (["AZURE_FEDERATED_TOKEN_FILE"], True),
+            (["AZURE_CLIENT_ID", "AZURE_CLIENT_SECRET", "AZURE_TENANT_ID"], 
False),
+        ],
+        ids=[
+            "bearer",
+            "account-key",
+            "access-key",
+            "master-key",
+            "federated-token-file",
+            "client-secret-triple",
+        ],
+    )
+    def test_get_credentials_azure_service_principal_env_precedence(
+        self, env_vars, should_raise, monkeypatch
+    ):
+        """A bearer token or any access-key spelling sits above client secret 
in object_store's
+        precedence order, and the connection's own client_id and tenant_id 
already complete two
+        of workload identity's three fields, so a federated token file alone 
is enough for env to
+        complete that tier -- but a client-secret triple in env is moot, since 
the connection's
+        own client_id/client_secret/tenant_id overwrite the same fields."""
+        for var in env_vars:
+            monkeypatch.setenv(var, "some-value")
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = "client-id"
+        mock_conn.password = "client-secret"
+        mock_conn.extra_dejson = {"tenant_id": "tenant-id"}
+        engine = DataFusionEngine()
+
+        if should_raise:
+            with pytest.raises(ValueError, match=", ".join(env_vars)):
+                engine._get_credentials(mock_conn)
+        else:
+            credentials, extra_config = engine._get_credentials(mock_conn)
+            assert credentials == {
+                "account": "client-id",
+                "client_id": "client-id",
+                "client_secret": "client-secret",
+                "tenant_id": "tenant-id",
+            }
+            assert extra_config == {}
+
+    def 
test_get_credentials_azure_env_precedence_guard_ignores_pure_ambient_auth(self, 
monkeypatch):
+        """The guard only fires for an explicit connection credential -- a 
connection with none
+        at all is meant to rely on the environment, so the same env vars must 
not raise here."""
+        monkeypatch.setenv("AZURE_STORAGE_ACCOUNT_KEY", "some-value")
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = "myaccount"
+        mock_conn.password = None
+        mock_conn.extra_dejson = {}
+        engine = DataFusionEngine()
+
+        credentials, extra_config = engine._get_credentials(mock_conn)
+
+        assert credentials == {"account": "myaccount"}
+        assert extra_config == {}
+
+    @pytest.mark.parametrize(
+        ("host", "login", "expected"),
+        [
+            ("myaccount.blob.core.windows.net", None, "myaccount"),
+            ("https://myaccount.blob.core.windows.net/container/path";, None, 
"myaccount"),
+            (None, "myaccount", "myaccount"),
+            ("justanadid", "myaccount", "myaccount"),
+            ("a" * 40 + ".blob.core.windows.net", "x", "a" * 24),
+            (None, None, None),
+        ],
+    )
+    def test_resolve_wasb_account(self, host, login, expected):
+        assert DataFusionEngine._resolve_wasb_account(host, login) == expected
+
+    @pytest.mark.parametrize(
+        "host",
+        [
+            "myaccount.blob.core.chinacloudapi.cn",
+            "myaccount.blob.core.usgovcloudapi.net",
+        ],
+    )
+    def test_resolve_wasb_account_rejects_non_public_cloud_hosts(self, host):
+        """The binding has no endpoint override, so a sovereign-cloud host 
would otherwise be
+        silently misrouted to the public *.blob.core.windows.net account of 
the same name."""
+        with pytest.raises(ValueError, match="does not resolve to the public"):
+            DataFusionEngine._resolve_wasb_account(host, None)
+
+    @pytest.mark.parametrize("env_var", ["AZURE_STORAGE_ENDPOINT", 
"AZURE_ENDPOINT"])
+    def 
test_resolve_wasb_account_allows_non_public_cloud_host_when_endpoint_set(self, 
env_var, monkeypatch):
+        """Once the worker sets an endpoint override, object_store uses it 
verbatim instead of
+        deriving the URL from the account name, so a real (dotted, portless) 
sovereign-cloud
+        hostname no longer needs to match *.blob.core.windows.net."""
+        monkeypatch.setenv(env_var, 
"https://myaccount.blob.core.chinacloudapi.cn";)
+
+        account = 
DataFusionEngine._resolve_wasb_account("myaccount.blob.core.chinacloudapi.cn", 
None)
+
+        assert account == "myaccount"
+
+    @pytest.mark.parametrize(
+        ("host", "login"),
+        [
+            ("http://127.0.0.1:10000/devstoreaccount1";, None),
+            ("azurite:10000", None),
+            ("azurite:10000", "devstoreaccount1"),
+        ],
+        ids=["ip-port-url", "dotless-host-port", 
"dotless-host-port-with-login"],
+    )
+    def test_resolve_wasb_account_rejects_emulator_style_hosts(self, host, 
login, monkeypatch):
+        """A host:port netloc (Azurite's shape) is never a real DNS hostname 
-- resolving the
+        account from it would otherwise silently produce the wrong value 
("127", or "None" once
+        the dotless fallback for an Active Directory ID kicks in), so this 
must always raise
+        regardless of login or of the endpoint override, which does not help 
this shape either."""
+        monkeypatch.setenv("AZURE_STORAGE_ENDPOINT", host.split("/")[0])
+
+        with pytest.raises(ValueError, match="looks like an emulator address"):
+            DataFusionEngine._resolve_wasb_account(host, login)
+
+    def test_resolve_wasb_account_rejects_dotless_host_without_login(self):
+        """The Active Directory ID fallback assumes login holds the real 
account name; without
+        one there's nothing to build a netloc from, and silently using the 
literal string "None"
+        is exactly the bug this function's own empty-connection case is meant 
to avoid."""
+        with pytest.raises(ValueError, match="no `login` was given"):
+            DataFusionEngine._resolve_wasb_account("azurite", None)
+
+    def 
test_resolve_wasb_account_azurite_workaround_uses_login_with_empty_host(self):
+        """The documented Azurite workaround -- the account name in login, 
host left empty --
+        already resolves correctly without needing the non-public-suffix check 
at all."""
+        assert DataFusionEngine._resolve_wasb_account(None, 
"devstoreaccount1") == "devstoreaccount1"
+
     def test_get_credentials_unknown_type(self):
         mock_conn = MagicMock()
         mock_conn.conn_type = "dummy"
diff --git 
a/providers/common/sql/tests/unit/common/sql/datafusion/test_object_storage_provider.py
 
b/providers/common/sql/tests/unit/common/sql/datafusion/test_object_storage_provider.py
index 4370d670dde..308879489f7 100644
--- 
a/providers/common/sql/tests/unit/common/sql/datafusion/test_object_storage_provider.py
+++ 
b/providers/common/sql/tests/unit/common/sql/datafusion/test_object_storage_provider.py
@@ -25,6 +25,7 @@ import pytest
 from airflow.providers.common.sql.config import ConnectionConfig, StorageType
 from airflow.providers.common.sql.datafusion.exceptions import 
ObjectStoreCreationException
 from airflow.providers.common.sql.datafusion.object_storage_provider import (
+    AzureObjectStorageProvider,
     GCSObjectStorageProvider,
     LocalObjectStorageProvider,
     S3ObjectStorageProvider,
@@ -141,6 +142,52 @@ class TestObjectStorageProvider:
         with pytest.raises(ValueError, match="connection_config must be 
provided for gcs"):
             provider.create_object_store("gs://demo-data/path")
 
+    
@patch("airflow.providers.common.sql.datafusion.object_storage_provider.MicrosoftAzure")
+    def test_azure_provider_success(self, mock_azure):
+        provider = AzureObjectStorageProvider()
+        connection_config = ConnectionConfig(
+            conn_id="wasb_default",
+            credentials={"account": "myaccount", "access_key": "fake_key"},
+        )
+
+        store = provider.create_object_store("az://demo-container/path", 
connection_config)
+
+        mock_azure.assert_called_once_with(
+            container_name="demo-container", account="myaccount", 
access_key="fake_key"
+        )
+        assert store == mock_azure.return_value
+        assert provider.get_storage_type == StorageType.AZURE
+        assert provider.get_scheme() == "az://"
+
+    def test_azure_provider_failure(self):
+        provider = AzureObjectStorageProvider()
+        connection_config = ConnectionConfig(conn_id="wasb_default")
+
+        with patch(
+            
"airflow.providers.common.sql.datafusion.object_storage_provider.MicrosoftAzure",
+            side_effect=Exception("Error"),
+        ):
+            with pytest.raises(ObjectStoreCreationException, match="Failed to 
create Azure object store"):
+                provider.create_object_store("az://demo-container/path", 
connection_config)
+
+    def test_azure_provider_requires_connection_config(self):
+        provider = AzureObjectStorageProvider()
+
+        with pytest.raises(ValueError, match="connection_config must be 
provided for azure"):
+            provider.create_object_store("az://demo-container/path")
+
+    def test_azure_provider_partial_service_principal_raises_clear_error(self):
+        """Uses the real MicrosoftAzure binding, not a mock, since it's the 
one that panics
+        on a partial client_id/client_secret/tenant_id combination."""
+        provider = AzureObjectStorageProvider()
+        connection_config = ConnectionConfig(
+            conn_id="wasb_default",
+            credentials={"client_id": "only-this-one-set"},
+        )
+
+        with pytest.raises(ObjectStoreCreationException, match="Failed to 
create Azure object store"):
+            provider.create_object_store("az://demo-container/path", 
connection_config)
+
     
@patch("airflow.providers.common.sql.datafusion.object_storage_provider.LocalFileSystem")
     def test_local_provider(self, mock_local):
         provider = LocalObjectStorageProvider()
@@ -152,6 +199,7 @@ class TestObjectStorageProvider:
     def test_get_object_storage_provider(self):
         assert isinstance(get_object_storage_provider(StorageType.S3), 
S3ObjectStorageProvider)
         assert isinstance(get_object_storage_provider(StorageType.GCS), 
GCSObjectStorageProvider)
+        assert isinstance(get_object_storage_provider(StorageType.AZURE), 
AzureObjectStorageProvider)
         assert isinstance(get_object_storage_provider(StorageType.LOCAL), 
LocalObjectStorageProvider)
 
         with pytest.raises(ValueError, match="Unsupported storage type"):
diff --git a/providers/common/sql/tests/unit/common/sql/test_config.py 
b/providers/common/sql/tests/unit/common/sql/test_config.py
index 14f3827c6a0..359721681fd 100644
--- a/providers/common/sql/tests/unit/common/sql/test_config.py
+++ b/providers/common/sql/tests/unit/common/sql/test_config.py
@@ -34,6 +34,7 @@ class TestDataSourceConfig:
         [
             ("s3://bucket/path", StorageType.S3),
             ("gs://bucket/path", StorageType.GCS),
+            ("az://container/path", StorageType.AZURE),
             ("file:///path/to/file", StorageType.LOCAL),
         ],
     )

Reply via email to