kaxil commented on code in PR #74228:
URL: https://github.com/apache/airflow/pull/74228#discussion_r4195350446


##########
providers/common/sql/src/airflow/providers/common/sql/datafusion/engine.py:
##########
@@ -80,9 +80,13 @@ def _register_object_store(
             object_store = storage_provider.create_object_store(
                 datasource_config.uri, connection_config=connection_config
             )
-            schema = storage_provider.get_scheme()
-            self.session_context.register_object_store(schema=schema, 
store=object_store)
-            self.log.info("Registered object store for schema: %s", schema)
+            schema = storage_provider.get_scheme(datasource_config.uri)
+            # DataFusion's object-store registry keys on (schema, host); 
omitting host only
+            # matches URIs with an empty authority (e.g. file:///path), so a 
bucket/container
+            # URI's netloc must be passed explicitly or lookup fails at query 
time.
+            host = urlsplit(datasource_config.uri).netloc

Review Comment:
   DataFusion's registry key drops the userinfo part of the URL, so 
`abfss://[email protected]` and 
`abfss://[email protected]` both register under 
`abfss://acct.dfs.core.windows.net` and the second store replaces the first. 
Each `MicrosoftAzure` store is bound to one container, so with two abfs 
datasources on the same account (bronze and silver in one `AnalyticsOperator`, 
or `LLMSchemaCompareOperator`) the first table silently reads from the second 
container. I reproduced it on datafusion 50 and 51 with two 
`LocalFileSystem(prefix=...)` stores: `select * from t1` returns cont2's rows. 
`az://` doesn't hit this because its host is the container.
   
   Could we rewrite `abfs(s)://<container>@<account>.<host>/<path>` to 
`az://<container>/<path>` after the account check, and register both the store 
and the table under that URI? A two-container test would pin it.



##########
providers/common/sql/src/airflow/providers/common/sql/datafusion/engine.py:
##########
@@ -80,9 +80,13 @@ def _register_object_store(
             object_store = storage_provider.create_object_store(
                 datasource_config.uri, connection_config=connection_config
             )
-            schema = storage_provider.get_scheme()
-            self.session_context.register_object_store(schema=schema, 
store=object_store)
-            self.log.info("Registered object store for schema: %s", schema)
+            schema = storage_provider.get_scheme(datasource_config.uri)

Review Comment:
   This breaks an explicit `storage_type=StorageType.LOCAL` with a bare path 
like `/data/x.csv`. In 2.2.0 the zero-arg `get_scheme()` always returned 
`file://` for the local provider, so that config registered and queried fine. 
Now `get_scheme("/data/x.csv")` matches nothing in `SCHEMES` and registration 
fails with `does not match any known scheme`. `_extract_storage_type` rejects a 
bare path, so the explicit override is the only way users reach this. Could 
`LocalObjectStorageProvider` keep returning `"file://"` regardless of the URI?



##########
providers/common/sql/src/airflow/providers/common/sql/datafusion/object_storage_provider.py:
##########
@@ -124,6 +144,16 @@ def create_object_store(self, path: str, 
connection_config: ConnectionConfig | N
             credentials = connection_config.credentials
             container = self.get_bucket(path)
 
+            uri_account = self._get_uri_account(path)
+            resolved_account = credentials.get("account")
+            if uri_account and resolved_account and uri_account.lower() != 
resolved_account.lower():

Review Comment:
   When the wasb connection has no host or login, `_resolve_wasb_account` 
returns None and `_remove_none_values` drops `account`, so this check is 
skipped and `MicrosoftAzure` falls back to `AZURE_STORAGE_ACCOUNT_NAME`. With 
that set to `envacct`, `abfss://[email protected]/data.parquet` 
sends its request to `https://envacct.blob.core.windows.net/cont/data.parquet`, 
which is the "silently picking one" case the docs say raises. Passing 
`account=uri_account` when the connection resolves none would fix it (an 
explicit `account` kwarg wins over the env var, I checked). Could you add a 
test with an account-less connection as well?



##########
providers/common/sql/tests/unit/common/sql/datafusion/test_engine.py:
##########
@@ -235,6 +238,44 @@ def test_execute_query_with_local_csv(self, mock_get_conn):
         finally:
             os.unlink(csv_path)
 
+    @patch.object(DataFusionEngine, "_get_connection_config")
+    def test_execute_query_with_bucket_style_uri_matches_real_registry(self, 
mock_get_conn):

Review Comment:
   This fails without `host` only because a `LocalFileSystem` has no bucket for 
DataFusion to default the host from; a real `AmazonS3` store under `s3://` 
resolves without it. So the docstring describes a failure production didn't 
have, and nothing in the suite sends an abfs URI through the real registry, 
which is the one scheme that needs `host`. Could this use an 
`abfss://[email protected]{csv_path}` URI instead (a 
`LocalFileSystem` registered with that netloc reads it fine), plus the 
two-container case? Small one: the `patch(...)` below wants `autospec=True`.



##########
providers/common/sql/src/airflow/providers/common/sql/datafusion/engine.py:
##########
@@ -80,9 +80,13 @@ def _register_object_store(
             object_store = storage_provider.create_object_store(
                 datasource_config.uri, connection_config=connection_config
             )
-            schema = storage_provider.get_scheme()
-            self.session_context.register_object_store(schema=schema, 
store=object_store)
-            self.log.info("Registered object store for schema: %s", schema)
+            schema = storage_provider.get_scheme(datasource_config.uri)
+            # DataFusion's object-store registry keys on (schema, host); 
omitting host only

Review Comment:
   I don't think this holds for s3/gs/az. When `host` is omitted, 
datafusion-python defaults it to the store's own bucket or container name, so 
`s3://bucket/...` and `az://container/...` already resolved before this change 
(a real `AmazonS3` or `MicrosoftAzure` store with no host finds the store and 
fails later on credentials). Only abfs(s) needs it, because its netloc is 
`container@account.<host>` rather than the container. Could the comment say 
that?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to