Skip to content

RedisBroker

faststream.redis.broker.RedisBroker #

RedisBroker(
    url: str = "redis://localhost:6379",
    **kwargs: Unpack[RedisBrokerParams],
)

Bases: RedisRegistrator, BrokerUsecase[UnifyRedisDict, 'Redis[bytes]', RedisBrokerConfig]

Redis broker.

Initialized the RedisBroker.

Source code in faststream/redis/broker/broker.py
64
65
66
67
68
69
70
def __init__(
    self,
    url: str = "redis://localhost:6379",
    **kwargs: Unpack["RedisBrokerParams"],
) -> None:
    """Initialized the RedisBroker."""
    self._init_broker(url, dict(kwargs))

middlewares property #

middlewares: Sequence[BrokerMiddleware[MsgType]]

context property #

context: ContextRepo

config instance-attribute #

config: ConfigComposition[BrokerConfigType] = (
    ConfigComposition(config)
)

routers instance-attribute #

routers: list[Registrator[MsgType, Any]] = []

subscribers property #

subscribers: list[SubscriberUsecase[MsgType]]

publishers property #

publishers: list[PublisherUsecase]

parent property writable #

parent: Registrator[MsgType, Any] | None

specification instance-attribute #

specification = specification

running instance-attribute #

running = False

provider property #

provider: Provider

add_middleware #

add_middleware(
    middleware: BrokerMiddleware[Any, Any],
) -> None

Append BrokerMiddleware to the end of middlewares list.

Current middleware will be used as a most inner of the stack.

Source code in faststream/_internal/broker/registrator.py
52
53
54
55
56
57
def add_middleware(self, middleware: "BrokerMiddleware[Any, Any]") -> None:
    """Append BrokerMiddleware to the end of middlewares list.

    Current middleware will be used as a most inner of the stack.
    """
    self.config.add_middleware(middleware)

insert_middleware #

insert_middleware(
    middleware: BrokerMiddleware[Any, Any],
) -> None

Insert BrokerMiddleware to the start of middlewares list.

Current middleware will be used as a most outer of the stack.

Source code in faststream/_internal/broker/registrator.py
59
60
61
62
63
64
def insert_middleware(self, middleware: "BrokerMiddleware[Any, Any]") -> None:
    """Insert BrokerMiddleware to the start of middlewares list.

    Current middleware will be used as a most outer of the stack.
    """
    self.config.insert_middleware(middleware)

subscriber #

