Merhabalar,
Büyük ölçekli projelerde hayat kurtaran Kafka’yı Python ile nasıl kullanabileceğimizi bu yazımda göreceğiz;
Apache Kafka, dağıtık olarak veri akışı (streaming) sağlayabilmemize imkan tanıyan, Bu veri akış sistemini yatayda ölçeklendirebildiğimiz, veri iletiminde düşük gecikme ile gerçek zamanlı veri-alışverişi sağlayan teknolojidir.
Bu teknolojinin mimarisine, nasıl çalıştığına dair Youtube’da Barış Dere’nin güzel bir eğitim serisine denk, geldim. Eğer sadece adını duyduğunuz bu teknoloji ile derinlemesine tanışmak isterseniz, video serisini dikkatle takip etmenizi öneririm:
Lokal Ortamda Kafka ve Zookeper Ortamlarını Hazırlama
Ancak bu makalenin konusu, Python kullanarak nasıl consumer ve procuder oluşturacağımız. Öncelikle test ortamını oluşturabilmek için, Docker’ın nimetlerinden faydalanarak başlıyoruz, compose dosyasını şuradaki repodan bilgisayarınıza klonlayabilirsiniz.
Bizim işimize yarayacak konfigürasyona sahip docker-compose.yml dosyamız ise şöyle:
version: "3"
services:
zookeeper:
image: 'bitnami/zookeeper:latest'
ports:
- '2181:2181'
environment:
- ALLOW_ANONYMOUS_LOGIN=yes
kafka:
image: 'bitnami/kafka:latest'
ports:
- '9092:9092'
environment:
- KAFKA_BROKER_ID=1
- KAFKA_CFG_LISTENERS=PLAINTEXT://:9092
- KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://127.0.0.1:9092
- KAFKA_CFG_ZOOKEEPER_CONNECT=zookeeper:2181
- ALLOW_PLAINTEXT_LISTENER=yes
depends_on:
- zookeeper
docker-compose up -d
İşlemler başarıyla tamamlandığında aşağıdakine benzer bir çıktı alabiliyor olmalısınız;
~ docker ps --format "{{.Names}}: {{.Command}}"
kafkamq_kafka_1: "/opt/bitnami/script…"
kafkamq_zookeeper_1: "/opt/bitnami/script…"
Artık 1 brokar ve 1 zookeper’ımız ile Python tarafına başlayabiliriz.
Yaygın olarak kullanılan 3 Kafka kütüphanesi mevcut ve bunların birbirlerine göre avantajları ve dezavanjları var tabii.
Kafka için Yaygın Python Kütüphaneleri
kafka-python
kafka-python, C uzantısı olmayan saf bir Python kütüphanesidir. API, iyi tasarlanmış ve yeni başlayanlar için kullanımı kolaydır. Aynı zamanda aktif olarak geliştirilmeye devam eden bir projedir. Python-kafka’nın dezavantajı hızıdır. Performansa önem veriyorsanız, confluent-kafka’ya geçmenizi tavsiye ederim.
Confluent-kafka
Confluent-kafka, 3 kütüphane arasında hiç şüphesiz en iyi performansa sahip. API iyi tasarlanmış ve parametreler orijinal Apache Kafka ile aynı adı ve aynı varsayılan parametreler ile kullanılıyor. Bu sayede kodlarınızı orijinal fonksiyon isimlerine kolayca bağlayabilirsiniz. Şahsen, kullanıcı davranışını özelleştirme esnekliğini seviyorum. Ayrıca Confluent tarafından aktif olarak geliştirilmekte ve desteklenmektedir.
Bir dezavantaj, Windows kullanıcılarının çalışması için biraz uğraşması gerekebilmesidir.
pykafka
Kafka-python ve conflunet-kafka ile karşılaştırıldığında, pykafka’nın gelişimi daha az aktiftir. Github repo geçmişi, Kasım 2018’den bu yana güncellenmediğini gösteriyor. Ayrıca, pykafka’nın farklı API tasarımları var ve ilk kez basit olmayabilecek farklı varsayılan parametreler kullanıyor.
Kafka-python Kullanarak Consumer ve Producer Oluşturma
Bu anlatımda Kafka-python üzerinden ilerlemeyi en yaygın olan kütüphane olması sebebiyle tercih ettim.
Producer.py
# Created by Sezer BOZKIR<admin@sezerbozkir.com> at 1.06.2021
from kafka import KafkaProducer
from time import sleep
def kafka_producer_sync(procuder: KafkaProducer, msg, size):
for _ in range(size):
print("Mesaj gonderimi basladi.")
future = procuder.send("sample_topic", msg)
result = future.get(timeout=60)
print(result)
procuder.flush()
sleep(1)
def success(metadata):
print(metadata.topic)
def error(exception):
print(exception)
def kafka_producer_async(producer: KafkaProducer, msg, size):
for _ in range(size):
producer.send("sample_topic", msg).add_callback(success).add_errback(error)
producer.flush()
sleep(5)
if __name__ == '__main__':
main_producer = KafkaProducer(bootstrap_servers="localhost:9092")
main_msg = ('Bu Ornek Bir Kafka Mesaj Gonderimidir' * 20).encode()[:100]
kafka_producer_sync(main_producer, main_msg, 1000000)
Kodu çalıştırdığımızda, herşey yolundaysa, aşağıdaki gibi bir çıktı bizi karşılıyor olmalı:
Senkron mesaj gonderimi basladi. RecordMetadata(topic='sample_topic', partition=0, topic_partition=TopicPartition(topic='sample_topic', partition=0), offset=112, timestamp=1622705225229, log_start_offset=0, checksum=None, serialized_key_size=-1, serialized_value_size=100, serialized_header_size=-1) RecordMetadata(topic='sample_topic', partition=0, topic_partition=TopicPartition(topic='sample_topic', partition=0), offset=113, timestamp=1622705226340, log_start_offset=0, checksum=None, serialized_key_size=-1, serialized_value_size=100, serialized_header_size=-1) RecordMetadata(topic='sample_topic', partition=0, topic_partition=TopicPartition(topic='sample_topic', partition=0), offset=120, timestamp=1622705233382, log_start_offset=0, checksum=None, serialized_key_size=-1, serialized_value_size=100, serialized_header_size=-1) RecordMetadata(topic='sample_topic', partition=0, topic_partition=TopicPartition(topic='sample_topic', partition=0), offset=121, timestamp=1622705234389, log_start_offset=0, checksum=None, serialized_key_size=-1, serialized_value_size=100, serialized_header_size=-1) Asenkron mesaj gonderimi basladi. sample_topic sample_topic sample_topic sample_topic sample_topic sample_topic sample_topic sample_topic sample_topic sample_topic
Kodun içerisinde, çıktıyı takip edebilmek için özellikle koyduğum “sleep(1)” satırı her saniye bir mesaj gönderimi için eklenmiş durumda. Bir projenizde test etmek için bu kısımları kaldırmayı unutmayın.
Consumer.py
# Created by Sezer BOZKIR<admin@sezerbozkir.com> at 1.06.2021
from kafka import KafkaConsumer
from pprint import pprint
if __name__ == '__main__':
consumer = KafkaConsumer('sample_topic', bootstrap_servers="localhost:9092",
enable_auto_commit=False, auto_offset_reset="earliest")
pprint(consumer.metrics())
for msg in consumer:
pprint(msg)
Consumer çıktısı:
{'consumer-coordinator-metrics': {'assigned-partitions': 0.0,
'commit-latency-avg': 0.0,
'commit-latency-max': -inf,
'commit-rate': 0.0,
'heartbeat-rate': 0.0,
'heartbeat-response-time-max': -inf,
'join-rate': 0.0,
'join-time-avg': 0.0,
'join-time-max': -inf,
'last-heartbeat-seconds-ago': inf,
'sync-rate': 0.0,
'sync-time-avg': 0.0,
'sync-time-max': -inf},
'consumer-fetch-manager-metrics': {'bytes-consumed-rate': 0.0,
'fetch-latency-avg': 0.0,
'fetch-latency-max': -inf,
'fetch-rate': 0.0,
'fetch-size-avg': 0.0,
'fetch-size-max': -inf,
'fetch-throttle-time-avg': 0.0,
'fetch-throttle-time-max': -inf,
'records-consumed-rate': 0.0,
'records-lag-max': -inf,
'records-per-request-avg': 0.0},
'consumer-metrics': {'connection-close-rate': 0.0,
'connection-count': 1.0,
'connection-creation-rate': 0.03321440047958481,
'incoming-byte-rate': 14.365008050665505,
'io-ratio': 0.0,
'io-time-ns-avg': 0.0,
'io-wait-ratio': 0.0,
'io-wait-time-ns-avg': 0.0,
'network-io-rate': 0.1328585144317845,
'outgoing-byte-rate': 2.258594159259351,
'request-latency-avg': 53.37703227996826,
'request-latency-max': 104.21919822692871,
'request-rate': 0.0664291403230769,
'request-size-avg': 34.0,
'request-size-max': 36.0,
'response-rate': 0.06665893644332424,
'select-rate': 0.0},
'consumer-node-metrics.node-bootstrap-0': {'incoming-byte-rate': 14.364983153300857,
'outgoing-byte-rate': 2.2585833534456787,
'request-latency-avg': 53.37703227996826,
'request-latency-max': 104.21919822692871,
'request-rate': 0.0664289065386803,
'request-size-avg': 34.0,
'request-size-max': 36.0,
'response-rate': 0.06665883935227039},
'kafka-metrics-count': {'count': 50.0}}
ConsumerRecord(topic='sample_topic', partition=0, offset=132, timestamp=1622705358347, timestamp_type=0, key=None, value=b'Bu Ornek Bir Kafka Mesaj GonderimidirBu Ornek Bir Kafka Mesaj GonderimidirBu Ornek Bir Kafka Mesaj G', headers=[], checksum=None, serialized_key_size=-1, serialized_value_size=100, serialized_header_size=-1)
Burada da öncelikle takip etmeye başladığınız topic’e dair bilgiler sonrasında da gönderdiğimiz mesajları görüntülüyoruz. Docker-compose loglarını dinlerseniz, aynı zamanda topic içeriğini buradan da görüntüleyebilirsiniz:
kafka_1 | [2021-06-03 07:18:35,444] INFO [broker-1-to-controller-send-thread]: Recorded new controller, from now on will use broker 127.0.0.1:9092 (id: 1 rack: null) (kafka.server.BrokerToControllerRequestThread)
kafka_1 | [2021-06-03 07:21:23,197] INFO Creating topic sample_topic with configuration {} and initial partition assignment Map(0 -> ArrayBuffer(1)) (kafka.zk.AdminZkClient)
kafka_1 | [2021-06-03 07:21:23,358] INFO [ReplicaFetcherManager on broker 1] Removed fetcher for partitions Set(sample_topic-0) (kafka.server.ReplicaFetcherManager)
kafka_1 | [2021-06-03 07:21:23,415] INFO [Log partition=sample_topic-0, dir=/bitnami/kafka/data] Loading producer state till offset 0 with message format version 2 (kafka.log.Log)
kafka_1 | [2021-06-03 07:21:23,418] INFO Created log for partition sample_topic-0 in /bitnami/kafka/data/sample_topic-0 with properties {compression.type -> producer, min.insync.replicas -> 1, message.downconversion.enable -> true, segment.jitter.ms -> 0, cleanup.policy -> [delete], flush.ms -> 9223372036854775807, retention.ms -> 604800000, segment.bytes -> 1073741824, flush.messages -> 9223372036854775807, message.format.version -> 2.8-IV1, max.compaction.lag.ms -> 9223372036854775807, file.delete.delay.ms -> 60000, max.message.bytes -> 1048588, min.compaction.lag.ms -> 0, message.timestamp.type -> CreateTime, preallocate -> false, index.interval.bytes -> 4096, min.cleanable.dirty.ratio -> 0.5, unclean.leader.election.enable -> false, retention.bytes -> -1, delete.retention.ms -> 86400000, segment.ms -> 604800000, message.timestamp.difference.max.ms -> 9223372036854775807, segment.index.bytes -> 10485760}. (kafka.log.LogManager)
Elbette buradan yola çıkarak gerçek projelerde çoklu broker’lar üzerinden, birden fazla zookeper ile projelere nasıl dahil olabileceğinize dair yolculuk edebilirsiniz. Makaledeki örneğin kodlarına GitHub sayfamdan erişebilirsiniz:
https://github.com/Natgho/apache-kafka-with-python
Bir başka yazımda görüşmek üzere…
Kaynaklar:
https://kafka-python.readthedocs.io/en/master/usage.html
https://stackoverflow.com/questions/64729122/how-does-apache-kafka-work-with-multiple-brokers-and-single-broker
https://towardsdatascience.com/kafka-python-explained-in-10-lines-of-code-800e3e07dad1