Python 操作 Kafka 实战指南
一、说明Kafka 作为高吞吐分布式消息队列广泛用于日志收集、异步解耦、数据流处理、事件推送。Python 生态主流使用confluent-kafka官方推荐高性能底层 librdkafka对比老旧的kafka-python内存占用更低、吞吐量更强生产环境优先选用。环境说明 Python 3.8 Kafka 2.8/3.x 依赖pip install confluent-kafka二、核心概念快速回顾Brokerkafka 服务节点Topic消息主题消息分类载体Partition分区实现水平扩展、并发消费Producer生产者推送消息Consumer消费者拉取消息Consumer Group消费组同一组内一条消息只能被一个消费者消费Offset消息在分区内唯一序号三、基础配置封装统一配置文件kafka_config.pyfrom confluent_kafka import Producer, Consumer # kafka集群地址多个节点逗号分隔 BOOTSTRAP_SERVERS 127.0.0.1:9092四、生产者同步发送 异步回调 批量发送4.1 基础生产者 消息发送回调回调函数用于确认消息是否投递成功记录失败消息是生产环境必备。from confluent_kafka import Producer from kafka_config import BOOTSTRAP_SERVERS producer_conf { bootstrap.servers: BOOTSTRAP_SERVERS, # 确认机制1表示leader写入成功即返回 acks: 1, # 消息超时时间 message.timeout.ms: 5000 } p Producer(producer_conf) def delivery_report(err, msg): 消息投递回调 if err is not None: print(f消息发送失败: {err}) else: print(f消息发送成功,topic:{msg.topic()},partition:{msg.partition()},offset:{msg.offset()}) def send_message(topic: str, data: str, key: str None): # 发送消息key用于决定消息分配到哪个分区 p.produce( topictopic, keykey.encode(utf-8) if key else None, valuedata.encode(utf-8), on_deliverydelivery_report ) # 轮询触发回调 p.poll(0) if __name__ __main__: topic_name demo-topic for i in range(10): send_message(topic_name, f测试消息{i}, keyfkey_{i}) # flush等待所有消息发送完成退出前必须调用 p.flush()4.2 批量发送优化高频场景不要频繁调用 produce积攒消息批量推送提升吞吐量messages [] batch_size 20 topic_name demo-topic for i in range(100): messages.append(f批量消息{i}) if len(messages) batch_size: for msg in messages: p.produce(topic_name, valuemsg.encode(utf-8), on_deliverydelivery_report) p.flush() messages.clear() # 发送剩余消息 if messages: for msg in messages: p.produce(topic_name, valuemsg.encode(utf-8), on_deliverydelivery_report) p.flush()五、消费者持续拉取、手动提交 offset重点自动提交 offset 存在丢消息风险生产环境推荐手动提交 offsetfrom confluent_kafka import Consumer, KafkaError from kafka_config import BOOTSTRAP_SERVERS consumer_conf { bootstrap.servers: BOOTSTRAP_SERVERS, group.id: demo-consumer-group, # 首次启动消费策略latest 最新消息 / earliest从头消费 auto.offset.reset: earliest, # 关闭自动提交offset enable.auto.commit: False, fetch.min.bytes: 1, fetch.max.wait.ms: 500 } c Consumer(consumer_conf) def consume_topic(topic: str): c.subscribe([topic]) try: while True: # 阻塞等待消息超时时间ms msg c.consume(timeout1000) if msg is None: continue # 处理kafka服务端消息 if msg.error(): if msg.error().code() KafkaError._PARTITION_EOF: continue else: raise msg.error() # 业务处理消息 msg_key msg.key().decode(utf-8) if msg.key() else None msg_value msg.value().decode(utf-8) print(f收到消息 key{msg_key}, data{msg_value}) # 业务逻辑执行完成后手动提交offset c.commit(asynchronousFalse) except KeyboardInterrupt: pass finally: # 关闭消费者 c.close() if __name__ __main__: consume_topic(demo-topic)六、JSON 消息收发实际项目绝大部分传递 JSON 数据封装通用工具方法import json from confluent_kafka import Producer, Consumer # 发送json def send_json(producer: Producer, topic: str, payload: dict, keyNone): data json.dumps(payload, ensure_asciiFalse) producer.produce( topictopic, keykey.encode(utf-8) if key else None, valuedata.encode(utf-8), on_deliverydelivery_report ) producer.poll(0) # 消费解析json payload json.loads(msg.value().decode(utf-8)) print(payload[user_name])七、异步方案适配 FastAPI 异步项目confluent-kafka本身是同步库不能直接在 async 函数阻塞调用。 两种解决方案使用threading将消费者放到独立线程推荐 FastAPI 项目aiokafka 纯异步库适合全异步架构aiokafka 异步示例纯异步 Python安装pip install aiokafkaimport asyncio from aiokafka import AIOKafkaProducer, AIOKafkaConsumer BOOTSTRAP_SERVERS 127.0.0.1:9092 # 异步生产者 async def async_producer_demo(): producer AIOKafkaProducer(bootstrap_serversBOOTSTRAP_SERVERS) await producer.start() try: await producer.send_and_wait(demo-topic, basync kafka message) finally: await producer.stop() # 异步消费者 async def async_consumer_demo(): consumer AIOKafkaConsumer( demo-topic, bootstrap_serversBOOTSTRAP_SERVERS, group_idasync-group, auto_offset_resetearliest ) await consumer.start() try: async for msg in consumer: print(收到消息, msg.value.decode()) finally: await consumer.stop() if __name__ __main__: asyncio.run(async_consumer_demo())八、生产环境高频问题 最佳实践offset 自动提交风险消息还未处理完成offset 提前提交程序崩溃导致消息丢失业务处理成功后手动提交 offset。消息丢失场景生产者未调用 flush、acks 配置为 0、网络波动消息未投递务必实现 delivery_report 日志记录失败消息。消息重复消费kafka 不保证 Exactly Once仅保证 At-Least Once业务代码必须实现幂等唯一业务编号去重。分区数量规划消费者并发上限 topic 分区总数想要提升消费并发需要增加分区。kafka-python vs confluent-kafkakafka-python 纯 Python 实现性能差不再推荐新项目生产统一使用 confluent-kafka。序列化规范统一使用 JSON/Protobuf 传递数据不要直接传递复杂对象。九、拓展方向消息重试队列、死信队列失败消息转发 DLQKafka 监控、消息延迟告警FastAPI 集成 kafka项目启动时创建消费者后台任务消息压缩配置lz4 压缩减少网络流量