Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions docs/docs/pypaimon/catalogs.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -170,6 +170,17 @@ PyPaimon supports filesystem, JDBC, and REST catalogs. See [Catalog](../concepts

Use this catalog for the database and table operations below.

## OSS and custom S3 checksums

OSS rejects PyArrow 22+ optional checksum trailers (checksums sent after the data).
For PyArrow-backed OSS and explicit S3 endpoints (including AWS), PyPaimon sets
`AWS_REQUEST_CHECKSUM_CALCULATION=WHEN_REQUIRED` to disable optional checksums.
This also affects later AWS clients in the process. Explicit environment values take
precedence; required checksums and older PyArrow versions are unchanged.

Set `fs.s3.checksum-compatibility.auto-configure=false` on all relevant catalogs before
creating clients to disable this setup. Existing process settings are not cleared.

## Create Database

Tables belong to a database. Create the database before creating its tables.
Expand Down
3 changes: 3 additions & 0 deletions paimon-python/pypaimon/common/options/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,9 @@ class OssOptions:


class S3Options:
CHECKSUM_COMPATIBILITY_AUTO_CONFIGURE = (
ConfigOptions.key("fs.s3.checksum-compatibility.auto-configure").boolean_type().default_value(True)
.with_description("Apply a process-wide optional-checksum default for PyArrow OSS/custom S3 clients."))
S3_ACCESS_KEY_ID = ConfigOptions.key("fs.s3.accessKeyId").string_type().no_default_value().with_description(
"S3 access key ID")
S3_ACCESS_KEY_SECRET = ConfigOptions.key("fs.s3.accessKeySecret").string_type().no_default_value().with_description(
Expand Down
53 changes: 46 additions & 7 deletions paimon-python/pypaimon/filesystem/pyarrow_file_io.py
Original file line number Diff line number Diff line change
Expand Up @@ -54,15 +54,15 @@ class PyArrowFileIO(FileIO):
def __init__(self, path: str, catalog_options: Options):
self.properties = catalog_options
self.logger = logging.getLogger(__name__)
self._pyarrow_gte_8 = parse(pyarrow.__version__) >= parse("8.0.0")
# force_virtual_addressing landed in PyArrow 16; below it the OSS bucket
# goes into endpoint_override, so keys must omit it (init + path share
# this flag so they can't drift).
self._pyarrow_gte_16 = parse(pyarrow.__version__) >= parse("16.0.0")
self._oss_bucket_in_endpoint = not self._pyarrow_gte_16
self._set_pyarrow_version()
scheme, netloc, _ = self.parse_location(path)
self.uri_reader_factory = UriReaderFactory(catalog_options)
self._is_oss = scheme in {"oss"}
self._is_s3 = scheme in {"s3", "s3a", "s3n"}
self._s3_endpoint = (
self._get_s3_property("endpoint", S3Options.S3_ENDPOINT.key())
if self._is_s3 else None
)
self._oss_bucket = None
_oss_impl = self.properties.get(OssOptions.OSS_IMPL)
self._use_jindo = False
Expand All @@ -86,7 +86,7 @@ def __init__(self, path: str, catalog_options: Options):
"Falling back to legacy PyArrow S3FileSystem implementation. "
"Install pyjindosdk for better performance: pip install pyjindosdk")
self.filesystem = self._initialize_oss_fs(path)
elif scheme in {"s3", "s3a", "s3n"}:
elif self._is_s3:
self.filesystem = self._initialize_s3_fs()
elif scheme in {"hdfs", "viewfs"}:
self.filesystem = self._initialize_hdfs_fs(scheme, netloc)
Expand All @@ -95,15 +95,52 @@ def __init__(self, path: str, catalog_options: Options):
else:
raise ValueError(f"Unrecognized filesystem type in URI: {scheme}")

def _set_pyarrow_version(self):
self._pyarrow_gte_8 = parse(pyarrow.__version__) >= parse("8.0.0")
# force_virtual_addressing landed in PyArrow 16; below it the OSS bucket
# goes into endpoint_override, so keys must omit it (init + path share
# this flag so they can't drift).
self._pyarrow_gte_16 = parse(pyarrow.__version__) >= parse("16.0.0")
self._oss_bucket_in_endpoint = not self._pyarrow_gte_16

def _uses_s3_compatibility(self) -> bool:
return (not self._use_jindo
and (self._is_oss or bool(self._s3_endpoint)))

def _configure_s3_checksums(self):
if parse(pyarrow.__version__) < parse("22.0.0"):
return
if self._uses_s3_compatibility() and self.properties.get(S3Options.CHECKSUM_COMPATIBILITY_AUTO_CONFIGURE):
# Process-wide default; preserve explicit settings and do not restore it.
os.environ.setdefault("AWS_REQUEST_CHECKSUM_CALCULATION", "WHEN_REQUIRED")

def __getstate__(self):
state = self.__dict__.copy()
# threading.Lock cannot be pickled; recreated in __setstate__.
state.pop("_legacy_bucket_lock", None)
state.pop("logger", None)
# Recreate S3-compatible clients with the worker's AWS SDK settings.
if self._uses_s3_compatibility():
state.pop("filesystem", None)
return state

def __setstate__(self, state):
self.__dict__.update(state)
self.logger = logging.getLogger(__name__)
self._set_pyarrow_version()
if "_is_s3" not in state:
self._is_s3 = (not self._is_oss
and isinstance(self.filesystem, pafs.S3FileSystem))
if "_s3_endpoint" not in state:
self._s3_endpoint = (
self._get_s3_property("endpoint", S3Options.S3_ENDPOINT.key())
if self._is_s3 else None)
self._legacy_bucket_lock = threading.Lock()
if self._uses_s3_compatibility():
self.filesystem = (
self._initialize_oss_fs(None)
if self._is_oss else self._initialize_s3_fs()
)

@staticmethod
def parse_location(location: str):
Expand Down Expand Up @@ -193,6 +230,7 @@ def _initialize_jindo_fs(self, path) -> FileSystem:
return pafs.PyFileSystem(fs_handler)

def _initialize_oss_fs(self, path) -> FileSystem:
self._configure_s3_checksums()
if self.properties.get(OssOptions.OSS_ACCESS_KEY_ID):
# When explicit credentials are provided, disable the EC2 Instance Metadata
# Service (IMDS) probe to avoid multi-second timeouts in non-AWS environments.
Expand Down Expand Up @@ -222,6 +260,7 @@ def _initialize_oss_fs(self, path) -> FileSystem:
return pafs.S3FileSystem(**client_kwargs)

def _initialize_s3_fs(self) -> FileSystem:
self._configure_s3_checksums()
access_key = self._get_property(
S3Options.S3_ACCESS_KEY_ID.key(),
*self._s3_key_variants("access-key", "access.key"))
Expand Down
Loading
Loading