From beab721733db03ff31fac5f5f1afb2c83377a580 Mon Sep 17 00:00:00 2001 From: Mrigank Bhatnagar Date: Fri, 9 Oct 2026 12:13:03 +0530 Subject: [PATCH 1/3] Give the two FFI example crates a runnable entry point (#1727) --- examples/datafusion-ffi-example/README.md | 4 ++ .../python/tests/_test_run_demo.py | 19 +++++++ examples/datafusion-ffi-example/run_demo.py | 57 +++++++++++++++++++ .../README.md | 4 ++ .../python/tests/_test_run_demo.py | 17 ++++++ .../run_demo.py | 47 +++++++++++++++ 6 files changed, 148 insertions(+) create mode 100644 examples/datafusion-ffi-example/python/tests/_test_run_demo.py create mode 100644 examples/datafusion-ffi-example/run_demo.py create mode 100644 examples/datafusion-ffi-query-planner-example/python/tests/_test_run_demo.py create mode 100644 examples/datafusion-ffi-query-planner-example/run_demo.py 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..0037b919b --- /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( + [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..630b85c46 --- /dev/null +++ b/examples/datafusion-ffi-example/run_demo.py @@ -0,0 +1,57 @@ +""" +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 + +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.") + +from datafusion import LogicalPlan, SessionConfig, SessionContext, udf + + + + +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..d5ef87b25 --- /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( + [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..b788327c5 --- /dev/null +++ b/examples/datafusion-ffi-query-planner-example/run_demo.py @@ -0,0 +1,47 @@ +""" +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 + +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.") + +from datafusion import SessionConfig, SessionContext + + + + +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()) From 30596749aec97cc3467866e8fbe4dc8961c257e3 Mon Sep 17 00:00:00 2001 From: Mrigank Bhatnagar Date: Fri, 9 Oct 2026 14:06:30 +0530 Subject: [PATCH 2/3] Fix ruff lints for example scripts --- .../python/tests/_test_run_demo.py | 2 +- examples/datafusion-ffi-example/run_demo.py | 9 +++++---- .../python/tests/_test_run_demo.py | 2 +- .../datafusion-ffi-query-planner-example/run_demo.py | 9 +++++---- 4 files changed, 12 insertions(+), 10 deletions(-) diff --git a/examples/datafusion-ffi-example/python/tests/_test_run_demo.py b/examples/datafusion-ffi-example/python/tests/_test_run_demo.py index 0037b919b..0d076bf20 100644 --- a/examples/datafusion-ffi-example/python/tests/_test_run_demo.py +++ b/examples/datafusion-ffi-example/python/tests/_test_run_demo.py @@ -5,7 +5,7 @@ def test_run_demo(): script = Path(__file__).parent.parent.parent / "run_demo.py" - result = subprocess.run( + result = subprocess.run( # noqa: S603 [sys.executable, str(script)], capture_output=True, text=True, diff --git a/examples/datafusion-ffi-example/run_demo.py b/examples/datafusion-ffi-example/run_demo.py index 630b85c46..6dd949dc9 100644 --- a/examples/datafusion-ffi-example/run_demo.py +++ b/examples/datafusion-ffi-example/run_demo.py @@ -5,6 +5,7 @@ catalog provider, config extension, codec round-trip, then the same bytes decoded a second time. """ + import sys try: @@ -18,9 +19,7 @@ except ImportError: sys.exit("build the extension first:\n uv run maturin develop\nSee README.md.") -from datafusion import LogicalPlan, SessionConfig, SessionContext, udf - - +from datafusion import LogicalPlan, SessionConfig, SessionContext, udf # noqa: I001, E402 print("1. table provider") @@ -38,7 +37,9 @@ print("\n4. config extension") config = MyConfig() -config = SessionConfig({"datafusion.catalog.information_schema": "true"}).with_extension(config) +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() 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 index d5ef87b25..b467e4399 100644 --- 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 @@ -5,7 +5,7 @@ def test_run_demo(): script = Path(__file__).parent.parent.parent / "run_demo.py" - result = subprocess.run( + result = subprocess.run( # noqa: S603 [sys.executable, str(script)], capture_output=True, text=True, diff --git a/examples/datafusion-ffi-query-planner-example/run_demo.py b/examples/datafusion-ffi-query-planner-example/run_demo.py index b788327c5..d78747a2a 100644 --- a/examples/datafusion-ffi-query-planner-example/run_demo.py +++ b/examples/datafusion-ffi-query-planner-example/run_demo.py @@ -5,6 +5,7 @@ plan returned, the effect of SET ffi_query_planner.max_rows, and two planners nesting. """ + import sys try: @@ -15,16 +16,16 @@ except ImportError: sys.exit("build the extension first:\n uv run maturin develop\nSee README.md.") -from datafusion import SessionConfig, SessionContext - - +from datafusion import SessionConfig, SessionContext # noqa: I001, E402 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)") +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()) From 0082932b7e62fed7701dc56698ecae6fab2f88ac Mon Sep 17 00:00:00 2001 From: Mrigank Bhatnagar Date: Fri, 9 Oct 2026 14:09:42 +0530 Subject: [PATCH 3/3] Fix: rearrange imports instead of suppressing lint --- examples/datafusion-ffi-example/run_demo.py | 4 ++-- examples/datafusion-ffi-query-planner-example/run_demo.py | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/examples/datafusion-ffi-example/run_demo.py b/examples/datafusion-ffi-example/run_demo.py index 6dd949dc9..5f53e80f3 100644 --- a/examples/datafusion-ffi-example/run_demo.py +++ b/examples/datafusion-ffi-example/run_demo.py @@ -8,6 +8,8 @@ import sys +from datafusion import LogicalPlan, SessionConfig, SessionContext, udf + try: from datafusion_ffi_example import ( IsNullUDF, @@ -19,8 +21,6 @@ except ImportError: sys.exit("build the extension first:\n uv run maturin develop\nSee README.md.") -from datafusion import LogicalPlan, SessionConfig, SessionContext, udf # noqa: I001, E402 - print("1. table provider") ctx = SessionContext() diff --git a/examples/datafusion-ffi-query-planner-example/run_demo.py b/examples/datafusion-ffi-query-planner-example/run_demo.py index d78747a2a..03efc1159 100644 --- a/examples/datafusion-ffi-query-planner-example/run_demo.py +++ b/examples/datafusion-ffi-query-planner-example/run_demo.py @@ -8,6 +8,8 @@ import sys +from datafusion import SessionConfig, SessionContext + try: from datafusion_ffi_query_planner_example import ( MyPlannerConfig, @@ -16,8 +18,6 @@ except ImportError: sys.exit("build the extension first:\n uv run maturin develop\nSee README.md.") -from datafusion import SessionConfig, SessionContext # noqa: I001, E402 - print("1. logical plan") config = SessionConfig().with_extension(MyPlannerConfig(max_rows=5))