Takip et

RabbitMQ ve Python’ın Puka Kütüphanesi ile Birden Fazla Tüketiciye Mesaj İletimi

RabbitMQ ve Python’ın Puka Kütüphanesi ile Birden Fazla Tüketiciye Mesaj İletimi Giriş: Mesaj Kuyrukları ve Dağıtık Sistemler Günümüzün

RabbitMQ ve Python’ın Puka Kütüphanesi ile Birden Fazla Tüketiciye Mesaj İletimi

Giriş: Mesaj Kuyrukları ve Dağıtık Sistemler

Günümüzün modern yazılım mimarilerinde, uygulamalar genellikle bağımsız servislerden oluşan dağıtık sistemler şeklinde tasarlanır. Bu servislerin birbiriyle güvenli, esnek ve ölçeklenebilir bir şekilde iletişim kurması kritik öneme sahiptir. Monolitik uygulamaların yerini alan mikroservis mimarileri, bu iletişim ihtiyacını daha da belirgin hale getirmiştir. İşte tam bu noktada mesaj kuyrukları (message queues) devreye girer.

Mesaj kuyrukları, farklı uygulama bileşenlerinin asenkron olarak iletişim kurmasını sağlayan bir aracı (broker) görevi görür. Bir bileşen (üretici/producer), mesajı kuyruğa bırakır ve diğer bileşenler (tüketici/consumer) bu mesajları kuyruktan alır. Bu modelin en büyük avantajları arasında sistem bileşenleri arasındaki bağımlılığı azaltması, ölçeklenebilirliği artırması, yük dengelemesi sağlaması ve sistemin genel güvenilirliğini yükseltmesi yer alır. Örneğin, bir web uygulaması bir kullanıcının fotoğrafını işlemek gibi zaman alıcı bir görevi direkt olarak yerine getirmek yerine, bu görevi bir mesaj kuyruğuna bırakabilir. Arka plandaki bir işleyici (worker) servisi bu mesajı kuyruktan alıp fotoğrafı işlerken, web uygulaması kullanıcıya anında yanıt verebilir.

Mesaj Kuyrukları Neden Önemli?

* Bağımsızlık (Decoupling): Üreticiler ve tüketiciler birbirlerinin varlığından haberdar olmak zorunda değildir. Sadece belirli bir mesaj formatını ve kuyruğu bilmeleri yeterlidir. Bu, sistemin farklı parçalarının bağımsız olarak geliştirilmesine, dağıtılmasına ve ölçeklenmesine olanak tanır.
* Asenkron İletişim: Mesajlar anında işlenmek zorunda değildir. Üretici mesajı bırakır ve kendi işine devam ederken, tüketici mesajı uygun olduğunda işler. Bu, özellikle yavaş veya yoğun işlemler için sistem yanıt sürelerini iyileştirir.
* Yük Dengeleme ve Ölçeklenebilirlik: Birden fazla tüketici aynı kuyruktan mesaj alarak iş yükünü paylaşabilir. Yeni bir tüketici eklemek veya mevcut tüketicileri kaldırmak, sistemin genel performansını etkilemeden kolayca yapılabilir.
* Güvenilirlik ve Dayanıklılık: Mesaj kuyrukları, mesajların işlenene kadar kalıcı olarak depolanmasını sağlayabilir. Bir tüketici çökerse veya bağlantısı kesilirse, mesaj kaybolmaz ve başka bir tüketici tarafından işlenebilir.
* Hata Toleransı: Bir bileşen arızalandığında, diğer bileşenler çalışmaya devam edebilir. Arızalanan bileşen tekrar ayağa kalktığında, biriken mesajları işlemeye başlayabilir.

RabbitMQ’ya Kısa Bir Bakış

RabbitMQ, bu alandaki en popüler ve güçlü mesaj brokerlarından biridir. Advanced Message Queuing Protocol (AMQP) standardını uygulayan açık kaynaklı bir mesaj kuyruğu yazılımıdır. Erlang ile yazılmış olup, yüksek performans, güvenilirlik ve esneklik sunar. Çeşitli programlama dilleri için istemci kütüphaneleri (client libraries) bulunur ve bu da onu farklı teknolojilerle entegrasyon için ideal bir seçim haline getirir.

Bu makalede, RabbitMQ’nun temel prensiplerini ve Python için geliştirilmiş Puka kütüphanesini kullanarak birden fazla tüketiciye mesajların nasıl iletileceğini detaylı bir şekilde inceleyeceğiz. Puka, Twisted çatısı ile entegre olabilen asenkron bir RabbitMQ istemcisidir ve özellikle yüksek eşzamanlılık gerektiren uygulamalarda tercih edilebilir.

RabbitMQ Temelleri: Kavramlar ve Kurulum

RabbitMQ’yu etkili bir şekilde kullanabilmek için bazı temel kavramları anlamak önemlidir. Bu kavramlar, mesajların nasıl üretildiği, yönlendirildiği ve tüketildiği konusunda bir çerçeve sunar.