subscriber(
    channel: Union[PubSub, str] = ...,
    *,
    list: None = None,
    stream: None = None,
    dependencies: Sequence[Dependant] = (),
    parser: Optional[CustomCallable] = None,
    decoder: Optional[CustomCallable] = None,
    codec: Optional[CodecProto] = None,
    ack_policy: AckPolicy = EMPTY,
    no_reply: bool = False,
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    include_in_schema: bool = True,
    max_workers: None = None,
) -> ChannelSubscriber
subscriber(
    channel: Union[PubSub, str] = ...,
    *,
    list: None = None,
    stream: None = None,
    dependencies: Sequence[Dependant] = (),
    parser: Optional[CustomCallable] = None,
    decoder: Optional[CustomCallable] = None,
    codec: Optional[CodecProto] = None,
    ack_policy: AckPolicy = EMPTY,
    no_reply: bool = False,
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    include_in_schema: bool = True,
    max_workers: int = ...,
) -> ChannelConcurrentSubscriber
subscriber(
    channel: None = None,
    *,
    list: Union[str, ListSub[Literal[False]]] = ...,
    stream: None = None,
    dependencies: Sequence[Dependant] = (),
    parser: Optional[CustomCallable] = None,
    decoder: Optional[CustomCallable] = None,
    codec: Optional[CodecProto] = None,
    ack_policy: AckPolicy = EMPTY,
    no_reply: bool = False,
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    include_in_schema: bool = True,
    max_workers: None = None,
) -> ListSubscriber
subscriber(
    channel: None = None,
    *,
    list: ListSub[Literal[True]] = ...,
    stream: None = None,
    dependencies: Sequence[Dependant] = (),
    parser: Optional[CustomCallable] = None,
    decoder: Optional[CustomCallable] = None,
    codec: Optional[CodecProto] = None,
    ack_policy: AckPolicy = EMPTY,
    no_reply: bool = False,
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    include_in_schema: bool = True,
    max_workers: None = None,
) -> ListBatchSubscriber
subscriber(
    channel: None = None,
    *,
    list: Union[ListSub, str] = ...,
    stream: None = None,
    dependencies: Sequence[Dependant] = (),
    parser: Optional[CustomCallable] = None,
    decoder: Optional[CustomCallable] = None,
    codec: Optional[CodecProto] = None,
    ack_policy: AckPolicy = EMPTY,
    no_reply: bool = False,
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    include_in_schema: bool = True,
    max_workers: None = None,
) -> Union[ListSubscriber, ListBatchSubscriber]
subscriber(
    channel: None = None,
    *,
    list: Union[ListSub, str] = ...,
    stream: None = None,
    dependencies: Sequence[Dependant] = (),
    parser: Optional[CustomCallable] = None,
    decoder: Optional[CustomCallable] = None,
    codec: Optional[CodecProto] = None,
    ack_policy: AckPolicy = EMPTY,
    no_reply: bool = False,
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    include_in_schema: bool = True,
    max_workers: int = ...,
) -> ListConcurrentSubscriber
subscriber(
    channel: None = None,
    *,
    list: None = None,
    stream: Union[str, StreamSub[Literal[False]]] = ...,
    dependencies: Sequence[Dependant] = (),
    parser: Optional[CustomCallable] = None,
    decoder: Optional[CustomCallable] = None,
    codec: Optional[CodecProto] = None,
    ack_policy: AckPolicy = EMPTY,
    no_reply: bool = False,
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    include_in_schema: bool = True,
    max_workers: None = None,
) -> StreamSubscriber
subscriber(
    channel: None = None,
    *,
    list: None = None,
    stream: StreamSub[Literal[True]] = ...,
    dependencies: Sequence[Dependant] = (),
    parser: Optional[CustomCallable] = None,
    decoder: Optional[CustomCallable] = None,
    codec: Optional[CodecProto] = None,
    ack_policy: AckPolicy = EMPTY,
    no_reply: bool = False,
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    include_in_schema: bool = True,
    max_workers: None = None,
) -> StreamBatchSubscriber
subscriber(
    channel: None = None,
    *,
    list: None = None,
    stream: Union[StreamSub, str] = ...,
    dependencies: Sequence[Dependant] = (),
    parser: Optional[CustomCallable] = None,
    decoder: Optional[CustomCallable] = None,
    codec: Optional[CodecProto] = None,
    ack_policy: AckPolicy = EMPTY,
    no_reply: bool = False,
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    include_in_schema: bool = True,
    max_workers: None = None,
) -> Union[StreamSubscriber, StreamBatchSubscriber]
subscriber(
    channel: None = None,
    *,
    list: None = None,
    stream: Union[StreamSub, str] = ...,
    dependencies: Sequence[Dependant] = (),
    parser: Optional[CustomCallable] = None,
    decoder: Optional[CustomCallable] = None,
    codec: Optional[CodecProto] = None,
    ack_policy: AckPolicy = EMPTY,
    no_reply: bool = False,
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    include_in_schema: bool = True,
    max_workers: int = ...,
) -> StreamConcurrentSubscriber
subscriber(
    channel: Union[PubSub, str, None] = None,
    *,
    list: Union[ListSub, str, None] = None,
    stream: Union[StreamSub, str, None] = None,
    dependencies: Sequence[Dependant] = (),
    parser: Optional[CustomCallable] = None,
    decoder: Optional[CustomCallable] = None,
    codec: Optional[CodecProto] = None,
    ack_policy: AckPolicy = EMPTY,
    no_reply: bool = False,
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    include_in_schema: bool = True,
    max_workers: int | None = None,
) -> LogicSubscriber
subscriber(
    channel: Union[PubSub, str, None] = None,
    *,
    list: Union[ListSub, str, None] = None,
    stream: Union[StreamSub, str, None] = None,
    dependencies: Sequence[Dependant] = (),
    parser: Optional[CustomCallable] = None,
    decoder: Optional[CustomCallable] = None,
    codec: Optional[CodecProto] = None,
    ack_policy: AckPolicy = EMPTY,
    no_reply: bool = False,
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    include_in_schema: bool = True,
    max_workers: int | None = None,
) -> LogicSubscriber

