feat: rebalance b/w paritions example

This commit is contained in:
2026-01-25 02:24:37 +00:00
parent f8d7806ae1
commit be773bd77c
5 changed files with 124 additions and 2 deletions
+25
View File
@@ -0,0 +1,25 @@
from kafka import KafkaConsumer
topic_name = "rebal-topic"
group_id = "rebal-group"
consumer = KafkaConsumer(
topic_name,
bootstrap_servers='localhost:9092', # advertised listener
auto_offset_reset='earliest', # start from the beginning if no offset is committed
enable_auto_commit=True, # read = commit, this is True by default
# who am I in the consumer group, without it you'll always read from beginning
# as kafka does not know if you have read before as you had not name for yourself
group_id=group_id
)
try:
for message in consumer:
# message is of type ConsumerRecord and data is in value which is bytes
text = message.value.decode('utf-8')
print(f"Received: {text} from partition: {message.partition}")
except KeyboardInterrupt:
print("Stopping consumer...")
finally:
consumer.close()