- Raise an error when trying to use a consumer after it has been closed
- Add
Consumer#stopmethod to stop aneachloop - Add crystal versions 1.15.1 and 1.16.3 to test matrix
- Use ::sleep(Time::Span) instead of ::sleep(Number) to fix deprecations warnings with Crystal >= 1.14
- Add
Kafka.version_infoandKafka.librdkafka_versionmethods.
- Update
Kafka::Consumer#pollandKafka::Consumer#eachto automatically raise aKafka::ConsumerExceptionif the message is an error. Passraise_on_error: falseto maintain the previous behaviour.
- Add
#topicmethod toKafka::Message. Thanks @oozzal!
- Upgrade Crystal version to v1.11.2
- Fix to prevent exception when Delivery Report string is null pointer
- Call
rd_kafka_pollautomatically inKafka::Producer
- Rename main src file
- Save statistics option on
Kafka::Producer
- Remove topic + partition name from rebalance log
- Integration & unit tests
- Documentation for all key methods and examples in the README
- Format all files using
crystal tool format - Add
Fiber.yieldat the start of each loop inKafka::Consumer#eachto allow other Fibers to run in between each iteration - Fix
Invalid memory accesserror and raise exception when unknown or invalid config passed toKafka::Consumer.new - Fix
Invalid memory accesserror and raise exception when unknown or invalid config passed toKafka::Producer.new - Fix
Invalid memory accesserror and raise exception whenLibRdKafka.kafka_newfails to create consumer - Fix
Invalid memory accesserror and raise exception whenKafka::Consumer#subscribefails to subscribe to topics - Call
rd_kafka_destroy()after closing consumer as advised in the librdkafka documentation
- Refactor setting rebalancing callback into separate class
- Refactor building of config for Producer/Consumer into separate class
- Improve logging around consumer partition assignment and producer delivery reports
- Forked from https://github.com/CloudKarafka/kafka.cr
- Added new
Kafka::Producer#producemethod without key argument