diff options
Diffstat (limited to 'kafka')
-rw-r--r-- | kafka/client_async.py | 1 | ||||
-rw-r--r-- | kafka/conn.py | 1 | ||||
-rw-r--r-- | kafka/consumer/fetcher.py | 3 | ||||
-rw-r--r-- | kafka/consumer/group.py | 3 | ||||
-rw-r--r-- | kafka/consumer/subscription_state.py | 3 | ||||
-rw-r--r-- | kafka/coordinator/base.py | 1 | ||||
-rw-r--r-- | kafka/coordinator/consumer.py | 1 | ||||
-rw-r--r-- | kafka/metrics/metrics.py | 3 |
8 files changed, 16 insertions, 0 deletions
diff --git a/kafka/client_async.py b/kafka/client_async.py index a9704fa..a86817d 100644 --- a/kafka/client_async.py +++ b/kafka/client_async.py @@ -414,6 +414,7 @@ class KafkaClient(object): return def __del__(self): + log.debug('%s: __del__', self) self._close() def is_disconnected(self, node_id): diff --git a/kafka/conn.py b/kafka/conn.py index a2d5ee6..fa14761 100644 --- a/kafka/conn.py +++ b/kafka/conn.py @@ -689,6 +689,7 @@ class BrokerConnection(object): self._sock = None def __del__(self): + log.debug('%s: __del__', self) self._close_socket() def close(self, error=None): diff --git a/kafka/consumer/fetcher.py b/kafka/consumer/fetcher.py index 6ec1b71..01deb5c 100644 --- a/kafka/consumer/fetcher.py +++ b/kafka/consumer/fetcher.py @@ -120,6 +120,9 @@ class Fetcher(six.Iterator): self._sensors = FetchManagerMetrics(metrics, self.config['metric_group_prefix']) self._isolation_level = READ_UNCOMMITTED + def __del__(self): + log.debug('%s: __del__', self) + def send_fetches(self): """Send FetchRequests for all assigned partitions that do not already have an in-flight fetch or pending fetch data. diff --git a/kafka/consumer/group.py b/kafka/consumer/group.py index 9abf15e..543c75d 100644 --- a/kafka/consumer/group.py +++ b/kafka/consumer/group.py @@ -421,6 +421,9 @@ class KafkaConsumer(six.Iterator): """ return self._subscription.assigned_partitions() + def __del__(self): + log.debug('%s: __del__', self) + def close(self, autocommit=True): """Close the consumer, waiting indefinitely for any needed cleanup. diff --git a/kafka/consumer/subscription_state.py b/kafka/consumer/subscription_state.py index 10d722e..c1cb020 100644 --- a/kafka/consumer/subscription_state.py +++ b/kafka/consumer/subscription_state.py @@ -73,6 +73,9 @@ class SubscriptionState(object): # initialize to true for the consumers to fetch offset upon starting up self.needs_fetch_committed_offsets = True + def __del__(self): + log.debug('%s: __del__', self) + def subscribe(self, topics=(), pattern=None, listener=None): """Subscribe to a list of topics, or a topic regex pattern. diff --git a/kafka/coordinator/base.py b/kafka/coordinator/base.py index 7deeaf0..506ab54 100644 --- a/kafka/coordinator/base.py +++ b/kafka/coordinator/base.py @@ -735,6 +735,7 @@ class BaseCoordinator(object): self._heartbeat_thread = None def __del__(self): + log.debug('BaseCoordinator: __del__') self._close_heartbeat_thread() def close(self): diff --git a/kafka/coordinator/consumer.py b/kafka/coordinator/consumer.py index cb1de0d..bb992e4 100644 --- a/kafka/coordinator/consumer.py +++ b/kafka/coordinator/consumer.py @@ -125,6 +125,7 @@ class ConsumerCoordinator(BaseCoordinator): self._cluster.add_listener(WeakMethod(self._handle_metadata_update)) def __del__(self): + log.debug('ConsumerCoordinator: __del__') if hasattr(self, '_cluster') and self._cluster: self._cluster.remove_listener(WeakMethod(self._handle_metadata_update)) super(ConsumerCoordinator, self).__del__() diff --git a/kafka/metrics/metrics.py b/kafka/metrics/metrics.py index e9c465d..33b50b5 100644 --- a/kafka/metrics/metrics.py +++ b/kafka/metrics/metrics.py @@ -71,6 +71,9 @@ class Metrics(object): 'total number of registered metrics'), AnonMeasurable(lambda config, now: len(self._metrics))) + def __del__(self): + logger.debug('%s: __del__', self) + @property def config(self): return self._config |