Skip to content

Commit b6a691b

Browse files
committed
Finish 0.1.2
2 parents 0aa9cdf + 2283921 commit b6a691b

File tree

4 files changed

+3
-15
lines changed

4 files changed

+3
-15
lines changed

README.md

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -38,7 +38,5 @@ AioKafkaBroker parameters:
3838
* `kafka_topic` - custom topic in kafka.
3939
* `result_backend` - custom result backend.
4040
* `task_id_generator` - custom task_id genertaor.
41-
* `aiokafka_producer` - custom `aiokafka` producer.
42-
* `aiokafka_consumer` - custom `aiokafka` consumer.
4341
* `kafka_admin_client` - custom `kafka` admin client.
4442
* `delete_topic_on_shutdown` - flag to delete topic on broker shutdown.

pyproject.toml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@ name = "taskiq-aio-kafka"
33
description = "Kafka broker for taskiq"
44
authors = ["Taskiq team <[email protected]>"]
55
maintainers = ["Taskiq team <[email protected]>"]
6-
version = "0.1.1"
6+
version = "0.1.2"
77
readme = "README.md"
88
license = "LICENSE"
99
classifiers = [

taskiq_aio_kafka/broker.py

Lines changed: 2 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -46,8 +46,6 @@ def __init__( # noqa: WPS211
4646
kafka_topic: Optional[NewTopic] = None,
4747
result_backend: Optional[AsyncResultBackend[_T]] = None,
4848
task_id_generator: Optional[Callable[[], str]] = None,
49-
aiokafka_producer: Optional[AIOKafkaProducer] = None,
50-
aiokafka_consumer: Optional[AIOKafkaConsumer] = None,
5149
kafka_admin_client: Optional[KafkaAdminClient] = None,
5250
loop: Optional[asyncio.AbstractEventLoop] = None,
5351
delete_topic_on_shutdown: bool = False,
@@ -58,8 +56,6 @@ def __init__( # noqa: WPS211
5856
:param kafka_topic: kafka topic.
5957
:param result_backend: custom result backend.
6058
:param task_id_generator: custom task_id generator.
61-
:param aiokafka_producer: configured AIOKafkaProducer.
62-
:param aiokafka_consumer: configured AIOKafkaConsumer.
6359
:param kafka_admin_client: configured KafkaAdminClient.
6460
:param loop: specific even loop.
6561
:param delete_topic_on_shutdown: delete or don't delete topic on shutdown.
@@ -69,10 +65,10 @@ def __init__( # noqa: WPS211
6965
"""
7066
super().__init__(result_backend, task_id_generator)
7167

72-
if (aiokafka_producer or aiokafka_consumer) and not bootstrap_servers:
68+
if kafka_admin_client and not bootstrap_servers:
7369
raise WrongAioKafkaBrokerParametersError(
7470
(
75-
"If you specify `aiokafka_producer` and/or `aiokafka_consumer`, "
71+
"If you specify `kafka_admin_client`, "
7672
"you must specify `bootstrap_servers`."
7773
),
7874
)

tests/conftest.py

Lines changed: 0 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -114,8 +114,6 @@ async def broker_without_arguments(
114114
@pytest.fixture()
115115
async def broker(
116116
kafka_url: str,
117-
test_kafka_producer: AIOKafkaProducer,
118-
test_kafka_consumer: AIOKafkaConsumer,
119117
base_topic: NewTopic,
120118
) -> AsyncGenerator[AioKafkaBroker, None]:
121119
"""Yield new broker instance.
@@ -125,16 +123,12 @@ async def broker(
125123
and shutdown after test.
126124
127125
:param kafka_url: url to kafka.
128-
:param test_kafka_producer: custom AIOKafkaProducer.
129-
:param test_kafka_consumer: custom AIOKafkaConsumer.
130126
:param base_topic: base topic.
131127
132128
:yields: broker.
133129
"""
134130
broker = AioKafkaBroker(
135131
bootstrap_servers=kafka_url,
136-
aiokafka_producer=test_kafka_producer,
137-
aiokafka_consumer=test_kafka_consumer,
138132
kafka_topic=base_topic,
139133
delete_topic_on_shutdown=True,
140134
)

0 commit comments

Comments
 (0)