Skip to content

Commit 53e2b21

Browse files
committed
Fix Windows CI: release parquet/IPC file handles before tempdir cleanup
Signed-off-by: Arham Chopra <arham.chopra@cubistsystematic.com>
1 parent 4f95ec6 commit 53e2b21

1 file changed

Lines changed: 35 additions & 25 deletions

File tree

csp/tests/adapters/test_parquet_output.py

Lines changed: 35 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,18 @@
2323
START = datetime(2022, 1, 1, tzinfo=pytz.utc)
2424

2525

26+
def _read_ipc(path):
27+
"""Read an Arrow IPC stream fully into memory, leaving no open file handle.
28+
29+
``pyarrow.memory_map`` keeps an OS handle open until GC, which blocks
30+
``tempfile.TemporaryDirectory`` cleanup on Windows ("file in use by another
31+
process"). Reading the bytes up front and parsing from an in-memory buffer
32+
avoids holding the file open.
33+
"""
34+
with open(path, "rb") as fh:
35+
return pyarrow.ipc.open_stream(pyarrow.py_buffer(fh.read()))
36+
37+
2638
class TestOutputScalarTypes(unittest.TestCase):
2739
"""Test writing all scalar types to parquet and reading them back."""
2840

@@ -410,7 +422,7 @@ def g():
410422
writer.publish("y", csp.curve(float, [(timedelta(seconds=i + 1), i * 1.5) for i in range(5)]))
411423

412424
csp.run(g, starttime=START, endtime=timedelta(seconds=10))
413-
reader = pyarrow.ipc.open_stream(pyarrow.memory_map(fname))
425+
reader = _read_ipc(fname)
414426
table = reader.read_all()
415427
self.assertEqual(table.column("x").to_pylist(), [0, 10, 20, 30, 40])
416428
self.assertEqual(table.column("y").to_pylist(), [0.0, 1.5, 3.0, 4.5, 6.0])
@@ -426,7 +438,7 @@ def g():
426438
writer.publish("x", csp.curve(int, [(timedelta(seconds=1), 42)]))
427439

428440
csp.run(g, starttime=START, endtime=timedelta(seconds=5))
429-
reader = pyarrow.ipc.open_stream(pyarrow.memory_map(fname))
441+
reader = _read_ipc(fname)
430442
table = reader.read_all()
431443
self.assertEqual(table.column("x").to_pylist(), [42])
432444

@@ -772,9 +784,8 @@ def g():
772784
writer.publish("x", csp.curve(int, [(timedelta(seconds=i + 1), i) for i in range(n_rows)]))
773785

774786
csp.run(g, starttime=START, endtime=timedelta(seconds=n_rows + 5))
775-
pf = pyarrow.parquet.ParquetFile(fname)
776-
# Check actual codec in row group metadata
777-
actual_codec = pf.metadata.row_group(0).column(0).compression
787+
# Read footer eagerly (no lingering file handle: Windows can't delete an open file)
788+
actual_codec = pyarrow.parquet.read_metadata(fname).row_group(0).column(0).compression
778789
self.assertEqual(actual_codec, expected_codec)
779790
# Verify data integrity
780791
df = pandas.read_parquet(fname)
@@ -804,7 +815,7 @@ def g():
804815
writer.publish("x", csp.curve(int, [(timedelta(seconds=i + 1), i) for i in range(10)]))
805816

806817
csp.run(g, starttime=START, endtime=timedelta(seconds=15))
807-
codec = pyarrow.parquet.ParquetFile(fname).metadata.row_group(0).column(0).compression
818+
codec = pyarrow.parquet.read_metadata(fname).row_group(0).column(0).compression
808819
self.assertEqual(codec, "ZSTD")
809820

810821
def test_invalid_compression_raises_clear_error(self):
@@ -839,9 +850,8 @@ def g():
839850
writer.publish("x", csp.curve(int, [(timedelta(seconds=i + 1), i) for i in range(n_rows)]))
840851

841852
csp.run(g, starttime=START, endtime=timedelta(seconds=n_rows + 5))
842-
pf = pyarrow.parquet.ParquetFile(fname)
843853
expected_row_groups = math.ceil(n_rows / batch_size)
844-
self.assertEqual(pf.metadata.num_row_groups, expected_row_groups)
854+
self.assertEqual(pyarrow.parquet.read_metadata(fname).num_row_groups, expected_row_groups)
845855
# Verify data integrity
846856
df = pandas.read_parquet(fname)
847857
self.assertEqual(df["x"].tolist(), list(range(n_rows)))
@@ -881,8 +891,7 @@ def g():
881891
writer.publish("x", csp.curve(int, [(timedelta(seconds=1), 1)]))
882892

