Skip to content

Commit 4a8cb2c

Browse files
timsaucerclaude
andcommitted
Clear the small stuff off the distributed example
Six unrelated tidy-ups, no behaviour change to any query. `driver.py` reads a worker's row count off the last line of its stdout rather than parsing the whole buffer. A worker shares stdout with everything loaded into it, so one stray print turns `json.loads` into a failure a long way from its cause; an unreadable or absent report now says which stage and partition produced it. `dfx_storage` exports `BundledLogicalCodec`. It was constructed with `module = "dfx_storage"` but never added to the module, so that claim was untrue and the physical half was exported while the logical half was not. `write_partition` writes the file with the schema the stream declares instead of the one the child's stream reports. They agree for any well-behaved child, and pinning it to the declaration makes a disagreement loud: `StreamWriter::write` rejects a batch whose schema differs, so a child contradicting its own `schema()` fails there rather than publishing a file readers were told to expect something else from. `run_tpch.py` raises instead of asserting. That comparison is the only thing making the script a check rather than a demo, and `python -O` drops an `assert`, which would leave it printing a table it never verified. The tolerance rationale moves to the new function's docstring, where it was otherwise duplicated. The CI artifact is `test-example-wheels-x86_64`. It stopped being manylinux when these builds moved to the host, and it carries five projects rather than the two the old name and step titles implied. `dfx_engine` documents why it declares no dependency on the sibling libraries. I tried declaring them -- `dfx_engine.session` imports both, so `import dfx_engine` fails without them -- and it breaks the build outright, because neither is published: Because dfx-storage was not found in the package registry and your project depends on dfx-storage, we can conclude that your project's requirements are unsatisfiable. So the note records the dead end next to the field someone will otherwise fill in again. The requirement itself stays in the README, beside the install command that satisfies it. Verified: all three wheels build, install together the way test.yml does, and their suites pass 13 / 11 / 19. Both workflow files parse as YAML; actionlint could not run locally, as it needs Docker. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
1 parent 9f9f688 commit 4a8cb2c

7 files changed

Lines changed: 96 additions & 29 deletions

File tree

‎.github/workflows/build.yml‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -186,7 +186,7 @@ jobs:
186186
manylinux: "2_28"
187187

