Skip to content
This repository was archived by the owner on Sep 9, 2026. It is now read-only.

Commit 9a6b1e6

Browse files
author
Joan Fontanals
authored
test: move the pydantic check inside test (#1812)
Signed-off-by: Joan Fontanals Martinez <joan.martinez@jina.ai>
1 parent 83d2236 commit 9a6b1e6

2 files changed

Lines changed: 44 additions & 38 deletions

File tree

‎docarray/store/helpers.py‎

Lines changed: 29 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -174,32 +174,36 @@ def _from_binary_stream(
174174
compress: Optional[str] = None,
175175
show_progress: bool = False,
176176
) -> Iterator['T']:
177-
if show_progress:
178-
pbar, t = _get_progressbar(
179-
'Deserializing', disable=not show_progress, total=total
180-
)
181-
else:
182-
pbar = nullcontext()
183-
184-
with pbar:
177+
try:
185178
if show_progress:
186-
_total_size = 0
187-
pbar.start_task(t)
188-
while True:
189-
len_bytes = stream.read(4)
190-
if len(len_bytes) < 4:
191-
raise ValueError('Unexpected end of stream')
192-
len_item = int.from_bytes(len_bytes, 'big', signed=False)
193-
if len_item == 0:
194-
break
195-
item_bytes = stream.read(len_item)
196-
if len(item_bytes) < len_item:
197-
raise ValueError('Unexpected end of stream')
198-
item = cls.from_bytes(item_bytes, protocol=protocol, compress=compress)
199-
200-
yield item
179+
pbar, t = _get_progressbar(
180+
'Deserializing', disable=not show_progress, total=total
181+
)
182+
else:
183+
pbar = nullcontext()
201184

185+
with pbar:
202186
if show_progress:
203-
_total_size += len_item + 4
204-
pbar.update(t, advance=1, total_size=str(filesize.decimal(_total_size)))
187+
_total_size = 0
188+
pbar.start_task(t)
189+
while True:
190+
len_bytes = stream.read(4)
191+
if len(len_bytes) < 4:
192+
raise ValueError('Unexpected end of stream')
193+
len_item = int.from_bytes(len_bytes, 'big', signed=False)
194+
if len_item == 0:
195+
break
196+
item_bytes = stream.read(len_item)
197+
if len(item_bytes) < len_item:
198+
raise ValueError('Unexpected end of stream')
199+
item = cls.from_bytes(item_bytes, protocol=protocol, compress=compress)
200+
201+
yield item
202+
203+
if show_progress:
204+
_total_size += len_item + 4
205+
pbar.update(
206+
t, advance=1, total_size=str(filesize.decimal(_total_size))
207+
)
208+
finally:
205209
stream.close()

‎tests/integrations/store/test_file.py‎

Lines changed: 15 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -6,12 +6,12 @@
66
from docarray import DocList
77
from docarray.documents import TextDoc
88
from docarray.store.file import ConcurrentPushException, FileDocStore
9-
from docarray.utils._internal.cache import _get_cache_path
109
from docarray.utils._internal.pydantic import is_pydantic_v2
10+
from docarray.utils._internal.cache import _get_cache_path
1111
from tests.integrations.store import gen_text_docs, get_test_da, profile_memory
1212

1313
DA_LEN: int = 2**10
14-
TOLERANCE_RATIO = 0.1 # Percentage of difference allowed in stream vs non-stream test
14+
TOLERANCE_RATIO = 0.1 # Percentage of difference allowed when streaming between a long and a shorter DA
1515

1616

1717
def test_path_resolution():
@@ -23,7 +23,6 @@ def test_path_resolution():
2323

2424

2525
def test_pushpull_correct(capsys, tmp_path: Path):
26-
tmp_path.mkdir(parents=True, exist_ok=True)
2726
namespace_dir = tmp_path
2827
da1 = get_test_da(DA_LEN)
2928

@@ -51,7 +50,6 @@ def test_pushpull_correct(capsys, tmp_path: Path):
5150

5251

5352
def test_pushpull_stream_correct(capsys, tmp_path: Path):
54-
tmp_path.mkdir(parents=True, exist_ok=True)
5553
namespace_dir = tmp_path
5654
da1 = get_test_da(DA_LEN)
5755

@@ -85,10 +83,8 @@ def test_pushpull_stream_correct(capsys, tmp_path: Path):
8583

8684

8785
# for some reason this test is failing with pydantic v2
88-
@pytest.mark.skipif(is_pydantic_v2, reason="Not working with pydantic v2 for now")
8986
@pytest.mark.slow
9087
def test_pull_stream_vs_pull_full(tmp_path: Path):
91-
tmp_path.mkdir(parents=True, exist_ok=True)
9288
namespace_dir = tmp_path
9389
DocList[TextDoc].push_stream(
9490
gen_text_docs(DA_LEN * 1),
@@ -136,15 +132,23 @@ def get_total_full(url: str):
136132
), 'Streamed and non-streamed pull should have similar statistics'
137133

138134
assert (
139-
abs(long_stream_peak - short_stream_peak) / short_stream_peak < TOLERANCE_RATIO
140-
), 'Streamed memory usage should not be dependent on the size of the data'
135+
long_full_peak > long_stream_peak
136+
), 'Peak of memory using full should be larger than when streaming'
141137
assert (
142-
abs(long_full_peak - short_full_peak) / short_full_peak > TOLERANCE_RATIO
143-
), 'Full pull memory usage should be dependent on the size of the data'
138+
short_full_peak > short_stream_peak
139+
), 'Peak of memory using full should be larger than when streaming'
140+
if not is_pydantic_v2:
141+
# I bet there is some memory that Pydantic is leaking
142+
assert (
143+
abs(long_stream_peak - short_stream_peak) / short_stream_peak
144+
< TOLERANCE_RATIO
145+
), 'Streamed memory usage should not be dependent on the size of the data'
146+
assert (
147+
abs(long_full_peak - short_full_peak) / short_full_peak > TOLERANCE_RATIO
148+
), 'Full pull memory usage should be dependent on the size of the data'
144149

145150

146151
def test_list_and_delete(tmp_path: Path):
147-
tmp_path.mkdir(parents=True, exist_ok=True)
148152
namespace_dir = str(tmp_path)
149153

150154
da_names = FileDocStore.list(namespace_dir, show_table=False)
@@ -179,7 +183,6 @@ def test_list_and_delete(tmp_path: Path):
179183

180184
def test_concurrent_push_pull(tmp_path: Path):
181185
# Push to DA that is being pulled should not mess up the pull
182-
tmp_path.mkdir(parents=True, exist_ok=True)
183186
namespace_dir = tmp_path
184187

185188
DocList[TextDoc].push_stream(
@@ -214,7 +217,6 @@ def test_concurrent_push(tmp_path: Path):
214217
# Double push should fail the second push
215218
import time
216219

217-
tmp_path.mkdir(parents=True, exist_ok=True)
218220
namespace_dir = tmp_path
219221

220222
DocList[TextDoc].push_stream(

0 commit comments

Comments
 (0)