This is an automated email from the ASF dual-hosted git repository.
kevinjqliu pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iceberg-python.git
The following commit(s) were added to refs/heads/main by this push:
new 84c91f07 Support ADLS with Pyarrow file IO (#2111)
84c91f07 is described below
commit 84c91f078a9e165acb20350463fe42bb98b7a65a
Author: NikitaMatskevich <[email protected]>
AuthorDate: Fri Jun 20 17:11:08 2025 +0200
Support ADLS with Pyarrow file IO (#2111)
<!--
Thanks for opening a pull request!
-->
<!-- In the case this PR will resolve an issue, please replace
${GITHUB_ISSUE_ID} below with the actual Github issue id. -->
<!-- Closes #${GITHUB_ISSUE_ID} -->
# Rationale for this change
Starting from version 20, PyArrow supports ADLS filesystem. This PR adds
Pyarrow Azure support to Pyiceberg.
PyArrow is the [default
IO](https://github.com/apache/iceberg-python/blob/main/pyiceberg/io/__init__.py#L366-L369)
for Pyiceberg catalogs. In Azure environment it handles wider spectrum
of auth strategies then Fsspec, including, for instance, [Managed
Identities](https://learn.microsoft.com/en-us/entra/identity/managed-identities-azure-resources/overview).
Also, prior to this PR
https://github.com/apache/iceberg-python/pull/1663 (that is not merged
yet) there was no support for wasb(s) with Fsspec.
See the corresponding issue for more details:
https://github.com/apache/iceberg-python/issues/2112
# Are these changes tested?
Tests are added under tests/io/test_pyarrow.py.
# Are there any user-facing changes?
There are no API breaking changes. Direct impact of the PR: Pyarrow
FileIO in Pyiceberg supports Azure cloud environment. Examples of impact
for final users:
- Pyiceberg is usable in services with Managed Identities auth strategy.
- Pyiceberg is usable with wasb(s) schemes in Azure.
<!-- In the case of user-facing changes, please add the changelog label.
-->
---------
Co-authored-by: Kevin Liu <[email protected]>
Co-authored-by: Kevin Liu <[email protected]>
---
mkdocs/docs/configuration.md | 24 +++---
pyiceberg/io/__init__.py | 4 +
pyiceberg/io/pyarrow.py | 47 ++++++++++
tests/conftest.py | 79 ++++++++++++-----
tests/io/test_pyarrow.py | 201 ++++++++++++++++++++++++++++++++++++++++++-
5 files changed, 321 insertions(+), 34 deletions(-)
diff --git a/mkdocs/docs/configuration.md b/mkdocs/docs/configuration.md
index 3037e0f5..539fa482 100644
--- a/mkdocs/docs/configuration.md
+++ b/mkdocs/docs/configuration.md
@@ -146,16 +146,20 @@ For the FileIO there are several configuration options
available:
<!-- markdown-link-check-disable -->
-| Key | Example
| Description
|
-| ---------------------- |
-----------------------------------------------------------------------------------------
|
--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
|
-| adls.connection-string |
AccountName=devstoreaccount1;AccountKey=Eby8vdM02xNOcqF...;BlobEndpoint=<http://localhost/>
| A [connection
string](https://learn.microsoft.com/en-us/azure/storage/common/storage-configure-connection-string).
This could be used to use FileIO with any adls-compatible object storage
service that has a different endpoint (like
[azurite](https://github.com/azure/azurite)). |
-| adls.account-name | devstoreaccount1
| The account that you want to connect to
|
-| adls.account-key | Eby8vdM02xNOcqF...
| The key to authentication against the
account.
|
-| adls.sas-token | NuHOuuzdQN7VRM%2FOpOeqBlawRCA845IY05h9eu1Yte4%3D
| The shared access signature
|
-| adls.tenant-id | ad667be4-b811-11ed-afa1-0242ac120002
| The tenant-id
|
-| adls.client-id | ad667be4-b811-11ed-afa1-0242ac120002
| The client-id
|
-| adls.client-secret | oCA3R6P\*ka#oa1Sms2J74z...
| The client-secret
|
-| adls.account-host | accountname1.blob.core.windows.net
| The storage account host. See
[AzureBlobFileSystem](https://github.com/fsspec/adlfs/blob/adb9c53b74a0d420625b86dd00fbe615b43201d2/adlfs/spec.py#L125)
for reference
|
+| Key | Example
| Description
|
+|------------------------------|---------------------------------------------------------------------------------------------|---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|
+| adls.connection-string |
AccountName=devstoreaccount1;AccountKey=Eby8vdM02xNOcqF...;BlobEndpoint=<http://localhost/>
| A [connection
string](https://learn.microsoft.com/en-us/azure/storage/common/storage-configure-connection-string).
This could be used to use FileIO with any adls-compatible object storage
service that has a different endpoint (like
[azurite](https://github.com/azure/azurite)). |
+| adls.account-name | devstoreaccount1
| The account that you want to
connect to
|
+| adls.account-key | Eby8vdM02xNOcqF...
| The key to authentication
against the account.
|
+| adls.sas-token |
NuHOuuzdQN7VRM%2FOpOeqBlawRCA845IY05h9eu1Yte4%3D
| The shared access signature
|
+| adls.tenant-id | ad667be4-b811-11ed-afa1-0242ac120002
| The tenant-id
|
+| adls.client-id | ad667be4-b811-11ed-afa1-0242ac120002
| The client-id
|
+| adls.client-secret | oCA3R6P\*ka#oa1Sms2J74z...
| The client-secret
|
+| adls.account-host | accountname1.blob.core.windows.net
| The storage account host. See
[AzureBlobFileSystem](https://github.com/fsspec/adlfs/blob/adb9c53b74a0d420625b86dd00fbe615b43201d2/adlfs/spec.py#L125)
for reference
|
+| adls.blob-storage-authority | .blob.core.windows.net
| The hostname[:port] of the Blob
Service. Defaults to `.blob.core.windows.net`. Useful for connecting to a local
emulator, like [azurite](https://github.com/azure/azurite). See
[AzureFileSystem](https://arrow.apache.org/docs/python/filesystems.html#azure-storage-file-system)
for reference |
+| adls.dfs-storage-authority | .dfs.core.windows.net
| The hostname[:port] of the Data
Lake Gen 2 Service. Defaults to `.dfs.core.windows.net`. Useful for connecting
to a local emulator, like [azurite](https://github.com/azure/azurite). See
[AzureFileSystem](https://arrow.apache.org/docs/python/filesystems.html#azure-storage-file-system)
for reference |
+| adls.blob-storage-scheme | https
| Either `http` or `https`.
Defaults to `https`. Useful for connecting to a local emulator, like
[azurite](https://github.com/azure/azurite). See
[AzureFileSystem](https://arrow.apache.org/docs/python/filesystems.html#azure-storage-file-system)
for reference |
+| adls.dfs-storage-scheme | https
| Either `http` or `https`.
Defaults to `https`. Useful for connecting to a local emulator, like
[azurite](https://github.com/azure/azurite). See
[AzureFileSystem](https://arrow.apache.org/docs/python/filesystems.html#azure-storage-file-system)
for reference |
<!-- markdown-link-check-enable-->
diff --git a/pyiceberg/io/__init__.py b/pyiceberg/io/__init__.py
index ba973807..2a73aba7 100644
--- a/pyiceberg/io/__init__.py
+++ b/pyiceberg/io/__init__.py
@@ -80,6 +80,10 @@ ADLS_TENANT_ID = "adls.tenant-id"
ADLS_CLIENT_ID = "adls.client-id"
ADLS_CLIENT_SECRET = "adls.client-secret"
ADLS_ACCOUNT_HOST = "adls.account-host"
+ADLS_BLOB_STORAGE_AUTHORITY = "adls.blob-storage-authority"
+ADLS_DFS_STORAGE_AUTHORITY = "adls.dfs-storage-authority"
+ADLS_BLOB_STORAGE_SCHEME = "adls.blob-storage-scheme"
+ADLS_DFS_STORAGE_SCHEME = "adls.dfs-storage-scheme"
GCS_TOKEN = "gcs.oauth2.token"
GCS_TOKEN_EXPIRES_AT_MS = "gcs.oauth2.token-expires-at"
GCS_PROJECT_ID = "gcs.project-id"
diff --git a/pyiceberg/io/pyarrow.py b/pyiceberg/io/pyarrow.py
index eeacaccf..c23be93e 100644
--- a/pyiceberg/io/pyarrow.py
+++ b/pyiceberg/io/pyarrow.py
@@ -84,6 +84,13 @@ from pyiceberg.expressions.visitors import (
)
from pyiceberg.expressions.visitors import visit as boolean_expression_visit
from pyiceberg.io import (
+ ADLS_ACCOUNT_KEY,
+ ADLS_ACCOUNT_NAME,
+ ADLS_BLOB_STORAGE_AUTHORITY,
+ ADLS_BLOB_STORAGE_SCHEME,
+ ADLS_DFS_STORAGE_AUTHORITY,
+ ADLS_DFS_STORAGE_SCHEME,
+ ADLS_SAS_TOKEN,
AWS_ACCESS_KEY_ID,
AWS_REGION,
AWS_ROLE_ARN,
@@ -390,6 +397,9 @@ class PyArrowFileIO(FileIO):
elif scheme in {"gs", "gcs"}:
return self._initialize_gcs_fs()
+ elif scheme in {"abfs", "abfss", "wasb", "wasbs"}:
+ return self._initialize_azure_fs()
+
elif scheme in {"file"}:
return self._initialize_local_fs()
@@ -471,6 +481,43 @@ class PyArrowFileIO(FileIO):
return S3FileSystem(**client_kwargs)
+ def _initialize_azure_fs(self) -> FileSystem:
+ from packaging import version
+
+ MIN_PYARROW_VERSION_SUPPORTING_AZURE_FS = "20.0.0"
+ if version.parse(pyarrow.__version__) <
version.parse(MIN_PYARROW_VERSION_SUPPORTING_AZURE_FS):
+ raise ImportError(
+ f"pyarrow version >= {MIN_PYARROW_VERSION_SUPPORTING_AZURE_FS}
required for AzureFileSystem support, "
+ f"but found version {pyarrow.__version__}."
+ )
+
+ from pyarrow.fs import AzureFileSystem
+
+ client_kwargs: Dict[str, str] = {}
+
+ if account_name := self.properties.get(ADLS_ACCOUNT_NAME):
+ client_kwargs["account_name"] = account_name
+
+ if account_key := self.properties.get(ADLS_ACCOUNT_KEY):
+ client_kwargs["account_key"] = account_key
+
+ if blob_storage_authority :=
self.properties.get(ADLS_BLOB_STORAGE_AUTHORITY):
+ client_kwargs["blob_storage_authority"] = blob_storage_authority
+
+ if dfs_storage_authority :=
self.properties.get(ADLS_DFS_STORAGE_AUTHORITY):
+ client_kwargs["dfs_storage_authority"] = dfs_storage_authority
+
+ if blob_storage_scheme :=
self.properties.get(ADLS_BLOB_STORAGE_SCHEME):
+ client_kwargs["blob_storage_scheme"] = blob_storage_scheme
+
+ if dfs_storage_scheme := self.properties.get(ADLS_DFS_STORAGE_SCHEME):
+ client_kwargs["dfs_storage_scheme"] = dfs_storage_scheme
+
+ if sas_token := self.properties.get(ADLS_SAS_TOKEN):
+ client_kwargs["sas_token"] = sas_token
+
+ return AzureFileSystem(**client_kwargs)
+
def _initialize_hdfs_fs(self, scheme: str, netloc: Optional[str]) ->
FileSystem:
from pyarrow.fs import HadoopFileSystem
diff --git a/tests/conftest.py b/tests/conftest.py
index 729e29cb..a584f98c 100644
--- a/tests/conftest.py
+++ b/tests/conftest.py
@@ -52,6 +52,12 @@ from pyiceberg.catalog import Catalog, load_catalog
from pyiceberg.catalog.noop import NoopCatalog
from pyiceberg.expressions import BoundReference
from pyiceberg.io import (
+ ADLS_ACCOUNT_KEY,
+ ADLS_ACCOUNT_NAME,
+ ADLS_BLOB_STORAGE_AUTHORITY,
+ ADLS_BLOB_STORAGE_SCHEME,
+ ADLS_DFS_STORAGE_AUTHORITY,
+ ADLS_DFS_STORAGE_SCHEME,
GCS_PROJECT_ID,
GCS_SERVICE_HOST,
GCS_TOKEN,
@@ -348,6 +354,11 @@ def table_schema_with_all_types() -> Schema:
)
[email protected](params=["abfs", "abfss", "wasb", "wasbs"])
+def adls_scheme(request: pytest.FixtureRequest) -> str:
+ return request.param
+
+
@pytest.fixture(scope="session")
def pyarrow_schema_simple_without_ids() -> "pa.Schema":
import pyarrow as pa
@@ -2088,6 +2099,26 @@ def fsspec_fileio_gcs(request: pytest.FixtureRequest) ->
FsspecFileIO:
return fsspec.FsspecFileIO(properties=properties)
[email protected]
+def adls_fsspec_fileio(request: pytest.FixtureRequest) ->
Generator[FsspecFileIO, None, None]:
+ from azure.storage.blob import BlobServiceClient
+
+ azurite_url = request.config.getoption("--adls.endpoint")
+ azurite_account_name = request.config.getoption("--adls.account-name")
+ azurite_account_key = request.config.getoption("--adls.account-key")
+ azurite_connection_string =
f"DefaultEndpointsProtocol=http;AccountName={azurite_account_name};AccountKey={azurite_account_key};BlobEndpoint={azurite_url}/{azurite_account_name};"
+ properties = {
+ "adls.connection-string": azurite_connection_string,
+ "adls.account-name": azurite_account_name,
+ }
+
+ bbs =
BlobServiceClient.from_connection_string(conn_str=azurite_connection_string)
+ bbs.create_container("tests")
+ yield fsspec.FsspecFileIO(properties=properties)
+ bbs.delete_container("tests")
+ bbs.close()
+
+
@pytest.fixture
def pyarrow_fileio_gcs(request: pytest.FixtureRequest) -> "PyArrowFileIO":
from pyiceberg.io.pyarrow import PyArrowFileIO
@@ -2101,6 +2132,34 @@ def pyarrow_fileio_gcs(request: pytest.FixtureRequest)
-> "PyArrowFileIO":
return PyArrowFileIO(properties=properties)
[email protected]
+def pyarrow_fileio_adls(request: pytest.FixtureRequest) -> Generator[Any,
None, None]:
+ from azure.storage.blob import BlobServiceClient
+
+ from pyiceberg.io.pyarrow import PyArrowFileIO
+
+ azurite_url = request.config.getoption("--adls.endpoint")
+ azurite_scheme, azurite_authority = azurite_url.split("://", 1)
+
+ azurite_account_name = request.config.getoption("--adls.account-name")
+ azurite_account_key = request.config.getoption("--adls.account-key")
+ azurite_connection_string =
f"DefaultEndpointsProtocol=http;AccountName={azurite_account_name};AccountKey={azurite_account_key};BlobEndpoint={azurite_url}/{azurite_account_name};"
+ properties = {
+ ADLS_ACCOUNT_NAME: azurite_account_name,
+ ADLS_ACCOUNT_KEY: azurite_account_key,
+ ADLS_BLOB_STORAGE_AUTHORITY: azurite_authority,
+ ADLS_DFS_STORAGE_AUTHORITY: azurite_authority,
+ ADLS_BLOB_STORAGE_SCHEME: azurite_scheme,
+ ADLS_DFS_STORAGE_SCHEME: azurite_scheme,
+ }
+
+ bbs =
BlobServiceClient.from_connection_string(conn_str=azurite_connection_string)
+ bbs.create_container("warehouse")
+ yield PyArrowFileIO(properties=properties)
+ bbs.delete_container("warehouse")
+ bbs.close()
+
+
def aws_credentials() -> None:
os.environ["AWS_ACCESS_KEY_ID"] = "testing"
os.environ["AWS_SECRET_ACCESS_KEY"] = "testing"
@@ -2162,26 +2221,6 @@ def fixture_dynamodb(_aws_credentials: None) ->
Generator[boto3.client, None, No
yield boto3.client("dynamodb", region_name="us-east-1")
[email protected]
-def adls_fsspec_fileio(request: pytest.FixtureRequest) ->
Generator[FsspecFileIO, None, None]:
- from azure.storage.blob import BlobServiceClient
-
- azurite_url = request.config.getoption("--adls.endpoint")
- azurite_account_name = request.config.getoption("--adls.account-name")
- azurite_account_key = request.config.getoption("--adls.account-key")
- azurite_connection_string =
f"DefaultEndpointsProtocol=http;AccountName={azurite_account_name};AccountKey={azurite_account_key};BlobEndpoint={azurite_url}/{azurite_account_name};"
- properties = {
- "adls.connection-string": azurite_connection_string,
- "adls.account-name": azurite_account_name,
- }
-
- bbs =
BlobServiceClient.from_connection_string(conn_str=azurite_connection_string)
- bbs.create_container("tests")
- yield fsspec.FsspecFileIO(properties=properties)
- bbs.delete_container("tests")
- bbs.close()
-
-
@pytest.fixture(scope="session")
def empty_home_dir_path(tmp_path_factory: pytest.TempPathFactory) -> str:
home_path = str(tmp_path_factory.mktemp("home"))
diff --git a/tests/io/test_pyarrow.py b/tests/io/test_pyarrow.py
index 9daefacd..850b1292 100644
--- a/tests/io/test_pyarrow.py
+++ b/tests/io/test_pyarrow.py
@@ -24,9 +24,11 @@ from typing import Any, List, Optional
from unittest.mock import MagicMock, patch
from uuid import uuid4
+import pyarrow
import pyarrow as pa
import pyarrow.parquet as pq
import pytest
+from packaging import version
from pyarrow.fs import FileType, LocalFileSystem, S3FileSystem
from pyiceberg.exceptions import ResolveError
@@ -106,6 +108,11 @@ from pyiceberg.types import (
from tests.catalog.test_base import InMemoryCatalog
from tests.conftest import UNIFIED_AWS_SESSION_PROPERTIES
+skip_if_pyarrow_too_old = pytest.mark.skipif(
+ version.parse(pyarrow.__version__) < version.parse("20.0.0"),
+ reason="Requires pyarrow version >= 20.0.0",
+)
+
def test_pyarrow_infer_local_fs_from_path() -> None:
"""Test path with `file` scheme and no scheme both use LocalFileSystem"""
@@ -1672,7 +1679,7 @@ def test_new_output_file_gcs(pyarrow_fileio_gcs:
PyArrowFileIO) -> None:
@pytest.mark.gcs
@pytest.mark.skip(reason="Open issue on Arrow:
https://github.com/apache/arrow/issues/36993")
def test_write_and_read_file_gcs(pyarrow_fileio_gcs: PyArrowFileIO) -> None:
- """Test writing and reading a file using FsspecInputFile and
FsspecOutputFile"""
+ """Test writing and reading a file using PyArrowFile"""
location = f"gs://warehouse/{uuid4()}.txt"
output_file = pyarrow_fileio_gcs.new_output(location=location)
with output_file.create() as f:
@@ -1689,7 +1696,7 @@ def test_write_and_read_file_gcs(pyarrow_fileio_gcs:
PyArrowFileIO) -> None:
@pytest.mark.gcs
def test_getting_length_of_file_gcs(pyarrow_fileio_gcs: PyArrowFileIO) -> None:
- """Test getting the length of an FsspecInputFile and FsspecOutputFile"""
+ """Test getting the length of PyArrowFile"""
filename = str(uuid4())
output_file =
pyarrow_fileio_gcs.new_output(location=f"gs://warehouse/{filename}")
@@ -1753,7 +1760,7 @@ def
test_read_specified_bytes_for_file_gcs(pyarrow_fileio_gcs: PyArrowFileIO) ->
@pytest.mark.gcs
@pytest.mark.skip(reason="Open issue on Arrow:
https://github.com/apache/arrow/issues/36993")
def test_raise_on_opening_file_not_found_gcs(pyarrow_fileio_gcs:
PyArrowFileIO) -> None:
- """Test that an fsspec input file raises appropriately when the gcs file
is not found"""
+ """Test that PyArrowFile raises appropriately when the gcs file is not
found"""
filename = str(uuid4())
input_file =
pyarrow_fileio_gcs.new_input(location=f"gs://warehouse/{filename}")
@@ -1815,7 +1822,7 @@ def
test_converting_an_outputfile_to_an_inputfile_gcs(pyarrow_fileio_gcs: PyArro
@pytest.mark.gcs
@pytest.mark.skip(reason="Open issue on Arrow:
https://github.com/apache/arrow/issues/36993")
def test_writing_avro_file_gcs(generated_manifest_entry_file: str,
pyarrow_fileio_gcs: PyArrowFileIO) -> None:
- """Test that bytes match when reading a local avro file, writing it using
fsspec file-io, and then reading it again"""
+ """Test that bytes match when reading a local avro file, writing it using
pyarrow file-io, and then reading it again"""
filename = str(uuid4())
with
PyArrowFileIO().new_input(location=generated_manifest_entry_file).open() as f:
b1 = f.read()
@@ -1828,6 +1835,192 @@ def
test_writing_avro_file_gcs(generated_manifest_entry_file: str, pyarrow_filei
pyarrow_fileio_gcs.delete(f"gs://warehouse/{filename}")
[email protected]
+@skip_if_pyarrow_too_old
+def test_new_input_file_adls(pyarrow_fileio_adls: PyArrowFileIO, adls_scheme:
str) -> None:
+ """Test creating a new input file from pyarrow file-io"""
+ filename = str(uuid4())
+
+ input_file =
pyarrow_fileio_adls.new_input(f"{adls_scheme}://warehouse/{filename}")
+
+ assert isinstance(input_file, PyArrowFile)
+ assert input_file.location == f"{adls_scheme}://warehouse/{filename}"
+
+
[email protected]
+@skip_if_pyarrow_too_old
+def test_new_output_file_adls(pyarrow_fileio_adls: PyArrowFileIO, adls_scheme:
str) -> None:
+ """Test creating a new output file from pyarrow file-io"""
+ filename = str(uuid4())
+
+ output_file =
pyarrow_fileio_adls.new_output(f"{adls_scheme}://warehouse/{filename}")
+
+ assert isinstance(output_file, PyArrowFile)
+ assert output_file.location == f"{adls_scheme}://warehouse/{filename}"
+
+
[email protected]
+@skip_if_pyarrow_too_old
+def test_write_and_read_file_adls(pyarrow_fileio_adls: PyArrowFileIO,
adls_scheme: str) -> None:
+ """Test writing and reading a file using PyArrowFile"""
+ location = f"{adls_scheme}://warehouse/{uuid4()}.txt"
+ output_file = pyarrow_fileio_adls.new_output(location=location)
+ with output_file.create() as f:
+ assert f.write(b"foo") == 3
+
+ assert output_file.exists()
+
+ input_file = pyarrow_fileio_adls.new_input(location=location)
+ with input_file.open() as f:
+ assert f.read() == b"foo"
+
+ pyarrow_fileio_adls.delete(input_file)
+
+
[email protected]
+@skip_if_pyarrow_too_old
+def test_getting_length_of_file_adls(pyarrow_fileio_adls: PyArrowFileIO,
adls_scheme: str) -> None:
+ """Test getting the length of PyArrowFile"""
+ filename = str(uuid4())
+
+ output_file =
pyarrow_fileio_adls.new_output(location=f"{adls_scheme}://warehouse/{filename}")
+ with output_file.create() as f:
+ f.write(b"foobar")
+
+ assert len(output_file) == 6
+
+ input_file =
pyarrow_fileio_adls.new_input(location=f"{adls_scheme}://warehouse/{filename}")
+ assert len(input_file) == 6
+
+ pyarrow_fileio_adls.delete(output_file)
+
+
[email protected]
+@skip_if_pyarrow_too_old
+def test_file_tell_adls(pyarrow_fileio_adls: PyArrowFileIO, adls_scheme: str)
-> None:
+ location = f"{adls_scheme}://warehouse/{uuid4()}"
+
+ output_file = pyarrow_fileio_adls.new_output(location=location)
+ with output_file.create() as write_file:
+ write_file.write(b"foobar")
+
+ input_file = pyarrow_fileio_adls.new_input(location=location)
+ with input_file.open() as f:
+ f.seek(0)
+ assert f.tell() == 0
+ f.seek(1)
+ assert f.tell() == 1
+ f.seek(3)
+ assert f.tell() == 3
+ f.seek(0)
+ assert f.tell() == 0
+
+
[email protected]
+@skip_if_pyarrow_too_old
+def test_read_specified_bytes_for_file_adls(pyarrow_fileio_adls:
PyArrowFileIO) -> None:
+ location = f"abfss://warehouse/{uuid4()}"
+
+ output_file = pyarrow_fileio_adls.new_output(location=location)
+ with output_file.create() as write_file:
+ write_file.write(b"foo")
+
+ input_file = pyarrow_fileio_adls.new_input(location=location)
+ with input_file.open() as f:
+ f.seek(0)
+ assert b"f" == f.read(1)
+ f.seek(0)
+ assert b"fo" == f.read(2)
+ f.seek(1)
+ assert b"o" == f.read(1)
+ f.seek(1)
+ assert b"oo" == f.read(2)
+ f.seek(0)
+ assert b"foo" == f.read(999) # test reading amount larger than entire
content length
+
+ pyarrow_fileio_adls.delete(input_file)
+
+
[email protected]
+@skip_if_pyarrow_too_old
+def test_raise_on_opening_file_not_found_adls(pyarrow_fileio_adls:
PyArrowFileIO, adls_scheme: str) -> None:
+ """Test that PyArrowFile raises appropriately when the adls file is not
found"""
+
+ filename = str(uuid4())
+ input_file =
pyarrow_fileio_adls.new_input(location=f"{adls_scheme}://warehouse/{filename}")
+ with pytest.raises(FileNotFoundError) as exc_info:
+ input_file.open().read()
+
+ assert filename in str(exc_info.value)
+
+
[email protected]
+@skip_if_pyarrow_too_old
+def test_checking_if_a_file_exists_adls(pyarrow_fileio_adls: PyArrowFileIO,
adls_scheme: str) -> None:
+ """Test checking if a file exists"""
+ non_existent_file =
pyarrow_fileio_adls.new_input(location=f"{adls_scheme}://warehouse/does-not-exist.txt")
+ assert not non_existent_file.exists()
+
+ location = f"{adls_scheme}://warehouse/{uuid4()}"
+ output_file = pyarrow_fileio_adls.new_output(location=location)
+ assert not output_file.exists()
+ with output_file.create() as f:
+ f.write(b"foo")
+
+ existing_input_file = pyarrow_fileio_adls.new_input(location=location)
+ assert existing_input_file.exists()
+
+ existing_output_file = pyarrow_fileio_adls.new_output(location=location)
+ assert existing_output_file.exists()
+
+ pyarrow_fileio_adls.delete(existing_output_file)
+
+
[email protected]
+@skip_if_pyarrow_too_old
+def test_closing_a_file_adls(pyarrow_fileio_adls: PyArrowFileIO, adls_scheme:
str) -> None:
+ """Test closing an output file and input file"""
+ filename = str(uuid4())
+ output_file =
pyarrow_fileio_adls.new_output(location=f"{adls_scheme}://warehouse/{filename}")
+ with output_file.create() as write_file:
+ write_file.write(b"foo")
+ assert not write_file.closed # type: ignore
+ assert write_file.closed # type: ignore
+
+ input_file =
pyarrow_fileio_adls.new_input(location=f"{adls_scheme}://warehouse/{filename}")
+ with input_file.open() as f:
+ assert not f.closed # type: ignore
+ assert f.closed # type: ignore
+
+ pyarrow_fileio_adls.delete(f"{adls_scheme}://warehouse/{filename}")
+
+
[email protected]
+@skip_if_pyarrow_too_old
+def test_converting_an_outputfile_to_an_inputfile_adls(pyarrow_fileio_adls:
PyArrowFileIO, adls_scheme: str) -> None:
+ """Test converting an output file to an input file"""
+ filename = str(uuid4())
+ output_file =
pyarrow_fileio_adls.new_output(location=f"{adls_scheme}://warehouse/{filename}")
+ input_file = output_file.to_input_file()
+ assert input_file.location == output_file.location
+
+
[email protected]
+@skip_if_pyarrow_too_old
+def test_writing_avro_file_adls(generated_manifest_entry_file: str,
pyarrow_fileio_adls: PyArrowFileIO, adls_scheme: str) -> None:
+ """Test that bytes match when reading a local avro file, writing it using
pyarrow file-io, and then reading it again"""
+ filename = str(uuid4())
+ with
PyArrowFileIO().new_input(location=generated_manifest_entry_file).open() as f:
+ b1 = f.read()
+ with
pyarrow_fileio_adls.new_output(location=f"{adls_scheme}://warehouse/{filename}").create()
as out_f:
+ out_f.write(b1)
+ with
pyarrow_fileio_adls.new_input(location=f"{adls_scheme}://warehouse/{filename}").open()
as in_f:
+ b2 = in_f.read()
+ assert b1 == b2 # Check that bytes of read from local avro file
match bytes written to s3
+
+ pyarrow_fileio_adls.delete(f"{adls_scheme}://warehouse/{filename}")
+
+
def test_parse_location() -> None:
def check_results(location: str, expected_schema: str, expected_netloc:
str, expected_uri: str) -> None:
schema, netloc, uri = PyArrowFileIO.parse_location(location)