Temel RabbitMQ Kavramları

*

Producer (Üretici)

Mesajları oluşturan ve RabbitMQ’ya gönderen uygulamadır. Üretici, mesajları doğrudan kuyruklara göndermez; bunun yerine bir Exchange’e gönderir.

*

Consumer (Tüketici)

Mesajları RabbitMQ’dan alan ve işleyen uygulamadır. Tüketiciler, belirli bir kuyruğa abone olur ve oradan mesajları çeker.

*

Queue (Kuyruk)

Mesajların depolandığı yerdir. Tüketiciler, mesajları bu kuyruklardan alır. Kuyruklar, isimleri ile tanımlanır ve mesajları FIFO (First-In, First-Out) prensibiyle tutar.

*

Exchange (Değiştirici)

Üreticiden gelen mesajları alır ve belirli kurallara göre bir veya daha fazla kuyruğa yönlendirir. RabbitMQ’da farklı Exchange türleri bulunur:
* Direct Exchange: Mesajı, mesajın routing key‘i ile kuyruğun binding key‘i tam olarak eşleşen kuyruklara yönlendirir.
* Fanout Exchange: Mesajı, kendisine bağlı olan tüm kuyruklara yönlendirir. Yayınlama (publish/subscribe) senaryoları için idealdir.
* Topic Exchange: Mesajı, mesajın routing key‘i ile kuyruğun binding key‘i belirli bir kalıba (pattern) göre eşleşen kuyruklara yönlendirir. Daha esnek yönlendirme sağlar.
* Headers Exchange: Mesajı, mesajın başlık (header) özelliklerine göre yönlendirir. Daha az kullanılır.

*

Binding (Bağlama)

Bir Exchange ile bir Kuyruk arasındaki bağlantıdır. Binding, Exchange’e hangi kuyruklara mesaj göndereceğini söyler. Binding’ler genellikle bir binding key ile tanımlanır.

RabbitMQ Kurulumu

RabbitMQ’yu kullanmaya başlamadan önce onu bir sunucuya kurmanız gerekir. En yaygın kurulum yöntemleri şunlardır:

* Docker ile Kurulum: En hızlı ve pratik yöntemlerden biridir.

docker run -it --rm --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:3-management

Bu komut, RabbitMQ’yu yönetim arayüzü ile birlikte bir Docker konteynerinde başlatır. 5672 portu AMQP bağlantıları için, 15672 portu ise yönetim arayüzü için kullanılır.

* İşletim Sistemi Üzerine Kurulum: İşletim sisteminize (Ubuntu, CentOS, Windows vb.) özgü paket yöneticileri veya kurulum dosyaları kullanılarak kurulabilir. Örneğin, Ubuntu’da:

sudo apt update
    sudo apt install rabbitmq-server

Kurulumdan sonra RabbitMQ servisini başlatmanız gerekebilir: sudo systemctl start rabbitmq-server.

RabbitMQ Yönetim Arayüzü

RabbitMQ ile birlikte gelen yönetim arayüzü (Management UI), kuyrukları, exchange’leri, bağlantıları ve mesajları izlemek için çok değerli bir araçtır. Docker ile kurduğunuzda http://localhost:15672 adresinden erişebilirsiniz (varsayılan kullanıcı adı: guest, şifre: guest). Bu arayüz, sistemin durumunu anlamak ve sorun gidermek için oldukça faydalıdır.

Python ve Puka: Neden Puka?

Python için RabbitMQ ile etkileşim kurmak üzere birkaç farklı istemci kütüphanesi bulunmaktadır. En bilinenleri Pika, aio-pika ve Puka’dır. Bu makale, özellikle Puka kütüphanesine odaklanmaktadır.

Puka Nedir ve Avantajları Nelerdir?

Puka, Python için Twisted çatısı ile uyumlu, asenkron bir AMQP istemcisidir. Asenkron yapısı sayesinde, I/O işlemlerini beklerken uygulamanın bloke olmamasını sağlar. Bu, özellikle aynı anda birden fazla bağlantıyı yönetmesi gereken veya düşük gecikmeli, yüksek eşzamanlılık gerektiren uygulamalar için büyük bir avantajdır.

Puka’nın temel avantajları şunlardır:
* Asenkron Yapı: Twisted event loop üzerine inşa edildiği için, ağ işlemlerini bloke etmeden aynı anda birçok mesajı gönderebilir veya alabilir.
* Twisted Entegrasyonu: Eğer uygulamanız zaten Twisted kullanıyorsa, Puka doğal bir seçim olacaktır. Twisted’ın güçlü ağ ve eşzamanlılık yeteneklerinden faydalanır.
* Basit API: Puka, AMQP protokolünün karmaşıklığını soyutlayarak, nispeten basit ve kullanımı kolay bir API sunar.
* Geri Çağrı Tabanlı Model: Asenkron işlemleri yönetmek için geri çağrı (callback) fonksiyonlarını kullanır.

