State API
quixstreams.state.base.state
State
Primary interface for working with key-value state data from StreamingDataFrame
State.get
Get the value for key if key is present in the state, else default
Arguments:
key: keydefault: default value to return if the key is not found
Returns:
value or None if the key is not found and default is not provided
State.get_bytes
Get the value for key if key is present in the state, else default
Arguments:
key: keydefault: default value to return if the key is not found
Returns:
value as bytes or None if the key is not found and default is not provided
State.set
Set value for the key, optionally with a per-write expiry.
Arguments:
key: keyvalue: valuettl: optional event-time TTL. When set, the entry expiresttlafter the current record's event-time and is filtered from subsequent reads.None(default) writes a sentinel stamp meaning "never expires", overwriting any prior TTL on the same key.
State.set_bytes
Set bytes value for the key, optionally with a per-write expiry.
Arguments:
key: keyvalue: value as bytesttl: see :meth:set.
State.delete
Delete value for the key.
This function always returns None, even if value is not found.
Arguments:
key: key
State.exists
Check if the key exists in state.
Arguments:
key: key
Returns:
True if key exists, False otherwise
TransactionState
TransactionState.__init__
Simple key-value state to be provided into StreamingDataFrame functions
Arguments:
transaction: instance ofPartitionTransactionprefix: serialized key prefix shared across callstimestamp: optional event-time of the current record (ms). Used by TTL-aware partitions to stamp values onset()withrecord.timestamp + ttland to filter expired entries onget(). The framework injects this on every record via theStreamingDataFramestateful wrapper.
TransactionState.get
Get the value for key if key is present in the state, else default
Arguments:
key: keydefault: default value to return if the key is not found
Returns:
value or None if the key is not found and default is not provided
TransactionState.get_bytes
Get the bytes value for key if key is present in the state, else default
Arguments:
key: keydefault: default value to return if the key is not found
Returns:
value or None if the key is not found and default is not provided
TransactionState.set
Set value for the key, optionally with a per-write expiry.
Arguments:
key: keyvalue: valuettl: optional event-time TTL. See :class:State.set.
TransactionState.set_bytes
Set bytes value for the key, optionally with a per-write expiry.
Arguments:
key: keyvalue: value as bytesttl: optional event-time TTL. See :class:State.set.
TransactionState.delete
Delete value for the key.
This function always returns None, even if value is not found.
Arguments:
key: key
TransactionState.exists
Check if the key exists in state.
Arguments:
key: key
Returns:
True if key exists, False otherwise
quixstreams.state.rocksdb.options
RocksDBOptions
RocksDB database options.
Arguments:
dumps: function to dump data to JSONloads: function to load data from JSONopen_max_retries: number of times to retry opening the database if it's locked by another process. To disable retrying, pass 0open_retry_backoff: number of seconds to wait between each retry.on_corrupted_recreate: when True, the corrupted DB will be destroyed if theuse_changelog_topics=Trueis also set on the Application. If this option is True, butuse_changelog_topics=False, the DB won't be destroyed. Note: risk of data loss! Make sure that the changelog topics are up-to-date before disabling it in production. Default -True.block_cache_size: size of the RocksDB block cache, in bytes. This is an aggregate ceiling for the whole process: one cache of this size is shared by every store partition, following RocksDB's own guidance that a single block cache be shared across databases. It is not a per-partition budget, so a 32-partition application does not reserve 32 times this value. A cache is filled lazily, so a small deployment never realizes the full capacity; size it against total available memory rather than per store. Note that index and bloom-filter blocks are held outside this budget (cache_index_and_filter_blocksis not enabled), so real usage is this value plus roughlybloom_filter_bits_per_key / 8bytes per key. Default -1 GiB.max_evictions_per_flush: cap on TTL-driven evictions performed during a singleflush()for stores with TTL enabled. Larger values increase per-flush latency but let the sweep keep up with higher steady-state expiration rates. Only meaningful for TTL-enabled stores; ignored otherwise.
This is the sweep's throughput dial, and it is the one to reach for:
the drain rate is max_evictions_per_flush / commit_interval, but a
checkpoint's cost is mostly fixed (a producer flush barrier plus an
offset commit, milliseconds against a remote broker), so shrinking
commit_interval buys no drain speed and starves message processing.
Raise this instead.
Interaction with the producer queue. Each eviction produces one
changelog tombstone (see ttl_changelog_tombstones), so a sweep can
enqueue far more records than librdkafka's
queue.buffering.max.messages (default 100_000) would appear to
allow. In practice it does not, because Producer.produce() polls
after every produce, so against a live broker the queue is a rolling
window rather than an accumulator: enqueueing 150_000 tombstones was
measured to peak at a queue depth of ~6,800 (6.8% of the default),
with the checkpoint's producer flush completing in ~4ms. The drain rate
governs, not the queue depth.
Two situations still bound it: the peak depth scales with
partitions-per-process (roughly ~15 partitions sweeping concurrently at
that depth would approach the default cap, so raise
queue.buffering.max.messages for high partition counts), and a broker
that is not draining at all will raise BufferError once the queue
genuinely fills — though by then the application has larger problems.
queue.buffering.max.kbytes (~1 GB default) is a second cap that binds
first for multi-KB records.
Default - 150_000. Measured on one partition against a live broker: a
full 150_000-eviction sweep completed in 1.12s (index scan 0.31s,
produces 0.53s, flush 0.004s, commit 0.27s) — a ~268x margin under a 300s
max.poll.interval.ms. The original 10_000 sustained only ~300
evictions/s once checkpoints grew, which loses the race against expiry
and lets the store grow without bound.
- legacy_records_ttl: expiry for pre-existing records when enabling
TTL on a populated legacy store that already holds un-stamped
records. When None (the default), the migration still completes:
the pre-existing records are backfilled using the ttl the service
itself uses (the max ttl= in the triggering flush) and a WARNING
names the implicit value. When set to a strictly positive
timedelta, that value is used instead: the partition backfills
its pre-existing un-stamped records with a uniform expiry of
high_water + legacy_records_ttl (event-time high-water at the
enable moment) and flips into TTL mode in place — no state deletion.
New records keep getting their true event-time expiry. The backfill
runs exactly once; a redeploy / restart never re-runs it. Ignored for
windowed / timestamped stores (they opt out of the TTL stamp
machinery at the class level). Must be strictly positive if set;
<= 0 raises ValueError at construction.
Default - None.
- legacy_backfill_chunk_size: number of pre-existing records re-stamped
per write-batch during the one-time legacy backfill (see
legacy_records_ttl). The backfill iterates the populated default CF
in chunks of this size; each chunk is re-stamped, produced to the
changelog, flushed, and committed before the next chunk is read, so peak
transient memory is bounded to one chunk regardless of total store size.
Lower it on memory-constrained deployments. Only meaningful on the single
backfilling flush; ignored otherwise and on windowed / timestamped
stores. Must be strictly positive; <= 0 raises ValueError at
construction.
Default - 150_000. Raising it mainly reduces the number of confirming
flushes (one broker round-trip each), whose fixed cost otherwise dominates
a large backfill; the producer-queue note under
max_evictions_per_flush applies here too.
- ttl_changelog_tombstones: when True (the default), TTL-driven
evictions are also produced to the changelog as tombstones
(value=None) so log compaction physically reclaims expired keys in
step with the local store — cleanup.policy=compact alone then shrinks
the changelog as keys expire (no delete policy / retention tuning
needed to reclaim). When False, evictions are local-only (the
pre-change behavior): the changelog keeps each expired key's last record
until compacted by other means, and rebuilds rely on the read-time
expiry filter. Read-time consistency is identical either way. Only
meaningful for TTL-enabled stores; ignored for windowed / timestamped
stores and for no-ttl= workloads.
Default - True.
- ttl_rollback: operational lever that UNDOES an automatic adoption of
a v3.24.0-stamped store, putting it back in legacy mode.
When a store opened by an older release wrote TTL-stamped values without
recording that it had TTL enabled, this release recognises the stamps
and adopts the store into TTL mode automatically. The adoption is
reversible until a live state.set(..., ttl=...) write confirms it.
Turn this lever on and restart if that adoption was wrong for your
store: values that were already adopted on this volume are restored
byte-identical from their backup, and a store rebuilt from its changelog
is left in legacy mode instead of being adopted. A store whose TTL mode
was established by this release's own migration, or one whose adoption
has already been confirmed, is unaffected.
Set it either here or as the environment variable
QUIXSTREAMS_STATE_TTL_ROLLBACK=1; the lever is on when either is
set, and the partition logs at open which of the two turned it on.
Prefer this option in Quix Cloud, where a deployment environment
variable that is not declared in the app's app.yaml is silently
dropped on redeploy — a lever that can vanish between runs is not a
lever you can reason about afterwards. Intended as a temporary measure:
clear it once the store is in the mode you want.
Mutually exclusive with ttl_force_flip (setting both raises).
Default - False (off).
- ttl_force_flip: operational REPAIR lever, the inverse of
ttl_rollback: it forces a store into TTL mode and records the flip,
then lets recovery finish any leftover migration.
Use it when a store holds TTL-stamped values but has lost the on-disk
evidence the automatic open-time repair needs to recognise them — for
example after the state directory was rebuilt, or after a rollback that
removed the TTL flag while leaving the values themselves stamped. The
symptom is a crash loop in which every read of a stamped value fails in
the value deserializer, and the raised StateMigrationError names
this lever. Read the error before setting it: it is a manual assertion
that the values in the store REALLY are TTL stamps, and on a genuinely
legacy store it would make every read strip eight bytes of real data.
Set it either here or as the environment variable
QUIXSTREAMS_STATE_TTL_FORCE_FLIP=1, with the same precedence and
open-time logging as ttl_rollback. It is a no-op on a store that is
already in TTL mode, so it is safe to leave set for one restart and then
clear — which is how it is meant to be used.
Mutually exclusive with ttl_rollback (setting both raises).
Default - False (off).
Please see rocksdict.Options for a complete description of other options.
RocksDBOptions.to_options
Convert parameters to rocksdict.Options
Returns:
instance of rocksdict.Options