Subscribe a handler to a RabbitMQ queue.

PARAMETER DESCRIPTION
channel

Redis PubSub object name to send message.

TYPE: Union[PubSub, str, None] DEFAULT: None

list

Redis List object name to send message.

TYPE: Union[ListSub, str, None] DEFAULT: None

stream

Redis Stream object name to send message.

TYPE: Union[StreamSub, str, None] DEFAULT: None

ack_policy

Acknowledgement policy for message processing.

TYPE: AckPolicy DEFAULT: EMPTY

dependencies

Dependencies list ([Depends(),]) to apply to the subscriber.

TYPE: Sequence[Dependant] DEFAULT: ()

parser

Parser to map original IncomingMessage Msg to FastStream one.

TYPE: Optional[CustomCallable] DEFAULT: None

decoder

Function to decode FastStream msg bytes body to python objects.

TYPE: Optional[CustomCallable] DEFAULT: None

codec

Custom codec object.

TYPE: Optional[CodecProto] DEFAULT: None

no_reply

Whether to disable FastStream RPC and Reply To auto responses or not.

TYPE: bool DEFAULT: False

message_format

Which format to use when parsing messages.

TYPE: type[MessageFormat] | None DEFAULT: None

persistent

Whether to make the subscriber persistent or not.

TYPE: bool DEFAULT: True

max_workers

Number of workers to process messages concurrently.

TYPE: int | None DEFAULT: None

title

AsyncAPI subscriber object title.

TYPE: str | None DEFAULT: None

description

AsyncAPI subscriber object description. Uses decorated docstring as default.

TYPE: str | None DEFAULT: None

include_in_schema

Whether to include operation in AsyncAPI schema or not.

TYPE: bool DEFAULT: True

RETURNS DESCRIPTION
SubscriberType

The subscriber object.

TYPE: LogicSubscriber

Source code in faststream/redis/broker/registrator.py
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
@override
def subscriber(
    self,
    channel: Union["PubSub", str, None] = None,
    *,
    list: Union["ListSub", str, None] = None,
    stream: Union["StreamSub", str, None] = None,
    # broker arguments
    dependencies: Sequence["Dependant"] = (),
    parser: Optional["CustomCallable"] = None,
    decoder: Optional["CustomCallable"] = None,
    codec: Optional["CodecProto"] = None,
    ack_policy: AckPolicy = EMPTY,
    no_reply: bool = False,
    message_format: type["MessageFormat"] | None = None,
    persistent: bool = True,
    # AsyncAPI information
    title: str | None = None,
    description: str | None = None,
    include_in_schema: bool = True,
    max_workers: int | None = None,
) -> "LogicSubscriber":
    """Subscribe a handler to a RabbitMQ queue.

    Args:
        channel: Redis PubSub object name to send message.
        list: Redis List object name to send message.
        stream: Redis Stream object name to send message.
        ack_policy: Acknowledgement policy for message processing.
        dependencies: Dependencies list (`[Depends(),]`) to apply to the subscriber.
        parser: Parser to map original **IncomingMessage** Msg to FastStream one.
        decoder: Function to decode FastStream msg bytes body to python objects.
        codec: Custom codec object.
        no_reply: Whether to disable **FastStream** RPC and Reply To auto responses or not.
        message_format: Which format to use when parsing messages.
        persistent: Whether to make the subscriber persistent or not.
        max_workers: Number of workers to process messages concurrently.
        title: AsyncAPI subscriber object title.
        description: AsyncAPI subscriber object description. Uses decorated docstring as default.
        include_in_schema: Whether to include operation in AsyncAPI schema or not.

    Returns:
        SubscriberType: The subscriber object.
    """
    subscriber = create_subscriber(
        channel=channel,
        list=list,
        stream=stream,
        # subscriber args
        max_workers=max_workers or 1,
        no_reply=no_reply,
        ack_policy=ack_policy,
        message_format=message_format,
        config=cast("RedisBrokerConfig", self.config),
        # AsyncAPI
        title_=title,
        description_=description,
        include_in_schema=include_in_schema,
    )

    super().subscriber(subscriber, persistent=persistent)

    return subscriber.add_call(
        parser_=parser or self._parser,
        decoder_=decoder or self._decoder,
        codec_=codec,
        dependencies_=dependencies,
    )