188188
# The example libraries below are test fixtures, not release artifacts:
189-
# the `test-ffi-manylinux-x86_64` artifact they feed is consumed only by
189+
# the `test-example-wheels-x86_64` artifact they feed is consumed only by
190190
# test.yml, which installs them on a runner like this one. They therefore
191191
# build on the host (`container: off`) rather than in the manylinux
192192
# container, which drops a container start and an in-container rustup
@@ -261,11 +261,11 @@ jobs:
261261
name: dist-manylinux-x86_64-${{ matrix.python-tag }}
262262
path: dist/*
263263

264-
- name: Archive FFI test wheel
264+
- name: Archive example test wheels
265265
if: matrix.python-tag == 'abi3'
266266
uses: actions/upload-artifact@v7
267267
with:
268-
name: test-ffi-manylinux-x86_64
268+
name: test-example-wheels-x86_64
269269
path: |
270270
examples/datafusion-ffi-example/dist/*
271271
examples/datafusion-ffi-query-planner-example/dist/*

‎.github/workflows/test.yml‎

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -73,11 +73,11 @@ jobs:
7373
path: wheels/
7474

7575
# FFI test wheel only built once (under the abi3 matrix entry in build.yml).
76-
- name: Download pre-built FFI test wheel
76+
- name: Download pre-built example test wheels
7777
if: matrix.wheel-tag == 'abi3'
7878
uses: actions/download-artifact@v8
7979
with:
80-
name: test-ffi-manylinux-x86_64
80+
name: test-example-wheels-x86_64
8181
path: wheels/
8282

8383
- name: Install from pre-built wheels
@@ -93,9 +93,9 @@ jobs:
9393
uv venv --python "${{ steps.setup-python.outputs.python-path }}"
9494
VENV_PY="$PWD/.venv/bin/python"
9595
uv sync --python "$VENV_PY" --dev --no-install-package datafusion
96-
# Search recursively: the FFI artifact bundles more than one
97-
# project, so upload-artifact keeps a `<project>/dist/` prefix
98-
# and the wheels are not all at the top of wheels/.
96+
# Search recursively: the example artifact bundles five projects,
97+
# so upload-artifact keeps a `<project>/dist/` prefix and the
98+
# wheels are not all at the top of wheels/.
9999
WHEELS=$(find wheels/ -name "*.whl")
100100
if [ -n "$WHEELS" ]; then
101101
echo "Installing wheels:"

‎examples/distributed/engine-library/pyproject.toml‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,18 @@ classifiers = [
2727
"Programming Language :: Python :: Implementation :: CPython",
2828
]
2929
dynamic = ["version"]
30+
# No `dependencies`, deliberately, though `dfx_engine.session` imports
31+
# `dfx_storage` and `dfx_udfs` at module scope and `import dfx_engine` fails
32+
# without them. Declaring them does not work: neither is published, so
33+
# resolution fails before a wheel is built at all --
34+
#
35+
# Because dfx-storage was not found in the package registry and your
36+
# project depends on dfx-storage, we can conclude that your project's
37+
# requirements are unsatisfiable.
38+
#
39+
# which would break the build rather than document the requirement. The
40+
# requirement is in the README instead, next to the install command that
41+
# satisfies it.
3042

3143
[tool.maturin]
3244
features = ["pyo3/extension-module"]

‎examples/distributed/engine-library/python/dfx_engine/driver.py‎

Lines changed: 25 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -142,6 +142,30 @@ def _dispatch(
142142
)
143143

144144

145+
def _report(task: tuple[int, int], stdout: str) -> int:
146+
"""Read a worker's row count off its last line of output.
147+
148+
The last line, not the whole stream: a worker's stdout is shared with
149+
everything loaded into it, and one stray `print` from a library -- or a
150+
warning some future dependency decides to write there -- would turn
151+
`json.loads` on the whole buffer into a confusing failure a long way from
152+
its cause.
153+
"""
154+
stage_id, partition = task
155+
lines = [line for line in stdout.splitlines() if line.strip()]
156+
if not lines:
157+
message = f"stage {stage_id} partition {partition} printed no report"
158+
raise RuntimeError(message)
159+
try:
160+
return json.loads(lines[-1])["rows"]
161+
except (ValueError, KeyError) as err:
162+
message = (
163+
f"stage {stage_id} partition {partition} printed an unreadable "
164+
f"report {lines[-1]!r}"
165+
)
166+
raise RuntimeError(message) from err
167+
168+
145169
def run_distributed(
146170
sql: str, spec: SessionSpec, extra_udfs: list[ScalarUDF] | None = None
147171
) -> DistributedResult:
@@ -216,7 +240,7 @@ def run_distributed(
216240
stage_id, partition = task
217241
failures.append(f"stage {stage_id} partition {partition} failed:\n{stderr}")
218242
continue
219-
task_rows[task] = json.loads(stdout)["rows"]
243+
task_rows[task] = _report(task, stdout)
220244

221245
if failures:
222246
raise RuntimeError("\n".join(failures))

‎examples/distributed/engine-library/src/stage.rs‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -125,6 +125,14 @@ impl ShuffleStageExec {
125125
let final_path = partition_path(&self.shuffle_dir, self.stage_id, partition);
126126
let temp_path = temp_partition_path(&self.shuffle_dir, self.stage_id, partition);
127127
let shuffle_dir = self.shuffle_dir.clone();
128+
// The schema the file is written with is the one this stream declares,
129+
// not the one the child's stream happens to report. They agree for any
130+
// well-behaved child, and pinning it to the declaration is what makes
131+
// a disagreement loud: `StreamWriter::write` rejects a batch whose
132+
// schema differs, so a child that contradicts its own `schema()` fails
133+
// here instead of publishing a file readers were told to expect
134+
// something else from.
135+
let written_schema = Arc::clone(&schema);
128136

129137
let collected = async move {
130138
let mut stream = input.execute(partition, context)?;
@@ -139,7 +147,7 @@ impl ShuffleStageExec {
139147
let file = fs::File::create(&temp_path).map_err(|err| {
140148
exec_datafusion_err!("dfx_engine: creating {}: {err}", temp_path.display())
141149
})?;
142-
let mut writer = StreamWriter::try_new(file, stream.schema().as_ref())
150+
let mut writer = StreamWriter::try_new(file, written_schema.as_ref())
143151
.map_err(|err| exec_datafusion_err!("dfx_engine: ipc writer: {err}"))?;
144152
for batch in &batches {
145153
writer

‎examples/distributed/run_tpch.py‎

Lines changed: 38 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -92,6 +92,39 @@ def reshard(
9292
return written
9393

9494

95+
def compare(table: pa.Table, reference: pa.Table) -> None:
96+
"""Raise unless `table` matches `reference`, floats to 1e-6 relative.
97+
98+
Floats get a tolerance rather than equality. Splitting a `sum` across
99+
partitions changes the order the additions happen in, and floating point
100+
addition is not associative, so the last bits of `sum_charge` legitimately
101+
differ between the two runs. Every distributed engine has this property;
102+
it is worth knowing before someone diffs two runs and concludes the split
103+
is broken.
104+
"""
105+
if table.column_names != reference.column_names:
106+
message = (
107+
f"column names differ: {table.column_names} vs {reference.column_names}"
108+
)
109+
raise ValueError(message)
110+
111+
for name in table.column_names:
112+
got = table.column(name).to_pylist()
113+
want = reference.column(name).to_pylist()
114+
if len(got) != len(want):
115+
message = f"{name}: {len(got)} rows distributed, {len(want)} local"
116+
raise ValueError(message)
117+
for lhs, rhs in zip(got, want, strict=True):
118+
close = (
119+
abs(lhs - rhs) <= 1e-6 * max(1.0, abs(rhs))
120+
if isinstance(lhs, float)
121+
else lhs == rhs
122+
)
123+
if not close:
124+
message = f"{name}: {lhs!r} distributed, {rhs!r} local"
125+
raise ValueError(message)
126+
127+
95128
def main(argv: list[str] | None = None) -> int:
96129
parser = argparse.ArgumentParser(description=__doc__)
97130
parser.add_argument(
@@ -150,24 +183,11 @@ def main(argv: list[str] | None = None) -> int:
150183
table = pa.Table.from_batches(result.batches)
151184
reference = pa.Table.from_batches(local)
152185

153-
# Compared with a tolerance, not for equality. Splitting a `sum` across
154-
# partitions changes the order the additions happen in, and floating
155-
# point addition is not associative -- so the last bits of `sum_charge`
156-
# legitimately differ between the two runs. Any distributed engine has
157-
# this property; it is worth knowing before someone diffs two runs and
158-
# concludes the split is broken.
159-
assert table.column_names == reference.column_names
160-
for name in table.column_names:
161-
got, want = (
162-
table.column(name).to_pylist(),
163-
reference.column(name).to_pylist(),
164-
)
165-
assert len(got) == len(want), name
166-
for lhs, rhs in zip(got, want, strict=True):
167-
if isinstance(lhs, float):
168-
assert abs(lhs - rhs) <= 1e-6 * max(1.0, abs(rhs)), (name, lhs, rhs)
169-
else:
170-
assert lhs == rhs, (name, lhs, rhs)
186+
# Raises rather than asserts. This comparison is the only thing that
187+
# makes the script a check rather than a demo, and `python -O` removes
188+
# an `assert` -- which would leave it printing a table it never
189+
# verified.
190+
compare(table, reference)
171191

172192
print("\nsame answer both ways (floats to within 1e-6 relative):\n")
173193
names = table.column_names

‎examples/distributed/storage-library/src/lib.rs‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@
2222
2323
use pyo3::prelude::*;
2424

25-
use crate::extension::{BundledPhysicalCodec, DfxStorageExtension};
25+
use crate::extension::{BundledLogicalCodec, BundledPhysicalCodec, DfxStorageExtension};
2626
use crate::table_provider::PyPartitionedParquetTable;
2727

2828
mod codec;
@@ -33,6 +33,9 @@ mod table_provider;
3333
#[pymodule]
3434
fn dfx_storage(m: &Bound<'_, PyModule>) -> PyResult<()> {
3535
pyo3_log::init();
36+
// Both bundled codecs, so that `module = "dfx_storage"` on each is true
37+
// and a caller inspecting a session's codecs sees a type it can look up.
38+
m.add_class::<BundledLogicalCodec>()?;
3639
m.add_class::<BundledPhysicalCodec>()?;
3740
m.add_class::<DfxStorageExtension>()?;
3841
m.add_class::<PyPartitionedParquetTable>()?;

0 commit comments

Comments
 (0)