Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions examples/datafusion-ffi-example/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
19 changes: 19 additions & 0 deletions examples/datafusion-ffi-example/python/tests/_test_run_demo.py
Original file line number Diff line number Diff line change
@@ -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
58 changes: 58 additions & 0 deletions examples/datafusion-ffi-example/run_demo.py
Original file line number Diff line number Diff line change
@@ -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()
4 changes: 4 additions & 0 deletions examples/datafusion-ffi-query-planner-example/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Original file line number Diff line number Diff line change
@@ -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
48 changes: 48 additions & 0 deletions examples/datafusion-ffi-query-planner-example/run_demo.py
Original file line number Diff line number Diff line change
@@ -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())