Ancak, Puka’nın dezavantajı, Pika veya aio-pika kadar aktif olarak geliştirilmemesi ve topluluk desteğinin daha az olmasıdır. Yine de, belirli senaryolarda ve Twisted tabanlı projelerde hala geçerli bir seçenek olabilir.

Puka Kurulumu

Puka’yı Python projenize dahil etmek oldukça basittir. pip paket yöneticisini kullanarak kurabilirsiniz:

pip install puka

Bu komut, Puka kütüphanesini ve bağımlılıklarını sisteminize yükleyecektir.

Puka ile RabbitMQ Bağlantısı ve Temel İşlemler

Şimdi, Puka kullanarak RabbitMQ ile nasıl bağlantı kurulacağını ve temel mesajlaşma işlemlerinin nasıl gerçekleştirileceğini adım adım inceleyelim.

Puka İstemcisi ile Bağlantı Kurma

Puka ile RabbitMQ’ya bağlanmak için puka.Client sınıfını kullanırız. Bağlantı asenkron olduğu için, connect() metodunun döndürdüğü defer.Deferred nesnesini dinlememiz gerekir.

import puka
from twisted.internet import reactor, defer

RabbitMQ bağlantı URL'si

RABBITMQ_URL = "amqp://guest:guest@localhost:5672/" @defer.inlineCallbacks def connect_to_rabbitmq(): client = puka.Client(RABBITMQ_URL) print("RabbitMQ'ya bağlanılıyor...") try: yield client.connect() print("RabbitMQ'ya başarıyla bağlandı.") # Diğer işlemleri burada yapabiliriz yield client.close() print("Bağlantı kapatıldı.") except Exception as e: print(f"Bağlantı hatası: {e}") finally: reactor.stop() if __name__ == "__main__": reactor.callWhenRunning(connect_to_rabbitmq) reactor.run()

Bu kod parçası, Twisted reactor‘ını kullanarak RabbitMQ’ya bağlanır ve bağlantı başarılı olursa bir mesaj yazdırır. defer.inlineCallbacks dekoratörü, asenkron kodları senkron gibi yazmamızı sağlar (yield anahtar kelimesi ile).

Kuyruk ve Exchange Bildirimi

Mesaj göndermeden veya almadan önce, kullanacağımız kuyrukları ve exchange’leri RabbitMQ’ya bildirmemiz (declare) gerekir. Bu işlem, eğer mevcut değillerse onları oluşturur.

import puka
from twisted.internet import reactor, defer

RABBITMQ_URL = "amqp://guest:guest@localhost:5672/"

@defer.inlineCallbacks
def declare_entities():
    client = puka.Client(RABBITMQ_URL)
    yield client.connect()
    print("Bağlantı başarılı.")

    # Exchange bildirimi (fanout türünde)
    exchange_name = 'my_fanout_exchange'
    yield client.exchange_declare(exchange=exchange_name, type='fanout')
    print(f"'{exchange_name}' exchange bildirildi.")

    # Kuyruk bildirimi
    queue_name = 'my_queue'
    yield client.queue_declare(queue=queue_name)
    print(f"'{queue_name}' kuyruğu bildirildi.")

    yield client.close()
    reactor.stop()

if __name__ == "__main__":
    reactor.callWhenRunning(declare_entities)
    reactor.run()

exchange_declare ve queue_declare metodları, ilgili varlıkları RabbitMQ üzerinde tanımlar. type='fanout' parametresi, exchange’in türünü belirtir.

Binding Oluşturma

Bir Exchange’den bir kuyruğa mesajların yönlendirilmesi için bir binding oluşturulmalıdır.

import puka
from twisted.internet import reactor, defer

RABBITMQ_URL = "amqp://guest:guest@localhost:5672/"

@defer.inlineCallbacks
def create_binding():
    client = puka.Client(RABBITMQ_URL)
    yield client.connect()
    print("Bağlantı başarılı.")

    exchange_name = 'my_direct_exchange'
    queue_name = 'my_direct_queue'
    routing_key = 'info' # Direct exchange için routing key

    # Exchange ve kuyruk bildirimleri (gerekirse)
    yield client.exchange_declare(exchange=exchange_name, type='direct')
    yield client.queue_declare(queue=queue_name)

    # Binding oluşturma
    yield client.queue_bind(exchange=exchange_name, queue=queue_name, routing_key=routing_key)
    print(f"'{exchange_name}' exchange ile '{queue_name}' kuyruğu '{routing_key}' anahtarıyla bağlandı.")

    yield client.close()
    reactor.stop()

if __name__ == "__main__":
    reactor.callWhenRunning(create_binding)
    reactor.run()

queue_bind metodu, bir exchange’i bir kuyruğa bağlar. routing_key parametresi, direct ve topic exchange’ler için mesaj yönlendirmede kullanılır. fanout exchange’lerde routing_key genellikle boş bırakılır veya göz ardı edilir.

Mesaj Gönderme (Producer)

