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]