883893
csp.run(g, starttime=START, endtime=timedelta(seconds=5))
884-
pf = pyarrow.parquet.ParquetFile(fname)
885-
metadata = pf.schema_arrow.metadata
894+
metadata = pyarrow.parquet.read_schema(fname).metadata
886895
self.assertEqual(metadata[b"created_by"], b"test")
887896
self.assertEqual(metadata[b"version"], b"1.0")
888897

@@ -901,8 +910,7 @@ def g():
901910
writer.publish("x", csp.curve(int, [(timedelta(seconds=1), 1)]))
902911

903912
csp.run(g, starttime=START, endtime=timedelta(seconds=5))
904-
pf = pyarrow.parquet.ParquetFile(fname)
905-
x_field = pf.schema_arrow.field("x")
913+
x_field = pyarrow.parquet.read_schema(fname).field("x")
906914
self.assertEqual(x_field.metadata[b"units"], b"meters")
907915
self.assertEqual(x_field.metadata[b"source"], b"sensor_1")
908916

@@ -931,17 +939,17 @@ def g():
931939
self.assertTrue(os.path.isfile(x_file))
932940
self.assertTrue(os.path.isfile(ts_file))
933941

934-
pf_x = pyarrow.parquet.ParquetFile(x_file)
935-
pf_ts = pyarrow.parquet.ParquetFile(ts_file)
942+
schema_x = pyarrow.parquet.read_schema(x_file)
943+
schema_ts = pyarrow.parquet.read_schema(ts_file)
936944

937945
# File-level metadata on both files
938-
self.assertEqual(pf_x.schema_arrow.metadata[b"author"], b"test_suite")
939-
self.assertEqual(pf_x.schema_arrow.metadata[b"version"], b"2.0")
940-
self.assertEqual(pf_ts.schema_arrow.metadata[b"author"], b"test_suite")
941-
self.assertEqual(pf_ts.schema_arrow.metadata[b"version"], b"2.0")
946+
self.assertEqual(schema_x.metadata[b"author"], b"test_suite")
947+
self.assertEqual(schema_x.metadata[b"version"], b"2.0")
948+
self.assertEqual(schema_ts.metadata[b"author"], b"test_suite")
949+
self.assertEqual(schema_ts.metadata[b"version"], b"2.0")
942950

943951
# Column-level metadata preserved on x
944-
x_field = pf_x.schema_arrow.field("x")
952+
x_field = schema_x.field("x")
945953
self.assertEqual(x_field.metadata[b"units"], b"kg")
946954

947955
def test_file_metadata_in_split_ipc_mode(self):
@@ -965,7 +973,7 @@ def g():
965973

966974
val_file = os.path.join(outdir, "val.arrow")
967975
self.assertTrue(os.path.isfile(val_file))
968-
reader = pyarrow.ipc.open_stream(pyarrow.memory_map(val_file))
976+
reader = _read_ipc(val_file)
969977
self.assertEqual(reader.schema.metadata[b"source"], b"ipc_test")
970978

971979

@@ -1276,12 +1284,14 @@ def g():
12761284
self.assertEqual(df0["value"].tolist(), list(range(10)))
12771285
self.assertEqual(df1["value"].tolist(), list(range(10, 20)))
12781286
# Verify batch_size=4 → correct row group counts
1279-
pf0 = pyarrow.parquet.ParquetFile(os.path.join(d, "p0.parquet"))
1280-
pf1 = pyarrow.parquet.ParquetFile(os.path.join(d, "p1.parquet"))
12811287
import math
12821288

1283-
self.assertEqual(pf0.metadata.num_row_groups, math.ceil(10 / 4)) # 3
1284-
self.assertEqual(pf1.metadata.num_row_groups, math.ceil(10 / 4)) # 3
1289+
self.assertEqual(
1290+
pyarrow.parquet.read_metadata(os.path.join(d, "p0.parquet")).num_row_groups, math.ceil(10 / 4)
1291+
) # 3
1292+
self.assertEqual(
1293+
pyarrow.parquet.read_metadata(os.path.join(d, "p1.parquet")).num_row_groups, math.ceil(10 / 4)
1294+
) # 3
12851295

12861296
def test_file_visitor_exact_contract(self):
12871297
"""file_visitor called once per closed file, in order, including final at stop()."""
@@ -1335,7 +1345,7 @@ def g():
13351345
writer.publish("x", csp.curve(int, [(timedelta(seconds=1), 1)]))
13361346

13371347
csp.run(g, starttime=START, endtime=timedelta(seconds=5))
1338-
reader = pyarrow.ipc.open_stream(pyarrow.memory_map(fname))
1348+
reader = _read_ipc(fname)
13391349
schema = reader.schema
13401350
self.assertEqual(schema.metadata[b"author"], b"test")
13411351
x_field = schema.field("x")

0 commit comments

Comments
 (0)