summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorcharsyam <charsyam@naver.com>2017-03-03 07:15:01 +0900
committerDana Powers <dana.powers@gmail.com>2017-03-02 14:15:01 -0800
commit3a630f2f886d9182bc6fe593d3659b0f3986fb4b (patch)
tree6b2f539a1227065cc7dabd1e06acc74c10e49cf3
parenta22ea165649b3510d770243f6f3809d598cb4f81 (diff)
downloadkafka-python-3a630f2f886d9182bc6fe593d3659b0f3986fb4b.tar.gz
Add send_list_offset_request for searching offset by timestamp (#1001)
-rw-r--r--kafka/client.py10
-rw-r--r--kafka/protocol/legacy.py29
-rw-r--r--kafka/structs.py6
3 files changed, 45 insertions, 0 deletions
diff --git a/kafka/client.py b/kafka/client.py
index ff0169b..9df5bd9 100644
--- a/kafka/client.py
+++ b/kafka/client.py
@@ -686,6 +686,16 @@ class SimpleClient(object):
return [resp if not callback else callback(resp) for resp in resps
if not fail_on_error or not self._raise_on_response_error(resp)]
+ def send_list_offset_request(self, payloads=[], fail_on_error=True,
+ callback=None):
+ resps = self._send_broker_aware_request(
+ payloads,
+ KafkaProtocol.encode_list_offset_request,
+ KafkaProtocol.decode_list_offset_response)
+
+ return [resp if not callback else callback(resp) for resp in resps
+ if not fail_on_error or not self._raise_on_response_error(resp)]
+
def send_offset_commit_request(self, group, payloads=[],
fail_on_error=True, callback=None):
encoder = functools.partial(KafkaProtocol.encode_offset_commit_request,
diff --git a/kafka/protocol/legacy.py b/kafka/protocol/legacy.py
index 6d9329d..c855d05 100644
--- a/kafka/protocol/legacy.py
+++ b/kafka/protocol/legacy.py
@@ -249,6 +249,35 @@ class KafkaProtocol(object):
]
@classmethod
+ def encode_list_offset_request(cls, payloads=()):
+ return kafka.protocol.offset.OffsetRequest[1](
+ replica_id=-1,
+ topics=[(
+ topic,
+ [(
+ partition,
+ payload.time)
+ for partition, payload in six.iteritems(topic_payloads)])
+ for topic, topic_payloads in six.iteritems(group_by_topic_and_partition(payloads))])
+
+ @classmethod
+ def decode_list_offset_response(cls, response):
+ """
+ Decode OffsetResponse_v2 into ListOffsetResponsePayloads
+
+ Arguments:
+ response: OffsetResponse_v2
+
+ Returns: list of ListOffsetResponsePayloads
+ """
+ return [
+ kafka.structs.ListOffsetResponsePayload(topic, partition, error, timestamp, offset)
+ for topic, partitions in response.topics
+ for partition, error, timestamp, offset in partitions
+ ]
+
+
+ @classmethod
def encode_metadata_request(cls, topics=(), payloads=None):
"""
Encode a MetadataRequest
diff --git a/kafka/structs.py b/kafka/structs.py
index 7d1d96a..48321e7 100644
--- a/kafka/structs.py
+++ b/kafka/structs.py
@@ -37,9 +37,15 @@ FetchResponsePayload = namedtuple("FetchResponsePayload",
OffsetRequestPayload = namedtuple("OffsetRequestPayload",
["topic", "partition", "time", "max_offsets"])
+ListOffsetRequestPayload = namedtuple("ListOffsetRequestPayload",
+ ["topic", "partition", "time"])
+
OffsetResponsePayload = namedtuple("OffsetResponsePayload",
["topic", "partition", "error", "offsets"])
+ListOffsetResponsePayload = namedtuple("ListOffsetResponsePayload",
+ ["topic", "partition", "error", "timestamp", "offset"])
+
# https://cwiki.apache.org/confluence/display/KAFKA/A+Guide+To+The+Kafka+Protocol#AGuideToTheKafkaProtocol-OffsetCommit/FetchAPI
OffsetCommitRequestPayload = namedtuple("OffsetCommitRequestPayload",
["topic", "partition", "offset", "metadata"])