datafusion.extensions#

Protocols and value types for installing extensions on a session context.

An extension is a reusable configuration object — typically shipped by a separate compiled library — that contributes components to a SessionContext. It implements SessionComponentsExportable by returning a SessionExtensionComponents describing what it contributes, and is installed with with_extensions():

ctx = SessionContext().with_extensions(MyLibraryExtension())

Codecs and planners install in two phases: every SessionComponentsExportable runs first and its codecs are installed, then every SessionPlannerExportable runs in argument order. A bundle implements either hook or both. Bundle order is significant for planners, which nest, and irrelevant for codecs, which accumulate.

Of the five names here, only the two bundle hooks are @runtime_checkable, because with_extensions() dispatches on them from Python. QueryPlannerExportable and PhysicalOptimizerRuleExportable are type hints only, matching the other capsule-getter protocols in datafusion.user_defined and datafusion.catalog.

See Extension bundles in the online documentation for why the phases are split and for a worked implementation.

Classes#

PhysicalOptimizerRuleExportable

Type hint for object that has __datafusion_physical_optimizer_rule__ PyCapsule.

QueryPlannerExportable

Type hint for object that has a __datafusion_query_planner__ PyCapsule.

SessionComponentsExportable

Type hint for extension bundles installable via with_extensions.

SessionExtensionComponents

Components an extension contributes to a session context.

SessionPlannerExportable

Type hint for extension bundles that contribute a query planner.

Module Contents#

class datafusion.extensions.PhysicalOptimizerRuleExportable#

Bases: Protocol

Type hint for object that has __datafusion_physical_optimizer_rule__ PyCapsule.

The method returns a PyCapsule wrapping an FFI_PhysicalOptimizerRule, typically produced by a separate compiled extension. It takes no argument: a rule needs neither a codec nor a task-context provider, so there is nothing session-scoped to hand it.

Rules accumulate rather than replace. Install one with add_physical_optimizer_rule() — see Optimizer rules and configuration.

Examples

The getter is the whole protocol, and a capsule is what it must return — anything else is refused where it is installed rather than at plan time:

>>> from datafusion import SessionContext
>>> ctx = SessionContext()
>>> ctx.add_physical_optimizer_rule(object())
Traceback (most recent call last):
    ...
RuntimeError: "Invalid datafusion_physical_optimizer_rule...

Real usage. Skipped here (needs a built extension library); parsed out of this docstring and run for real by test_physical_optimizer_rule_docstring_example_still_runs in datafusion-ffi-example.

>>> from datafusion_ffi_example import MyPhysicalOptimizerRule  
>>> ctx.add_physical_optimizer_rule(MyPhysicalOptimizerRule())  
__datafusion_physical_optimizer_rule__() → object#
class datafusion.extensions.QueryPlannerExportable#

Bases: Protocol

Type hint for object that has a __datafusion_query_planner__ PyCapsule.

The method returns a PyCapsule wrapping an FFI_QueryPlanner, typically produced by a separate compiled extension. session is a handle on the session the planner is being installed on; take the extension codecs from it rather than building your own, and duck-type it — see LogicalExtensionCodecExportable for session.

Unlike the two bundle hooks in this module, this protocol is a type hint only: it is not @runtime_checkable, so isinstance against it raises TypeError.

Examples

A SessionContext satisfies this protocol, which is what lets a foreign planner wrap the one a session already has:

>>> from datafusion import SessionContext
>>> ctx = SessionContext()
>>> type(ctx.__datafusion_query_planner__(ctx)).__name__
'PyCapsule'

The protocol itself is not runtime-checkable:

>>> from datafusion.extensions import QueryPlannerExportable
>>> try:
...     isinstance(ctx, QueryPlannerExportable)
... except TypeError as e:
...     print("runtime_checkable" in str(e))
True
__datafusion_query_planner__(session: Any) → object#
class datafusion.extensions.SessionComponentsExportable#

Bases: Protocol

Type hint for extension bundles installable via with_extensions.

Runtime-checkable, so isinstance answers whether an object implements the protocol. Only the presence of the method is checked, which is the same question with_extensions() asks before calling it.

Implementations are reusable configuration objects: they must create fresh components on every call using the context supplied by with_extensions(), and must not retain that context or cache the components they bound to it, since the next call may install onto a different session. They should also avoid mutating the context they are handed — a registration made during binding is not rolled back if a later extension fails. See Failure and rollback.

A bundle that also contributes a query planner implements SessionPlannerExportable alongside this protocol.

Parameters:

ctx – The session the components will run on. Take the task-context provider off it. Do not read its codec chains expecting to find this call’s codecs, including your own: this hook runs before anything is installed, so ctx still carries whatever chains the receiver had. SessionPlannerExportable is the hook that sees the completed chains — see Two phases, because codecs and planners compose differently.

Returns:

The codecs this bundle contributes.

Examples

>>> from datafusion import SessionExtensionComponents
>>> from datafusion.extensions import SessionComponentsExportable
>>> class MyLibraryExtension:
...     def __datafusion_session_components__(self, ctx):
...         return SessionExtensionComponents()
>>> isinstance(MyLibraryExtension(), SessionComponentsExportable)
True
>>> isinstance(object(), SessionComponentsExportable)
False
__datafusion_session_components__(ctx: datafusion.context.SessionContext) → SessionExtensionComponents#
class datafusion.extensions.SessionExtensionComponents#

