Kafka integration with asyncio.
Project description
asyncio client for kafka
AIOKafkaProducer
AIOKafkaProducer is a high-level, asynchronous message producer.
Example of AIOKafkaProducer usage:
import asyncio
from aiokafka import AIOKafkaProducer
@asyncio.coroutine
def produce(loop):
# Just adds message to sending queue
future = yield from producer.send('foobar', b'some_message_bytes')
# waiting for message to be delivered
resp = yield from future
print("Message produced: partition {}; offset {}".format(
resp.partition, resp.offset))
# Also can use a helper to send and wait in 1 call
resp = yield from producer.send_and_wait(
'foobar', key=b'foo', value=b'bar')
resp = yield from producer.send_and_wait(
'foobar', b'message for partition 1', partition=1)
loop = asyncio.get_event_loop()
producer = AIOKafkaProducer(loop=loop, bootstrap_servers='localhost:9092')
loop.run_until_complete(producer.start())
loop.run_until_complete(produce(loop))
loop.run_until_complete(producer.stop())
loop.close()
AIOKafkaConsumer
AIOKafkaConsumer is a high-level, asynchronous message consumer. It interacts with the assigned kafka Group Coordinator node to allow multiple consumers to load balance consumption of topics (requires kafka >= 0.9.0.0).
Example of AIOKafkaConsumer usage:
import asyncio
from kafka.common import KafkaError
from aiokafka import AIOKafkaConsumer
@asyncio.coroutine
def consume_task(consumer):
while True:
try:
msg = yield from consumer.getone()
print("consumed: ", msg.topic, msg.partition, msg.offset, msg.value)
except KafkaError as err:
print("error while consuming message: ", err)
loop = asyncio.get_event_loop()
consumer = AIOKafkaConsumer(
'topic1', 'topic2', loop=loop, bootstrap_servers='localhost:1234')
loop.run_until_complete(consumer.start())
c_task = loop.create_task(consume_task(consumer))
try:
loop.run_forever()
finally:
loop.run_until_complete(consumer.stop())
c_task.cancel()
loop.close()
Running tests
Docker is required to run tests. See https://docs.docker.com/engine/installation for installation notes.
Setting up tests requirements (assuming you’re within virtualenv on ubuntu 14.04+):
sudo apt-get install -y libsnappy-dev pip install flake8 pytest pytest-cov pytest-catchlog docker-py python-snappy coveralls .
Running tests:
make cov
To run tests with a specific version of Kafka (default one is 0.9.0.1) use KAFKA_VERSION variable:
make cov KAFKA_VERSION=0.8.2.1
CHANGES
0.1.2 (2016-04-30)
Added Python3.5 usage example to docs
Don’t raise retriable exceptions in 3.5’s async for iterator
Fix Cancellation issue with producer’s send_and_wait method
0.1.1 (2016-04-15)
Fix packaging issues. Removed unneded files from package.
0.1.0 (2016-04-15)
Initial release
Added full support for Kafka 9.0. Older Kafka versions are not tested.
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
File details
Details for the file aiokafka-0.1.2.tar.gz.
File metadata
- Download URL: aiokafka-0.1.2.tar.gz
- Upload date:
- Size: 39.6 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
8f9c544d85be2a83bb8f2b52b93cf133778eea667e0cf25f6074b0a344e2f7d8
|
|
| MD5 |
bb61352a635b39ae984857f65e468ad3
|
|
| BLAKE2b-256 |
fc35bed3eb04190421a9398dc7946a0077b1600e6e569045d96e6ffe39909630
|