publisher #

publisher(
    channel: None = None,
    *,
    list: None = None,
    stream: Union[StreamSub, str] = ...,
    headers: dict[str, Any] | None = None,
    reply_to: str = "",
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    schema: Any | None = None,
    include_in_schema: bool = True,
) -> StreamPublisher
publisher(
    channel: None = None,
    *,
    list: Union[str, ListSub[Literal[False]]] = ...,
    stream: None = None,
    headers: dict[str, Any] | None = None,
    reply_to: str = "",
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    schema: Any | None = None,
    include_in_schema: bool = True,
) -> ListPublisher
publisher(
    channel: None = None,
    *,
    list: ListSub[Literal[True]] = ...,
    stream: None = None,
    headers: dict[str, Any] | None = None,
    reply_to: str = "",
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    schema: Any | None = None,
    include_in_schema: bool = True,
) -> ListBatchPublisher
publisher(
    channel: None = None,
    *,
    list: Union[ListSub, str] = ...,
    stream: None = None,
    headers: dict[str, Any] | None = None,
    reply_to: str = "",
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    schema: Any | None = None,
    include_in_schema: bool = True,
) -> Union[ListPublisher, ListBatchPublisher]
publisher(
    channel: Union[PubSub, str] = ...,
    *,
    list: None = None,
    stream: None = None,
    headers: dict[str, Any] | None = None,
    reply_to: str = "",
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    schema: Any | None = None,
    include_in_schema: bool = True,
) -> ChannelPublisher
publisher(
    channel: Union[PubSub, str, None] = None,
    *,
    list: Union[ListSub, str, None] = None,
    stream: Union[StreamSub, str, None] = None,
    headers: dict[str, Any] | None = None,
    reply_to: str = "",
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    schema: Any | None = None,
    include_in_schema: bool = True,
) -> LogicPublisher
publisher(
    channel: Union[PubSub, str, None] = None,
    *,
    list: Union[ListSub, str, None] = None,
    stream: Union[StreamSub, str, None] = None,
    headers: dict[str, Any] | None = None,
    reply_to: str = "",
    message_format: type[MessageFormat] | None = None,
    persistent: bool = True,
    title: str | None = None,
    description: str | None = None,
    schema: Any | None = None,
    include_in_schema: bool = True,
) -> LogicPublisher

Creates long-living and AsyncAPI-documented publisher object.

You can use it as a handler decorator (handler should be decorated by @broker.subscriber(...) too) - @broker.publisher(...). In such case publisher will publish your handler return value.

Or you can create a publisher object to call it lately - broker.publisher(...).publish(...).

PARAMETER DESCRIPTION
channel

Redis PubSub object name to send message.

TYPE: Union[PubSub, str, None] DEFAULT: None

list

Redis List object name to send message.

TYPE: Union[ListSub, str, None] DEFAULT: None

stream

Redis Stream object name to send message.

TYPE: Union[StreamSub, str, None] DEFAULT: None

headers

Message headers to store meta-information. Can be overridden by publish.headers if specified.

TYPE: dict[str, Any] | None DEFAULT: None

reply_to

Reply message destination PubSub object name.

TYPE: str DEFAULT: ''

message_format

Which format to use when parsing messages.

TYPE: type[MessageFormat] | None DEFAULT: None

title

AsyncAPI publisher object title.

TYPE: str | None DEFAULT: None

description

AsyncAPI publisher object description.

TYPE: str | None DEFAULT: None

schema

AsyncAPI publishing message type. Should be any python-native object annotation or pydantic.BaseModel.

TYPE: Any | None DEFAULT: None

include_in_schema

Whether to include operation in AsyncAPI schema or not.

TYPE: bool DEFAULT: True

