<?xml version="1.0" encoding="UTF-8"?>
<rss version="2.0" xmlns:atom="http://www.w3.org/2005/Atom" xmlns:dc="http://purl.org/dc/elements/1.1/">
  <channel>
    <title>DEV Community: Amrou Bellalouna</title>
    <description>The latest articles on DEV Community by Amrou Bellalouna (@shtlrs).</description>
    <link>https://dev.to/shtlrs</link>
    <image>
      <url>https://media2.dev.to/dynamic/image/width=90,height=90,fit=cover,gravity=auto,format=auto/https:%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Fuser%2Fprofile_image%2F4111705%2F58cc1ba8-7bf0-4edb-b0d4-61fb2ed716fc.jpg</url>
      <title>DEV Community: Amrou Bellalouna</title>
      <link>https://dev.to/shtlrs</link>
    </image>
    <atom:link rel="self" type="application/rss+xml" href="https://dev.to/feed/shtlrs"/>
    <language>en</language>
    <item>
      <title>How We Upgraded RabbitMQ to v4 Without Breaking 8M Daily Celery Tasks</title>
      <dc:creator>Amrou Bellalouna</dc:creator>
      <pubDate>Tue, 08 Sep 2026 14:00:00 +0000</pubDate>
      <link>https://dev.to/shtlrs/how-we-upgraded-rabbitmq-to-v4-without-breaking-8m-daily-celery-tasks-3mf5</link>
      <guid>https://dev.to/shtlrs/how-we-upgraded-rabbitmq-to-v4-without-breaking-8m-daily-celery-tasks-3mf5</guid>
      <description>&lt;p&gt;Upgrading to RabbitMQ v4 threatened to break our entire usage of Celery, more specifically tasks with ETAs.&lt;/p&gt;

&lt;p&gt;At 8M messages/day with zero downtime tolerance, we needed a migration strategy that preserves delayed task execution while switching from classic to quorum queues.&lt;/p&gt;

&lt;h2&gt;
  
  
  Prerequisites
&lt;/h2&gt;

&lt;p&gt;This post assumes familiarity with:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;a href="https://docs.celeryq.dev/en/stable/getting-started/introduction.html" rel="noopener noreferrer"&gt;Celery task queues&lt;/a&gt; (workers, tasks, brokers)&lt;/li&gt;
&lt;li&gt;
&lt;a href="https://www.rabbitmq.com/tutorials/tutorial-one-python.html" rel="noopener noreferrer"&gt;RabbitMQ fundamentals&lt;/a&gt; (exchanges, queues, routing)&lt;/li&gt;
&lt;li&gt;
&lt;a href="https://www.cloudamqp.com/blog/part1-rabbitmq-for-beginners-what-is-rabbitmq.html" rel="noopener noreferrer"&gt;Message queue concepts&lt;/a&gt; (producers, consumers, brokers)&lt;/li&gt;
&lt;li&gt;
&lt;a href="https://www.cloudamqp.com/blog/what-is-a-rabbitmq-vhost.html" rel="noopener noreferrer"&gt;Virtual hosts&lt;/a&gt; (multi-tenancy in RabbitMQ)&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;If you're comfortable with distributed task processing in Python, you're good to go.&lt;/p&gt;

&lt;h2&gt;
  
  
  Context
&lt;/h2&gt;

&lt;p&gt;At &lt;a href="https://kraken.tech" rel="noopener noreferrer"&gt;Kraken&lt;/a&gt;, we use &lt;a href="https://docs.celeryq.dev/" rel="noopener noreferrer"&gt;Celery&lt;/a&gt; to offload long-running tasks to workers via &lt;a href="https://www.rabbitmq.com/" rel="noopener noreferrer"&gt;RabbitMQ&lt;/a&gt; (v3.13.7.1). Our platform team needed to upgrade to RabbitMQ v4.2.2, which introduces breaking changes:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;a href="https://www.rabbitmq.com/blog/2021/08/21/4.0-deprecation-announcements#removal-of-global-qos" rel="noopener noreferrer"&gt;Global QoS removal&lt;/a&gt; - ETA/countdown tasks now &lt;a href="https://docs.celeryq.dev/en/v5.6.0/getting-started/backends-and-brokers/rabbitmq.html#limitations" rel="noopener noreferrer"&gt;block workers until execution time&lt;/a&gt;
&lt;/li&gt;
&lt;li&gt;
&lt;a href="https://www.rabbitmq.com/blog/2021/08/21/4.0-deprecation-announcements#removal-of-classic-queue-mirroring" rel="noopener noreferrer"&gt;Classic Queue Mirroring removal&lt;/a&gt; - forced migration to &lt;a href="https://www.rabbitmq.com/docs/quorum-queues" rel="noopener noreferrer"&gt;Quorum Queues&lt;/a&gt;
&lt;/li&gt;
&lt;/ul&gt;