Mesaj göndermek için basic_publish metodunu kullanırız.

import puka
from twisted.internet import reactor, defer
import time

RABBITMQ_URL = "amqp://guest:guest@localhost:5672/"

@defer.inlineCallbacks
def send_message():
    client = puka.Client(RABBITMQ_URL)
    yield client.connect()
    print("Producer: Bağlantı başarılı.")

    exchange_name = 'my_fanout_exchange'
    # Fanout exchange'e mesaj gönderirken routing_key boş bırakılır
    yield client.exchange_declare(exchange=exchange_name, type='fanout')
    print(f"Producer: '{exchange_name}' exchange bildirildi.")

    message_count = 5
    for i in range(message_count):
        message_body = f"Hello from producer! Message {i+1}"
        yield client.basic_publish(exchange=exchange_name, routing_key='', body=message_body)
        print(f"Producer: Mesaj gönderildi: '{message_body}'")
        time.sleep(0.1) # Küçük bir gecikme

    print(f"Producer: Toplam {message_count} mesaj gönderildi.")
    yield client.close()
    reactor.stop()

if __name__ == "__main__":
    reactor.callWhenRunning(send_message)
    reactor.run()

basic_publish metoduna exchange adı, routing_key (eğer varsa) ve mesajın body‘si verilir.

Mesaj Alma (Consumer) ve Onaylama (Acknowledgement)

Tüketiciler, basic_consume metodu ile bir kuyruktan mesajları almaya başlar. Mesaj işlendikten sonra, RabbitMQ’ya mesajın başarıyla alındığını ve işlendiğini bildirmek için bir onay (acknowledgement – ACK) göndermek önemlidir. Bu, mesajın kaybolmamasını sağlar.

import puka
from twisted.internet import reactor, defer

RABBITMQ_URL = "amqp://guest:guest@localhost:5672/"

@defer.inlineCallbacks
def receive_messages():
    client = puka.Client(RABBITMQ_URL)
    yield client.connect()
    print("Consumer: Bağlantı başarılı.")

    exchange_name = 'my_fanout_exchange'
    # Geçici, özel ve otomatik silinen bir kuyruk oluşturalım
    result = yield client.queue_declare(exclusive=True, auto_delete=True)
    queue_name = result['queue']
    print(f"Consumer: Geçici kuyruk oluşturuldu: '{queue_name}'")

    yield client.exchange_declare(exchange=exchange_name, type='fanout')
    yield client.queue_bind(exchange=exchange_name, queue=queue_name, routing_key='')
    print(f"Consumer: '{exchange_name}' exchange ile '{queue_name}' kuyruğu bağlandı.")

    # Mesajları tüketmeye başla
    consume_result = yield client.basic_consume(queue=queue_name, no_ack=False)
    print("Consumer: Mesajlar bekleniyor. Çıkmak için Ctrl+C.")

    try:
        while True:
            message = yield consume_result.wait()
            print(f"Consumer: Mesaj alındı: '{message['body']}' (Delivery Tag: {message['delivery_tag']})")
            # Mesajı işleme simülasyonu
            # time.sleep(1)
            
            # Mesajı onaylama
            yield client.basic_ack(message['delivery_tag'])
            print(f"Consumer: Mesaj onaylandı: {message['delivery_tag']}")
    except KeyboardInterrupt:
        print("Consumer: Kapatılıyor...")
    finally:
        yield client.close()
        reactor.stop()

if __name__ == "__main__":
    reactor.callWhenRunning(receive_messages)
    reactor.run()

queue_declare(exclusive=True, auto_delete=True) ile oluşturulan kuyruk, sadece bu bağlantı tarafından kullanılabilir (exclusive) ve bağlantı kapandığında otomatik olarak silinir (auto_delete). Bu, özellikle fanout exchange ile yayınlama senaryolarında her tüketicinin kendi geçici kuyruğuna sahip olması istendiğinde kullanışlıdır.

basic_consume metoduna no_ack=False parametresi verilerek manuel onaylama etkinleştirilir. Mesaj işlendikten sonra client.basic_ack(message['delivery_tag']) ile RabbitMQ’ya mesajın işlendiği bildirilir. Eğer bir mesaj onaylanmazsa ve tüketici bağlantısı kesilirse, RabbitMQ bu mesajı başka bir tüketiciye yeniden teslim eder.

Birden Fazla Tüketiciye Mesaj İletimi Senaryoları

RabbitMQ ve Puka’yı kullanarak birden fazla tüketiciye mesaj iletmenin iki ana senaryosu vardır: iş yükünü paylaşma (work queues) ve mesajları tüm tüketicilere yayınlama (publish/subscribe).

Senaryo 1: İş Yükünü Paylaştırma (Work Queues)

Bu senaryoda, aynı işi yapacak birden fazla tüketici vardır ve her mesaj sadece bir tüketici tarafından işlenmelidir. RabbitMQ, mesajları kuyruktan tüketicilere round-robin (sırayla) dağıtır. Bu, iş yükünü dengelemek ve paralel işlemeyi sağlamak için kullanılır.

