Skip to content

Commit bb903d0

Browse files
committed
[python] Open OSS Vortex files via vortex.open(path, store=...)
The reader built an object store with vortex.store.from_url(...) and then called vortex_store.open(). The pinned Vortex 0.70.0 S3Store (and the other object stores) have no .open() method, so the OSS read workflow enabled by the endpoint normalization raised AttributeError: 'builtins.S3Store' object has no attribute 'open' before reading a single row. Switch to the supported entry point, vortex.open(path, store=...), passing the store so it carries the endpoint/credentials while vortex resolves the object. The local-file branch is unchanged. Add a reader test driving the remote branch with the real OSS store kwargs from to_vortex_specified and a stand-in vortex module: it asserts the reader opens through vortex.open(path, store=...) and never store.open(). The SDK boundary is faked so it runs on the Python lane without the native package; the end-to-end OSS read is validated against the pinned SDK in the native lane.
1 parent 18159ad commit bb903d0

2 files changed

Lines changed: 128 additions & 1 deletion

File tree

‎paimon-python/pypaimon/read/reader/format_vortex_reader.py‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -47,7 +47,11 @@ def __init__(self, file_io: FileIO, file_path: str, read_fields: List[DataField]
4747
if store_kwargs:
4848
from vortex import store
4949
vortex_store = store.from_url(file_path_for_vortex, **store_kwargs)
50-
vortex_file = vortex_store.open()
50+
# vortex 0.70.0 S3Store (and the other object stores) have no
51+
# ``.open()``; the supported entry point is ``vortex.open(path,
52+
# store=...)``. Passing the store carries the endpoint/credentials
53+
# the OSS path needs while vortex resolves the object.
54+
vortex_file = vortex.open(file_path_for_vortex, store=vortex_store)
5155
else:
5256
vortex_file = vortex.open(file_path_for_vortex)
5357

Lines changed: 123 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,123 @@
1+
# Licensed to the Apache Software Foundation (ASF) under one
2+
# or more contributor license agreements. See the NOTICE file
3+
# distributed with this work for additional information
4+
# regarding copyright ownership. The ASF licenses this file
5+
# to you under the Apache License, Version 2.0 (the
6+
# "License"); you may not use this file except in compliance
7+
# with the License. You may obtain a copy of the License at
8+
#
9+
# http://www.apache.org/licenses/LICENSE-2.0
10+
#
11+
# Unless required by applicable law or agreed to in writing,
12+
# software distributed under the License is distributed on an
13+
# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14+
# KIND, either express or implied. See the License for the
15+
# specific language governing permissions and limitations
16+
# under the License.
17+
18+
"""Reader-side coverage for the remote (object-store) branch of
19+
``FormatVortexReader``.
20+
21+
The pinned Vortex 0.70.0 object stores (``S3Store`` et al.) have no
22+
``.open()`` method; the supported entry point is ``vortex.open(path,
23+
store=...)``. This test drives the reader with the OSS store kwargs produced
24+
by the real ``to_vortex_specified`` and a stand-in ``vortex`` module, and
25+
asserts the reader opens the file through ``vortex.open(path, store=...)`` and
26+
never calls ``store.open()`` -- the pre-existing mismatch this change fixes.
27+
28+
It fakes the SDK boundary so it runs on the normal Python lane without the
29+
native ``vortex`` package; the real end-to-end OSS read is validated against
30+
the pinned SDK in the native CI lane.
31+
"""
32+
33+
import sys
34+
import unittest
35+
from unittest import mock
36+
37+
from pypaimon.common.file_io import FileIO
38+
from pypaimon.common.options import Options
39+
from pypaimon.common.options.config import OssOptions
40+
from pypaimon.schema.data_types import AtomicType, DataField
41+
42+
43+
class _FakeArrowSchema:
44+
def __init__(self, names):
45+
self.names = names
46+
47+
48+
class _FakeDType:
49+
def __init__(self, names):
50+
self._names = names
51+
52+
def to_arrow_schema(self):
53+
return _FakeArrowSchema(self._names)
54+
55+
56+
class _FakeScan:
57+
def to_arrow(self):
58+
return iter(())
59+
60+
61+
class _FakeVortexFile:
62+
def __init__(self, names):
63+
self.dtype = _FakeDType(names)
64+
65+
def scan(self, *args, **kwargs):
66+
return _FakeScan()
67+
68+
69+
class _FakeStore:
70+
"""A vortex object store stand-in. ``.open()`` must never be called: 0.70.0
71+
stores do not expose it, so hitting it means the reader regressed."""
72+
73+
def open(self, *args, **kwargs):
74+
raise AssertionError(
75+
"store.open() must not be called; use vortex.open(path, store=...)")
76+
77+
78+
class FormatVortexReaderStoreBranchTest(unittest.TestCase):
79+
80+
def test_remote_store_opens_through_vortex_open_with_store(self):
81+
captured = {}
82+
fake_store_obj = _FakeStore()
83+
84+
fake_store_module = mock.Mock()
85+
fake_store_module.from_url = mock.Mock(return_value=fake_store_obj)
86+
87+
def fake_open(path, store=None):
88+
captured['path'] = path
89+
captured['store'] = store
90+
return _FakeVortexFile(['a'])
91+
92+
fake_vortex = mock.Mock()
93+
fake_vortex.open = mock.Mock(side_effect=fake_open)
94+
fake_vortex.store = fake_store_module
95+
96+
file_path = "oss://test-bucket/db.db/t/bucket-0/data.vortex"
97+
file_io = FileIO.get(file_path, Options({
98+
OssOptions.OSS_ENDPOINT.key(): "oss-region.example.com",
99+
OssOptions.OSS_ACCESS_KEY_ID.key(): "k",
100+
OssOptions.OSS_ACCESS_KEY_SECRET.key(): "s",
101+
}))
102+
read_fields = [DataField(0, 'a', AtomicType('INT'))]
103+
104+
from pypaimon.read.reader.format_vortex_reader import FormatVortexReader
105+
with mock.patch.dict(
106+
sys.modules,
107+
{'vortex': fake_vortex, 'vortex.store': fake_store_module}):
108+
FormatVortexReader(file_io, file_path, read_fields, None)
109+
110+
# Built the store from the generated OSS kwargs...
111+
fake_store_module.from_url.assert_called_once()
112+
_, from_url_kwargs = fake_store_module.from_url.call_args
113+
self.assertEqual(
114+
from_url_kwargs.get('endpoint'),
115+
"https://test-bucket.oss-region.example.com")
116+
# ...and opened via vortex.open(path, store=...), not store.open().
117+
fake_vortex.open.assert_called_once()
118+
self.assertIs(captured['store'], fake_store_obj)
119+
self.assertTrue(str(captured['path']).startswith('s3://test-bucket/'))
120+
121+
122+
if __name__ == '__main__':
123+
unittest.main()

0 commit comments

Comments
 (0)