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

guan404ming 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 832d0a5506a Support cloud URI globs in DocumentLoaderOperator (#70299)
832d0a5506a is described below

commit 832d0a5506a5ecb4dddb47fac8a5c2fb072b221a
Author: Guan-Ming Chiu <[email protected]>
AuthorDate: Fri Jul 24 10:05:45 2026 +0800

    Support cloud URI globs in DocumentLoaderOperator (#70299)
---
 .../common/ai/docs/operators/document_loader.rst   |  9 +++---
 .../common/ai/operators/document_loader.py         | 32 +++++++++++++++------
 .../common/ai/operators/test_document_loader.py    | 33 ++++++++++++++++++++++
 3 files changed, 61 insertions(+), 13 deletions(-)

diff --git a/providers/common/ai/docs/operators/document_loader.rst 
b/providers/common/ai/docs/operators/document_loader.rst
index 2aa32e6594d..843d059e282 100644
--- a/providers/common/ai/docs/operators/document_loader.rst
+++ b/providers/common/ai/docs/operators/document_loader.rst
@@ -195,8 +195,9 @@ Cloud storage URIs
 ``source_path`` accepts any URI that
 :class:`~airflow.sdk.ObjectStoragePath` resolves via fsspec
 (``s3://``, ``gs://``, ``azure://``, ``file://``, ...). Point it at a
-single object or a directory; cross-directory globs in cloud URIs are not
-supported in this version.
+single object, a directory, or a glob pattern such as
+``s3://bucket/logs/**/*.json``. Wildcards must not appear in the scheme or
+bucket segment.
 
 .. exampleinclude:: 
/../../ai/src/airflow/providers/common/ai/example_dags/example_document_loader.py
     :language: python
@@ -261,8 +262,8 @@ Parameters
      - Local file, directory, or glob pattern, **or** a storage URI
        (``s3://``, ``gs://``, ``azure://``, ``file://``) resolved via
        :class:`~airflow.sdk.ObjectStoragePath`. ``**`` is recursive for
-       local globs; cross-directory globs in cloud URIs are not supported.
-       Mutually exclusive with ``source_bytes``.
+       both local globs and cloud URIs; wildcards must not appear in the
+       scheme or bucket segment. Mutually exclusive with ``source_bytes``.
    * - ``source_conn_id``
      - Airflow connection ID for the cloud-storage credentials used by
        ``ObjectStoragePath`` (``aws_default``, ``google_cloud_default``,
diff --git 
a/providers/common/ai/src/airflow/providers/common/ai/operators/document_loader.py
 
b/providers/common/ai/src/airflow/providers/common/ai/operators/document_loader.py
index b452cab72ed..1c8c6a935b3 100644
--- 
a/providers/common/ai/src/airflow/providers/common/ai/operators/document_loader.py
+++ 
b/providers/common/ai/src/airflow/providers/common/ai/operators/document_loader.py
@@ -68,9 +68,9 @@ class DocumentLoaderOperator(BaseOperator):
     :param source_path: A local path, glob pattern, or storage URI
         (``s3://``, ``gs://``, ``azure://``, ``file://``, ...). Cloud URIs
         go through :class:`~airflow.sdk.ObjectStoragePath` / fsspec.
-        ``**`` enables recursive matching for local globs. Cloud URIs
-        accept a single file or a directory; cross-directory globs in a
-        cloud URI are not supported in this version.
+        ``**`` enables recursive matching, for both local paths and cloud
+        URIs. Wildcards must not appear in the scheme or bucket segment
+        of a cloud URI.
     :param source_conn_id: Airflow connection ID used by
         ``ObjectStoragePath`` for cloud URIs (``aws_default``,
         ``google_cloud_default``, ...). Ignored for local paths.
@@ -208,7 +208,10 @@ class DocumentLoaderOperator(BaseOperator):
         return self._filter_files([p for p in candidates if p.is_file()], 
is_directory_mode=is_directory_mode)
 
     def _resolve_remote_files(self, source_path: str) -> list[FilePathT]:
-        from airflow.sdk import ObjectStoragePath
+        from airflow.providers.common.compat.sdk import ObjectStoragePath
+
+        if any(char in source_path for char in "*?["):
+            return self._resolve_remote_glob(source_path)
 
         root = ObjectStoragePath(source_path, conn_id=self.source_conn_id)
         try:
@@ -219,11 +222,7 @@ class DocumentLoaderOperator(BaseOperator):
             pass
 
         if not root.is_dir():
-            raise FileNotFoundError(
-                f"Cloud URI '{source_path}' is neither a file nor a directory. 
"
-                "Cross-directory globs in cloud URIs aren't supported here; "
-                "point ``source_path`` at a single object or a directory."
-            )
+            raise FileNotFoundError(f"Cloud URI '{source_path}' is neither a 
file nor a directory.")
 
         candidates = sorted(
             (p for p in root.iterdir() if not p.name.startswith(".")),
@@ -231,6 +230,21 @@ class DocumentLoaderOperator(BaseOperator):
         )
         return self._filter_files([p for p in candidates if p.is_file()], 
is_directory_mode=True)
 
+    def _resolve_remote_glob(self, source_path: str) -> list[FilePathT]:
+        from airflow.providers.common.compat.sdk import ObjectStoragePath
+
+        segments = source_path.split("/")
+        magic_at = next(i for i, seg in enumerate(segments) if any(c in seg 
for c in "*?["))
+        # segments[:3] is ``scheme:``, ``""``, ``bucket``; a wildcard there 
has no fixed root.
+        if magic_at < 3:
+            raise ValueError(
+                f"Cloud URI '{source_path}' must not use wildcards in the 
scheme or bucket segment."
+            )
+
+        root = ObjectStoragePath("/".join(segments[:magic_at]), 
conn_id=self.source_conn_id)
+        candidates = sorted(root.glob("/".join(segments[magic_at:])), key=str)
+        return self._filter_files([p for p in candidates if p.is_file()], 
is_directory_mode=False)
+
     def _filter_files(self, results: list[FilePathT], *, is_directory_mode: 
bool) -> list[FilePathT]:
         if self.file_extensions:
             allowed = {(ext if ext.startswith(".") else f".{ext}").lower() for 
ext in self.file_extensions}
diff --git 
a/providers/common/ai/tests/unit/common/ai/operators/test_document_loader.py 
b/providers/common/ai/tests/unit/common/ai/operators/test_document_loader.py
index 6681c0f463d..0a949c19010 100644
--- a/providers/common/ai/tests/unit/common/ai/operators/test_document_loader.py
+++ b/providers/common/ai/tests/unit/common/ai/operators/test_document_loader.py
@@ -494,6 +494,39 @@ class TestCloudUriDispatch:
         with pytest.raises(FileNotFoundError, match="neither a file nor a 
directory"):
             op.execute(context=MagicMock())
 
+    @patch("airflow.sdk.ObjectStoragePath")
+    def test_glob_uri_matches_across_directories(self, mock_osp_cls):
+        def _mock_match(name: str, content: bytes):
+            match = MagicMock()
+            match.is_file.return_value = True
+            match.name = name
+            match.suffix = "." + name.rsplit(".", 1)[-1]
+            match.read_bytes.return_value = content
+            return match
+
+        root = MagicMock()
+        root.glob.return_value = [_mock_match("a.txt", b"alpha"), 
_mock_match("b.txt", b"beta")]
+        mock_osp_cls.return_value = root
+
+        op = DocumentLoaderOperator(
+            task_id="test",
+            source_path="s3://bucket/logs/**/*.txt",
+            source_conn_id="aws_default",
+        )
+        result = op.execute(context=MagicMock())
+
+        mock_osp_cls.assert_called_once_with("s3://bucket/logs", 
conn_id="aws_default")
+        root.glob.assert_called_once_with("**/*.txt")
+        assert {doc["text"] for doc in result} == {"alpha", "beta"}
+
+    @patch("airflow.sdk.ObjectStoragePath")
+    def test_glob_in_bucket_segment_raises(self, mock_osp_cls):
+        op = DocumentLoaderOperator(task_id="test", 
source_path="s3://bucket-*/dir/a.txt")
+        with pytest.raises(ValueError, match="scheme or bucket segment"):
+            op.execute(context=MagicMock())
+
+        mock_osp_cls.assert_not_called()
+
 
 class TestEncoding:
     def test_strict_utf8_default_raises_with_path_context(self, tmp_path):

Reply via email to