Direct Exchange Kullanımı

İş kuyrukları genellikle direct exchange veya varsayılan (default) exchange kullanılarak uygulanır. Varsayılan exchange, boş bir isimle temsil edilir ve gönderilen mesajın routing key‘ini doğrudan aynı isimdeki kuyruğa yönlendirir. Ancak, daha belirgin olmak için bir direct exchange kullanmak daha iyidir.

Producer (Gönderici) Kodu:

# producer_work_queue.py
import puka
from twisted.internet import reactor, defer
import time

RABBITMQ_URL = "amqp://guest:guest@localhost:5672/"
EXCHANGE_NAME = 'work_queue_exchange'
QUEUE_NAME = 'task_queue'
ROUTING_KEY = 'tasks'

@defer.inlineCallbacks
def send_tasks():
    client = puka.Client(RABBITMQ_URL)
    yield client.connect()
    print("Producer (Work Queue): Bağlantı başarılı.")

    yield client.exchange_declare(exchange=EXCHANGE_NAME, type='direct')
    yield client.queue_declare(queue=QUEUE_NAME, durable=True) # Kuyruğu kalıcı yap
    yield client.queue_bind(exchange=EXCHANGE_NAME, queue=QUEUE_NAME, routing_key=ROUTING_KEY)
    print(f"Producer: Exchange '{EXCHANGE_NAME}', Kuyruk '{QUEUE_NAME}' ve Binding '{ROUTING_KEY}' hazır.")

    for i in range(1, 11):
        message_body = f"Task {i}: Process this data."
        # Mesajı kalıcı yap
        yield client.basic_publish(exchange=EXCHANGE_NAME, routing_key=ROUTING_KEY, body=message_body,
                                   properties={'delivery_mode': 2}) # delivery_mode=2 kalıcı mesaj demektir
        print(f"Producer: Görev gönderildi: '{message_body}'")
        time.sleep(0.5)

    print("Producer: Tüm görevler gönderildi.")
    yield client.close()
    reactor.stop()

if __name__ == "__main__":
    reactor.callWhenRunning(send_tasks)
    reactor.run()

Burada durable=True ile kuyruğu ve delivery_mode=2 ile mesajları kalıcı hale getirdik. Bu, RabbitMQ sunucusu yeniden başlatıldığında bile kuyruğun ve içindeki mesajların kaybolmamasını sağlar.

Consumer (Tüketici) Kodu:

# consumer_work_queue.py
import puka
from twisted.internet import reactor, defer
import time

RABBITMQ_URL = "amqp://guest:guest@localhost:5672/"
QUEUE_NAME = 'task_queue'

@defer.inlineCallbacks
def process_tasks():
    client = puka.Client(RABBITMQ_URL)
    yield client.connect()
    print("Consumer (Work Queue): Bağlantı başarılı.")

    yield client.queue_declare(queue=QUEUE_NAME, durable=True)
    print(f"Consumer: Kuyruk '{QUEUE_NAME}' hazır. Görevler bekleniyor...")

    # QoS (Quality of Service) ayarı: Her tüketiciye aynı anda sadece 1 mesaj gönder
    yield client.basic_qos(prefetch_count=1)

    consume_result = yield client.basic_consume(queue=QUEUE_NAME, no_ack=False)

    try:
        while True:
            message = yield consume_result.wait()
            task_id = message['body'].split(':')[0]
            print(f"Consumer: Görev alındı: '{message['body']}'")
            # Görev işleme simülasyonu
            processing_time = len(message['body']) % 4 + 1 # Görev uzunluğuna göre rastgele süre
            print(f"Consumer: '{task_id}' işleniyor... ({processing_time} saniye)")
            time.sleep(processing_time)
            
            yield client.basic_ack(message['delivery_tag'])
            print(f"Consumer: Görev '{task_id}' tamamlandı ve onaylandı.")
    except KeyboardInterrupt:
        print("Consumer: Kapatılıyor...")
    finally:
        yield client.close()
        reactor.stop()

if __name__ == "__main__":
    reactor.callWhenRunning(process_tasks)
    reactor.run()

basic_qos(prefetch_count=1) ayarı çok önemlidir. Bu ayar, RabbitMQ’ya bir tüketiciye aynı anda sadece bir mesaj göndermesini söyler. Tüketici bu mesajı onaylamadan (ACK) önce başka bir mesaj almaz. Bu, tüketiciler arasında iş yükünün adil bir şekilde dağıtılmasını ve bir tüketicinin diğerlerinden daha yavaş çalışması durumunda bile mesajların yığılmamasını sağlar. Birden fazla consumer_work_queue.py örneği çalıştırdığınızda, görevlerin aralarında paylaşıldığını göreceksiniz.

Senaryo 2: Mesajları Tüm Tüketicilere Yayınlama (Publish/Subscribe)

