Skip to content

Commit 6603088

Browse files
Added keys to compressed messages (both gzip and snappy).
1 parent 9c5216a commit 6603088

File tree

1 file changed

+2
-2
lines changed

1 file changed

+2
-2
lines changed

kafka/protocol.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -568,7 +568,7 @@ def create_gzip_message(payloads, key=None):
568568
key: bytes, a key used for partition routing (optional)
569569
"""
570570
message_set = KafkaProtocol._encode_message_set(
571-
[create_message(payload) for payload in payloads])
571+
[create_message(payload, key) for payload in payloads])
572572

573573
gzipped = gzip_encode(message_set)
574574
codec = ATTRIBUTE_CODEC_MASK & CODEC_GZIP
@@ -589,7 +589,7 @@ def create_snappy_message(payloads, key=None):
589589
key: bytes, a key used for partition routing (optional)
590590
"""
591591
message_set = KafkaProtocol._encode_message_set(
592-
[create_message(payload) for payload in payloads])
592+
[create_message(payload, key) for payload in payloads])
593593

594594
snapped = snappy_encode(message_set)
595595
codec = ATTRIBUTE_CODEC_MASK & CODEC_SNAPPY

0 commit comments

Comments
 (0)