userver: testsuite.databases.kafka.classes.KafkaConsumer Class Reference
Loading...
Searching...
No Matches
testsuite.databases.kafka.classes.KafkaConsumer Class Reference

Detailed Description

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[ConsumedMessagereceive_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 = []

Constructor & Destructor Documentation

◆ __init__()

testsuite.databases.kafka.classes.KafkaConsumer.__init__ ( self,
bool enabled,
bootstrap_servers )

Definition at line 163 of file classes.py.

Member Function Documentation

◆ _commit()

testsuite.databases.kafka.classes.KafkaConsumer._commit ( self)
protected

Definition at line 192 of file classes.py.

◆ _subscribe()

testsuite.databases.kafka.classes.KafkaConsumer._subscribe ( self,
list[str] topics )
protected

Definition at line 178 of file classes.py.

◆ _unsubscribe()

testsuite.databases.kafka.classes.KafkaConsumer._unsubscribe ( self)
protected

Definition at line 199 of file classes.py.

◆ aclose()

testsuite.databases.kafka.classes.KafkaConsumer.aclose ( self)

Definition at line 258 of file classes.py.

◆ receive_batch()

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.

◆ receive_one()

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.

◆ start()

testsuite.databases.kafka.classes.KafkaConsumer.start ( self)

Definition at line 168 of file classes.py.

Member Data Documentation

◆ _bootstrap_servers

testsuite.databases.kafka.classes.KafkaConsumer._bootstrap_servers = bootstrap_servers
protected

Definition at line 165 of file classes.py.

◆ _enabled

testsuite.databases.kafka.classes.KafkaConsumer._enabled = enabled
protected

Definition at line 164 of file classes.py.

◆ _subscribed_topics

list testsuite.databases.kafka.classes.KafkaConsumer._subscribed_topics = []
protected

Definition at line 166 of file classes.py.

◆ consumer

testsuite.databases.kafka.classes.KafkaConsumer.consumer
Initial value:
= aiokafka.AIOKafkaConsumer(
group_id='Test-group',
bootstrap_servers=self._bootstrap_servers,
auto_offset_reset='earliest',
enable_auto_commit=False,
)

Definition at line 170 of file classes.py.


The documentation for this class was generated from the following file: