RedisClusterBroker
faststream.redis.broker.cluster_broker.RedisClusterBroker #
RedisClusterBroker(
url: str = "redis://localhost:6379",
**kwargs: Unpack[RedisClusterParams],
)
Bases: RedisBroker
A Redis Cluster broker.
Source code in faststream/redis/broker/cluster_broker.py
38 39 40 41 42 43 | |
config instance-attribute #
config: ConfigComposition[BrokerConfigType] = (
ConfigComposition(config)
)
subscriber #
subscriber(
channel: Union[PubSub, str, None] = None,
*,
list: Union[ListSub, str, None] = None,
stream: Union[StreamSub, str, None] = None,
**kwargs: Any,
) -> LogicSubscriber
Source code in faststream/redis/broker/cluster_broker.py
72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 | |
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: str | None = None,
maxlen: int | None = None,
pipeline: Optional[Pipeline[bytes]] = EMPTY,
) -> int | bytes
Source code in faststream/redis/broker/cluster_broker.py
115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 | |
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/cluster_broker.py
158 159 160 161 162 163 164 165 166 | |
start async #
start() -> None
Source code in faststream/redis/broker/cluster_broker.py
168 169 170 | |
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]] = EMPTY,
) -> int
Source code in faststream/redis/broker/cluster_broker.py
172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 | |
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 | |
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 | |
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 | |
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: |
list | Redis List object name to send message. TYPE: |
stream | Redis Stream object name to send message. TYPE: |
headers | Message headers to store meta-information. Can be overridden by TYPE: |
reply_to | Reply message destination PubSub object name. TYPE: |
message_format | Which format to use when parsing messages. TYPE: |
title | AsyncAPI publisher object title. TYPE: |
description | AsyncAPI publisher object description. TYPE: |
schema | AsyncAPI publishing message type. Should be any python-native object annotation or TYPE: |
include_in_schema | Whether to include operation in AsyncAPI schema or not. TYPE: |
persistent | Whether to make the publisher persistent or not. TYPE: |
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 | |
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 | |
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 | |
connect async #
connect() -> ConnectionType
Connect to a remote server.
Source code in faststream/_internal/broker/broker.py
108 109 110 111 112 113 114 | |
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 | |