kafka-python 是一个纯 Python 实现的 Kafka 客户端库(性能方面比confluent-kafka-python要差一些,毕竟 confluent-kafka-python 底层有用 c 写),开发者可以使用他 Apache Kafka 集群进行交互,发送和接收消息。
特点
- 易用性:简化 Kafka 的操作,易于上手和使用。
- 兼容性:与 Kafka 集群版本兼容性好。
- 异步性:支持异步消息发送,提高性能。
- 扩展性:可以根据需求扩展功能,如消费者组和分区管理。
- 稳定性:拥有较好的错误处理和异常管理机制。
下面用例测试的时候,记得要先把 kafka 服务启动起来!我自己用的 kafka-python 版本是 3.0.11
pip install kafka-python==3.0.11基础功能
kafka-python 提供了丰富的 API,覆盖 Kafka 生产者和消费者的核心操作,以下是几个常用的基础功能:
- 消息发送:通过
KafkaProducer向指定 Topic 发送消息,支持同步和异步两种发送方式。 - 消息消费:通过
KafkaConsumer订阅 Topic 并拉取消息,支持手动提交偏移量和自动提交偏移量。 - 消费者组:支持消费者组机制,多个消费者可以协同消费同一个 Topic 下的分区,实现负载均衡。
- 分区管理:支持手动指定分区发送消息,也支持通过分区器自动选择分区。
- 偏移量管理:支持手动提交偏移量,便于实现精确一次或至少一次的消费语义。
生产者
生产者负责向 Kafka 集群发送消息。
from kafka import KafkaProducer producer = KafkaProducer(bootstrap_servers='localhost:9092') # 发送消息 producer.send('test_topic', b'Hello, Kafka!')消费者
消费者用于从 Kafka 集群中读取消息。
from kafka import KafkaConsumer consumer = KafkaConsumer('test_topic', bootstrap_servers='localhost:9092') # 读取消息 for message in consumer: print(f"Received message: {message.value.decode()}") # Received message: Hello, Kafka!消费者组
消费者组允许多个消费者共同消费一个主题。(我自测发现这种消费会有很大延迟,后面再分析下具体原因)
from kafka import KafkaConsumer, TopicPartition from kafka.coordinator.assignors.roundrobin import RoundRobinPartitionAssignor consumer = KafkaConsumer( "test_topic", group_id='my-group', bootstrap_servers='localhost:9092', auto_offset_reset='earliest', enable_auto_commit=True, partition_assignment_strategy=[RoundRobinPartitionAssignor] ) # 消费消息 for message in consumer: print(f"Received message: {message.value.decode()}")消息确认
确保消息被正确处理后进行确认。
from kafka import KafkaConsumer consumer = KafkaConsumer('test_topic', group_id='my-group', bootstrap_servers='localhost:9092', auto_offset_reset='earliest') # 手动提交偏移量 for message in consumer: # 处理消息 print(f"Received message: {message.value.decode()}") consumer.commit() # 手动提交偏移量指定分区发送
kafka-python支持向特定分区发送消息。
from kafka import KafkaProducer producer = KafkaProducer(bootstrap_servers='localhost:9092') # 向特定分区发送消息 producer.send('test_topic', key=b'key1', value=b'Hello, Kafka!', partition=0) producer.flush()指定分区消费
from kafka import KafkaConsumer, TopicPartition # 创建 TopicPartition 对象 tp = TopicPartition('test_topic', 0) consumer = KafkaConsumer( bootstrap_servers='localhost:9092', auto_offset_reset='latest', enable_auto_commit=True ) # 指定分区消费 consumer.assign([tp]) for message in consumer: print(f"Received message: key: {message.key.decode()}; value: {message.value.decode()}") # Received message: key: key1; value: Hello, Kafka!高级分区消费
在高级分区消费中,可以更细致地控制消息的消费。
from kafka import KafkaConsumer, TopicPartition # 创建 TopicPartition 对象 tp = TopicPartition('test_topic', 0) consumer = KafkaConsumer( group_id="group_id", bootstrap_servers='localhost:9092', auto_offset_reset='latest', enable_auto_commit=False ) # 指定分区消费 consumer.assign([tp]) for message in consumer: print(f"Received message: key: {message.key.decode()}; value: {message.value.decode()}") consumer.commit() # Received message: key: key1; value: Hello, Kafka!消费者组管理
在kafka-python中,可以方便地管理消费者组。这允许多个消费者协调消费同一个主题的消息,确保消息不会被重复处理。
from kafka import KafkaConsumer, TopicPartition # 创建消费者实例,指定消费者组 consumer = KafkaConsumer(group_id='my-group', bootstrap_servers='localhost:9092') # 手动指定消费的分区和偏移量 tp = TopicPartition('my-topic', 0) consumer.assign([tp]) consumer.seek(tp, 10) # 从偏移量10开始消费 for message in consumer: print(f"Received message: {message.value.decode('utf-8')}")消费者偏移量
kafka-python允许开发者手动管理消费者偏移量,这在需要精确控制消费进度时非常有用。
from kafka import KafkaConsumer, TopicPartition # 创建消费者实例,指定消费者组 consumer = KafkaConsumer(group_id='my-group', bootstrap_servers='localhost:9092') # 手动指定消费的分区和偏移量 tp = TopicPartition('test_topic', 1) consumer.assign([tp]) consumer.seek(tp, 11) # 从偏移量11开始消费 # 手动提交偏移量 consumer.commit_async() # 获取当前偏移量 current_offset = consumer.position(tp) print(f"Current offset for partition {tp.partition}: {current_offset}") # Current offset for partition 1: 11生产者事务
在处理高可靠性消息时,使用事务可以确保消息的精确一次处理。
from kafka import KafkaProducer producer = KafkaProducer(bootstrap_servers='localhost:9092', transactional_id='my-transactional-id') # 初始化事务 producer.init_transactions() # 开始事务 producer.begin_transaction() # 发送消息 producer.send('test_topic', b'Hello, Kafka!') # 提交事务 producer.commit_transaction()消息重试
在发送消息时,如果遇到临时错误,可以使用重试机制来确保消息能够成功发送。
from kafka import KafkaProducer producer = KafkaProducer(bootstrap_servers='localhost:9092') from kafka.errors import KafkaError # 定义重试次数 retries = 3 for _ in range(retries): try: producer.send('test_topic', b'Hello, Kafka!') producer.flush() break except KafkaError as e: print(f"Error sending message: {e}") if _ == retries - 1: raise发送结果确认
异步发送可以提高生产者的吞吐量,因为它不需要等待每个消息的发送确认。如果要获取发送结果,则需要阻塞等待。
from kafka import KafkaProducer from kafka.errors import KafkaError producer = KafkaProducer(bootstrap_servers='localhost:9092') # 异步发送消息 future = producer.send('test_topic', b'Hello, Kafka!') # 获取发送结果 try: record_metadata = future.get(timeout=10) print(f"Message sent to {record_metadata.topic}, partition {record_metadata.partition}, offset {record_metadata.offset}") except KafkaError as e: print(f"Failed to send message: {e}")