persistent

Whether to make the publisher persistent or not.

TYPE: bool DEFAULT: True

Source code in faststream/redis/broker/registrator.py
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
@override
def publisher(
    self,
    channel: Union["PubSub", str, None] = None,
    *,
    list: Union["ListSub", str, None] = None,
    stream: Union["StreamSub", str, None] = None,
    headers: dict[str, Any] | None = None,
    reply_to: str = "",
    message_format: type["MessageFormat"] | None = None,
    persistent: bool = True,
    # AsyncAPI information
    title: str | None = None,
    description: str | None = None,
    schema: Any | None = None,
    include_in_schema: bool = True,
) -> "LogicPublisher":
    """Creates long-living and AsyncAPI-documented publisher object.

    You can use it as a handler decorator (handler should be decorated by `@broker.subscriber(...)` too) - `@broker.publisher(...)`.
    In such case publisher will publish your handler return value.

    Or you can create a publisher object to call it lately - `broker.publisher(...).publish(...)`.

    Args:
        channel: Redis PubSub object name to send message.
        list: Redis List object name to send message.
        stream: Redis Stream object name to send message.
        headers: Message headers to store meta-information. Can be overridden
            by `publish.headers` if specified.
        reply_to: Reply message destination PubSub object name.
        message_format: Which format to use when parsing messages.
        title: AsyncAPI publisher object title.
        description: AsyncAPI publisher object description.
        schema: AsyncAPI publishing message type. Should be any python-native
            object annotation or `pydantic.BaseModel`.
        include_in_schema: Whether to include operation in AsyncAPI schema or not.
        persistent: Whether to make the publisher persistent or not.
    """
    publisher = create_publisher(
        channel=channel,
        list=list,
        stream=stream,
        headers=headers,
        reply_to=reply_to,
        # Specific
        config=cast("RedisBrokerConfig", self.config),
        message_format=message_format,
        # AsyncAPI
        title_=title,
        description_=description,
        schema_=schema,
        include_in_schema=include_in_schema,
    )
    super().publisher(publisher, persistent=persistent)
    return publisher

include_router #

include_router(
    router: RedisRegistrator,
    *,
    prefix: str = "",
    dependencies: Sequence[Dependant] = (),
    middlewares: Sequence[BrokerMiddleware[Any, Any]] = (),
    include_in_schema: bool | None = None,
) -> None
Source code in faststream/redis/broker/registrator.py
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
@override
def include_router(
    self,
    router: "RedisRegistrator",  # type: ignore[override]
    *,
    prefix: str = "",
    dependencies: Sequence["Dependant"] = (),
    middlewares: Sequence["BrokerMiddleware[Any, Any]"] = (),
    include_in_schema: bool | None = None,
) -> None:
    if not isinstance(router, RedisRegistrator):
        msg = (
            f"Router must be an instance of RedisRegistrator, "
            f"got {type(router).__name__} instead"
        )
        raise SetupError(msg)

    super().include_router(
        router,
        prefix=prefix,
        dependencies=dependencies,
        middlewares=middlewares,
        include_in_schema=include_in_schema,
    )

include_routers #

include_routers(
    *routers: Registrator[MsgType, Any],
) -> None

Includes routers in the object.

Source code in faststream/_internal/broker/registrator.py
124
125
126
127
128
129
130
def include_routers(
    self,
    *routers: "Registrator[MsgType, Any]",
) -> None:
    """Includes routers in the object."""
    for r in routers:
        self.include_router(r)

connect async #

connect() -> ConnectionType

Connect to a remote server.

Source code in faststream/_internal/broker/broker.py
108
109
110
111
112
113
114
async def connect(self) -> ConnectionType:
    """Connect to a remote server."""
    if self._connection is None:
        self._connection = await self._connect()
        self._setup_logger()

    return self._connection

stop async #

stop(
    exc_type: type[BaseException] | None = None,
    exc_val: BaseException | None = None,
    exc_tb: Optional[TracebackType] = None,
) -> None
Source code in faststream/redis/broker/broker.py
181
182
183
184
185
186
187
188
189
async def stop(
    self,
    exc_type: type[BaseException] | None = None,
    exc_val: BaseException | None = None,
    exc_tb: Optional["TracebackType"] = None,
) -> None:
    await super().stop(exc_type, exc_val, exc_tb)
    await self.config.disconnect()
    self._connection = None