Bu senaryoda, bir üretici tarafından gönderilen mesajın, o exchange’e bağlı olan tüm aktif tüketicilere ulaşması istenir. Bu, bir olay yayınlama (event broadcasting) veya bildirim gönderme gibi durumlarda kullanılır.

Fanout Exchange Kullanımı

fanout exchange, bu tür yayınlama senaryoları için tasarlanmıştır. Bir fanout exchange’e gönderilen her mesaj, kendisine bağlı olan tüm kuyruklara kopyalanır ve bu kuyrukları dinleyen tüketicilere iletilir. routing key bu exchange türünde göz ardı edilir.

Geçici Kuyruklar ve Otomatik Silme

Publish/Subscribe modelinde, her tüketicinin genellikle kendi geçici kuyruğuna sahip olması istenir. Bu kuyruklar, tüketici bağlantısı kesildiğinde otomatik olarak silinmelidir. Puka’da queue_declare(exclusive=True, auto_delete=True) ile bu sağlanır. exclusive=True kuyruğu sadece bildiren bağlantının kullanabileceği anlamına gelirken, auto_delete=True son tüketici bağlantısı kesildiğinde kuyruğun otomatik olarak silinmesini sağlar.

Producer (Gönderici) Kodu:

# producer_fanout.py
import puka
from twisted.internet import reactor, defer
import time

RABBITMQ_URL = "amqp://guest:guest@localhost:5672/"
EXCHANGE_NAME = 'weather_updates'

@defer.inlineCallbacks
def send_weather_updates():
    client = puka.Client(RABBITMQ_URL)
    yield client.connect()
    print("Producer (Fanout): Bağlantı başarılı.")

    yield client.exchange_declare(exchange=EXCHANGE_NAME, type='fanout')
    print(f"Producer: Exchange '{EXCHANGE_NAME}' bildirildi.")

    updates = ["Güneşli ve sıcak", "Hafif yağmur bekleniyor", "Rüzgarlı ve soğuk", "Kar yağışı ihtimali"]
    for i, update in enumerate(updates):
        message_body = f"Hava Durumu Güncellemesi {i+1}: {update}"
        yield client.basic_publish(exchange=EXCHANGE_NAME, routing_key='', body=message_body)
        print(f"Producer: Hava durumu güncellendi: '{message_body}'")
        time.sleep(1)

    print("Producer: Tüm güncellemeler gönderildi.")
    yield client.close()
    reactor.stop()

if __name__ == "__main__":
    reactor.callWhenRunning(send_weather_updates)
    reactor.run()

fanout exchange kullandığımız için routing_key boş bırakılmıştır.

Consumer (Tüketici) Kodu:

# consumer_fanout.py
import puka
from twisted.internet import reactor, defer
import time

RABBITMQ_URL = "amqp://guest:guest@localhost:5672/"
EXCHANGE_NAME = 'weather_updates'

@defer.inlineCallbacks
def receive_weather_updates():
    client = puka.Client(RABBITMQ_URL)
    yield client.connect()
    print("Consumer (Fanout): Bağlantı başarılı.")

    yield client.exchange_declare(exchange=EXCHANGE_NAME, type='fanout')
    
    # Geçici, exclusive ve auto-delete kuyruk oluştur
    result = yield client.queue_declare(exclusive=True, auto_delete=True)
    queue_name = result['queue']
    print(f"Consumer: Geçici kuyruk oluşturuldu: '{queue_name}'")

    yield client.queue_bind(exchange=EXCHANGE_NAME, queue=queue_name, routing_key='')
    print(f"Consumer: Exchange '{EXCHANGE_NAME}' ile kuyruk '{queue_name}' bağlandı.")
    print("Consumer: Hava durumu güncellemeleri bekleniyor. Çıkmak için Ctrl+C.")

    consume_result = yield client.basic_consume(queue=queue_name, no_ack=False)

    try:
        while True:
            message = yield consume_result.wait()
            print(f"Consumer: Hava durumu alındı: '{message['body']}'")
            yield client.basic_ack(message['delivery_tag'])
    except KeyboardInterrupt:
        print("Consumer: Kapatılıyor...")
    finally:
        yield client.close()
        reactor.stop()

if __name__ == "__main__":
    reactor.callWhenRunning(receive_weather_updates)
    reactor.run()

Birden fazla consumer_fanout.py örneği çalıştırdığınızda, her bir tüketicinin aynı hava durumu güncellemelerini aldığını göreceksiniz. Her tüketici kendi özel geçici kuyruğuna sahip olduğundan, mesajlar her birine ayrı ayrı teslim edilir.

Gelişmiş Konular ve Güvenilirlik

Mesajlaşma sistemlerinde güvenilirlik, mesajların kaybolmamasını ve doğru bir şekilde işlenmesini sağlamak için kritik öneme sahiptir.

Mesaj Kalıcılığı (Message Durability)