&lt;h2&gt;
  
  
  Challenges
&lt;/h2&gt;

&lt;p&gt;&lt;strong&gt;Can't switch queue types on the fly&lt;/strong&gt;: RabbitMQ doesn't let you change a queue's type after it's created. Normally you'd just delete the queue and recreate it with the new type, but that wasn't an option for us. We're pushing 8M messages/day across &lt;a href="https://engineering.kraken.tech/news/2025/02/07/how-we-ship.html" rel="noopener noreferrer"&gt;many environments&lt;/a&gt;, and our SLAs don't allow for any downtime.&lt;/p&gt;

&lt;p&gt;&lt;strong&gt;Workers would block indefinitely&lt;/strong&gt;: Without &lt;a href="https://www.rabbitmq.com/docs/consumer-prefetch" rel="noopener noreferrer"&gt;global QoS&lt;/a&gt;, any task with an ETA would block a worker until that ETA arrived, completely defeating the purpose of async processing. Fortunately, Celery has a feature called &lt;a href="https://docs.celeryq.dev/en/v5.5.3/getting-started/backends-and-brokers/rabbitmq.html#native-delayed-delivery" rel="noopener noreferrer"&gt;Native Delayed Delivery&lt;/a&gt; that solves this—but it requires &lt;a href="https://www.rabbitmq.com/docs/quorum-queues" rel="noopener noreferrer"&gt;quorum queues&lt;/a&gt; bound to &lt;a href="https://www.rabbitmq.com/tutorials/tutorial-five-python.html" rel="noopener noreferrer"&gt;topic exchanges&lt;/a&gt;. Lucky for us, that's exactly what we needed to migrate to anyway.&lt;/p&gt;

&lt;h2&gt;
  
  
  Solution
&lt;/h2&gt;

&lt;p&gt;Our migration strategy:&lt;/p&gt;

&lt;ol&gt;
&lt;li&gt;Create a new &lt;a href="https://www.cloudamqp.com/blog/what-is-a-rabbitmq-vhost.html" rel="noopener noreferrer"&gt;vhost&lt;/a&gt; (&lt;code&gt;qhost&lt;/code&gt;), which will host quorum queues bound to topic exchanges&lt;/li&gt;
&lt;li&gt;Configure application code to support both queue types via a feature flag&lt;/li&gt;
&lt;li&gt;Transfer messages from old vhost (&lt;code&gt;chost&lt;/code&gt;) to new vhost without losing ETA information&lt;/li&gt;
&lt;li&gt;Decommission &lt;code&gt;chost&lt;/code&gt;
&lt;/li&gt;
&lt;li&gt;
&lt;a href="https://www.rabbitmq.com/docs/rolling-upgrade" rel="noopener noreferrer"&gt;Rolling upgrade&lt;/a&gt; to RabbitMQ v4&lt;/li&gt;
&lt;/ol&gt;

&lt;h3&gt;
  
  
  Phase 1: New Virtual Host
&lt;/h3&gt;

&lt;p&gt;The first step was straightforward: create a new virtual host (&lt;code&gt;qhost&lt;/code&gt;) to run alongside our existing one (&lt;code&gt;chost&lt;/code&gt;). For us, this meant updating some infrastructure manifests to provision the new vhost. Not the most exciting part, but essential for what comes next—we needed both vhosts running simultaneously to avoid any downtime.&lt;/p&gt;

&lt;h3&gt;
  
  
  Phase 2: Feature-Flagged Queue Configuration
