The optional Cython accelerators¶
Faust ships several hot code paths twice: a readable pure-Python implementation, and a Cython one used instead whenever the extension modules could be built. Nothing in Faust requires the extensions – every accelerated import falls back:
if not NO_CYTHON:
try:
from ._cython.streams import StreamIterator as _CStreamIterator
except ImportError:
_CStreamIterator = None
That fallback is what makes the accelerators optional, and it is also the single biggest hazard in maintaining them. This page is about the hazard.
The cython_optimizations opt-in¶
Two of the Cython fast paths never ran – each was guarded by a condition that could not become true (see Why parity tests exist). Repairing them activates code that has, by definition, never executed in production, so the repairs are behind a setting that defaults to off:
app = faust.App('myapp', cython_optimizations=True)
or CYTHON_OPTIMIZATIONS=1 in the environment (FAUST_CYTHON_OPTIMIZATIONS
when env_prefix is set). With it off, the extensions behave exactly
as the released versions do.
What it gates:
StreamIterator._try_get_quick_value– taking values already in the channel queue instead of always awaiting.ConductorHandlerevent reuse – decoding a message once and reusing the event across channels with matching key/value types, instead of deserializing once per subscribed channel.
What it does not gate: the on_topic_buffer_full argument fix. That one
was wrong in both implementations, is not a Cython-specific change, and
produces a metric that was simply incorrect before – so it applies
unconditionally.
One consequence to be aware of. While the setting is off, the Cython path and the pure-Python path genuinely differ. That is not new – it is what has shipped for years – and the flag does not introduce the divergence, it makes it selectable. The sharpest case is in the conductor: a reused event is never decoded a second time, so a channel whose payload would fail to deserialize raises no error when the event is reused, and raises one when it is not. That changes which channels receive a message, and how many acks it takes.
Consequently the parity suites run with the setting on – that is the configuration in which the two implementations are supposed to agree. A separate test in each suite pins the default-off behaviour, so the historical path stays covered too.
Retiring the setting¶
The setting is transitional. It exists to make adopting the repaired paths a decision rather than something that arrives in an upgrade, and it is meant to be removed, not kept. The intended sequence:
Now – ships defaulting to
False. Upgrading changes nothing.Default flipped to
Trueonce there is real-world evidence the repaired paths behave: the parity suites passing is necessary but not sufficient, since they only prove the two implementations agree under test. Record the flip asversion_changed={'<ver>': 'Enabled by default.'}.Deprecated – set
version_deprecatedanddeprecation_reasonon the setting. Users who set it explicitly get a warning; nobody else notices.Removed – delete the setting, both
bintattributes, the branches guarding the fast paths,faust.utils.optin, and the two default-off tests. At that point the fast paths are simply the behaviour, and the parity suites no longer need aconfmarker.
One thing to know before step 3.
__get__() emits a
UserWarning on every read of a deprecated setting, and faust reads
this one itself – once per Stream, once per assigned
partition. Deprecating it naively would make faust warn at itself, at a rate
that scales with the deployment, about a setting the user most likely never
set.
That is why the two extensions read it through
faust.utils.optin.cython_optimizations_enabled() rather than
app.conf.cython_optimizations: the helper takes the stored value and so
stays silent, while user-facing reads still warn, which is the point of
deprecating it. tests/unit/utils/test_optin.py pins both halves, so step 3
is genuinely a two-line change.
Testing the compiled code¶
The extensions have to be built in place, or the tests do not touch them.
pytest runs from the repository root, so import faust resolves
to the source tree – not to whatever pip install . compiled into
site-packages. With no .so next to the .pyx, every accelerated
import raises ImportError, the fallback engages, and the whole suite
tests pure Python. Silently: nothing warns, and the run is green either way.
$ USE_CYTHON=1 python setup.py build_ext --inplace
$ FAUST_REQUIRE_CYTHON=1 python -m pytest tests/unit tests/functional
FAUST_REQUIRE_CYTHON=1 asserts that the accelerators really were loaded,
turning the silent fallback into a failure. Set it whenever a run is supposed
to be testing the compiled code; the CI legs that build the extensions do.
Without it, a green run proves nothing about the Cython path, and any test that compares the two implementations degrades into comparing one implementation against itself.
Why parity tests exist¶
Two implementations of the same behaviour drift, and this pair has drifted repeatedly:
#608, “Fix cython stream_event_in to match python impl” – shipped, and fixed only after the fact.
Conductor’s full-queue path passed a channel toon_topic_buffer_fullwhere aTPwas expected, soMonitor.topic_buffer_full– aCounter[TP]– was keyed by channel from that path and byTPfrom the pressure-high path. The same partition accumulated under two keys, splitting its count and adding a second/statsentry for it.Both twins had it, so for a long time the comment in
faust/transport/conductor.pyrecorded the defect as deliberately left unfixed: correcting one alone would have made them disagree. The duplication turned a one-line bug into one nobody wanted to touch. It is fixed now – in both, together, which is what the parity suites make safe.Worth noting what did not catch it: the parity tests were green throughout, because both implementations were wrong in the same way. A differential test only finds divergence. Shared mistakes need an assertion about the behaviour itself, which is why the conductor suite now checks that the sensor is handed a
TPrather than only that both sides hand it the same thing.StreamIterator._try_get_quick_valuecarried two bugs that concealed each other.chan_queue_emptyholds the boundqueue.emptymethod:# streams.py # streams.pyx (before) if chan_queue_empty(): if self.chan_queue_empty:
A bound method is always truthy, so the extension always reported “queue empty” and took the awaiting path. That made the
elsebranch unreachable – which hid the fact that it returned the bare value fromget_nowait()instead of the(need_slow_get, value)pair the caller unpacks. Had the fast path ever run, it would have raisedTypeError, or silently mis-unpacked a two-element value.So the extension quietly did more work than the pure-Python code it was meant to accelerate, for as long as it has existed.
ConductorHandlerhad the same shape of fault, independently. The conductor deserializes a message once and reuses the event for every channel whose(key_type, value_type)pair matches. In the extension,event_keyidwas only ever assigned from_decode(), which returned it unchanged on the first pass – so it stayedNoneforever and the reuse branch was dead. Every subscribed channel re-deserialized the payload.That masked a second fault, again: had the keyid ever been set, a mismatched pair fell off the end of
_decodeand returned a bareNone, which unpacking into two names raisesTypeErroron. Fixing the reuse alone would have converted a silent inefficiency into a crash on any topic whose subscribers declare different key or value types.It was not only a performance difference. A channel whose event is reused never calls
decodeat all, so a channel that would have failed to deserialize raised no error under the pure-Python conductor and raised one under the extension – changing which channels got the message, and how many acks the message received.
None of these were caught by a test, because until recently no test ever imported the compiled modules.
The parity suites are tests/unit/test_cython_parity.py (windows, the
stream iterator’s queue fast path) and
tests/unit/transport/test_conductor_parity.py (the conductor’s
per-message fan-out, driven end to end through both implementations).
tests/unit/test_cython_parity.py covers both halves: it asserts the
accelerators are loaded when they are required, and compares the two
implementations where they can be driven directly.
Writing an accelerator¶
The conventions the existing modules follow:
Keep the pure-Python implementation. It is the reference, it is what PyPy and no-compiler installs use, and it is the other half of every parity test. Name it
_py_<name>or_Py<Name>and export both, so tests can reach the two independently.Mirror behaviour rather than approximating it. Anything the pure-Python version guarantees – iteration order, what happens when a mapping is mutated mid-pass, which exception comes out – is a guarantee of the accelerated one too.
Add parity tests in the same change, parametrised over both implementations. A differential test over randomised inputs is worth more than a handful of examples.
Measure first, and record it – time the accelerator against its twin in the same interpreter, and put the numbers in the commit message. An accelerator that does not clearly pay is a second implementation to keep in sync forever, in exchange for nothing. (A shared harness for this,
extra/tools/benchmark_cython.py, is proposed in #751.)
Not every hot path is worth compiling. The wins concentrate in code doing
real per-call arithmetic – the window types are ~4-5x faster compiled. Code
whose body is mostly await and calls back into Python gains much less,
because the time is in the awaiting, not the arithmetic.