diff --git a/.github/workflows/pr_lint_and_test.yaml b/.github/workflows/pr_lint_and_test.yaml index 2f18eb0..540217c 100644 --- a/.github/workflows/pr_lint_and_test.yaml +++ b/.github/workflows/pr_lint_and_test.yaml @@ -39,7 +39,7 @@ jobs: strategy: fail-fast: false matrix: - python-version: ["3.10", "3.11", "3.12", "3.13"] + python-version: ["3.10", "3.11", "3.12", "3.13", "3.14"] services: kafka: image: confluentinc/cp-kafka:8.0.0 diff --git a/docs/docs/sqlbroker/design.md b/docs/docs/sqlbroker/design.md index 1d92786..3952b2a 100644 --- a/docs/docs/sqlbroker/design.md +++ b/docs/docs/sqlbroker/design.md @@ -37,7 +37,7 @@ flowchart TD On start, the subscriber spawns four types of concurrent loops: -**1. Fetch loop** — Periodically fetches batches of `PENDING` or `RETRYABLE` messages from the database, simultaneously updating them: marking as `PROCESSING`, setting `acquired_at` to now, and incrementing `deliveries_count`. Only messages with `next_attempt_at <= now` are fetched, ordered by `next_attempt_at`. The fetched messages are placed into an internal queue. The fetch limit is the minimum of `fetch_batch_size` and the free buffer capacity (`fetch_batch_size * overfetch_factor` minus currently queued messages). If the last fetch was "full" (returned as many messages as the limit), the next fetch happens after `min_fetch_interval`; otherwise after `max_fetch_interval`. +**1. Fetch loop** — Periodically fetches batches of `PENDING` or `RETRYABLE` messages from the database, simultaneously updating them: marking as `PROCESSING`, setting `acquired_at` to now, and incrementing `deliveries_count`. Only messages with `next_attempt_at <= now` are fetched, ordered by `next_attempt_at`. The fetched messages are placed into an internal queue. The fetch limit is the minimum of `fetch_batch_size`, the free acquired-but-not-yet-processed capacity (`fetch_batch_size * max_not_processed_factor` minus currently unprocessed messages), and the free acquired-but-not-yet-persisted capacity (`fetch_batch_size * max_not_persisted_factor` minus currently unpersisted messages). If the last fetch was "full" (returned as many messages as the limit), the next fetch happens after `min_fetch_interval`; otherwise after `max_fetch_interval`. **2. Worker loops** (`max_workers` concurrent instances) — Each worker takes a message from the internal queue and first checks if `max_deliveries` has been exceeded; if so, the message is [Rejected](../sqlbroker/tutorial.md#reject){.internal-link} without processing. Otherwise, processing proceeds. Depending on the processing result, [`AckPolicy`](../getting-started/acknowledgement.md){.internal-link}, and manual [Ack](../sqlbroker/tutorial.md#ack){.internal-link}/[Nack](../sqlbroker/tutorial.md#nack){.internal-link}/[Reject](../sqlbroker/tutorial.md#reject){.internal-link}, the message is [Acked](../sqlbroker/tutorial.md#ack){.internal-link}, [Nacked](../sqlbroker/tutorial.md#nack){.internal-link}, or [Rejected](../sqlbroker/tutorial.md#reject){.internal-link}. For [Nacked](../sqlbroker/tutorial.md#nack){.internal-link} messages, the `retry_strategy` is consulted to determine if and when the message might be retried. If allowed to be retried, the message is marked as `RETRYABLE`; otherwise as `FAILED`. [Acked](../sqlbroker/tutorial.md#ack){.internal-link} messages are marked as `COMPLETED` and rejected messages are marked as `FAILED`. The message is then buffered for flushing. diff --git a/docs/docs/sqlbroker/tutorial.md b/docs/docs/sqlbroker/tutorial.md index d82cb03..04c7ec7 100644 --- a/docs/docs/sqlbroker/tutorial.md +++ b/docs/docs/sqlbroker/tutorial.md @@ -99,9 +99,10 @@ When `connection` is provided, the message insert participates in the same datab - **`queues`** — List of queue names to consume from. - **`max_workers`** (default: `1`) — Number of concurrent handler coroutines. - **`retry_strategy`** (default: `NoRetryStrategy()`) — Called to determine if and how soon a [Nacked](#nack){.internal-link} message is retried. -- **`fetch_batch_size`** — Maximum number of messages to fetch in a single batch. A fetch's actual limit might be lower if the free capacity of the acquired-but-not-yet-processed messages set is smaller. -- **`overfetch_factor`** (default: `1.5`) — Multiplier for `fetch_batch_size` to cap the size of the set of acquired-but-not-yet-processed messages. -- **`min_fetch_interval`** — Minimum interval between consecutive fetches. If the last fetch was full (returned as many messages as the fetch's limit), the next fetch happens after both (i) minimum fetch interval has passed, and (ii) capacity equal to the fetch batch size has freed up in the set of acquired-but-not-yet-processed messages. +- **`fetch_batch_size`** — Maximum number of messages to fetch in a single batch. A fetch's actual limit might be lower if either the acquired-but-not-yet-processed or acquired-but-not-yet-persisted set has less free capacity. +- **`max_not_processed_factor`** (default: `1.5`) — Multiplier for `fetch_batch_size` to cap the size of the set of acquired-but-not-yet-processed messages. +- **`max_not_persisted_factor`** (default: `2.0`) — Multiplier for `fetch_batch_size` to cap the size of the set of acquired messages whose state has not yet been persisted to the database. +- **`min_fetch_interval`** — Minimum interval between consecutive fetches. If the last fetch was full (returned as many messages as the fetch's limit), the next fetch happens after both (i) minimum fetch interval has passed, and (ii) capacity equal to the fetch batch size has freed up in both the acquired-but-not-yet-processed and acquired-but-not-yet-persisted sets. - **`max_fetch_interval`** — Maximum interval between consecutive fetches. - **`flush_interval`** — Interval between flushes of processed message state to the database. - **`release_stuck_interval`** (default: `60`) — Interval between checks for stuck [`PROCESSING`](#message-lifecycle){.internal-link} messages. diff --git a/faststream_sqlbroker/sqlbroker/broker/registrator.py b/faststream_sqlbroker/sqlbroker/broker/registrator.py index 28155c1..bd628e2 100644 --- a/faststream_sqlbroker/sqlbroker/broker/registrator.py +++ b/faststream_sqlbroker/sqlbroker/broker/registrator.py @@ -32,7 +32,8 @@ def subscriber( # type: ignore[override] max_fetch_interval: float, min_fetch_interval: float, fetch_batch_size: int, - overfetch_factor: float = 1.5, + max_not_processed_factor: float = 1.5, + max_not_persisted_factor: float = 2.0, flush_interval: float, release_stuck_interval: float = 60, release_stuck_timeout: float = 60 * 10, @@ -70,9 +71,12 @@ def subscriber( # type: ignore[override] Maximum number of messages to fetch in a single batch. A fetch's actual limit might be lower if the free capacity of the acquired-but-not-yet-processed messages set is smaller. - overfetch_factor: + max_not_processed_factor: Multiplier for `fetch_batch_size` to size the maximum size of the set of acquired-but-not-yet-processed messages. Defaults to `1.5`. + max_not_persisted_factor: + Multiplier for `fetch_batch_size` to cap the number of acquired + messages whose state has not yet been persisted. Defaults to `2.0`. flush_interval: Interval between flushes of processed message state to the database. release_stuck_interval: @@ -98,7 +102,8 @@ def subscriber( # type: ignore[override] max_fetch_interval=max_fetch_interval, min_fetch_interval=min_fetch_interval, fetch_batch_size=fetch_batch_size, - overfetch_factor=overfetch_factor, + max_not_processed_factor=max_not_processed_factor, + max_not_persisted_factor=max_not_persisted_factor, flush_interval=flush_interval, release_stuck_interval=release_stuck_interval, release_stuck_timeout=release_stuck_timeout, diff --git a/faststream_sqlbroker/sqlbroker/broker/router.py b/faststream_sqlbroker/sqlbroker/broker/router.py index 15428af..c8d715a 100644 --- a/faststream_sqlbroker/sqlbroker/broker/router.py +++ b/faststream_sqlbroker/sqlbroker/broker/router.py @@ -84,7 +84,8 @@ def __init__( max_fetch_interval: float, min_fetch_interval: float, fetch_batch_size: int, - overfetch_factor: float = 1.5, + max_not_processed_factor: float = 1.5, + max_not_persisted_factor: float = 1.5, flush_interval: float, release_stuck_interval: float = 60, release_stuck_timeout: float = 60 * 10, @@ -117,9 +118,13 @@ def __init__( The minimum allowed interval between consecutive fetches. fetch_batch_size: The maximum allowed number of messages to fetch in a single batch. - overfetch_factor: + max_not_processed_factor: The factor by which the fetch_batch_size is multiplied. Defaults to `1.5`. + max_not_persisted_factor: + The factor by which the fetch_batch_size is multiplied to cap + acquired messages whose state has not yet been persisted. + Defaults to `2.0`. flush_interval: The interval at which the state of messages is flushed to the database. release_stuck_interval: @@ -149,7 +154,8 @@ def __init__( max_fetch_interval=max_fetch_interval, min_fetch_interval=min_fetch_interval, fetch_batch_size=fetch_batch_size, - overfetch_factor=overfetch_factor, + max_not_processed_factor=max_not_processed_factor, + max_not_persisted_factor=max_not_persisted_factor, flush_interval=flush_interval, release_stuck_interval=release_stuck_interval, release_stuck_timeout=release_stuck_timeout, diff --git a/faststream_sqlbroker/sqlbroker/configs/subscriber.py b/faststream_sqlbroker/sqlbroker/configs/subscriber.py index 87c45b7..99e28c4 100644 --- a/faststream_sqlbroker/sqlbroker/configs/subscriber.py +++ b/faststream_sqlbroker/sqlbroker/configs/subscriber.py @@ -18,7 +18,8 @@ class SqlBrokerSubscriberConfig(SubscriberUsecaseConfig): max_fetch_interval: float min_fetch_interval: float fetch_batch_size: int - overfetch_factor: float + max_not_processed_factor: float + max_not_persisted_factor: float flush_interval: float release_stuck_interval: float release_stuck_timeout: float diff --git a/faststream_sqlbroker/sqlbroker/subscriber/factory.py b/faststream_sqlbroker/sqlbroker/subscriber/factory.py index 9218548..ef7e670 100644 --- a/faststream_sqlbroker/sqlbroker/subscriber/factory.py +++ b/faststream_sqlbroker/sqlbroker/subscriber/factory.py @@ -24,7 +24,8 @@ def create_subscriber( max_fetch_interval: float, min_fetch_interval: float, fetch_batch_size: int, - overfetch_factor: float, + max_not_processed_factor: float, + max_not_persisted_factor: float, flush_interval: float, release_stuck_interval: float, release_stuck_timeout: float, @@ -41,7 +42,7 @@ def create_subscriber( max_fetch_interval=max_fetch_interval, min_fetch_interval=min_fetch_interval, fetch_batch_size=fetch_batch_size, - overfetch_factor=overfetch_factor, + max_not_processed_factor=max_not_processed_factor, flush_interval=flush_interval, release_stuck_interval=release_stuck_interval, release_stuck_timeout=release_stuck_timeout, @@ -59,7 +60,8 @@ def create_subscriber( max_fetch_interval=max_fetch_interval, min_fetch_interval=min_fetch_interval, fetch_batch_size=fetch_batch_size, - overfetch_factor=overfetch_factor, + max_not_processed_factor=max_not_processed_factor, + max_not_persisted_factor=max_not_persisted_factor, flush_interval=flush_interval, release_stuck_interval=release_stuck_interval, release_stuck_timeout=release_stuck_timeout, @@ -84,7 +86,7 @@ def _validate_input_for_misconfiguration( max_fetch_interval: float, min_fetch_interval: float, fetch_batch_size: int, - overfetch_factor: float, + max_not_processed_factor: float, flush_interval: float, release_stuck_interval: float, release_stuck_timeout: float, diff --git a/faststream_sqlbroker/sqlbroker/subscriber/usecase.py b/faststream_sqlbroker/sqlbroker/subscriber/usecase.py index 1f0bf2c..eb8ae2f 100644 --- a/faststream_sqlbroker/sqlbroker/subscriber/usecase.py +++ b/faststream_sqlbroker/sqlbroker/subscriber/usecase.py @@ -66,8 +66,14 @@ def __init__( self._flush_interval = config.flush_interval self._fetch_batch_size = config.fetch_batch_size - self._max_not_processed = int(config.fetch_batch_size * config.overfetch_factor) + self._max_not_processed = int( + config.fetch_batch_size * config.max_not_processed_factor + ) + self._max_not_persisted = int( + config.fetch_batch_size * config.max_not_persisted_factor + ) self._not_processed_count = 0 + self._not_persisted_count = 0 self.graceful_timeout = self._outer_config.graceful_timeout self._release_stuck_timeout = config.release_stuck_timeout self._max_deliveries = config.max_deliveries @@ -157,8 +163,15 @@ async def _finalize_workers(self) -> None: task_flush = self.add_task(self._flush_results) await task_flush - def _check_if_may_fetch(self) -> None: - if self._max_not_processed - self._not_processed_count >= self._fetch_batch_size: + @property + def _free_slots(self) -> int: + return min( + self._max_not_processed - self._not_processed_count, + self._max_not_persisted - self._not_persisted_count, + ) + + def _check_if_may_fetch_eagerly(self) -> None: + if self._free_slots >= self._fetch_batch_size: self._may_fetch_event.set() async def _fetch_loop(self) -> None: @@ -167,9 +180,8 @@ async def _fetch_loop(self) -> None: break self._may_fetch_event.clear() - free_slots = self._max_not_processed - self._not_processed_count - if free_slots > 0: - limit = min(self._fetch_batch_size, free_slots) + if self._free_slots > 0: + limit = min(self._fetch_batch_size, self._free_slots) try: batch = await self._client.fetch(self._queues, limit=limit) @@ -180,15 +192,15 @@ async def _fetch_loop(self) -> None: for msg in batch: self._not_processed_count += 1 + self._not_persisted_count += 1 await self._pending_consume_queue.put(msg) - - self._check_if_may_fetch() - self._last_fetch_was_full = len(batch) == limit if not self._last_fetch_was_full: await self._sleep_until_stop_event(self._max_fetch_interval) continue + self._check_if_may_fetch_eagerly() + async with self._task_context( asyncio.sleep, func_args=(self._min_fetch_interval,) ) as min_fetch_interval_reached_task: @@ -224,7 +236,7 @@ async def _worker_loop(self) -> None: ) self._not_processed_count -= 1 - self._check_if_may_fetch() + self._check_if_may_fetch_eagerly() self._buffer_results(message) self._pending_consume_queue.task_done() @@ -298,12 +310,16 @@ async def _flush_results(self) -> None: try: await self._client.retry(to_update_in_primary) + self._not_persisted_count -= len(to_update_in_primary) + self._check_if_may_fetch_eagerly() except Exception: self._buffer_results(messages) raise try: await self._client.archive(to_persist_in_archive, to_delete_from_primary) + self._not_persisted_count -= len(to_delete_from_primary) + self._check_if_may_fetch_eagerly() except Exception: self._buffer_results(to_delete_from_primary) raise diff --git a/pyproject.toml b/pyproject.toml index d0d0433..266b33b 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -10,7 +10,7 @@ authors = [ { name = "Arseniy Popov", email = "arseniypopov@gmail.com" }, ] requires-python = ">=3.10" -version = "0.1.0a7" +version = "0.1.0a8" dependencies = [ "faststream>=0.7.0rc1", diff --git a/tests/basic.py b/tests/basic.py index 70db56e..ed054e8 100644 --- a/tests/basic.py +++ b/tests/basic.py @@ -53,7 +53,8 @@ def get_subscriber_params( kwargs.setdefault("max_fetch_interval", 0.1) kwargs.setdefault("min_fetch_interval", 0.01) kwargs.setdefault("fetch_batch_size", 5) - kwargs.setdefault("overfetch_factor", 1.5) + kwargs.setdefault("max_not_processed_factor", 1.5) + kwargs.setdefault("max_not_persisted_factor", 2.0) kwargs.setdefault("flush_interval", 0.01) kwargs.setdefault("release_stuck_interval", 60) kwargs.setdefault("release_stuck_timeout", 60 * 10) diff --git a/tests/test_ack_policy.py b/tests/test_ack_policy.py index 0e39993..a99deb5 100644 --- a/tests/test_ack_policy.py +++ b/tests/test_ack_policy.py @@ -42,7 +42,7 @@ async def test_consume_nack_on_error( max_fetch_interval=10, min_fetch_interval=10, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.1, release_stuck_interval=10, release_stuck_timeout=10, @@ -100,7 +100,7 @@ async def test_consume_reject_on_error( max_fetch_interval=10, min_fetch_interval=10, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.1, release_stuck_interval=10, release_stuck_timeout=10, @@ -160,7 +160,7 @@ async def test_consume_ack_and_ack_first( max_fetch_interval=10, min_fetch_interval=10, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.1, release_stuck_interval=10, release_stuck_timeout=10, @@ -220,7 +220,7 @@ async def test_consume_manual( max_fetch_interval=10, min_fetch_interval=10, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.1, release_stuck_interval=10, release_stuck_timeout=10, @@ -281,7 +281,7 @@ async def test_consume_manual_overrides_policy( max_fetch_interval=10, min_fetch_interval=10, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.1, release_stuck_interval=10, release_stuck_timeout=10, @@ -342,7 +342,7 @@ async def test_consume_manual_no_manual_ack( max_fetch_interval=10, min_fetch_interval=10, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.1, release_stuck_interval=10, release_stuck_timeout=10, @@ -387,7 +387,7 @@ async def test_consume_manual_nack_respects_policy( max_fetch_interval=0.01, min_fetch_interval=0, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.1, release_stuck_interval=10, release_stuck_timeout=10, diff --git a/tests/test_consume.py b/tests/test_consume.py index 2ab262d..a9a05fa 100644 --- a/tests/test_consume.py +++ b/tests/test_consume.py @@ -58,7 +58,7 @@ async def test_consume( max_fetch_interval=10, min_fetch_interval=10, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.01, release_stuck_interval=10, release_stuck_timeout=10, @@ -130,7 +130,7 @@ async def test_consume_nack_retry( max_fetch_interval=10, min_fetch_interval=10, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.1, release_stuck_interval=10, release_stuck_timeout=10, @@ -192,7 +192,7 @@ async def test_consume_nack_no_retry( max_fetch_interval=10, min_fetch_interval=10, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.1, release_stuck_interval=10, release_stuck_timeout=10, @@ -272,7 +272,7 @@ async def test_consume_max_deliveries( max_fetch_interval=0.1, min_fetch_interval=0.1, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.1, release_stuck_interval=10, release_stuck_timeout=1, @@ -360,7 +360,7 @@ async def test_consume_full_retry_flow( max_fetch_interval=0.01, min_fetch_interval=0.01, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.1, release_stuck_interval=10, release_stuck_timeout=10, @@ -428,7 +428,7 @@ async def test_consume_no_retry_strategy( max_fetch_interval=10, min_fetch_interval=10, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.1, release_stuck_interval=10, release_stuck_timeout=10, @@ -487,7 +487,7 @@ async def test_consume_by_queues( max_fetch_interval=10, min_fetch_interval=10, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.1, release_stuck_interval=10, release_stuck_timeout=10, @@ -526,7 +526,7 @@ async def test_consume_by_next_attempt_at( max_fetch_interval=0.01, min_fetch_interval=0.01, fetch_batch_size=1, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.01, release_stuck_interval=10, release_stuck_timeout=10, @@ -581,7 +581,7 @@ async def test_consume_current_messages_are_flushed_on_stop( max_fetch_interval=10, min_fetch_interval=10, fetch_batch_size=4, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.01, release_stuck_interval=10, release_stuck_timeout=10, @@ -668,7 +668,7 @@ async def test_consume_manual_ack_takes_precedence( max_fetch_interval=10, min_fetch_interval=10, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.01, release_stuck_interval=10, release_stuck_timeout=10, @@ -711,7 +711,7 @@ async def test_consume_manual_nack_takes_precedence( max_fetch_interval=10, min_fetch_interval=10, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.01, release_stuck_interval=10, release_stuck_timeout=10, @@ -779,7 +779,7 @@ async def test_consume_manual_reject_takes_precedence( max_fetch_interval=10, min_fetch_interval=10, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.01, release_stuck_interval=10, release_stuck_timeout=10, @@ -823,7 +823,7 @@ async def test_consume_context_fields( max_fetch_interval=10, min_fetch_interval=10, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.01, release_stuck_interval=10, release_stuck_timeout=10, @@ -878,7 +878,7 @@ async def test_consume_concurrency( max_fetch_interval=10, min_fetch_interval=0, fetch_batch_size=4, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.1, release_stuck_interval=10, release_stuck_timeout=10, @@ -930,7 +930,7 @@ async def test_consume_fetch_intervals_fetch_on_freed_capacity( max_fetch_interval=10, min_fetch_interval=0, fetch_batch_size=4, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.1, release_stuck_interval=10, release_stuck_timeout=10, @@ -940,13 +940,16 @@ async def test_consume_fetch_intervals_fetch_on_freed_capacity( async def handler(msg: Any) -> None: nonlocal attempted attempted.append(msg) + await asyncio.sleep(1) for _ in range(7): await broker.publish({"message": "hello"}, queue="default1") await broker.start() await asyncio.sleep(0.5) + assert len(attempted) == 4 + await asyncio.sleep(1) assert len(attempted) == 7 @pytest.mark.asyncio() @@ -956,6 +959,7 @@ async def test_consume_fetch_intervals_immediate_fetch_to_fill_capacity( """After first fetch, next fetch happened immediately to fill up capacity, because of the overfetch factor. """ + attempted = [] client = broker.config.broker_config.client client.fetch = AsyncMock(wraps=client.fetch) @@ -966,7 +970,7 @@ async def test_consume_fetch_intervals_immediate_fetch_to_fill_capacity( max_fetch_interval=10, min_fetch_interval=0, fetch_batch_size=4, - overfetch_factor=2, + max_not_processed_factor=2, flush_interval=0.1, release_stuck_interval=10, release_stuck_timeout=10, @@ -974,6 +978,8 @@ async def test_consume_fetch_intervals_immediate_fetch_to_fill_capacity( ack_policy=AckPolicy.NACK_ON_ERROR, ) async def handler(msg: Any) -> None: + nonlocal attempted + attempted.append(msg) await asyncio.sleep(4) for _ in range(7): @@ -982,6 +988,7 @@ async def handler(msg: Any) -> None: await asyncio.sleep(0.5) + assert len(attempted) == 4 assert client.fetch.await_count == 2 @pytest.mark.asyncio() @@ -1000,7 +1007,7 @@ async def test_consume_fetch_intervals_nonfull_fetch( max_fetch_interval=10, min_fetch_interval=0, fetch_batch_size=4, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=0.1, release_stuck_interval=10, release_stuck_timeout=10, @@ -1022,6 +1029,86 @@ async def handler(msg: Any) -> None: assert len(attempted) == 3 + @pytest.mark.asyncio() + async def test_consume_max_not_persisted_factor( + self, engine: AsyncEngine, recreate_tables: None, broker: SqlBroker + ) -> None: + """`max_not_persisted_factor` was saturated after 4 fetches.""" + client = broker.config.broker_config.client + client.fetch = AsyncMock(wraps=client.fetch) + fetched = [] + + @broker.subscriber( + queues=["default1"], + max_workers=4, + retry_strategy=NoRetryStrategy(), + max_fetch_interval=10, + min_fetch_interval=0, + fetch_batch_size=4, + max_not_processed_factor=2, + max_not_persisted_factor=4, + flush_interval=5, + release_stuck_interval=10, + release_stuck_timeout=10, + max_deliveries=20, + ack_policy=AckPolicy.NACK_ON_ERROR, + ) + async def handler(msg: Any) -> None: + nonlocal fetched + await asyncio.sleep(0) + fetched.append(msg) + + for _ in range(20): + await broker.publish({"message": "hello"}, queue="default1") + await broker.start() + + await asyncio.sleep(0.5) + + assert len(fetched) == 16 + assert client.fetch.await_count == 4 + + @pytest.mark.asyncio() + async def test_consume_max_not_persisted_factor_freed_capacity( + self, engine: AsyncEngine, recreate_tables: None, broker: SqlBroker + ) -> None: + """Top up freed `max_not_persisted_factor` capacity.""" + client = broker.config.broker_config.client + client.fetch = AsyncMock(wraps=client.fetch) + fetched = [] + + @broker.subscriber( + queues=["default1"], + max_workers=4, + retry_strategy=NoRetryStrategy(), + max_fetch_interval=10, + min_fetch_interval=0, + fetch_batch_size=4, + max_not_processed_factor=2, + max_not_persisted_factor=4, + flush_interval=2, + release_stuck_interval=10, + release_stuck_timeout=10, + max_deliveries=20, + ack_policy=AckPolicy.NACK_ON_ERROR, + ) + async def handler(msg: Any) -> None: + nonlocal fetched + await asyncio.sleep(0) + fetched.append(msg) + + for _ in range(20): + await broker.publish({"message": "hello"}, queue="default1") + await broker.start() + + await asyncio.sleep(0.5) + + assert len(fetched) == 16 + assert client.fetch.await_count == 4 + + await asyncio.sleep(2) + + assert len(fetched) == 20 + @pytest.mark.asyncio() async def test_consume_release_stuck( self, engine: AsyncEngine, recreate_tables: None, event: asyncio.Event @@ -1041,7 +1128,7 @@ async def test_consume_release_stuck( max_fetch_interval=0, min_fetch_interval=0, fetch_batch_size=5, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=10, release_stuck_interval=10, release_stuck_timeout=0.5, @@ -1093,7 +1180,7 @@ async def test_consume_work_sharing( max_fetch_interval=0, min_fetch_interval=0, fetch_batch_size=10, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=1, release_stuck_interval=10, release_stuck_timeout=10, @@ -1107,7 +1194,7 @@ async def test_consume_work_sharing( max_fetch_interval=0, min_fetch_interval=0, fetch_batch_size=10, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=1, release_stuck_interval=10, release_stuck_timeout=10, @@ -1121,7 +1208,7 @@ async def test_consume_work_sharing( max_fetch_interval=0, min_fetch_interval=0, fetch_batch_size=10, - overfetch_factor=1, + max_not_processed_factor=1, flush_interval=1, release_stuck_interval=10, release_stuck_timeout=10, diff --git a/tests/test_misconfigure.py b/tests/test_misconfigure.py index 4c7cf1e..2563652 100644 --- a/tests/test_misconfigure.py +++ b/tests/test_misconfigure.py @@ -41,7 +41,7 @@ async def test_warn_on_max_deliveries(broker: SqlBroker) -> None: max_fetch_interval=1.0, min_fetch_interval=0.1, fetch_batch_size=10, - overfetch_factor=1.0, + max_not_processed_factor=1.0, flush_interval=0.1, release_stuck_interval=10.0, release_stuck_timeout=10.0, @@ -63,7 +63,8 @@ async def test_subscriber_defaults(broker: SqlBroker) -> None: assert isinstance(subscriber.config.retry_strategy, NoRetryStrategy) assert subscriber.config.max_workers == 1 assert subscriber.config.ack_policy is AckPolicy.REJECT_ON_ERROR - assert subscriber.config.overfetch_factor == 1.5 + assert subscriber.config.max_not_processed_factor == 1.5 + assert subscriber.config.max_not_persisted_factor == 2.0 assert subscriber.config.max_deliveries is None assert subscriber.config.release_stuck_interval == 60 assert subscriber.config.release_stuck_timeout == 60 * 10 @@ -87,7 +88,7 @@ async def test_warn_when_retry_strategy_ignored(broker: SqlBroker) -> None: max_fetch_interval=1.0, min_fetch_interval=0.1, fetch_batch_size=10, - overfetch_factor=1.0, + max_not_processed_factor=1.0, flush_interval=0.1, release_stuck_interval=10.0, release_stuck_timeout=10.0, @@ -108,7 +109,7 @@ async def test_warn_when_nack_without_retry_strategy(broker: SqlBroker) -> None: max_fetch_interval=1.0, min_fetch_interval=0.1, fetch_batch_size=10, - overfetch_factor=1.0, + max_not_processed_factor=1.0, flush_interval=0.1, release_stuck_interval=10.0, release_stuck_timeout=10.0, @@ -134,7 +135,7 @@ async def test_fail_when_archiving_without_archive_table( max_fetch_interval=1.0, min_fetch_interval=0.1, fetch_batch_size=10, - overfetch_factor=1.0, + max_not_processed_factor=1.0, flush_interval=0.1, release_stuck_interval=10.0, release_stuck_timeout=10.0, @@ -154,7 +155,7 @@ async def test_no_fail_when_archiving_disabled_without_archive_table( max_fetch_interval=1.0, min_fetch_interval=0.1, fetch_batch_size=10, - overfetch_factor=1.0, + max_not_processed_factor=1.0, flush_interval=0.1, release_stuck_interval=10.0, release_stuck_timeout=10.0, @@ -176,7 +177,7 @@ async def test_warn_when_ack_first_used(broker: SqlBroker) -> None: max_fetch_interval=1.0, min_fetch_interval=0.1, fetch_batch_size=10, - overfetch_factor=1.0, + max_not_processed_factor=1.0, flush_interval=0.1, release_stuck_interval=10.0, release_stuck_timeout=10.0, diff --git a/uv.lock b/uv.lock index 4b3fa5f..9d9c25e 100644 --- a/uv.lock +++ b/uv.lock @@ -717,7 +717,7 @@ wheels = [ [[package]] name = "faststream-sqlbroker" -version = "0.1.0a7" +version = "0.1.0a8" source = { editable = "." } dependencies = [ { name = "faststream" },