Python 操作 Kafka 实战指南

📅 2026/8/3 3:25:41 👁️ 阅读次数 📝 编程学习
Python 操作 Kafka 实战指南

一、说明

Kafka 作为高吞吐分布式消息队列,广泛用于日志收集、异步解耦、数据流处理、事件推送。Python 生态主流使用confluent-kafka(官方推荐,高性能,底层 librdkafka),对比老旧的kafka-python,内存占用更低、吞吐量更强,生产环境优先选用。

环境说明 Python >=3.8 Kafka 2.8+/3.x 依赖:

pip install confluent-kafka

二、核心概念快速回顾

  1. Broker:kafka 服务节点
  2. Topic:消息主题,消息分类载体
  3. Partition:分区,实现水平扩展、并发消费
  4. Producer:生产者,推送消息
  5. Consumer:消费者,拉取消息
  6. Consumer Group:消费组,同一组内一条消息只能被一个消费者消费
  7. Offset:消息在分区内唯一序号

三、基础配置封装

统一配置文件kafka_config.py

from 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( topic=topic, key=key.encode("utf-8") if key else None, value=data.encode("utf-8"), on_delivery=delivery_report ) # 轮询,触发回调 p.poll(0) if __name__ == "__main__": topic_name = "demo-topic" for i in range(10): send_message(topic_name, f"测试消息{i}", key=f"key_{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, value=msg.encode("utf-8"), on_delivery=delivery_report) p.flush() messages.clear() # 发送剩余消息 if messages: for msg in messages: p.produce(topic_name, value=msg.encode("utf-8"), on_delivery=delivery_report) p.flush()

五、消费者(持续拉取、手动提交 offset)

重点:自动提交 offset 存在丢消息风险!生产环境推荐手动提交 offset

from 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(timeout=1000) 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(asynchronous=False) 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, key=None): data = json.dumps(payload, ensure_ascii=False) producer.produce( topic=topic, key=key.encode("utf-8") if key else None, value=data.encode("utf-8"), on_delivery=delivery_report ) producer.poll(0) # 消费解析json payload = json.loads(msg.value().decode("utf-8")) print(payload["user_name"])

七、异步方案(适配 FastAPI 异步项目)

confluent-kafka本身是同步库,不能直接在 async 函数阻塞调用。 两种解决方案:

  1. 使用threading将消费者放到独立线程(推荐 FastAPI 项目)
  2. aiokafka 纯异步库(适合全异步架构)

aiokafka 异步示例(纯异步 Python)

安装:

pip install aiokafka
import asyncio from aiokafka import AIOKafkaProducer, AIOKafkaConsumer BOOTSTRAP_SERVERS = "127.0.0.1:9092" # 异步生产者 async def async_producer_demo(): producer = AIOKafkaProducer(bootstrap_servers=BOOTSTRAP_SERVERS) await producer.start() try: await producer.send_and_wait("demo-topic", b"async kafka message") finally: await producer.stop() # 异步消费者 async def async_consumer_demo(): consumer = AIOKafkaConsumer( "demo-topic", bootstrap_servers=BOOTSTRAP_SERVERS, group_id="async-group", auto_offset_reset="earliest" ) 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())

八、生产环境高频问题 & 最佳实践

  1. offset 自动提交风险消息还未处理完成,offset 提前提交,程序崩溃导致消息丢失;业务处理成功后手动提交 offset。

  2. 消息丢失场景生产者未调用 flush、acks 配置为 0、网络波动消息未投递;务必实现 delivery_report 日志记录失败消息。

  3. 消息重复消费kafka 不保证 Exactly Once,仅保证 At-Least Once;业务代码必须实现幂等(唯一业务编号去重)。

  4. 分区数量规划消费者并发上限 = topic 分区总数,想要提升消费并发,需要增加分区。

  5. kafka-python vs confluent-kafkakafka-python 纯 Python 实现,性能差,不再推荐新项目;生产统一使用 confluent-kafka。

  6. 序列化规范统一使用 JSON/Protobuf 传递数据,不要直接传递复杂对象。

九、拓展方向

  1. 消息重试队列、死信队列(失败消息转发 DLQ)
  2. Kafka 监控、消息延迟告警
  3. FastAPI 集成 kafka,项目启动时创建消费者后台任务
  4. 消息压缩配置(lz4 压缩,减少网络流量)