Skip to content

Commit 41933f1

Browse files
Retry queue for inbound queue
1 parent a22a782 commit 41933f1

2 files changed

Lines changed: 42 additions & 2 deletions

File tree

app/config.py

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,10 @@ class Settings(BaseSettings):
4040
app_rabbit_inbound_exchange: str = "in-bound-ex"
4141
app_rabbit_inbound_routing_key: str = "in-bound-rk"
4242

43+
app_rabbit_inbound_retry_queue: str = "in-bound-retry-queue"
44+
app_rabbit_inbound_retry_exchange: str = "in-bound-retry-ex"
45+
app_rabbit_inbound_retry_routing_key: str = "in-bound-retry-rk"
46+
4347
app_rabbit_outbound_queue: str = "out-bound-queue"
4448
app_rabbit_outbound_exchange: str = "out-bound-ex"
4549
app_rabbit_outbound_routing_key: str = "out-bound-rk"

app/rabbit_mq_setup/rabbit_mq_initialisation.py

Lines changed: 38 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -37,7 +37,7 @@ async def configure_rabbitmq():
3737
)
3838

3939
await inbound_queue.bind(
40-
inbound_exchange, routing_key=settings.app_rabbit_inbound_exchange
40+
inbound_exchange, routing_key=settings.app_rabbit_inbound_routing_key
4141
)
4242

4343
inbound_dead_queue = await channel.declare_queue(
@@ -56,6 +56,42 @@ async def configure_rabbitmq():
5656
routing_key=f"{settings.app_rabbit_inbound_queue}.dlq.rk",
5757
)
5858

59+
inbound_retry_queue = await channel.declare_queue(
60+
settings.app_rabbit_inbound_retry_queue,
61+
durable=settings.app_rabbit_mq_durable,
62+
arguments={
63+
"x-dead-letter-exchange": f"{settings.app_rabbit_inbound_retry_queue}.dlx",
64+
"x-dead-letter-routing-key": f"{settings.app_rabbit_inbound_retry_queue}.dlq.rk",
65+
},
66+
)
67+
68+
inbound_retry_exchange = await channel.declare_exchange(
69+
settings.app_rabbit_inbound_retry_exchange,
70+
type=aio_pika.ExchangeType.TOPIC,
71+
durable=settings.app_rabbit_mq_durable,
72+
)
73+
74+
await inbound_retry_queue.bind(
75+
inbound_retry_exchange,
76+
routing_key=settings.app_rabbit_inbound_retry_exchange,
77+
)
78+
79+
inbound_retry_dead_queue = await channel.declare_queue(
80+
f"{settings.app_rabbit_inbound_retry_queue}.dlq",
81+
durable=settings.app_rabbit_mq_durable,
82+
)
83+
84+
inbound_retry_dead_exchange = await channel.declare_exchange(
85+
f"{settings.app_rabbit_inbound_retry_queue}.dlx",
86+
type=aio_pika.ExchangeType.TOPIC,
87+
durable=settings.app_rabbit_mq_durable,
88+
)
89+
90+
await inbound_retry_dead_queue.bind(
91+
inbound_retry_dead_exchange,
92+
routing_key=f"{settings.app_rabbit_inbound_retry_queue}.dlq.rk",
93+
)
94+
5995
outbound_queue = await channel.declare_queue(
6096
settings.app_rabbit_outbound_queue,
6197
durable=settings.app_rabbit_mq_durable,
@@ -72,7 +108,7 @@ async def configure_rabbitmq():
72108
)
73109

74110
await outbound_queue.bind(
75-
outbound_exchange, routing_key=settings.app_rabbit_outbound_exchange
111+
outbound_exchange, routing_key=settings.app_rabbit_outbound_routing_key
76112
)
77113

78114
outbound_dead_queue = await channel.declare_queue(

0 commit comments

Comments
 (0)