Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion kafka/admin/_cluster.py
Original file line number Diff line number Diff line change
Expand Up @@ -349,7 +349,8 @@ def update_features(self, feature_updates, validate_only=False, timeout_ms=60000
dict of {feature_name: 'OK' | error message}
"""
return self._manager.run(self._async_update_features,
feature_updates, validate_only, timeout_ms)
feature_updates, validate_only, timeout_ms,
timeout_ms=timeout_ms)


class UpdateFeatureType(EnumHelper, IntEnum):
Expand Down
10 changes: 7 additions & 3 deletions kafka/admin/_partitions.py
Original file line number Diff line number Diff line change
Expand Up @@ -212,7 +212,9 @@ def delete_records(self, records_to_delete, timeout_ms=None, partition_leader_id
Returns:
dict {topicPartition -> metadata}
"""
return self._manager.run(self._async_delete_records, records_to_delete, timeout_ms, partition_leader_id)
return self._manager.run(self._async_delete_records, records_to_delete,
timeout_ms, partition_leader_id,
timeout_ms=timeout_ms)

def _get_all_topic_partitions(self, topics=None):
return [
Expand Down Expand Up @@ -386,7 +388,8 @@ def list_partition_reassignments(self, topic_partitions=None, timeout_ms=None):
``'removing_replicas'`` (each a list of broker IDs).
"""
return self._manager.run(
self._async_list_partition_reassignments, topic_partitions, timeout_ms)
self._async_list_partition_reassignments, topic_partitions, timeout_ms,
timeout_ms=timeout_ms)

async def _async_describe_topic_partitions(self, topics, response_partition_limit, cursor):
_Topic = DescribeTopicPartitionsRequest.TopicRequest
Expand Down Expand Up @@ -573,7 +576,8 @@ def list_partition_offsets(self, topic_partition_specs, isolation_level=Isolatio
of ListOffsetsRequest compatible with the requested specs.
"""
return self._manager.run(
self._async_list_partition_offsets, topic_partition_specs, isolation_level, timeout_ms)
self._async_list_partition_offsets, topic_partition_specs, isolation_level, timeout_ms,
timeout_ms=timeout_ms)


class NewPartitions:
Expand Down
35 changes: 24 additions & 11 deletions kafka/consumer/fetcher.py
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,7 @@ class Fetcher:
'metrics': None,
'metric_group_prefix': 'consumer',
'request_timeout_ms': 30000,
'default_api_timeout_ms': 60000,
'retry_backoff_ms': 100,
'enable_incremental_fetch_sessions': True,
'isolation_level': 'read_uncommitted',
Expand Down Expand Up @@ -258,7 +259,7 @@ def _wake(_):
fut.add_both(_wake)

try:
self._net.run(self._manager.wait_for, wakeup, timeout_ms)
self._net.run(self._manager.wait_for, wakeup, timeout_ms, timeout_ms=timeout_ms)
except Errors.KafkaTimeoutError:
pass

Expand Down Expand Up @@ -395,7 +396,8 @@ def offsets_by_times(self, timestamps, timeout_ms=None):
timestamps: {TopicPartition: int} dict with timestamps to fetch
offsets by. -1 for the latest available, -2 for the earliest
available. Otherwise timestamp is treated as epoch milliseconds.
timeout_ms (int, optional): The maximum time in milliseconds to block.
timeout_ms (int, optional): Maximum time in milliseconds to block.
Defaults to default_api_timeout_ms.

Returns:
{TopicPartition: OffsetAndTimestamp or None}: Mapping of partition to
Expand All @@ -404,9 +406,12 @@ def offsets_by_times(self, timestamps, timeout_ms=None):
will be None.

Raises:
KafkaTimeoutError if timeout_ms provided
KafkaTimeoutError: if not completed within timeout_ms (default
default_api_timeout_ms).
"""
offsets = self._net.run(self._fetch_offsets_by_times_async, timestamps, timeout_ms)
if timeout_ms is None:
timeout_ms = self.config['default_api_timeout_ms']
offsets = self._net.run(self._fetch_offsets_by_times_async, timestamps, timeout_ms, timeout_ms=timeout_ms)
for tp in timestamps:
if tp not in offsets:
offsets[tp] = None
Expand Down Expand Up @@ -481,13 +486,15 @@ def beginning_offsets(self, partitions, timeout_ms=None):

Arguments:
partitions ([TopicPartition]): List of partitions for list offsets.
timeout_ms (int, optional): The maximum time in milliseconds to block.
timeout_ms (int, optional): Maximum time in milliseconds to block.
Defaults to default_api_timeout_ms.

Returns:
{TopicPartition: int}: Mapping of partition to retrieved offset.

Raises:
KafkaTimeoutError if timeout_ms provided.
KafkaTimeoutError: if not completed within timeout_ms (default
default_api_timeout_ms).
"""
return self.beginning_or_end_offset(
partitions, OffsetSpec.EARLIEST, timeout_ms)
Expand All @@ -500,13 +507,15 @@ def end_offsets(self, partitions, timeout_ms=None):

Arguments:
partitions ([TopicPartition]): List of partitions for list offsets.
timeout_ms (int, optional): The maximum time in milliseconds to block.
timeout_ms (int, optional): Maximum time in milliseconds to block.
Defaults to default_api_timeout_ms.

Returns:
{TopicPartition: int}: Mapping of partition to retrieved offset.

Raises:
KafkaTimeoutError if timeout_ms provided.
KafkaTimeoutError: if not completed within timeout_ms (default
default_api_timeout_ms).
"""
return self.beginning_or_end_offset(
partitions, OffsetSpec.LATEST, timeout_ms)
Expand All @@ -522,18 +531,22 @@ def beginning_or_end_offset(self, partitions, timestamp, timeout_ms=None):
timestamp (int or OffsetSpec): OffsetSpec.LATEST (-1) for the latest
available, OffsetSpec.EARLIEST (-2) for the earliest available.
Otherwise timestamp is treated as epoch milliseconds.
timeout_ms (int, optional): The maximum time in milliseconds to block.
timeout_ms (int, optional): Maximum time in milliseconds to block.
Defaults to default_api_timeout_ms.

Returns:
{TopicPartition: int}: Mapping of partition to retrieved offset.

Raises:
UnsupportedVersionError if broker does not support any compatible
ListOffsetsRequest api version.
KafkaTimeoutError if timeout_ms provided.
KafkaTimeoutError: if not completed within timeout_ms (default
default_api_timeout_ms).
"""
timestamps = dict([(tp, timestamp) for tp in partitions])
offsets = self._net.run(self._fetch_offsets_by_times_async, timestamps, timeout_ms)
if timeout_ms is None:
timeout_ms = self.config['default_api_timeout_ms']
offsets = self._net.run(self._fetch_offsets_by_times_async, timestamps, timeout_ms, timeout_ms=timeout_ms)
for tp in timestamps:
offsets[tp] = offsets[tp].offset
return offsets
Expand Down
82 changes: 57 additions & 25 deletions kafka/consumer/group.py
Original file line number Diff line number Diff line change
Expand Up @@ -679,19 +679,23 @@ def commit(self, offsets=None, timeout_ms=None):
message read if a consumer is restarted, the committed offset should be
the next message your application should consume, i.e.: last_offset + 1.

Blocks until either the commit succeeds or an unrecoverable error is
encountered (in which case it is thrown to the caller).
Blocks until the commit succeeds, an unrecoverable error is encountered
(in which case it is thrown to the caller), or the timeout expires.

Currently only supports kafka-topic offset storage (not zookeeper).

Arguments:
offsets (dict, optional): {TopicPartition: OffsetAndMetadata} dict
to commit with the configured group_id. Defaults to currently
consumed offsets for all subscribed partitions.
timeout_ms (numeric, optional): Maximum time in milliseconds to
block. Defaults to default_api_timeout_ms.

Raises:
IncompatibleBrokerVersion: if broker version < 0.8.1
IllegalStateError: if group_id is None
KafkaTimeoutError: if the commit does not complete within timeout_ms
(default default_api_timeout_ms).
"""
if self.config['api_version'] < (0, 8, 1):
raise Errors.IncompatibleBrokerVersion('Requires >= Kafka 0.8.1')
Expand Down Expand Up @@ -734,6 +738,8 @@ def committed(self, partition, metadata=False, timeout_ms=None):
partition (TopicPartition): The partition to check.
metadata (bool, optional): If True, return OffsetAndMetadata struct
instead of offset int. Default: False.
timeout_ms (numeric, optional): Maximum time in milliseconds to
block. Defaults to default_api_timeout_ms.

Returns:
The last committed offset (int or OffsetAndMetadata), or None if there was no prior commit.
Expand All @@ -742,7 +748,8 @@ def committed(self, partition, metadata=False, timeout_ms=None):
IncompatibleBrokerVersion: if broker version < 0.8.1
IllegalStateError: if group_id is None
TypeError: if partition is not TopicPartition
KafkaTimeoutError: if timeout_ms provided
KafkaTimeoutError: if the offsets are not fetched within timeout_ms
(default default_api_timeout_ms).
BrokerResponseError: if OffsetFetchRequest raises an error.
"""
if self.config['api_version'] < (0, 8, 1):
Expand All @@ -756,31 +763,44 @@ def committed(self, partition, metadata=False, timeout_ms=None):
return None
return committed[partition] if metadata else committed[partition].offset

def _fetch_all_topic_metadata(self):
def _fetch_all_topic_metadata(self, timeout_ms=None):
"""A blocking call that fetches topic metadata for all topics in the
cluster that the user is authorized to view.
"""
if timeout_ms is None:
timeout_ms = self.config['default_api_timeout_ms']
timer = Timer(timeout_ms)
if self._cluster.metadata_refresh_in_progress:
future = self._cluster.request_update()
self._net.run(self._manager.wait_for, future, None)
self._net.run(self._manager.wait_for, future, timer.timeout_ms, timeout_ms=timer.timeout_ms)
stash = self._cluster.need_all_topic_metadata
self._cluster.need_all_topic_metadata = True
future = self._cluster.request_update()
self._net.run(self._manager.wait_for, future, None)
self._cluster.need_all_topic_metadata = stash
try:
self._cluster.need_all_topic_metadata = True
future = self._cluster.request_update()
self._net.run(self._manager.wait_for, future, timer.timeout_ms, timeout_ms=timer.timeout_ms)
finally:
self._cluster.need_all_topic_metadata = stash

def topics(self):
def topics(self, timeout_ms=None):
"""Get all topics the user is authorized to view.
This will always issue a remote call to the cluster to fetch the latest
information.

Arguments:
timeout_ms (numeric, optional): Maximum time in milliseconds to
block. Defaults to default_api_timeout_ms.

Returns:
set: topics

Raises:
KafkaTimeoutError: if topic metadata is not fetched within
timeout_ms (default default_api_timeout_ms).
"""
self._fetch_all_topic_metadata()
self._fetch_all_topic_metadata(timeout_ms=timeout_ms)
return self._cluster.topics()

def partitions_for_topic(self, topic):
def partitions_for_topic(self, topic, timeout_ms=None):
"""This method first checks the local metadata cache for information
about the topic. If the topic is not found (either because the topic
does not exist, the user is not authorized to view the topic, or the
Expand All @@ -789,13 +809,19 @@ def partitions_for_topic(self, topic):

Arguments:
topic (str): Topic to check.
timeout_ms (numeric, optional): Maximum time in milliseconds to
block. Defaults to default_api_timeout_ms.

Returns:
set: Partition ids

Raises:
KafkaTimeoutError: if topic metadata is not fetched within
timeout_ms (default default_api_timeout_ms).
"""
partitions = self._cluster.partitions_for_topic(topic)
if partitions is None:
self._fetch_all_topic_metadata()
self._fetch_all_topic_metadata(timeout_ms=timeout_ms)
partitions = self._cluster.partitions_for_topic(topic)
return partitions or set()

Expand Down Expand Up @@ -1218,14 +1244,15 @@ def offsets_for_times(self, timestamps, timeout_ms=None):
partition. ``None`` will also be returned for the partition if there
are no messages in it.

Note: This method may block indefinitely if the partition does not exist
and no timeout_ms provided.
Note: This method blocks up to timeout_ms (default
default_api_timeout_ms) if the partition does not exist.

Arguments:
timestamps (dict): ``{TopicPartition: int}`` mapping from partition
to the timestamp to look up. Unit should be milliseconds since
beginning of the epoch (midnight Jan 1, 1970 (UTC))
timeout_ms (int, optional): Milliseconds to block fetching offsets.
Defaults to default_api_timeout_ms.

Returns:
``{TopicPartition: OffsetAndTimestamp}``: mapping from partition
Expand All @@ -1236,9 +1263,10 @@ def offsets_for_times(self, timestamps, timeout_ms=None):
ValueError: If the target timestamp is negative
UnsupportedVersionError: If the broker does not support looking
up the offsets by timestamp.
KafkaTimeoutError: If fetch failed in request_timeout_ms
KafkaTimeoutError: if the offsets are not fetched within timeout_ms
(default default_api_timeout_ms)
"""
timeout_ms = self.config['request_timeout_ms'] if timeout_ms is None else timeout_ms
timeout_ms = self.config['default_api_timeout_ms'] if timeout_ms is None else timeout_ms
for tp, ts in timestamps.items():
timestamps[tp] = int(ts)
if ts < 0:
Expand All @@ -1253,13 +1281,14 @@ def beginning_offsets(self, partitions, timeout_ms=None):
This method does not change the current consumer position of the
partitions.

Note: This method may block indefinitely if the partition does not exist
and no timeout_ms provided.
Note: This method blocks up to timeout_ms (default
default_api_timeout_ms) if the partition does not exist.

Arguments:
partitions (list): List of TopicPartition instances to fetch
offsets for.
timeout_ms (int, optional): Milliseconds to block fetching offsets.
Defaults to default_api_timeout_ms.

Returns:
``{TopicPartition: int}``: The earliest available offsets for the
Expand All @@ -1268,9 +1297,10 @@ def beginning_offsets(self, partitions, timeout_ms=None):
Raises:
UnsupportedVersionError: If the broker does not support looking
up the offsets by timestamp.
KafkaTimeoutError: If fetch failed in timeout_ms.
KafkaTimeoutError: if not completed within timeout_ms (default
default_api_timeout_ms).
"""
timeout_ms = self.config['request_timeout_ms'] if timeout_ms is None else timeout_ms
timeout_ms = self.config['default_api_timeout_ms'] if timeout_ms is None else timeout_ms
offsets = self._fetcher.beginning_offsets(partitions, timeout_ms)
return offsets

Expand All @@ -1282,23 +1312,25 @@ def end_offsets(self, partitions, timeout_ms=None):
This method does not change the current consumer position of the
partitions.

Note: This method may block indefinitely if the partition does not exist
and no timeout_ms provided.
Note: This method blocks up to timeout_ms (default
default_api_timeout_ms) if the partition does not exist.

Arguments:
partitions (list): List of TopicPartition instances to fetch
offsets for.
timeout_ms (int, optional): Milliseconds to block fetching offsets.
Defaults to default_api_timeout_ms.

Returns:
``{TopicPartition: int}``: The end offsets for the given partitions.

Raises:
UnsupportedVersionError: If the broker does not support looking
up the offsets by timestamp.
KafkaTimeoutError: If fetch failed in timeout_ms
KafkaTimeoutError: if not completed within timeout_ms (default
default_api_timeout_ms).
"""
timeout_ms = self.config['request_timeout_ms'] if timeout_ms is None else timeout_ms
timeout_ms = self.config['default_api_timeout_ms'] if timeout_ms is None else timeout_ms
offsets = self._fetcher.end_offsets(partitions, timeout_ms)
return offsets

Expand Down
Loading