&lt;/h3&gt;

&lt;p&gt;Next, we needed to make our application code flexible enough to handle both the old and new queue configurations. We introduced a &lt;code&gt;USE_QUORUM_QUEUES&lt;/code&gt; &lt;a href="https://12factor.net/config" rel="noopener noreferrer"&gt;environment variable&lt;/a&gt; (a common &lt;a href="https://martinfowler.com/articles/feature-toggles.html" rel="noopener noreferrer"&gt;feature flag&lt;/a&gt; pattern) to control which type of queues to create:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight python"&gt;&lt;code&gt;&lt;span class="kn"&gt;from&lt;/span&gt; &lt;span class="n"&gt;kombu&lt;/span&gt; &lt;span class="kn"&gt;import&lt;/span&gt; &lt;span class="n"&gt;Queue&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;Exchange&lt;/span&gt;
&lt;span class="kn"&gt;from&lt;/span&gt; &lt;span class="n"&gt;django.conf&lt;/span&gt; &lt;span class="kn"&gt;import&lt;/span&gt; &lt;span class="n"&gt;settings&lt;/span&gt;

&lt;span class="k"&gt;def&lt;/span&gt; &lt;span class="nf"&gt;build_queue&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;queue_name&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nb"&gt;str&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt; &lt;span class="o"&gt;-&amp;gt;&lt;/span&gt; &lt;span class="n"&gt;Queue&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
    &lt;span class="n"&gt;queue_type&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;quorum&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt; &lt;span class="k"&gt;if&lt;/span&gt; &lt;span class="n"&gt;settings&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;USE_QUORUM_QUEUES&lt;/span&gt; &lt;span class="k"&gt;else&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;classic&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;
    &lt;span class="n"&gt;exchange_type&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;topic&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt; &lt;span class="k"&gt;if&lt;/span&gt; &lt;span class="n"&gt;settings&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;USE_QUORUM_QUEUES&lt;/span&gt; &lt;span class="k"&gt;else&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;direct&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;

    &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="nc"&gt;Queue&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;
        &lt;span class="n"&gt;name&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="n"&gt;queue_name&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;
        &lt;span class="n"&gt;exchange&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="nc"&gt;Exchange&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;queue_type&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="nb"&gt;type&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="n"&gt;exchange_type&lt;/span&gt;&lt;span class="p"&gt;),&lt;/span&gt;
        &lt;span class="n"&gt;queue_arguments&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;x-queue-type&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="n"&gt;queue_type&lt;/span&gt;&lt;span class="p"&gt;}&lt;/span&gt;
    &lt;span class="p"&gt;)&lt;/span&gt;

&lt;span class="n"&gt;task_queues&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="nf"&gt;build_queue&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;first_queue&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;),&lt;/span&gt; &lt;span class="p"&gt;...,&lt;/span&gt; &lt;span class="nf"&gt;build_queue&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;last_queue&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)]&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;When we deployed this change, Kubernetes did its usual &lt;a href="https://kubernetes.io/docs/tutorials/kubernetes-basics/update/update-intro/" rel="noopener noreferrer"&gt;rolling update&lt;/a&gt;, and pods gradually switched from &lt;code&gt;chost&lt;/code&gt; to &lt;code&gt;qhost&lt;/code&gt;. This meant we temporarily had some pods on the old vhost and some on the new one, which was totally fine. Any stragglers would get caught in Phase 3.&lt;/p&gt;

&lt;h3&gt;
  
  
  Phase 3: Message Transfer with ETA Transformation
&lt;/h3&gt;

&lt;p&gt;Here's where things got interesting. We needed to transfer potentially millions of messages from &lt;code&gt;chost&lt;/code&gt; to &lt;code&gt;qhost&lt;/code&gt; without losing any data or causing downtime.&lt;/p&gt;

&lt;p&gt;At first glance, this sounds like a perfect job for RabbitMQ's &lt;a href="https://www.rabbitmq.com/docs/shovel" rel="noopener noreferrer"&gt;Shovel plugin&lt;/a&gt;, right? Just copy messages from one vhost to another. Unfortunately, a shovel copies messages as-is, including the &lt;code&gt;eta&lt;/code&gt; header. That's exactly what we're trying to avoid—those headers would cause workers to block, bringing us right back to square one.&lt;/p&gt;

&lt;p&gt;Instead, we had to transform messages during transfer. The idea was to replicate what &lt;a href="https://github.com/celery/celery/blob/0527296acb1f1790788301d4395ba6d5ce2a9704/celery/app/base.py#L854-L876" rel="noopener noreferrer"&gt;Celery does internally&lt;/a&gt; when Native Delayed Delivery is enabled: extract the ETA header, calculate the appropriate delay-based routing key, and route to the right exchange. If you want to understand how this works under the hood, &lt;a href="https://docs.particular.net/transports/rabbitmq/delayed-delivery" rel="noopener noreferrer"&gt;this guide&lt;/a&gt; explains the mechanics well.&lt;/p&gt;

&lt;p&gt;Here's the core logic (simplified for clarity):&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight python"&gt;&lt;code&gt;&lt;span class="kn"&gt;import&lt;/span&gt; &lt;span class="n"&gt;pika&lt;/span&gt;  &lt;span class="c1"&gt;# RabbitMQ Python client: https://pika.readthedocs.io/
&lt;/span&gt;&lt;span class="kn"&gt;from&lt;/span&gt; &lt;span class="n"&gt;kombu.transport&lt;/span&gt; &lt;span class="kn"&gt;import&lt;/span&gt; &lt;span class="n"&gt;native_delayed_delivery&lt;/span&gt; &lt;span class="k"&gt;as&lt;/span&gt; &lt;span class="n"&gt;kombu_utils&lt;/span&gt;

&lt;span class="k"&gt;def&lt;/span&gt; &lt;span class="nf"&gt;get_routing_details&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;method&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;properties&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;queue_name&lt;/span&gt;&lt;span class="p"&gt;):&lt;/span&gt;
    &lt;span class="n"&gt;target_exchange_name&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;method&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;exchange&lt;/span&gt; &lt;span class="ow"&gt;or&lt;/span&gt; &lt;span class="n"&gt;queue_name&lt;/span&gt;
    &lt;span class="n"&gt;target_routing_key&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;method&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;routing_key&lt;/span&gt; &lt;span class="ow"&gt;or&lt;/span&gt; &lt;span class="n"&gt;queue_name&lt;/span&gt;

    &lt;span class="n"&gt;eta_str&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nf"&gt;str&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;properties&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;headers&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;pop&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;eta&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="sh"&gt;""&lt;/span&gt;&lt;span class="p"&gt;))&lt;/span&gt;
    &lt;span class="n"&gt;countdown_in_seconds&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nf"&gt;compute_countdown&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;eta_str&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;

    &lt;span class="k"&gt;if&lt;/span&gt; &lt;span class="n"&gt;countdown_in_seconds&lt;/span&gt; &lt;span class="ow"&gt;and&lt;/span&gt; &lt;span class="n"&gt;countdown_in_seconds&lt;/span&gt; &lt;span class="o"&gt;&amp;gt;&lt;/span&gt; &lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
        &lt;span class="n"&gt;target_routing_key&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;kombu_utils&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;calculate_routing_key&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="nf"&gt;int&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;countdown_in_seconds&lt;/span&gt;&lt;span class="p"&gt;),&lt;/span&gt; &lt;span class="n"&gt;target_routing_key&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
        &lt;span class="n"&gt;target_exchange_name&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;celery_delayed_27&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;

    &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="n"&gt;target_routing_key&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;target_exchange_name&lt;/span&gt;

&lt;span class="n"&gt;chost_connection_string&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nf"&gt;read_from_env&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;CHOST_CONNECTION_STRING&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
&lt;span class="n"&gt;qhost_connection_string&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nf"&gt;read_from_env&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;QHOST_CONNECTION_STRING&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;