RabbitMQ, varsayılan olarak mesajları ve kuyrukları hafızada tutar. Bu, sunucu yeniden başlatıldığında mesajların veya kuyrukların kaybolabileceği anlamına gelir. Kalıcılığı sağlamak için:
* Kuyrukları Kalıcı Yapma: client.queue_declare(queue='my_queue', durable=True)
* Mesajları Kalıcı Yapma: client.basic_publish(..., properties={'delivery_mode': 2}). delivery_mode=2 mesajın diske yazılmasını sağlar.

Bu ayarlar, RabbitMQ sunucusu beklenmedik bir şekilde kapandığında bile mesajların ve kuyrukların korunmasına yardımcı olur.

Publisher Confirms

Üreticinin bir mesajı RabbitMQ’ya gönderdiğinden ve RabbitMQ’nun mesajı başarıyla kabul ettiğinden emin olmanın yolu publisher confirms kullanmaktır. Bu özellik etkinleştirildiğinde, RabbitMQ mesajı kabul ettiğinde üreticiye bir onay (ACK) gönderir. Eğer bir sorun olursa (örneğin, mesaj bir exchange’e yönlendirilemezse), bir NACK gönderilir. Puka’da client.confirm_select() ile bu modu etkinleştirebilir ve basic_publish metodunun döndürdüğü Deferred nesnesini bekleyerek onayı alabilirsiniz.

# Publisher Confirms örneği
@defer.inlineCallbacks
def send_with_confirm():
    client = puka.Client(RABBITMQ_URL)
    yield client.connect()
    yield client.confirm_select() # Publisher Confirms modunu etkinleştir

    exchange_name = 'my_direct_exchange'
    routing_key = 'test'
    yield client.exchange_declare(exchange=exchange_name, type='direct')

    message_body = "Mesaj onaylanacak mı?"
    confirm_deferred = yield client.basic_publish(exchange=exchange_name, routing_key=routing_key, body=message_body)
    
    try:
        # Onay bekleyelim
        confirm_result = yield confirm_deferred
        if confirm_result.get('ack'):
            print(f"Mesaj başarıyla onaylandı: '{message_body}'")
        else:
            print(f"Mesaj onaylanmadı (NACK): '{message_body}'")
    except Exception as e:
        print(f"Mesaj onayı alınamadı: {e}")
    finally:
        yield client.close()
        reactor.stop()

Tüketici Tarafında Hata Yönetimi: NACK ve Requeue

Bir tüketici mesajı işlerken bir hata ile karşılaşırsa, mesajı onaylamak (ACK) yerine reddedebilir (NACK – Negative Acknowledge). basic_nack metodu ile bu yapılır. requeue=True parametresi ile mesajın kuyruğa geri gönderilerek başka bir tüketici tarafından tekrar işlenmesi sağlanabilir.

# Consumer'da NACK örneği

... (bağlantı ve kuyruk bildirimleri)

consume_result = yield client.basic_consume(queue=queue_name, no_ack=False) try: while True: message = yield consume_result.wait() try: # Hata oluşabilecek bir işlem if "hata" in message['body']: raise ValueError("Mesajda hata var!") print(f"Mesaj işleniyor: '{message['body']}'") # Başarılı işlem sonrası ACK yield client.basic_ack(message['delivery_tag']) print(f"Mesaj onaylandı: {message['delivery_tag']}") except ValueError as e: print(f"Hata oluştu: {e}. Mesaj yeniden kuyruğa gönderiliyor: '{message['body']}'") # Hata durumunda NACK ve yeniden kuyruğa gönderme yield client.basic_nack(message['delivery_tag'], requeue=True) except KeyboardInterrupt: print("Kapatılıyor...") finally: yield client.close() reactor.stop()

basic_reject metodu da benzer şekilde mesajı reddetmek için kullanılır, ancak basic_nack‘ten farklı olarak toplu NACK yapamaz.

Önceden Getirme Sayısı (Prefetch Count / QoS)

Daha önce de bahsettiğimiz basic_qos(prefetch_count=N) ayarı, bir tüketicinin aynı anda kaç mesajı işleyebileceğini belirtir. Bu, iş yükünü adil bir şekilde dağıtmak ve yavaş tüketicilerin mesaj yığılmasına neden olmasını önlemek için çok önemlidir. prefetch_count=1 genellikle en güvenli ayardır, ancak daha yüksek değerler belirli senaryolarda performansı artırabilir (eğer tüketiciler mesajları çok hızlı işleyebiliyorsa).

Performans ve Ölçeklenebilirlik İpuçları

RabbitMQ ve Puka tabanlı bir sistemin performansını ve ölçeklenebilirliğini artırmak için bazı stratejiler:

Tüketici Sayısını Artırma

İş yükünü paylaştırma senaryosunda, işlem gücünü artırmak için aynı tüketici uygulamasından birden fazla örnek çalıştırabilirsiniz. RabbitMQ, mesajları bu tüketiciler arasında otomatik olarak dağıtacaktır.

Mesaj Boyutları ve Serileştirme

