From bff6cae885e9750b468d475dac3ec712bf69e853 Mon Sep 17 00:00:00 2001 From: David Arthur Date: Tue, 2 Apr 2013 20:35:22 -0400 Subject: A few fixes for offset APIs in 0.8.1 --- kafka/client.py | 2 +- kafka/protocol.py | 2 -- 2 files changed, 1 insertion(+), 3 deletions(-) diff --git a/kafka/client.py b/kafka/client.py index 8498a08..9893737 100644 --- a/kafka/client.py +++ b/kafka/client.py @@ -231,7 +231,7 @@ class KafkaClient(object): def send_offset_fetch_request(self, group, payloads=[], fail_on_error=True, callback=None): resps = self._send_broker_aware_request(payloads, - partial(KafkaProtocol.encode_offset_commit_fetch, group=group), + partial(KafkaProtocol.encode_offset_fetch_request, group=group), KafkaProtocol.decode_offset_fetch_response) out = [] for resp in resps: diff --git a/kafka/protocol.py b/kafka/protocol.py index fc866ad..94a7f2a 100644 --- a/kafka/protocol.py +++ b/kafka/protocol.py @@ -354,7 +354,6 @@ class KafkaProtocol(object): ====== data: bytes to decode """ - data = data[2:] # TODO remove me when versionId is removed ((correlation_id,), cur) = relative_unpack('>i', data, 0) (client_id, cur) = read_short_string(data, cur) ((num_topics,), cur) = relative_unpack('>i', data, cur) @@ -398,7 +397,6 @@ class KafkaProtocol(object): data: bytes to decode """ - data = data[2:] # TODO remove me when versionId is removed ((correlation_id,), cur) = relative_unpack('>i', data, 0) (client_id, cur) = read_short_string(data, cur) ((num_topics,), cur) = relative_unpack('>i', data, cur) -- cgit v1.2.1