Skip to content

Commit bd43d97

Browse files
timsaucerclaude
andcommitted
refactor: commit the optimizer rules through _commit_extensions
The rules keep their separate resolve step -- importing the capsules must stay fallible-and-early -- but the install half was a second private pymethod whose only caller was the line after `_commit_extensions`. Make it a parameter instead: the commit stays a single call, the rules go on in the same one-rebuild block, and the boundary does not grow an installer per component kind. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
1 parent b22387f commit bd43d97

3 files changed

Lines changed: 32 additions & 40 deletions

File tree

‎crates/core/src/context.rs‎

Lines changed: 27 additions & 36 deletions
Original file line numberDiff line numberDiff line change
@@ -1752,6 +1752,7 @@ impl PySessionContext {
17521752
/// hook never sees this call's functions in the registry — the
17531753
/// registrations have no fall-through, and a name is free to shadow one
17541754
/// the session already had.
1755+
#[allow(clippy::too_many_arguments)]
17551756
pub fn _commit_extensions<'py>(
17561757
slf: &Bound<'py, Self>,
17571758
extensions: Vec<Bound<'py, PyAny>>,
@@ -1760,6 +1761,7 @@ impl PySessionContext {
17601761
udfs: Vec<PyScalarUDF>,
17611762
udafs: Vec<PyAggregateUDF>,
17621763
udwfs: Vec<PyWindowUDF>,
1764+
rules: PyRef<'_, PyPhysicalOptimizerRules>,
17631765
) -> PyDataFusionResult<()> {
17641766
let py = slf.py();
17651767
// Nest the planners, outermost last. `planner` stays `None` when no
@@ -1798,17 +1800,38 @@ impl PySessionContext {
17981800
for udwf in udwfs {
17991801
this.ctx.register_udwf(udwf.function);
18001802
}
1803+
// Rules accumulate rather than replace, so unlike a planner there is
1804+
// no composition order to get right and no collision to refuse. All
1805+
// of them go on in **one** `SessionState` rebuild.
1806+
// [`Self::add_physical_optimizer_rule`] rebuilds per call, which for
1807+
// a bundle contributing several would clone the whole state that many
1808+
// times. Nothing here can fail: the capsules were imported by
1809+
// [`Self::_resolve_extension_physical_optimizer_rules`].
1810+
if !rules.rules.is_empty() {
1811+
let state_ref = this.ctx.state_ref();
1812+
let mut guard = state_ref.write();
1813+
// The session id has to be carried over for the same reason
1814+
// `add_physical_optimizer_rule` carries it: the builder mints a
1815+
// fresh one, and losing it leaves `session_id()` disagreeing with
1816+
// every `TaskContext` the session has already handed out.
1817+
let mut builder = SessionStateBuilder::new_from_existing(guard.clone())
1818+
.with_session_id(guard.session_id().to_string());
1819+
for rule in rules.rules.iter().cloned() {
1820+
builder = builder.with_physical_optimizer_rule(rule);
1821+
}
1822+
*guard = builder.build();
1823+
}
18011824
Ok(())
18021825
}
18031826

18041827
/// Import the physical optimizer rules a `with_extensions` call declared.
18051828
///
18061829
/// The fallible half of installing them, run while the call can still fail
18071830
/// harmlessly. Every capsule is imported here so that
1808-
/// [`Self::_install_extension_physical_optimizer_rules`] has nothing left
1809-
/// that can raise — a rule that failed to import after the planner was
1810-
/// bound would leave the session half-installed, and there is no derived
1811-
/// context to roll back to.
1831+
/// [`Self::_commit_extensions`] has nothing left that can raise — a rule
1832+
/// that failed to import after the planner was bound would leave the
1833+
/// session half-installed, and there is no derived context to roll back
1834+
/// to.
18121835
///
18131836
/// **Writes nothing.**
18141837
pub fn _resolve_extension_physical_optimizer_rules(
@@ -1821,38 +1844,6 @@ impl PySessionContext {
18211844
.collect::<Result<Vec<_>, _>>()?;
18221845
Ok(PyPhysicalOptimizerRules { rules })
18231846
}
1824-
1825-
/// Commit the physical optimizer rules for a `with_extensions` call.
1826-
///
1827-
/// Rules accumulate rather than replace, so unlike a planner there is no
1828-
/// composition order to get right and no collision to refuse.
1829-
///
1830-
/// All of them go on in **one** `SessionState` rebuild.
1831-
/// [`Self::add_physical_optimizer_rule`] rebuilds per call, which for a
1832-
/// bundle contributing several would clone the whole state that many times
1833-
/// and, worse, leave the earlier rules installed if a later one failed.
1834-
/// Nothing here can fail: the capsules were imported by
1835-
/// [`Self::_resolve_extension_physical_optimizer_rules`].
1836-
pub fn _install_extension_physical_optimizer_rules(
1837-
&self,
1838-
resolved: PyRef<'_, PyPhysicalOptimizerRules>,
1839-
) {
1840-
if resolved.rules.is_empty() {
1841-
return;
1842-
}
1843-
let state_ref = self.ctx.state_ref();
1844-
let mut guard = state_ref.write();
1845-
// The session id has to be carried over for the same reason
1846-
// `add_physical_optimizer_rule` carries it: the builder mints a fresh
1847-
// one, and losing it leaves `session_id()` disagreeing with every
1848-
// `TaskContext` the session has already handed out.
1849-
let mut builder = SessionStateBuilder::new_from_existing(guard.clone())
1850-
.with_session_id(guard.session_id().to_string());
1851-
for rule in resolved.rules.iter().cloned() {
1852-
builder = builder.with_physical_optimizer_rule(rule);
1853-
}
1854-
*guard = builder.build();
1855-
}
18561847
}
18571848

18581849
/// Physical optimizer rules imported for a `with_extensions` call.

‎python/datafusion/context.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2267,8 +2267,8 @@ def with_extensions(
22672267
[function._udf for function in resolved["udfs"]],
22682268
[function._udaf for function in resolved["udafs"]],
22692269
[function._udwf for function in resolved["udwfs"]],
2270+
resolved_rules,
22702271
)
2271-
new.ctx._install_extension_physical_optimizer_rules(resolved_rules)
22722272
return new
22732273

22742274
def table_provider(self, name: str) -> Table:

‎python/tests/test_wrapper_coverage.py‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -38,10 +38,11 @@
3838
# bundles declared.
3939
"_install_extension_codecs",
4040
"_commit_extensions",
41-
# Physical optimizer rules, split so the capsule import happens while
42-
# the call can still fail without leaving the session half-installed.
41+
# Physical optimizer rules import their capsules in a separate step,
42+
# so the call can still fail without leaving the session
43+
# half-installed; the commit itself is a `_commit_extensions`
44+
# parameter.
4345
"_resolve_extension_physical_optimizer_rules",
44-
"_install_extension_physical_optimizer_rules",
4546
}
4647
)
4748

0 commit comments

Comments
 (0)