Küçük mesajlar, genellikle büyük mesajlardan daha hızlı işlenir. Eğer çok büyük veriler göndermeniz gerekiyorsa, bu verileri bir depolama servisine (örn. S3, MinIO) kaydedip, mesaj kuyruğunda sadece bu verinin referansını (URL veya ID) göndermek daha verimli olabilir. Mesaj içeriği genellikle JSON, Protobuf veya MessagePack gibi hafif serileştirme formatlarında olmalıdır.

Exchange ve Kuyruk Tasarımı

* Doğru Exchange Türünü Seçin: İhtiyacınıza uygun exchange türünü seçmek (direct, fanout, topic) performansı ve yönlendirme mantığını optimize eder.
* Kuyruk Sayısı: Çok fazla küçük kuyruk yerine, benzer mesajları işleyen daha az sayıda kuyruk kullanmak daha verimli olabilir.
* Kuyruk Parametreleri: durable, exclusive, auto_delete, message-ttl gibi kuyruk parametrelerini dikkatli kullanın. Özellikle message-ttl (mesaj ömrü) ve dead-letter-exchange (işlenemeyen mesajlar için) gibi özellikler, hata yönetimi ve sistem dayanıklılığı için önemlidir.

Puka’nın Geleceği ve Alternatifler

Puka’nın Durumu

Puka, Twisted ekosistemi içinde asenkron RabbitMQ istemcisi olarak değerli bir rol oynamıştır. Ancak, projenin aktif geliştirme ve bakımının Pika veya aio-pika gibi diğer kütüphanelere kıyasla daha yavaş olduğu gözlemlenmektedir. Bu, yeni Python sürümleriyle veya AMQP protokolündeki değişikliklerle uyumluluk sorunlarına yol açabilir. Proje seçimi yaparken bu faktörü göz önünde bulundurmak önemlidir.

Alternatif Python Kütüphaneleri

* Pika: Python için en popüler RabbitMQ istemcisidir. Hem senkron hem de asenkron (asyncio ile uyumlu) bağlantı seçenekleri sunar. Geniş topluluk desteği ve aktif geliştirme sürecine sahiptir.
* aio-pika: Python’ın yerleşik asyncio kütüphanesi üzerine inşa edilmiş, tamamen asenkron bir RabbitMQ istemcisidir. Modern Python asenkron programlama paradigmalarına daha uygun bir yaklaşımdır.

Eğer yeni bir projeye başlıyorsanız ve Twisted kullanmıyorsanız, Pika veya aio-pika genellikle daha güncel ve desteklenen seçenekler olarak öne çıkar. Ancak, mevcut bir Twisted projesindeyseniz, Puka hala entegrasyon kolaylığı sağlayabilir.

Sonuç

RabbitMQ ve Python’ın Puka kütüphanesi, dağıtık sistemlerde güvenilir ve ölçeklenebilir mesaj iletimi için güçlü bir kombinasyon sunar. Mesaj kuyrukları, uygulama bileşenleri arasındaki bağımlılığı azaltarak, sistemin esnekliğini ve dayanıklılığını artırır. Puka’nın asenkron yapısı ve Twisted entegrasyonu, özellikle yüksek eşzamanlılık gerektiren Python uygulamalarında mesajlaşma yeteneklerini güçlendirir.

Bu makalede, RabbitMQ’nun temel kavramlarından (Exchange, Kuyruk, Binding) Puka ile bağlantı kurmaya, mesaj göndermeye ve almaya kadar adım adım ilerledik. Özellikle birden fazla tüketiciye mesaj iletmek için kullanılan iki ana senaryoyu (iş yükünü paylaşma için direct exchange ve yayınlama için fanout exchange) detaylı kod örnekleriyle inceledik. Mesaj kalıcılığı, publisher confirms ve prefetch count gibi güvenilirlik ve performans ipuçlarına da değindik.

Her ne kadar Puka’nın aktif geliştirme durumu diğer kütüphanelere göre biraz daha yavaş olsa da, Twisted ekosistemi içinde hala geçerli bir seçenektir. Mesaj kuyruklarının gücünü anlamak ve doğru araçlarla uygulamak, modern, dağıtık sistemlerin başarılı bir şekilde inşa edilmesinde temel bir adımdır. Bu bilgiler ışığında, projelerinizde RabbitMQ ve Puka’yı kullanarak sağlam ve ölçeklenebilir iletişim altyapıları oluşturabilirsiniz.

Yorumlar
İçeriği beğendiniz mi? Bir tartışma başlatın veya görüşlerinizi paylaşın.
Yorum Yaz

Bir yanıt yazın

E-posta adresiniz yayınlanmayacak. Gerekli alanlar * ile işaretlenmişlerdir

E-posta Bülteni
Yazılım Topluluğuna Katılın
En son güncellemeleri, yaratıcı ipuçlarını ve özel kaynakları doğrudan e-posta kutunuza alın. Tasarım ve inovasyonun geleceğini birlikte keşfedelim.