&lt;span class="n"&gt;source_channel&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;pika&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nc"&gt;BlockingConnection&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;pika&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nc"&gt;URLParameters&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;chost_connection_string&lt;/span&gt;&lt;span class="p"&gt;)).&lt;/span&gt;&lt;span class="nf"&gt;channel&lt;/span&gt;&lt;span class="p"&gt;()&lt;/span&gt;
&lt;span class="n"&gt;dest_channel&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;pika&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nc"&gt;BlockingConnection&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;pika&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nc"&gt;URLParameters&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;amqps://user:pwd@host:port/qhost&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)).&lt;/span&gt;&lt;span class="nf"&gt;channel&lt;/span&gt;&lt;span class="p"&gt;()&lt;/span&gt;

&lt;span class="k"&gt;for&lt;/span&gt; &lt;span class="n"&gt;method&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;properties&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;body&lt;/span&gt; &lt;span class="ow"&gt;in&lt;/span&gt; &lt;span class="n"&gt;source_channel&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;consume&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;target_queue&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;):&lt;/span&gt;
    &lt;span class="k"&gt;try&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
        &lt;span class="n"&gt;routing_key&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;exchange&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nf"&gt;get_routing_details&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;method&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;properties&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;target_queue&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
        &lt;span class="n"&gt;dest_channel&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;basic_publish&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;exchange&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="n"&gt;exchange&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;routing_key&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="n"&gt;routing_key&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;
                                          &lt;span class="n"&gt;body&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="n"&gt;body&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;properties&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="n"&gt;properties&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
        &lt;span class="n"&gt;source_conn&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;channel&lt;/span&gt;&lt;span class="p"&gt;().&lt;/span&gt;&lt;span class="nf"&gt;basic_ack&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;delivery_tag&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="n"&gt;method&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;delivery_tag&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
    &lt;span class="k"&gt;except&lt;/span&gt; &lt;span class="nb"&gt;Exception&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
        &lt;span class="n"&gt;source_conn&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;channel&lt;/span&gt;&lt;span class="p"&gt;().&lt;/span&gt;&lt;span class="nf"&gt;basic_nack&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;delivery_tag&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="n"&gt;method&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;delivery_tag&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;requeue&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="bp"&gt;True&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h4&gt;
  
  
  Some notes on the shoveling part:
&lt;/h4&gt;

&lt;ul&gt;
&lt;li&gt;The above is pseudocode to illustrate the concept. Our production version uses &lt;a href="https://aio-pika.readthedocs.io/" rel="noopener noreferrer"&gt;&lt;code&gt;aio_pika&lt;/code&gt;&lt;/a&gt; for &lt;a href="https://docs.python.org/3/library/asyncio.html" rel="noopener noreferrer"&gt;async I/O&lt;/a&gt; (to avoid blocking), &lt;a href="https://docs.python.org/3/library/multiprocessing.html" rel="noopener noreferrer"&gt;multiprocessing&lt;/a&gt; to handle high throughput, message backups for disaster recovery, and extensive logging to track everything.&lt;/li&gt;
&lt;li&gt;The script was deployed to run as a daemon, and would only work if &lt;code&gt;USE_QUORUM_QUEUES&lt;/code&gt; was set to &lt;code&gt;True&lt;/code&gt;.&lt;/li&gt;
&lt;li&gt;More mechanics were in place to determine if there were any messages to transfer in the first place, do transfers in batches, etc.&lt;/li&gt;
&lt;/ul&gt;

&lt;h3&gt;
  
  
  Phase 4 &amp;amp; 5: Cleanup and Upgrade
&lt;/h3&gt;

&lt;p&gt;Once all messages were safely transferred to &lt;code&gt;qhost&lt;/code&gt;, our platform team took over for the final steps. They deleted the &lt;code&gt;chost&lt;/code&gt; virtual host (which removed all the classic queues that are incompatible with v4), and then performed a &lt;a href="https://www.rabbitmq.com/docs/rolling-upgrade" rel="noopener noreferrer"&gt;rolling upgrade&lt;/a&gt; to RabbitMQ v4.2.2.&lt;/p&gt;

&lt;p&gt;And that was it! We successfully migrated from RabbitMQ v3 to v4 with zero downtime, while preserving ETA task behavior and handling 8M messages/day across all our environments.&lt;/p&gt;

</description>
      <category>architecture</category>
      <category>backend</category>
      <category>devops</category>
      <category>python</category>
    </item>
  </channel>
</rss>
