Kafka balanced consumer wrapper.
All consumers are created with the same group.id, after each test consumer commits offsets for all consumed messages. This is needed to make tests independent.
Definition at line 155 of file classes.py.
Public Member Functions | |
| __init__ (self, bool enabled, bootstrap_servers) | |
| start (self) | |
| ConsumedMessage | receive_one (self, list[str] topics, float timeout=20.0) |
| Waits until one message are consumed. | |
| list[ConsumedMessage] | receive_batch (self, list[str] topics, int|None max_batch_size, float timeout=3.0) |
| Waits until either max_batch_size messages are consumed or timeout expired. | |
| aclose (self) | |
Public Attributes | |
| consumer | |
Protected Member Functions | |
| _subscribe (self, list[str] topics) | |
| _commit (self) | |
| _unsubscribe (self) | |
Protected Attributes | |
| _enabled = enabled | |
| _bootstrap_servers = bootstrap_servers | |
| list | _subscribed_topics = [] |
| testsuite.databases.kafka.classes.KafkaConsumer.__init__ | ( | self, | |
| bool | enabled, | ||
| bootstrap_servers ) |
Definition at line 163 of file classes.py.
|
protected |
Definition at line 192 of file classes.py.
|
protected |
Definition at line 178 of file classes.py.
|
protected |
Definition at line 199 of file classes.py.
| testsuite.databases.kafka.classes.KafkaConsumer.aclose | ( | self | ) |
Definition at line 258 of file classes.py.
| list[ConsumedMessage] testsuite.databases.kafka.classes.KafkaConsumer.receive_batch | ( | self, | |
| list[str] | topics, | ||
| int | None | max_batch_size, | ||
| float | timeout = 3.0 ) |
Waits until either max_batch_size messages are consumed or timeout expired.
:param topics: list of topics to read messages from. :max_batch_size: maximum number of consumed messages. :param timeout: timeout to stop waiting. Default is 3 seconds.
:returns: :py:class:List[ConsumedMessage]
Definition at line 229 of file classes.py.
| ConsumedMessage testsuite.databases.kafka.classes.KafkaConsumer.receive_one | ( | self, | |
| list[str] | topics, | ||
| float | timeout = 20.0 ) |
Waits until one message are consumed.
:param topics: list of topics to read messages from. :param timeout: timeout to stop waiting. Default is 20 seconds.
:returns: :py:class:ConsumedMessage
Definition at line 207 of file classes.py.
| testsuite.databases.kafka.classes.KafkaConsumer.start | ( | self | ) |
Definition at line 168 of file classes.py.
|
protected |
Definition at line 165 of file classes.py.
|
protected |
Definition at line 164 of file classes.py.
|
protected |
Definition at line 166 of file classes.py.
| testsuite.databases.kafka.classes.KafkaConsumer.consumer |
Definition at line 170 of file classes.py.