start async #

start() -> None
Source code in faststream/redis/broker/broker.py
191
192
193
async def start(self) -> None:
    _ = await self.connect()
    await super().start()

publish async #

publish(
    message: SendableMessage = None,
    channel: str | None = None,
    *,
    reply_to: str = "",
    headers: dict[str, Any] | None = None,
    correlation_id: str | None = None,
    list: str | None = None,
    stream: None = None,
    maxlen: int | None = None,
    pipeline: Optional[Pipeline[bytes]] = None,
) -> int
publish(
    message: SendableMessage = None,
    channel: str | None = None,
    *,
    reply_to: str = "",
    headers: dict[str, Any] | None = None,
    correlation_id: str | None = None,
    list: str | None = None,
    stream: str = ...,
    maxlen: int | None = None,
    pipeline: Optional[Pipeline[bytes]] = None,
) -> bytes
publish(
    message: SendableMessage = None,
    channel: str | None = None,
    *,
    reply_to: str = "",
    headers: dict[str, Any] | None = None,
    correlation_id: str | None = None,
    list: str | None = None,
    stream: str | None = None,
    maxlen: int | None = None,
    pipeline: Optional[Pipeline[bytes]] = None,
) -> int | bytes

Publish message directly.

This method allows you to publish a message in a non-AsyncAPI-documented way. It can be used in other frameworks or to publish messages at specific intervals.

PARAMETER DESCRIPTION
message

Message body to send.

TYPE: SendableMessage DEFAULT: None

channel

Redis PubSub object name to send message.

TYPE: str | None DEFAULT: None

reply_to

Reply message destination PubSub object name.

TYPE: str DEFAULT: ''

headers

Message headers to store metainformation.

TYPE: dict[str, Any] | None DEFAULT: None

correlation_id

Manual message correlation_id setter. correlation_id is a useful option to trace messages.

TYPE: str | None DEFAULT: None

list

Redis List object name to send message.

TYPE: str | None DEFAULT: None

stream

Redis Stream object name to send message.

TYPE: str | None DEFAULT: None

maxlen

Redis Stream maxlen publish option. Remove eldest message if maxlen exceeded.

TYPE: int | None DEFAULT: None

pipeline

Redis pipeline to use for publishing messages.

TYPE: Optional[Pipeline[bytes]] DEFAULT: None

RETURNS DESCRIPTION
int

The result of the publish operation, typically the number of messages published.

TYPE: int | bytes

Source code in faststream/redis/broker/broker.py
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
@override
async def publish(
    self,
    message: "SendableMessage" = None,
    channel: str | None = None,
    *,
    reply_to: str = "",
    headers: dict[str, Any] | None = None,
    correlation_id: str | None = None,
    list: str | None = None,
    stream: str | None = None,
    maxlen: int | None = None,
    pipeline: Optional["Pipeline[bytes]"] = None,
) -> int | bytes:
    """Publish message directly.

    This method allows you to publish a message in a non-AsyncAPI-documented way.
    It can be used in other frameworks or to publish messages at specific intervals.

    Args:
        message:
            Message body to send.
        channel:
            Redis PubSub object name to send message.
        reply_to:
            Reply message destination PubSub object name.
        headers:
            Message headers to store metainformation.
        correlation_id:
            Manual message correlation_id setter. correlation_id is a useful option to trace messages.
        list:
            Redis List object name to send message.
        stream:
            Redis Stream object name to send message.
        maxlen:
            Redis Stream maxlen publish option. Remove eldest message if maxlen exceeded.
        pipeline:
            Redis pipeline to use for publishing messages.

    Returns:
        int: The result of the publish operation, typically the number of messages published.
    """
    cmd = RedisPublishCommand(
        message,
        correlation_id=correlation_id or self.config.id_generator(),
        channel=channel,
        list=list,
        stream=stream,
        maxlen=maxlen,
        reply_to=reply_to,
        headers=headers,
        pipeline=pipeline,
        _publish_type=PublishType.PUBLISH,
        message_format=self.message_format,
    )

    result: int | bytes = await super()._basic_publish(
        cmd,
        producer=self.config.producer,
    )
    return result

