diff options
| author | Dana Powers <dana.powers@rd.io> | 2015-12-10 16:24:32 -0800 | 
|---|---|---|
| committer | Dana Powers <dana.powers@rd.io> | 2015-12-10 18:37:03 -0800 | 
| commit | d54980a2cd918f243e30ecc23a588fb597957e41 (patch) | |
| tree | 6d666d4832418863cf31ff186dab79805ba2b497 /kafka/client.py | |
| parent | 7a804224949315251b9183fbfa56282ced881244 (diff) | |
| download | kafka-python-d54980a2cd918f243e30ecc23a588fb597957e41.tar.gz | |
Drop kafka_bytestring
Diffstat (limited to 'kafka/client.py')
| -rw-r--r-- | kafka/client.py | 7 | 
1 files changed, 2 insertions, 5 deletions
| diff --git a/kafka/client.py b/kafka/client.py index cb60d98..ca737c4 100644 --- a/kafka/client.py +++ b/kafka/client.py @@ -17,7 +17,6 @@ from kafka.common import (TopicAndPartition, BrokerMetadata, UnknownError,  from kafka.conn import collect_hosts, BrokerConnection, DEFAULT_SOCKET_TIMEOUT_SECONDS  from kafka.protocol import KafkaProtocol -from kafka.util import kafka_bytestring  log = logging.getLogger(__name__) @@ -212,7 +211,7 @@ class KafkaClient(object):                  failed_payloads(broker_payloads)                  continue -            conn = self._get_conn(broker.host.decode('utf-8'), broker.port) +            conn = self._get_conn(broker.host, broker.port)              request = encoder_fn(payloads=broker_payloads)              # decoder_fn=None signal that the server is expected to not              # send a response.  This probably only applies to @@ -305,7 +304,7 @@ class KafkaClient(object):          # Send the request, recv the response          try: -            conn = self._get_conn(broker.host.decode('utf-8'), broker.port) +            conn = self._get_conn(broker.host, broker.port)              conn.send(requestId, request)          except ConnectionError as e: @@ -410,14 +409,12 @@ class KafkaClient(object):          self.topic_partitions.clear()      def has_metadata_for_topic(self, topic): -        topic = kafka_bytestring(topic)          return (            topic in self.topic_partitions            and len(self.topic_partitions[topic]) > 0          )      def get_partition_ids_for_topic(self, topic): -        topic = kafka_bytestring(topic)          if topic not in self.topic_partitions:              return [] | 
