From fc6b96bb2b62682e1c62774c132d6fce4ecf4827 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Thu, 8 Oct 2026 04:49:26 -0700 Subject: [PATCH 01/10] [python] Configure compatible S3 checksums when rebuilding Arrow clients --- .../pypaimon/filesystem/pyarrow_file_io.py | 69 +++- .../tests/s3_client_compatibility_test.py | 311 ++++++++++++++++++ 2 files changed, 371 insertions(+), 9 deletions(-) create mode 100644 paimon-python/pypaimon/tests/s3_client_compatibility_test.py diff --git a/paimon-python/pypaimon/filesystem/pyarrow_file_io.py b/paimon-python/pypaimon/filesystem/pyarrow_file_io.py index ad5d3c8a4ea9..53a3d860c395 100644 --- a/paimon-python/pypaimon/filesystem/pyarrow_file_io.py +++ b/paimon-python/pypaimon/filesystem/pyarrow_file_io.py @@ -46,6 +46,10 @@ def _pyarrow_lt_7(): return parse(pyarrow.__version__) < parse("7.0.0") +_S3_CHECKSUM_LOCK = threading.Lock() +_S3_CHECKSUM_ENV = "AWS_REQUEST_CHECKSUM_CALCULATION" + + class LegacyOssDirectoryListingError(RuntimeError): """Raised when legacy PyArrow OSS cannot enumerate a directory.""" @@ -54,15 +58,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 @@ -86,7 +90,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) @@ -95,15 +99,61 @@ 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))) + + @staticmethod + def _create_s3_filesystem(client_kwargs, compatible: bool) -> FileSystem: + with _S3_CHECKSUM_LOCK: + if not compatible: + return pafs.S3FileSystem(**client_kwargs) + # PyArrow has no per-client checksum option; AWS reads this at construction. + previous = os.environ.get(_S3_CHECKSUM_ENV) + os.environ[_S3_CHECKSUM_ENV] = "WHEN_REQUIRED" + try: + return pafs.S3FileSystem(**client_kwargs) + finally: + if previous is None: + os.environ.pop(_S3_CHECKSUM_ENV, None) + else: + os.environ[_S3_CHECKSUM_ENV] = previous + 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): @@ -219,7 +269,7 @@ def _initialize_oss_fs(self, path) -> FileSystem: retry_config = self._create_s3_retry_config() client_kwargs.update(retry_config) - return pafs.S3FileSystem(**client_kwargs) + return self._create_s3_filesystem(client_kwargs, compatible=True) def _initialize_s3_fs(self) -> FileSystem: access_key = self._get_property( @@ -259,7 +309,8 @@ def _initialize_s3_fs(self) -> FileSystem: retry_config = self._create_s3_retry_config() client_kwargs.update(retry_config) - return pafs.S3FileSystem(**client_kwargs) + return self._create_s3_filesystem( + client_kwargs, compatible=bool(self._s3_endpoint)) def _initialize_hdfs_fs(self, scheme: str, netloc: Optional[str]) -> FileSystem: if 'HADOOP_HOME' not in os.environ: diff --git a/paimon-python/pypaimon/tests/s3_client_compatibility_test.py b/paimon-python/pypaimon/tests/s3_client_compatibility_test.py new file mode 100644 index 000000000000..99923e5c0644 --- /dev/null +++ b/paimon-python/pypaimon/tests/s3_client_compatibility_test.py @@ -0,0 +1,311 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +"""Unit tests for PyArrow-backed OSS and S3-compatible storage. + +No real OSS access is required. +""" + +import multiprocessing +import os +import pickle +import unittest +import threading +from http.server import BaseHTTPRequestHandler, HTTPServer +from socketserver import ThreadingMixIn +from unittest import mock + +import pyarrow +import pyarrow.fs as pafs +from packaging.version import parse + +from pypaimon.common.options import Options +from pypaimon.common.options.config import OssOptions, S3Options +from pypaimon.filesystem.pyarrow_file_io import PyArrowFileIO + + +def _restore_s3_file_io(connection): + os.environ.pop("AWS_REQUEST_CHECKSUM_CALCULATION", None) + connection.send("ready") + payload = connection.recv_bytes() + client = object() + settings = [] + + def create_client(**kwargs): + settings.append(os.environ.get("AWS_REQUEST_CHECKSUM_CALCULATION")) + return client + + with mock.patch("pyarrow.fs.S3FileSystem", side_effect=create_client) as s3: + restored = pickle.loads(payload) + connection.send(( + os.environ.get("AWS_REQUEST_CHECKSUM_CALCULATION"), + settings, + s3.call_count, + restored.filesystem is client, + )) + connection.close() + + +def _write_in_worker(connection): + os.environ["AWS_REQUEST_CHECKSUM_CALCULATION"] = "WHEN_SUPPORTED" + connection.send("ready") + try: + file_io = pickle.loads(connection.recv_bytes()) + with file_io.filesystem.open_output_stream("test-bucket/file") as stream: + stream.write(b"x") + connection.send("written") + except Exception as error: + connection.send(repr(error)) + finally: + connection.close() + + +class _HTTPServer(ThreadingMixIn, HTTPServer): + daemon_threads = True + + +class _UploadHandler(BaseHTTPRequestHandler): + def _respond(self, body=b""): + self.send_response(200) + self.send_header("Content-Length", str(len(body))) + self.send_header("ETag", '"etag"') + self.send_header("Connection", "close") + self.end_headers() + self.wfile.write(body) + self.close_connection = True + + def do_POST(self): + if "uploads" in self.path: + self._respond(b"test-bucket" + b"filetest-id" + b"") + else: + self._respond(b"test-bucket" + b"fileetag") + + def do_PUT(self): + self.server.put_headers.append(dict(self.headers)) + self._respond() + + def log_message(self, *args): + pass + + +class S3ClientCompatibilityTest(unittest.TestCase): + def _new_file_io(self, scheme="s3"): + if scheme == "oss": + options = Options({ + OssOptions.OSS_IMPL.key(): "legacy", + OssOptions.OSS_ENDPOINT.key(): "oss-cn-test.example.com", + OssOptions.OSS_ACCESS_KEY_ID.key(): "ak", + OssOptions.OSS_ACCESS_KEY_SECRET.key(): "sk", + }) + else: + options = Options({S3Options.S3_ENDPOINT.key(): "http://minio:9000"}) + with mock.patch("pyarrow.fs.S3FileSystem", return_value=pafs.LocalFileSystem()): + file_io = PyArrowFileIO(scheme + "://test-bucket/", options) + return file_io + + def test_environment_restored_after_client_creation_failure(self): + for previous in (None, "WHEN_SUPPORTED"): + with self.subTest(previous=previous), mock.patch.dict(os.environ, {}, clear=True): + if previous is not None: + os.environ["AWS_REQUEST_CHECKSUM_CALCULATION"] = previous + with mock.patch("pyarrow.fs.S3FileSystem", side_effect=RuntimeError("failed")): + with self.assertRaisesRegex(RuntimeError, "failed"): + self._new_file_io_failure() + self.assertEqual(previous, os.environ.get("AWS_REQUEST_CHECKSUM_CALCULATION")) + + def _new_file_io_failure(self): + PyArrowFileIO("s3://test-bucket/", Options({"fs.s3.endpoint": "http://minio:9000"})) + + def test_pickle_rebuilds_client_in_started_worker(self): + for scheme in ("oss", "s3", "s3a", "s3n"): + with self.subTest(scheme=scheme): + context = multiprocessing.get_context("spawn") + parent, child = context.Pipe() + process = context.Process(target=_restore_s3_file_io, args=(child,)) + process.start() + child.close() + try: + self.assertTrue(parent.poll(20)) + self.assertEqual("ready", parent.recv()) + file_io = self._new_file_io(scheme) + self.assertNotIn("filesystem", file_io.__getstate__()) + parent.send_bytes(pickle.dumps(file_io)) + self.assertTrue(parent.poll(20)) + self.assertEqual((None, ["WHEN_REQUIRED"], 1, True), parent.recv()) + finally: + parent.close() + process.join(20) + if process.is_alive(): + process.terminate() + process.join() + self.assertEqual(0, process.exitcode) + + def test_oss_initialization_disables_optional_checksum_trailers(self): + options = Options({ + OssOptions.OSS_ACCESS_KEY_ID.key(): "ak", + OssOptions.OSS_ACCESS_KEY_SECRET.key(): "sk", + OssOptions.OSS_ENDPOINT.key(): "oss-cn-test.example.com", + OssOptions.OSS_REGION.key(): "cn-test", + OssOptions.OSS_IMPL.key(): "legacy", + }) + settings = [] + + def create_client(**kwargs): + settings.append(os.environ.get("AWS_REQUEST_CHECKSUM_CALCULATION")) + return mock.Mock() + + with mock.patch.dict("os.environ", { + "AWS_REQUEST_CHECKSUM_CALCULATION": "WHEN_SUPPORTED"}, clear=True), \ + mock.patch("pyarrow.fs.S3FileSystem", side_effect=create_client): + PyArrowFileIO("oss://test-bucket/", options) + self.assertEqual( + "WHEN_SUPPORTED", + os.environ["AWS_REQUEST_CHECKSUM_CALCULATION"]) + self.assertEqual(["WHEN_REQUIRED"], settings) + + def test_initialization_configures_all_s3_schemes(self): + options = Options({ + S3Options.S3_ENDPOINT.key(): "http://minio:9000", + }) + for scheme in ("s3", "s3a", "s3n"): + settings = [] + + def create_client(**kwargs): + settings.append(os.environ.get("AWS_REQUEST_CHECKSUM_CALCULATION")) + return mock.Mock() + + with self.subTest(scheme=scheme), \ + mock.patch.dict("os.environ", {}, clear=True), \ + mock.patch("pyarrow.fs.S3FileSystem", side_effect=create_client): + PyArrowFileIO( + "{}://test-bucket/warehouse".format(scheme), options) + self.assertNotIn("AWS_REQUEST_CHECKSUM_CALCULATION", os.environ) + self.assertEqual(["WHEN_REQUIRED"], settings) + + def test_native_s3_does_not_change_checksum_setting(self): + with mock.patch.dict("os.environ", {}, clear=True), \ + mock.patch("pyarrow.fs.S3FileSystem", return_value=mock.Mock()): + PyArrowFileIO("s3://test-bucket/warehouse", Options({})) + self.assertNotIn( + "AWS_REQUEST_CHECKSUM_CALCULATION", os.environ) + + def test_worker_recomputes_pyarrow_version(self): + state = self._new_file_io().__getstate__() + state.update({ + "_pyarrow_gte_8": False, + "_pyarrow_gte_16": False, + "_oss_bucket_in_endpoint": True, + }) + restored = object.__new__(PyArrowFileIO) + with mock.patch.object( + PyArrowFileIO, "_initialize_s3_fs", return_value=object()): + restored.__setstate__(state) + + version = parse(pyarrow.__version__) + self.assertEqual(version >= parse("8.0.0"), restored._pyarrow_gte_8) + self.assertEqual(version >= parse("16.0.0"), restored._pyarrow_gte_16) + self.assertEqual(version < parse("16.0.0"), restored._oss_bucket_in_endpoint) + + def test_worker_accepts_old_non_oss_pickle(self): + state = self._new_file_io().__dict__.copy() + for key in ("_legacy_bucket_lock", "_is_s3", + "_s3_endpoint"): + state.pop(key) + state["filesystem"] = mock.Mock(spec=pafs.S3FileSystem) + restored = object.__new__(PyArrowFileIO) + with mock.patch.object( + PyArrowFileIO, "_initialize_s3_fs", return_value=object() + ) as initialize: + restored.__setstate__(state) + self.assertTrue(restored._is_s3) + self.assertEqual("http://minio:9000", restored._s3_endpoint) + initialize.assert_called_once() + + state["filesystem"] = pafs.LocalFileSystem() + restored = object.__new__(PyArrowFileIO) + restored.__setstate__(state) + self.assertFalse(restored._is_s3) + self.assertIsNone(restored._s3_endpoint) + + @unittest.skipUnless( + parse(pyarrow.__version__) >= parse("22.0.0"), + "requires PyArrow 22+ optional request checksums", + ) + def test_checksum_setting_is_scoped_to_compatible_client(self): + server = _HTTPServer( + ("127.0.0.1", 0), _UploadHandler) + server.requests = [] + server.put_headers = [] + server.bucket_objects = {"test-bucket": set()} + server_thread = threading.Thread(target=server.serve_forever) + server_thread.start() + try: + endpoint = "http://127.0.0.1:{}".format(server.server_port) + options = Options({ + S3Options.S3_ACCESS_KEY_ID.key(): "ak", + S3Options.S3_ACCESS_KEY_SECRET.key(): "sk", + S3Options.S3_ENDPOINT.key(): endpoint, + S3Options.S3_REGION.key(): "us-east-1", + "fs.s3.path.style.access": "true", + }) + with mock.patch.dict(os.environ, { + "AWS_REQUEST_CHECKSUM_CALCULATION": "WHEN_SUPPORTED", + "NO_PROXY": "127.0.0.1,localhost", + "no_proxy": "127.0.0.1,localhost", + }): + native = pafs.S3FileSystem( + access_key="ak", secret_key="sk", region="us-east-1", + endpoint_override=endpoint) + compatible = PyArrowFileIO("s3://test-bucket/", options) + self.assertEqual("WHEN_SUPPORTED", + os.environ["AWS_REQUEST_CHECKSUM_CALCULATION"]) + with native.open_output_stream("test-bucket/file") as stream: + stream.write(b"x") + with compatible.filesystem.open_output_stream( + "test-bucket/file") as stream: + stream.write(b"x") + + context = multiprocessing.get_context("spawn") + parent, child = context.Pipe() + process = context.Process(target=_write_in_worker, args=(child,)) + process.start() + child.close() + try: + self.assertTrue(parent.poll(20)) + self.assertEqual("ready", parent.recv()) + parent.send_bytes(pickle.dumps(compatible)) + self.assertTrue(parent.poll(20)) + self.assertEqual("written", parent.recv()) + finally: + parent.close() + process.join(20) + if process.is_alive(): + process.terminate() + process.join() + self.assertEqual(0, process.exitcode) + self.assertEqual(3, len(server.put_headers)) + self.assertIn("x-amz-trailer", { + key.lower() for key in server.put_headers[0]}) + for headers in server.put_headers[1:]: + self.assertNotIn("x-amz-trailer", {key.lower() for key in headers}) + finally: + server.shutdown() + server.server_close() + server_thread.join() From 8e463a4ac346f22698dc344e978897b3d6ccd581 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Thu, 8 Oct 2026 05:25:32 -0700 Subject: [PATCH 02/10] Avoid process environment mutation during S3 client construction --- paimon-python/README.md | 15 ++++++++ .../pypaimon/filesystem/pyarrow_file_io.py | 25 +------------ .../tests/s3_client_compatibility_test.py | 37 ++++++++++--------- 3 files changed, 36 insertions(+), 41 deletions(-) diff --git a/paimon-python/README.md b/paimon-python/README.md index f4dd3c7fe828..aa2637b948b6 100644 --- a/paimon-python/README.md +++ b/paimon-python/README.md @@ -30,6 +30,21 @@ python -m build Both produce a source archive and wheel in `dist/`. +## PyArrow checksums with OSS and custom S3 endpoints + +For endpoints that do not support the optional request checksums used by newer +PyArrow AWS SDKs, configure the process before starting the application and its +workers: + +```bash +export AWS_REQUEST_CHECKSUM_CALCULATION=WHEN_REQUIRED +``` + +This setting applies to AWS SDK clients throughout the process, not just Paimon. +PyPaimon does not modify it. Existing workers and clients must be restarted after +changing it. This only addresses optional write checksums; required checksums on +batch deletion are a separate compatibility concern. + # OSS metadata commits Install `pypaimon[oss]` (legacy PyArrow data access) or `pypaimon[jindo]` diff --git a/paimon-python/pypaimon/filesystem/pyarrow_file_io.py b/paimon-python/pypaimon/filesystem/pyarrow_file_io.py index 53a3d860c395..825ff952509e 100644 --- a/paimon-python/pypaimon/filesystem/pyarrow_file_io.py +++ b/paimon-python/pypaimon/filesystem/pyarrow_file_io.py @@ -46,10 +46,6 @@ def _pyarrow_lt_7(): return parse(pyarrow.__version__) < parse("7.0.0") -_S3_CHECKSUM_LOCK = threading.Lock() -_S3_CHECKSUM_ENV = "AWS_REQUEST_CHECKSUM_CALCULATION" - - class LegacyOssDirectoryListingError(RuntimeError): """Raised when legacy PyArrow OSS cannot enumerate a directory.""" @@ -111,22 +107,6 @@ def _uses_s3_compatibility(self) -> bool: return (not self._use_jindo and (self._is_oss or bool(self._s3_endpoint))) - @staticmethod - def _create_s3_filesystem(client_kwargs, compatible: bool) -> FileSystem: - with _S3_CHECKSUM_LOCK: - if not compatible: - return pafs.S3FileSystem(**client_kwargs) - # PyArrow has no per-client checksum option; AWS reads this at construction. - previous = os.environ.get(_S3_CHECKSUM_ENV) - os.environ[_S3_CHECKSUM_ENV] = "WHEN_REQUIRED" - try: - return pafs.S3FileSystem(**client_kwargs) - finally: - if previous is None: - os.environ.pop(_S3_CHECKSUM_ENV, None) - else: - os.environ[_S3_CHECKSUM_ENV] = previous - def __getstate__(self): state = self.__dict__.copy() # threading.Lock cannot be pickled; recreated in __setstate__. @@ -269,7 +249,7 @@ def _initialize_oss_fs(self, path) -> FileSystem: retry_config = self._create_s3_retry_config() client_kwargs.update(retry_config) - return self._create_s3_filesystem(client_kwargs, compatible=True) + return pafs.S3FileSystem(**client_kwargs) def _initialize_s3_fs(self) -> FileSystem: access_key = self._get_property( @@ -309,8 +289,7 @@ def _initialize_s3_fs(self) -> FileSystem: retry_config = self._create_s3_retry_config() client_kwargs.update(retry_config) - return self._create_s3_filesystem( - client_kwargs, compatible=bool(self._s3_endpoint)) + return pafs.S3FileSystem(**client_kwargs) def _initialize_hdfs_fs(self, scheme: str, netloc: Optional[str]) -> FileSystem: if 'HADOOP_HOME' not in os.environ: diff --git a/paimon-python/pypaimon/tests/s3_client_compatibility_test.py b/paimon-python/pypaimon/tests/s3_client_compatibility_test.py index 99923e5c0644..5f981b4f0448 100644 --- a/paimon-python/pypaimon/tests/s3_client_compatibility_test.py +++ b/paimon-python/pypaimon/tests/s3_client_compatibility_test.py @@ -25,6 +25,7 @@ import pickle import unittest import threading +from itertools import product from http.server import BaseHTTPRequestHandler, HTTPServer from socketserver import ThreadingMixIn from unittest import mock @@ -39,7 +40,6 @@ def _restore_s3_file_io(connection): - os.environ.pop("AWS_REQUEST_CHECKSUM_CALCULATION", None) connection.send("ready") payload = connection.recv_bytes() client = object() @@ -61,7 +61,6 @@ def create_client(**kwargs): def _write_in_worker(connection): - os.environ["AWS_REQUEST_CHECKSUM_CALCULATION"] = "WHEN_SUPPORTED" connection.send("ready") try: file_io = pickle.loads(connection.recv_bytes()) @@ -120,7 +119,7 @@ def _new_file_io(self, scheme="s3"): file_io = PyArrowFileIO(scheme + "://test-bucket/", options) return file_io - def test_environment_restored_after_client_creation_failure(self): + def test_environment_unchanged_after_client_creation_failure(self): for previous in (None, "WHEN_SUPPORTED"): with self.subTest(previous=previous), mock.patch.dict(os.environ, {}, clear=True): if previous is not None: @@ -133,10 +132,13 @@ def test_environment_restored_after_client_creation_failure(self): def _new_file_io_failure(self): PyArrowFileIO("s3://test-bucket/", Options({"fs.s3.endpoint": "http://minio:9000"})) + @mock.patch.dict(os.environ, {"AWS_REQUEST_CHECKSUM_CALCULATION": "WHEN_REQUIRED"}) def test_pickle_rebuilds_client_in_started_worker(self): - for scheme in ("oss", "s3", "s3a", "s3n"): - with self.subTest(scheme=scheme): - context = multiprocessing.get_context("spawn") + methods = [name for name in ("spawn", "fork") + if name in multiprocessing.get_all_start_methods()] + for scheme, method in product(("oss", "s3", "s3a", "s3n"), methods): + with self.subTest(scheme=scheme, method=method): + context = multiprocessing.get_context(method) parent, child = context.Pipe() process = context.Process(target=_restore_s3_file_io, args=(child,)) process.start() @@ -148,7 +150,7 @@ def test_pickle_rebuilds_client_in_started_worker(self): self.assertNotIn("filesystem", file_io.__getstate__()) parent.send_bytes(pickle.dumps(file_io)) self.assertTrue(parent.poll(20)) - self.assertEqual((None, ["WHEN_REQUIRED"], 1, True), parent.recv()) + self.assertEqual(("WHEN_REQUIRED", ["WHEN_REQUIRED"], 1, True), parent.recv()) finally: parent.close() process.join(20) @@ -157,7 +159,7 @@ def test_pickle_rebuilds_client_in_started_worker(self): process.join() self.assertEqual(0, process.exitcode) - def test_oss_initialization_disables_optional_checksum_trailers(self): + def test_oss_initialization_preserves_process_setting(self): options = Options({ OssOptions.OSS_ACCESS_KEY_ID.key(): "ak", OssOptions.OSS_ACCESS_KEY_SECRET.key(): "sk", @@ -178,9 +180,9 @@ def create_client(**kwargs): self.assertEqual( "WHEN_SUPPORTED", os.environ["AWS_REQUEST_CHECKSUM_CALCULATION"]) - self.assertEqual(["WHEN_REQUIRED"], settings) + self.assertEqual(["WHEN_SUPPORTED"], settings) - def test_initialization_configures_all_s3_schemes(self): + def test_initialization_preserves_environment_for_all_s3_schemes(self): options = Options({ S3Options.S3_ENDPOINT.key(): "http://minio:9000", }) @@ -197,7 +199,7 @@ def create_client(**kwargs): PyArrowFileIO( "{}://test-bucket/warehouse".format(scheme), options) self.assertNotIn("AWS_REQUEST_CHECKSUM_CALCULATION", os.environ) - self.assertEqual(["WHEN_REQUIRED"], settings) + self.assertEqual([None], settings) def test_native_s3_does_not_change_checksum_setting(self): with mock.patch.dict("os.environ", {}, clear=True), \ @@ -248,7 +250,7 @@ def test_worker_accepts_old_non_oss_pickle(self): parse(pyarrow.__version__) >= parse("22.0.0"), "requires PyArrow 22+ optional request checksums", ) - def test_checksum_setting_is_scoped_to_compatible_client(self): + def test_explicit_process_setting_applies_to_parent_and_worker(self): server = _HTTPServer( ("127.0.0.1", 0), _UploadHandler) server.requests = [] @@ -266,7 +268,7 @@ def test_checksum_setting_is_scoped_to_compatible_client(self): "fs.s3.path.style.access": "true", }) with mock.patch.dict(os.environ, { - "AWS_REQUEST_CHECKSUM_CALCULATION": "WHEN_SUPPORTED", + "AWS_REQUEST_CHECKSUM_CALCULATION": "WHEN_REQUIRED", "NO_PROXY": "127.0.0.1,localhost", "no_proxy": "127.0.0.1,localhost", }): @@ -274,7 +276,7 @@ def test_checksum_setting_is_scoped_to_compatible_client(self): access_key="ak", secret_key="sk", region="us-east-1", endpoint_override=endpoint) compatible = PyArrowFileIO("s3://test-bucket/", options) - self.assertEqual("WHEN_SUPPORTED", + self.assertEqual("WHEN_REQUIRED", os.environ["AWS_REQUEST_CHECKSUM_CALCULATION"]) with native.open_output_stream("test-bucket/file") as stream: stream.write(b"x") @@ -285,7 +287,8 @@ def test_checksum_setting_is_scoped_to_compatible_client(self): context = multiprocessing.get_context("spawn") parent, child = context.Pipe() process = context.Process(target=_write_in_worker, args=(child,)) - process.start() + with mock.patch.dict(os.environ, {"AWS_REQUEST_CHECKSUM_CALCULATION": "WHEN_REQUIRED"}): + process.start() child.close() try: self.assertTrue(parent.poll(20)) @@ -301,9 +304,7 @@ def test_checksum_setting_is_scoped_to_compatible_client(self): process.join() self.assertEqual(0, process.exitcode) self.assertEqual(3, len(server.put_headers)) - self.assertIn("x-amz-trailer", { - key.lower() for key in server.put_headers[0]}) - for headers in server.put_headers[1:]: + for headers in server.put_headers: self.assertNotIn("x-amz-trailer", {key.lower() for key in headers}) finally: server.shutdown() From a1b07c1a811023ba0c19ec2695ee2151744de8d0 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Thu, 8 Oct 2026 06:13:47 -0700 Subject: [PATCH 03/10] Clarify checksum configuration documentation --- paimon-python/README.md | 11 ++++------- 1 file changed, 4 insertions(+), 7 deletions(-) diff --git a/paimon-python/README.md b/paimon-python/README.md index aa2637b948b6..564a0c4c78c6 100644 --- a/paimon-python/README.md +++ b/paimon-python/README.md @@ -32,18 +32,15 @@ Both produce a source archive and wheel in `dist/`. ## PyArrow checksums with OSS and custom S3 endpoints -For endpoints that do not support the optional request checksums used by newer -PyArrow AWS SDKs, configure the process before starting the application and its -workers: +If your endpoint rejects optional request checksums, set this before starting +Python and its workers: ```bash export AWS_REQUEST_CHECKSUM_CALCULATION=WHEN_REQUIRED ``` -This setting applies to AWS SDK clients throughout the process, not just Paimon. -PyPaimon does not modify it. Existing workers and clients must be restarted after -changing it. This only addresses optional write checksums; required checksums on -batch deletion are a separate compatibility concern. +This affects AWS SDK clients process-wide; PyPaimon does not set it automatically. +Required checksums, including those for batch deletion, remain enabled. # OSS metadata commits From bd3e8f43f7e1adfed3e79c1cce922755c4616473 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Thu, 8 Oct 2026 06:24:14 -0700 Subject: [PATCH 04/10] Enable checksum compatibility by default for S3-compatible clients --- paimon-python/README.md | 17 +++++---- .../pypaimon/common/options/config.py | 3 ++ .../pypaimon/filesystem/pyarrow_file_io.py | 7 ++++ .../tests/s3_client_compatibility_test.py | 35 +++++++++++++------ 4 files changed, 42 insertions(+), 20 deletions(-) diff --git a/paimon-python/README.md b/paimon-python/README.md index 564a0c4c78c6..a5f3f4ab2b4c 100644 --- a/paimon-python/README.md +++ b/paimon-python/README.md @@ -32,15 +32,14 @@ Both produce a source archive and wheel in `dist/`. ## PyArrow checksums with OSS and custom S3 endpoints -If your endpoint rejects optional request checksums, set this before starting -Python and its workers: - -```bash -export AWS_REQUEST_CHECKSUM_CALCULATION=WHEN_REQUIRED -``` - -This affects AWS SDK clients process-wide; PyPaimon does not set it automatically. -Required checksums, including those for batch deletion, remain enabled. +PyArrow-backed OSS and custom S3 clients default +`AWS_REQUEST_CHECKSUM_CALCULATION` to `WHEN_REQUIRED` for compatibility. +No user configuration is needed; explicit environment settings take precedence. +This affects subsequently created AWS SDK clients process-wide, including workers. +Required checksums remain enabled. + +Set `fs.s3.checksum-compatibility.enabled=false` to disable this default. +Disabling it does not undo a setting already applied to the process. # OSS metadata commits diff --git a/paimon-python/pypaimon/common/options/config.py b/paimon-python/pypaimon/common/options/config.py index e67030957954..049fd0e913bb 100644 --- a/paimon-python/pypaimon/common/options/config.py +++ b/paimon-python/pypaimon/common/options/config.py @@ -46,6 +46,9 @@ class OssOptions: class S3Options: + CHECKSUM_COMPATIBILITY_ENABLED = ( + ConfigOptions.key("fs.s3.checksum-compatibility.enabled").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( diff --git a/paimon-python/pypaimon/filesystem/pyarrow_file_io.py b/paimon-python/pypaimon/filesystem/pyarrow_file_io.py index 825ff952509e..ffc0ab0d5886 100644 --- a/paimon-python/pypaimon/filesystem/pyarrow_file_io.py +++ b/paimon-python/pypaimon/filesystem/pyarrow_file_io.py @@ -107,6 +107,11 @@ 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 self._uses_s3_compatibility() and self.properties.get(S3Options.CHECKSUM_COMPATIBILITY_ENABLED): + # 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__. @@ -223,6 +228,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. @@ -252,6 +258,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")) diff --git a/paimon-python/pypaimon/tests/s3_client_compatibility_test.py b/paimon-python/pypaimon/tests/s3_client_compatibility_test.py index 5f981b4f0448..7fc03dbcba2a 100644 --- a/paimon-python/pypaimon/tests/s3_client_compatibility_test.py +++ b/paimon-python/pypaimon/tests/s3_client_compatibility_test.py @@ -40,6 +40,7 @@ def _restore_s3_file_io(connection): + os.environ.pop("AWS_REQUEST_CHECKSUM_CALCULATION", None) connection.send("ready") payload = connection.recv_bytes() client = object() @@ -61,6 +62,7 @@ def create_client(**kwargs): def _write_in_worker(connection): + os.environ.pop("AWS_REQUEST_CHECKSUM_CALCULATION", None) connection.send("ready") try: file_io = pickle.loads(connection.recv_bytes()) @@ -105,6 +107,20 @@ def log_message(self, *args): class S3ClientCompatibilityTest(unittest.TestCase): + def setUp(self): + environment = mock.patch.dict(os.environ) + environment.start() + self.addCleanup(environment.stop) + os.environ.pop("AWS_REQUEST_CHECKSUM_CALCULATION", None) + + def test_checksum_compatibility_can_be_disabled(self): + with mock.patch("pyarrow.fs.S3FileSystem", return_value=mock.Mock()): + PyArrowFileIO("s3://test-bucket/", Options({ + "fs.s3.endpoint": "http://minio:9000", + "fs.s3.checksum-compatibility.enabled": "false", + })) + self.assertNotIn("AWS_REQUEST_CHECKSUM_CALCULATION", os.environ) + def _new_file_io(self, scheme="s3"): if scheme == "oss": options = Options({ @@ -119,7 +135,7 @@ def _new_file_io(self, scheme="s3"): file_io = PyArrowFileIO(scheme + "://test-bucket/", options) return file_io - def test_environment_unchanged_after_client_creation_failure(self): + def test_process_default_retained_after_client_creation_failure(self): for previous in (None, "WHEN_SUPPORTED"): with self.subTest(previous=previous), mock.patch.dict(os.environ, {}, clear=True): if previous is not None: @@ -127,12 +143,11 @@ def test_environment_unchanged_after_client_creation_failure(self): with mock.patch("pyarrow.fs.S3FileSystem", side_effect=RuntimeError("failed")): with self.assertRaisesRegex(RuntimeError, "failed"): self._new_file_io_failure() - self.assertEqual(previous, os.environ.get("AWS_REQUEST_CHECKSUM_CALCULATION")) + self.assertEqual(previous or "WHEN_REQUIRED", os.environ.get("AWS_REQUEST_CHECKSUM_CALCULATION")) def _new_file_io_failure(self): PyArrowFileIO("s3://test-bucket/", Options({"fs.s3.endpoint": "http://minio:9000"})) - @mock.patch.dict(os.environ, {"AWS_REQUEST_CHECKSUM_CALCULATION": "WHEN_REQUIRED"}) def test_pickle_rebuilds_client_in_started_worker(self): methods = [name for name in ("spawn", "fork") if name in multiprocessing.get_all_start_methods()] @@ -182,7 +197,7 @@ def create_client(**kwargs): os.environ["AWS_REQUEST_CHECKSUM_CALCULATION"]) self.assertEqual(["WHEN_SUPPORTED"], settings) - def test_initialization_preserves_environment_for_all_s3_schemes(self): + def test_initialization_sets_default_for_all_s3_schemes(self): options = Options({ S3Options.S3_ENDPOINT.key(): "http://minio:9000", }) @@ -198,8 +213,8 @@ def create_client(**kwargs): mock.patch("pyarrow.fs.S3FileSystem", side_effect=create_client): PyArrowFileIO( "{}://test-bucket/warehouse".format(scheme), options) - self.assertNotIn("AWS_REQUEST_CHECKSUM_CALCULATION", os.environ) - self.assertEqual([None], settings) + self.assertEqual("WHEN_REQUIRED", os.environ["AWS_REQUEST_CHECKSUM_CALCULATION"]) + self.assertEqual(["WHEN_REQUIRED"], settings) def test_native_s3_does_not_change_checksum_setting(self): with mock.patch.dict("os.environ", {}, clear=True), \ @@ -250,7 +265,7 @@ def test_worker_accepts_old_non_oss_pickle(self): parse(pyarrow.__version__) >= parse("22.0.0"), "requires PyArrow 22+ optional request checksums", ) - def test_explicit_process_setting_applies_to_parent_and_worker(self): + def test_default_applies_to_parent_and_worker(self): server = _HTTPServer( ("127.0.0.1", 0), _UploadHandler) server.requests = [] @@ -268,14 +283,13 @@ def test_explicit_process_setting_applies_to_parent_and_worker(self): "fs.s3.path.style.access": "true", }) with mock.patch.dict(os.environ, { - "AWS_REQUEST_CHECKSUM_CALCULATION": "WHEN_REQUIRED", "NO_PROXY": "127.0.0.1,localhost", "no_proxy": "127.0.0.1,localhost", }): + compatible = PyArrowFileIO("s3://test-bucket/", options) native = pafs.S3FileSystem( access_key="ak", secret_key="sk", region="us-east-1", endpoint_override=endpoint) - compatible = PyArrowFileIO("s3://test-bucket/", options) self.assertEqual("WHEN_REQUIRED", os.environ["AWS_REQUEST_CHECKSUM_CALCULATION"]) with native.open_output_stream("test-bucket/file") as stream: @@ -287,8 +301,7 @@ def test_explicit_process_setting_applies_to_parent_and_worker(self): context = multiprocessing.get_context("spawn") parent, child = context.Pipe() process = context.Process(target=_write_in_worker, args=(child,)) - with mock.patch.dict(os.environ, {"AWS_REQUEST_CHECKSUM_CALCULATION": "WHEN_REQUIRED"}): - process.start() + process.start() child.close() try: self.assertTrue(parent.poll(20)) From c286cf0b73eba91c17b7b5df52a3e7b7344a05c7 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Thu, 8 Oct 2026 06:32:21 -0700 Subject: [PATCH 05/10] Limit checksum defaults to PyArrow 22 and newer --- docs/docs/pypaimon/catalogs.mdx | 10 ++++++++ paimon-python/README.md | 11 -------- .../pypaimon/filesystem/pyarrow_file_io.py | 2 ++ .../tests/s3_client_compatibility_test.py | 25 ++++++++++++++++--- 4 files changed, 33 insertions(+), 15 deletions(-) diff --git a/docs/docs/pypaimon/catalogs.mdx b/docs/docs/pypaimon/catalogs.mdx index 45f035b23654..4188f2f41fa1 100644 --- a/docs/docs/pypaimon/catalogs.mdx +++ b/docs/docs/pypaimon/catalogs.mdx @@ -170,6 +170,16 @@ 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 + +With PyArrow 22+, PyArrow-backed OSS/custom-S3 clients automatically default +`AWS_REQUEST_CHECKSUM_CALCULATION` to `WHEN_REQUIRED`. Older PyArrow versions are unchanged. +Explicit environment settings take precedence; required checksums remain enabled. + +Set catalog option `fs.s3.checksum-compatibility.enabled=false` to opt out. +The default affects subsequently created AWS clients process-wide, including data-loading +subprocesses. Opting out does not undo a setting already applied to the process. + ## Create Database Tables belong to a database. Create the database before creating its tables. diff --git a/paimon-python/README.md b/paimon-python/README.md index a5f3f4ab2b4c..f4dd3c7fe828 100644 --- a/paimon-python/README.md +++ b/paimon-python/README.md @@ -30,17 +30,6 @@ python -m build Both produce a source archive and wheel in `dist/`. -## PyArrow checksums with OSS and custom S3 endpoints - -PyArrow-backed OSS and custom S3 clients default -`AWS_REQUEST_CHECKSUM_CALCULATION` to `WHEN_REQUIRED` for compatibility. -No user configuration is needed; explicit environment settings take precedence. -This affects subsequently created AWS SDK clients process-wide, including workers. -Required checksums remain enabled. - -Set `fs.s3.checksum-compatibility.enabled=false` to disable this default. -Disabling it does not undo a setting already applied to the process. - # OSS metadata commits Install `pypaimon[oss]` (legacy PyArrow data access) or `pypaimon[jindo]` diff --git a/paimon-python/pypaimon/filesystem/pyarrow_file_io.py b/paimon-python/pypaimon/filesystem/pyarrow_file_io.py index ffc0ab0d5886..8a1e6ae28deb 100644 --- a/paimon-python/pypaimon/filesystem/pyarrow_file_io.py +++ b/paimon-python/pypaimon/filesystem/pyarrow_file_io.py @@ -108,6 +108,8 @@ def _uses_s3_compatibility(self) -> bool: 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_ENABLED): # Process-wide default; preserve explicit settings and do not restore it. os.environ.setdefault("AWS_REQUEST_CHECKSUM_CALCULATION", "WHEN_REQUIRED") diff --git a/paimon-python/pypaimon/tests/s3_client_compatibility_test.py b/paimon-python/pypaimon/tests/s3_client_compatibility_test.py index 7fc03dbcba2a..c4cba6712513 100644 --- a/paimon-python/pypaimon/tests/s3_client_compatibility_test.py +++ b/paimon-python/pypaimon/tests/s3_client_compatibility_test.py @@ -113,6 +113,20 @@ def setUp(self): self.addCleanup(environment.stop) os.environ.pop("AWS_REQUEST_CHECKSUM_CALCULATION", None) + def _expected_checksum_default(self): + return "WHEN_REQUIRED" if parse(pyarrow.__version__) >= parse("22.0.0") else None + + def test_checksum_default_version_boundary(self): + for version in ("19.0.1", "20.0.0", "21.0.0", "22.0.0", "23.0.0"): + for scheme in ("oss", "s3", "s3a", "s3n"): + with self.subTest(version=version, scheme=scheme), \ + mock.patch.dict(os.environ), \ + mock.patch.object(pyarrow, "__version__", version): + os.environ.pop("AWS_REQUEST_CHECKSUM_CALCULATION", None) + self._new_file_io(scheme) + expected = "WHEN_REQUIRED" if version in ("22.0.0", "23.0.0") else None + self.assertEqual(expected, os.environ.get("AWS_REQUEST_CHECKSUM_CALCULATION")) + def test_checksum_compatibility_can_be_disabled(self): with mock.patch("pyarrow.fs.S3FileSystem", return_value=mock.Mock()): PyArrowFileIO("s3://test-bucket/", Options({ @@ -143,7 +157,8 @@ def test_process_default_retained_after_client_creation_failure(self): with mock.patch("pyarrow.fs.S3FileSystem", side_effect=RuntimeError("failed")): with self.assertRaisesRegex(RuntimeError, "failed"): self._new_file_io_failure() - self.assertEqual(previous or "WHEN_REQUIRED", os.environ.get("AWS_REQUEST_CHECKSUM_CALCULATION")) + self.assertEqual(previous or self._expected_checksum_default(), + os.environ.get("AWS_REQUEST_CHECKSUM_CALCULATION")) def _new_file_io_failure(self): PyArrowFileIO("s3://test-bucket/", Options({"fs.s3.endpoint": "http://minio:9000"})) @@ -165,7 +180,8 @@ def test_pickle_rebuilds_client_in_started_worker(self): self.assertNotIn("filesystem", file_io.__getstate__()) parent.send_bytes(pickle.dumps(file_io)) self.assertTrue(parent.poll(20)) - self.assertEqual(("WHEN_REQUIRED", ["WHEN_REQUIRED"], 1, True), parent.recv()) + expected = self._expected_checksum_default() + self.assertEqual((expected, [expected], 1, True), parent.recv()) finally: parent.close() process.join(20) @@ -213,8 +229,9 @@ def create_client(**kwargs): mock.patch("pyarrow.fs.S3FileSystem", side_effect=create_client): PyArrowFileIO( "{}://test-bucket/warehouse".format(scheme), options) - self.assertEqual("WHEN_REQUIRED", os.environ["AWS_REQUEST_CHECKSUM_CALCULATION"]) - self.assertEqual(["WHEN_REQUIRED"], settings) + self.assertEqual(self._expected_checksum_default(), + os.environ.get("AWS_REQUEST_CHECKSUM_CALCULATION")) + self.assertEqual([self._expected_checksum_default()], settings) def test_native_s3_does_not_change_checksum_setting(self): with mock.patch.dict("os.environ", {}, clear=True), \ From 497acdd2556f74e63b62657038374f1a302257fb Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Thu, 8 Oct 2026 06:33:46 -0700 Subject: [PATCH 06/10] Simplify checksum configuration documentation --- docs/docs/pypaimon/catalogs.mdx | 11 +++++------ 1 file changed, 5 insertions(+), 6 deletions(-) diff --git a/docs/docs/pypaimon/catalogs.mdx b/docs/docs/pypaimon/catalogs.mdx index 4188f2f41fa1..2fef6f91ce5f 100644 --- a/docs/docs/pypaimon/catalogs.mdx +++ b/docs/docs/pypaimon/catalogs.mdx @@ -172,13 +172,12 @@ Use this catalog for the database and table operations below. ## OSS and custom S3 checksums -With PyArrow 22+, PyArrow-backed OSS/custom-S3 clients automatically default -`AWS_REQUEST_CHECKSUM_CALCULATION` to `WHEN_REQUIRED`. Older PyArrow versions are unchanged. -Explicit environment settings take precedence; required checksums remain enabled. +For PyArrow 22+ OSS/custom-S3 clients, PyPaimon defaults +`AWS_REQUEST_CHECKSUM_CALCULATION=WHEN_REQUIRED` process-wide, preserving explicit values. +Older PyArrow versions and required checksums are unaffected. -Set catalog option `fs.s3.checksum-compatibility.enabled=false` to opt out. -The default affects subsequently created AWS clients process-wide, including data-loading -subprocesses. Opting out does not undo a setting already applied to the process. +Set `fs.s3.checksum-compatibility.enabled=false` to opt out; this does not undo an +existing process setting. ## Create Database From 17a97d3dacd04a26a092f2e7759f2b824aa5d4a6 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Thu, 8 Oct 2026 06:46:11 -0700 Subject: [PATCH 07/10] Clarify process-wide checksum auto-configuration semantics --- docs/docs/pypaimon/catalogs.mdx | 5 +++-- .../pypaimon/common/options/config.py | 4 ++-- .../pypaimon/filesystem/pyarrow_file_io.py | 2 +- .../tests/s3_client_compatibility_test.py | 19 ++++++++++++++++++- 4 files changed, 24 insertions(+), 6 deletions(-) diff --git a/docs/docs/pypaimon/catalogs.mdx b/docs/docs/pypaimon/catalogs.mdx index 2fef6f91ce5f..b88684f6616a 100644 --- a/docs/docs/pypaimon/catalogs.mdx +++ b/docs/docs/pypaimon/catalogs.mdx @@ -176,8 +176,9 @@ For PyArrow 22+ OSS/custom-S3 clients, PyPaimon defaults `AWS_REQUEST_CHECKSUM_CALCULATION=WHEN_REQUIRED` process-wide, preserving explicit values. Older PyArrow versions and required checksums are unaffected. -Set `fs.s3.checksum-compatibility.enabled=false` to opt out; this does not undo an -existing process setting. +Set `fs.s3.checksum-compatibility.auto-configure=false` on all relevant catalogs before +creating any S3 clients to disable automatic configuration. It neither clears an existing +process setting nor provides per-client checksum isolation. ## Create Database diff --git a/paimon-python/pypaimon/common/options/config.py b/paimon-python/pypaimon/common/options/config.py index 049fd0e913bb..f6d7026789b4 100644 --- a/paimon-python/pypaimon/common/options/config.py +++ b/paimon-python/pypaimon/common/options/config.py @@ -46,8 +46,8 @@ class OssOptions: class S3Options: - CHECKSUM_COMPATIBILITY_ENABLED = ( - ConfigOptions.key("fs.s3.checksum-compatibility.enabled").boolean_type().default_value(True) + 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") diff --git a/paimon-python/pypaimon/filesystem/pyarrow_file_io.py b/paimon-python/pypaimon/filesystem/pyarrow_file_io.py index 8a1e6ae28deb..446bb64ed5f2 100644 --- a/paimon-python/pypaimon/filesystem/pyarrow_file_io.py +++ b/paimon-python/pypaimon/filesystem/pyarrow_file_io.py @@ -110,7 +110,7 @@ def _uses_s3_compatibility(self) -> bool: 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_ENABLED): + 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") diff --git a/paimon-python/pypaimon/tests/s3_client_compatibility_test.py b/paimon-python/pypaimon/tests/s3_client_compatibility_test.py index c4cba6712513..3e5456779747 100644 --- a/paimon-python/pypaimon/tests/s3_client_compatibility_test.py +++ b/paimon-python/pypaimon/tests/s3_client_compatibility_test.py @@ -131,10 +131,27 @@ def test_checksum_compatibility_can_be_disabled(self): with mock.patch("pyarrow.fs.S3FileSystem", return_value=mock.Mock()): PyArrowFileIO("s3://test-bucket/", Options({ "fs.s3.endpoint": "http://minio:9000", - "fs.s3.checksum-compatibility.enabled": "false", + "fs.s3.checksum-compatibility.auto-configure": "false", })) self.assertNotIn("AWS_REQUEST_CHECKSUM_CALCULATION", os.environ) + @mock.patch.object(pyarrow, "__version__", "23.0.0") + def test_disabling_auto_configuration_preserves_existing_process_default(self): + settings = [] + + def create_client(**kwargs): + settings.append(os.environ.get("AWS_REQUEST_CHECKSUM_CALCULATION")) + return mock.Mock() + + with mock.patch("pyarrow.fs.S3FileSystem", side_effect=create_client): + for auto_configure in ("true", "false"): + PyArrowFileIO("s3://test-bucket/", Options({ + "fs.s3.endpoint": "http://minio:9000", + "fs.s3.checksum-compatibility.auto-configure": auto_configure, + })) + self.assertEqual(["WHEN_REQUIRED", "WHEN_REQUIRED"], settings) + self.assertEqual("WHEN_REQUIRED", os.environ["AWS_REQUEST_CHECKSUM_CALCULATION"]) + def _new_file_io(self, scheme="s3"): if scheme == "oss": options = Options({ From 467888d90bd16bf46c65ee5d08a1e97c84d8db0b Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Thu, 8 Oct 2026 07:16:21 -0700 Subject: [PATCH 08/10] Clarify checksum endpoint scope and explicit overrides --- docs/docs/pypaimon/catalogs.mdx | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/docs/docs/pypaimon/catalogs.mdx b/docs/docs/pypaimon/catalogs.mdx index b88684f6616a..43c189864ed4 100644 --- a/docs/docs/pypaimon/catalogs.mdx +++ b/docs/docs/pypaimon/catalogs.mdx @@ -172,9 +172,11 @@ Use this catalog for the database and table operations below. ## OSS and custom S3 checksums -For PyArrow 22+ OSS/custom-S3 clients, PyPaimon defaults -`AWS_REQUEST_CHECKSUM_CALCULATION=WHEN_REQUIRED` process-wide, preserving explicit values. -Older PyArrow versions and required checksums are unaffected. +With PyArrow 22+, PyArrow-backed OSS or S3 clients with an explicit endpoint +(including AWS endpoints) default `AWS_REQUEST_CHECKSUM_CALCULATION=WHEN_REQUIRED` +process-wide. Explicit values take precedence: `WHEN_SUPPORTED` retains optional checksums +and may remain incompatible with the endpoint. Older PyArrow versions and required +checksums are unaffected. Set `fs.s3.checksum-compatibility.auto-configure=false` on all relevant catalogs before creating any S3 clients to disable automatic configuration. It neither clears an existing From e02f12015f7dcad6cd481cf062544b3d11443fec Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Thu, 8 Oct 2026 07:28:56 -0700 Subject: [PATCH 09/10] Tighten checksum configuration documentation --- docs/docs/pypaimon/catalogs.mdx | 17 ++++++++--------- 1 file changed, 8 insertions(+), 9 deletions(-) diff --git a/docs/docs/pypaimon/catalogs.mdx b/docs/docs/pypaimon/catalogs.mdx index 43c189864ed4..2352ffe97e3e 100644 --- a/docs/docs/pypaimon/catalogs.mdx +++ b/docs/docs/pypaimon/catalogs.mdx @@ -172,15 +172,14 @@ Use this catalog for the database and table operations below. ## OSS and custom S3 checksums -With PyArrow 22+, PyArrow-backed OSS or S3 clients with an explicit endpoint -(including AWS endpoints) default `AWS_REQUEST_CHECKSUM_CALCULATION=WHEN_REQUIRED` -process-wide. Explicit values take precedence: `WHEN_SUPPORTED` retains optional checksums -and may remain incompatible with the endpoint. Older PyArrow versions and required -checksums are unaffected. - -Set `fs.s3.checksum-compatibility.auto-configure=false` on all relevant catalogs before -creating any S3 clients to disable automatic configuration. It neither clears an existing -process setting nor provides per-client checksum isolation. +On PyArrow 22+, PyArrow-backed OSS or S3 with an explicit endpoint (including AWS) +automatically defaults `AWS_REQUEST_CHECKSUM_CALCULATION=WHEN_REQUIRED` process-wide. +Explicit values win; `WHEN_SUPPORTED` may still cause incompatibility. Older PyArrow +versions and required checksums are unchanged. + +To disable auto-configuration, set `fs.s3.checksum-compatibility.auto-configure=false` +on all relevant catalogs before creating clients. This neither clears existing settings +nor isolates clients. ## Create Database From 478faa9c98ac425e93deef3bad885525e2249377 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Thu, 8 Oct 2026 08:46:03 -0700 Subject: [PATCH 10/10] Explain checksum compatibility in plain language --- docs/docs/pypaimon/catalogs.mdx | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/docs/docs/pypaimon/catalogs.mdx b/docs/docs/pypaimon/catalogs.mdx index 2352ffe97e3e..e5354c76968b 100644 --- a/docs/docs/pypaimon/catalogs.mdx +++ b/docs/docs/pypaimon/catalogs.mdx @@ -172,14 +172,14 @@ Use this catalog for the database and table operations below. ## OSS and custom S3 checksums -On PyArrow 22+, PyArrow-backed OSS or S3 with an explicit endpoint (including AWS) -automatically defaults `AWS_REQUEST_CHECKSUM_CALCULATION=WHEN_REQUIRED` process-wide. -Explicit values win; `WHEN_SUPPORTED` may still cause incompatibility. Older PyArrow -versions and required checksums are unchanged. - -To disable auto-configuration, set `fs.s3.checksum-compatibility.auto-configure=false` -on all relevant catalogs before creating clients. This neither clears existing settings -nor isolates clients. +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