diff --git a/examples/datafusion-ffi-example/README.md b/examples/datafusion-ffi-example/README.md index aadea909f..62e6b4ce1 100644 --- a/examples/datafusion-ffi-example/README.md +++ b/examples/datafusion-ffi-example/README.md @@ -19,6 +19,10 @@ # DataFusion Python FFI provider example +**What this is:** A testbed that exports table providers, catalogs, functions, and codecs to Python. +**If you are learning the protocol:** Read the [Extension Guide](https://datafusion.apache.org/python/extension-guide/index.html). +**To run the demo:** `uv run python examples/datafusion-ffi-example/run_demo.py` (after `uv run maturin develop` in this directory). + This crate is the **provider library** in the three-library query-planning example. It exports table providers, functions, and the logical and physical codecs needed to serialize objects owned by this library. The companion planner is in [`../datafusion-ffi-query-planner-example`](../datafusion-ffi-query-planner-example/). The example intentionally uses separate `cdylib` crates for these roles: diff --git a/examples/datafusion-ffi-example/python/tests/_test_run_demo.py b/examples/datafusion-ffi-example/python/tests/_test_run_demo.py new file mode 100644 index 000000000..0d076bf20 --- /dev/null +++ b/examples/datafusion-ffi-example/python/tests/_test_run_demo.py @@ -0,0 +1,19 @@ +import subprocess +import sys +from pathlib import Path + + +def test_run_demo(): + script = Path(__file__).parent.parent.parent / "run_demo.py" + result = subprocess.run( # noqa: S603 + [sys.executable, str(script)], + capture_output=True, + text=True, + check=True, + ) + assert "1. table provider" in result.stdout + assert "2. functions" in result.stdout + assert "3. catalog provider" in result.stdout + assert "4. config extension" in result.stdout + assert "5. codec round-trip" in result.stdout + assert "6. the same bytes decoded a second time" in result.stdout diff --git a/examples/datafusion-ffi-example/run_demo.py b/examples/datafusion-ffi-example/run_demo.py new file mode 100644 index 000000000..5f53e80f3 --- /dev/null +++ b/examples/datafusion-ffi-example/run_demo.py @@ -0,0 +1,58 @@ +""" +DataFusion Python FFI provider example. + +Walks the conformance matrix in numbered sections: table provider, functions, +catalog provider, config extension, codec round-trip, then the same bytes +decoded a second time. +""" + +import sys + +from datafusion import LogicalPlan, SessionConfig, SessionContext, udf + +try: + from datafusion_ffi_example import ( + IsNullUDF, + MyCatalogProvider, + MyConfig, + MyLogicalExtensionCodec, + MyTableProvider, + ) +except ImportError: + sys.exit("build the extension first:\n uv run maturin develop\nSee README.md.") + + +print("1. table provider") +ctx = SessionContext() +ctx.register_table("numbers", MyTableProvider(1, 6, 1)) +ctx.sql('SELECT "A" FROM numbers').show() + +print("\n2. functions") +ctx.register_udf(udf(IsNullUDF())) +ctx.sql('SELECT "A", my_custom_is_null("A") AS is_null FROM numbers').show() + +print("\n3. catalog provider") +ctx.register_catalog_provider("ffi_catalog", MyCatalogProvider()) +ctx.sql("SELECT * FROM ffi_catalog.my_schema.my_table").show() + +print("\n4. config extension") +config = MyConfig() +config = SessionConfig( + {"datafusion.catalog.information_schema": "true"} +).with_extension(config) +config.set("my_config.baz_count", "42") +ctx2 = SessionContext(config) +ctx2.sql("SHOW my_config.baz_count;").show() + +print("\n5. codec round-trip") +codec = MyLogicalExtensionCodec() +ctx3 = SessionContext().with_logical_extension_codec(codec) +ctx3.register_table("numbers", MyTableProvider(1, 4, 1)) +plan = ctx3.sql('SELECT "A" FROM numbers').logical_plan() +blob = plan.to_bytes(ctx3) +restored = LogicalPlan.from_bytes(ctx3, blob) +ctx3.create_dataframe_from_logical_plan(restored).show() + +print("\n6. the same bytes decoded a second time") +restored2 = LogicalPlan.from_bytes(ctx3, blob) +ctx3.create_dataframe_from_logical_plan(restored2).show() diff --git a/examples/datafusion-ffi-query-planner-example/README.md b/examples/datafusion-ffi-query-planner-example/README.md index af128886f..88ef66d1c 100644 --- a/examples/datafusion-ffi-query-planner-example/README.md +++ b/examples/datafusion-ffi-query-planner-example/README.md @@ -19,6 +19,10 @@ # DataFusion Python FFI query planner example +**What this is:** A query planner extension that adds a configurable global limit to query plans. +**If you are learning the protocol:** Read the [Extension Guide](https://datafusion.apache.org/python/extension-guide/index.html). +**To run the demo:** `uv run python examples/datafusion-ffi-query-planner-example/run_demo.py` (after `uv run maturin develop` in this directory). + This crate is an independent query-planner Python extension. Together with [`../datafusion-ffi-example`](../datafusion-ffi-example/) it demonstrates a real three-library plan exchange: - **A — `datafusion-python`:** owns the session and final execution. diff --git a/examples/datafusion-ffi-query-planner-example/python/tests/_test_run_demo.py b/examples/datafusion-ffi-query-planner-example/python/tests/_test_run_demo.py new file mode 100644 index 000000000..b467e4399 --- /dev/null +++ b/examples/datafusion-ffi-query-planner-example/python/tests/_test_run_demo.py @@ -0,0 +1,17 @@ +import subprocess +import sys +from pathlib import Path + + +def test_run_demo(): + script = Path(__file__).parent.parent.parent / "run_demo.py" + result = subprocess.run( # noqa: S603 + [sys.executable, str(script)], + capture_output=True, + text=True, + check=True, + ) + assert "1. logical plan" in result.stdout + assert "2. physical plan returned" in result.stdout + assert "3. effect of SET ffi_query_planner.max_rows" in result.stdout + assert "4. two planners nesting" in result.stdout diff --git a/examples/datafusion-ffi-query-planner-example/run_demo.py b/examples/datafusion-ffi-query-planner-example/run_demo.py new file mode 100644 index 000000000..03efc1159 --- /dev/null +++ b/examples/datafusion-ffi-query-planner-example/run_demo.py @@ -0,0 +1,48 @@ +""" +DataFusion Python FFI query planner example. + +Prints plans, showing the logical plan handed to the planner, the physical +plan returned, the effect of SET ffi_query_planner.max_rows, and two +planners nesting. +""" + +import sys + +from datafusion import SessionConfig, SessionContext + +try: + from datafusion_ffi_query_planner_example import ( + MyPlannerConfig, + MyQueryPlanner, + ) +except ImportError: + sys.exit("build the extension first:\n uv run maturin develop\nSee README.md.") + + +print("1. logical plan") +config = SessionConfig().with_extension(MyPlannerConfig(max_rows=5)) +ctx = SessionContext(config) + +ctx.sql( + "CREATE TABLE t AS SELECT * FROM (VALUES (1), (2), (3), (4), (5), (6), (7)) AS t(a)" +) +df = ctx.sql("SELECT * FROM t") +print(df.logical_plan().display_indent()) + +print("\n2. physical plan returned") +planner = MyQueryPlanner() +ctx.set_query_planner(planner) + +plan = df.execution_plan() +print(plan.display_indent()) + +print("\n3. effect of SET ffi_query_planner.max_rows") +ctx.sql("SET ffi_query_planner.max_rows = 2").collect() +plan2 = df.execution_plan() +print(plan2.display_indent()) + +print("\n4. two planners nesting") +outer_planner = MyQueryPlanner(fallback=planner) +ctx.set_query_planner(outer_planner) +plan3 = df.execution_plan() +print(plan3.display_indent())