Components an extension contributes to a session context.

Returned by SessionComponentsExportable.__datafusion_session_components__() and consumed by with_extensions(). Every component must be created against the context passed to that method; components bound to a different session hold a task-context provider for that other session and cannot be rebound.

Construction is keyword-only, so later releases can add component kinds without changing what an existing call means.

Query planners are not listed here. They install in a second phase so each can wrap the one before it — see SessionPlannerExportable.

Examples

A bundle that contributes no codecs is valid — a planner-only library returns this, or omits the hook entirely:

>>> from datafusion import SessionExtensionComponents
>>> components = SessionExtensionComponents()
>>> components.logical_extension_codecs
()

A bundle that contributes one kind of component names it, leaving the rest empty. Here the codec is a capsule wrapped in an object that declares the id its payloads will carry:

>>> from datafusion import SessionContext
>>> class NamedCodec:
...     __datafusion_codec_id__ = "my_library.v1"
...
...     def __init__(self, capsule):
...         self._capsule = capsule
...
...     def __datafusion_logical_extension_codec__(self, session=None):
...         return self._capsule
>>> ctx = SessionContext()
>>> components = SessionExtensionComponents(
...     logical_extension_codecs=(NamedCodec(
...         ctx.__datafusion_logical_extension_codec__()
...     ),)
... )
>>> components.logical_extension_codecs[0].__datafusion_codec_id__
'my_library.v1'
>>> components.physical_extension_codecs
()

A single codec is not an iterable of codecs, and forgetting the trailing comma is the easy way to write one by accident:

>>> SessionExtensionComponents(logical_extension_codecs=NamedCodec(ctx))
Traceback (most recent call last):
    ...
TypeError: logical_extension_codecs must be an iterable of codec objects...
__post_init__() → None#

Normalize each codec field to a tuple, rejecting what cannot become one.

logical_extension_codecs: tuple[datafusion.user_defined.LogicalExtensionCodecExportable, Ellipsis] = ()#

Logical codecs to add to the session’s codec chain, in declaration order.

Objects exposing __datafusion_logical_extension_codec__, never bare PyCapsule objects — a codec’s id is read off the object it is handed over as. Any iterable is accepted and stored as a tuple. See Codecs are objects, not capsules.

physical_extension_codecs: tuple[datafusion.user_defined.PhysicalExtensionCodecExportable, Ellipsis] = ()#

Physical codecs to add to the session’s codec chain, in declaration order.

As logical_extension_codecs, for __datafusion_physical_extension_codec__.

class datafusion.extensions.SessionPlannerExportable#

Bases: Protocol

Type hint for extension bundles that contribute a query planner.

A session holds exactly one query planner, so planners compose by nesting rather than by chaining: each wraps the one before it and delegates to it for the work it does not handle. with_extensions() runs this hook once per bundle that implements it, in argument order, handing each the planner built so far. Returning a planner that wraps fallback puts this bundle outside the previous one, so the last bundle listed ends up outermost and is consulted first.

The hook runs after every codec from every bundle is installed, and ctx is the context carrying those final chains. That ordering is the point: a planner captured here sees the complete codec set, so a nested planner is not left encoding through a chain that a later bundle has grown. See Two phases, because codecs and planners compose differently.

Return ``None`` to contribute no planner, leaving fallback in place. That is the no-op, and it is not the same as returning fallback: the capsule the first bundle receives wraps the session’s planner for export, so handing it back installs that planner as a foreign one and every later plan crosses an FFI boundary that was not there before. A bundle that decides at runtime it has nothing to contribute returns None.

Ignoring fallback and returning a planner that does not delegate to it is legal and means “replace” — but it discards every planner listed before this one, including any the session already had.

Parameters:
  • ctx – The session the planner will run on, carrying the final codec chains.

  • fallback – The planner built so far, as a PyCapsule. For the first bundle this is the session’s existing planner, which is the DataFusion default unless one was installed earlier.

Returns:

A planner wrapping fallback, or None to contribute none.

Examples

A real library returns its own planner wrapping fallback, e.g. my_library.Planner(fallback=fallback). The two degenerate cases are worth contrasting, because both plan queries successfully and only one of them is the no-op:

>>> from datafusion import SessionContext
>>> from datafusion.extensions import SessionPlannerExportable
>>> class Contributes:
...     def __datafusion_session_planner__(self, ctx, fallback):
...         return None  # the no-op: session keeps its own planner
>>> class Replaces:
...     def __datafusion_session_planner__(self, ctx, fallback):
...         return fallback  # installs it as a *foreign* planner
>>> for bundle in (Contributes(), Replaces()):
...     ctx = SessionContext().with_extensions(bundle)
...     ctx.sql("SELECT 1 AS n").collect()[0].column(0).to_pylist()
[1]
[1]
>>> isinstance(Contributes(), SessionPlannerExportable)
True
>>> isinstance(object(), SessionPlannerExportable)
False
__datafusion_session_planner__(ctx: datafusion.context.SessionContext, fallback: types.CapsuleType) → QueryPlannerExportable | types.CapsuleType | None#