request async #

request(
    message: SendableMessage,
    channel: str | None = None,
    *,
    list: str | None = None,
    stream: str | None = None,
    maxlen: int | None = None,
    correlation_id: str | None = None,
    headers: dict[str, Any] | None = None,
    timeout: float | None = 30.0,
) -> RedisChannelMessage
Source code in faststream/redis/broker/broker.py
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
@override
async def request(  # type: ignore[override]
    self,
    message: "SendableMessage",
    channel: str | None = None,
    *,
    list: str | None = None,
    stream: str | None = None,
    maxlen: int | None = None,
    correlation_id: str | None = None,
    headers: dict[str, Any] | None = None,
    timeout: float | None = 30.0,
) -> "RedisChannelMessage":
    cmd = RedisPublishCommand(
        message,
        correlation_id=correlation_id or self.config.id_generator(),
        channel=channel,
        list=list,
        stream=stream,
        maxlen=maxlen,
        headers=headers,
        timeout=timeout,
        _publish_type=PublishType.REQUEST,
        message_format=self.message_format,
    )
    msg: RedisChannelMessage = await super()._basic_request(
        cmd,
        producer=self.config.producer,
    )
    return msg

publish_batch async #

publish_batch(
    *messages: SendableMessage,
    list: str,
    correlation_id: str | None = None,
    reply_to: str = "",
    headers: dict[str, Any] | None = None,
    pipeline: Optional[Pipeline[bytes]] = None,
) -> int

Publish multiple messages to Redis List by one request.

PARAMETER DESCRIPTION
*messages

Messages bodies to send.

TYPE: SendableMessage DEFAULT: ()

list

Redis List object name to send messages.

TYPE: str

correlation_id

Manual message correlation_id setter. correlation_id is a useful option to trace messages.

TYPE: str | None DEFAULT: None

reply_to

Reply message destination PubSub object name.

TYPE: str DEFAULT: ''

headers

Message headers to store metainformation.

TYPE: dict[str, Any] | None DEFAULT: None

pipeline

Redis pipeline to use for publishing messages.

TYPE: Optional[Pipeline[bytes]] DEFAULT: None

RETURNS DESCRIPTION
int

The result of the batch publish operation.

TYPE: int

Source code in faststream/redis/broker/broker.py
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
@override
async def publish_batch(  # type: ignore[override]
    self,
    *messages: "SendableMessage",
    list: str,
    correlation_id: str | None = None,
    reply_to: str = "",
    headers: dict[str, Any] | None = None,
    pipeline: Optional["Pipeline[bytes]"] = None,
) -> int:
    """Publish multiple messages to Redis List by one request.

    Args:
        *messages: Messages bodies to send.
        list: Redis List object name to send messages.
        correlation_id: Manual message **correlation_id** setter. **correlation_id** is a useful option to trace messages.
        reply_to: Reply message destination PubSub object name.
        headers: Message headers to store metainformation.
        pipeline: Redis pipeline to use for publishing messages.

    Returns:
        int: The result of the batch publish operation.
    """
    cmd = RedisPublishCommand(
        *messages,
        list=list,
        reply_to=reply_to,
        headers=headers,
        correlation_id=correlation_id or self.config.id_generator(),
        pipeline=pipeline,
        _publish_type=PublishType.PUBLISH,
        message_format=self.message_format,
    )

    result: int = await self._basic_publish_batch(
        cmd,
        producer=self.config.producer,
    )
    return result

ping async #

ping(timeout: float | None = 3) -> bool
Source code in faststream/redis/broker/broker.py
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
@override
async def ping(self, timeout: float | None = 3) -> bool:
    sleep_time = (timeout or 10) / 10

    with move_on_after(timeout) as cancel_scope:
        if self._connection is None:
            return False

        while True:
            if cancel_scope.cancel_called:
                return False

            try:
                if await self._connection.ping():
                    return True

            except ConnectionError:
                pass

            await anyio.sleep(sleep_time)

    return False