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#
Type hint for object that has __datafusion_physical_optimizer_rule__ PyCapsule. |
|
Type hint for object that has a __datafusion_query_planner__ PyCapsule. |
|
Type hint for extension bundles installable via |
|
Components an extension contributes to a session context. |
|
Type hint for extension bundles that contribute a query planner. |
Module Contents#
- class datafusion.extensions.PhysicalOptimizerRuleExportable#
Bases:
ProtocolType 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_runsindatafusion-ffi-example.>>> from datafusion_ffi_example import MyPhysicalOptimizerRule >>> ctx.add_physical_optimizer_rule(MyPhysicalOptimizerRule())
- __datafusion_physical_optimizer_rule__() object#
- class datafusion.extensions.QueryPlannerExportable#
Bases:
ProtocolType 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.sessionis 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 — seeLogicalExtensionCodecExportableforsession.Unlike the two bundle hooks in this module, this protocol is a type hint only: it is not
@runtime_checkable, soisinstanceagainst it raisesTypeError.Examples
A
SessionContextsatisfies 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:
ProtocolType hint for extension bundles installable via
with_extensions.Runtime-checkable, so
isinstanceanswers whether an object implements the protocol. Only the presence of the method is checked, which is the same questionwith_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
SessionPlannerExportablealongside 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
ctxstill carries whatever chains the receiver had.SessionPlannerExportableis 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 bywith_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 barePyCapsuleobjects — 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:
ProtocolType 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 wrapsfallbackputs 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
ctxis 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
fallbackin place. That is the no-op, and it is not the same as returningfallback: 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 returnsNone.Ignoring
fallbackand 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, orNoneto 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#