Kafka集成测试首次读取失败,后续正常的问题排查
我有个自定义函数forward_to_kafka(list: List),能把列表里的所有事件发送到Kafka,用JB Big Data Tools插件验证过功能正常。但写pytest集成测试时遇到问题:清空Kafka镜像(无历史数据)后,首次推送消息的测试总是失败——测试里的KafkaConsumer接收的消息是空的,但用Big Data Tools能确认消息已经成功投递,后续的测试却能正常运行。
测试代码如下:
@pytest.mark.parametrize("topic, hits", [ ("topic1", []), ("topic3", ["ev1", "ev2", "ev3"]), ("topic2", ["event"]), ("topic4", ['{"event": ["arr_elem"]}', '{"event_num": 13}', '{"event": {"subev": "value"}}']), ]) @pytest.mark.integration def test_forward_to_kafka_integration(self, topic, hits, output): kafka_host = 'localhost:9094' producer = KafkaProducer(bootstrap_servers=[kafka_host], acks='all',) output.forward_to_kafka(producer, topic, [message.encode() for message in hits]) consumer = KafkaConsumer(topic, bootstrap_servers=[kafka_host], group_id=f"{topic}_grp", auto_offset_reset='earliest', consumer_timeout_ms=1000) received_messages = [message.value.decode() for message in consumer] print(received_messages) assert all([message in received_messages for message in hits])
测试失败日志:
FAILED [ 50%][] test_output.py:90 (TestForwardToKafka.test_forward_to_kafka_integration[topic3-hits1]) self = <tests.test_output.TestForwardToKafka object at 0x107b49520> topic = 'topic3', hits = ['ev1', 'ev2', 'ev3'] output = <h3ra.output.Output object at 0x107c14b80> @pytest.mark.parametrize("topic, hits", [ ("topic1", []), ("topic3", ["ev1", "ev2", "ev3"]), ("topic2", ["event"]), ("topic4", ['{"event": ["arr_elem"]}', '{"event_num": 13}', '{"event": {"subev": "value"}}']), ]) @pytest.mark.integration def test_forward_to_kafka_integration(self, topic, hits, output): kafka_host = 'localhost:9094' producer = KafkaProducer(bootstrap_servers=[kafka_host], acks='all',) output.forward_to_kafka(producer, topic, [message.encode() for message in hits]) consumer = KafkaConsumer(topic, bootstrap_servers=[kafka_host], group_id=f"{topic}_grp", auto_offset_reset='earliest', consumer_timeout_ms=1000) received_messages = [message.value.decode() for message in consumer] print(received_messages) > assert all([message in received_messages for message in hits]) E assert False E + where False = all([False, False, False]) test_output.py:107: AssertionError
替换output.forward_to_kafka为以下代码可复现问题:
for message in hits: producer.send(topic, message) producer.flush()
Docker Compose配置:
version: "3.9" services: kafka: # DNS-1035 hostname: kafka image: docker-proxy.artifactory.tcsbank.ru/bitnami/kafka:3.5 expose: - "9092" - "9093" - "9094" volumes: - "kafka_data:/bitnami" environment: - ALLOW_PLAINTEXT_LISTENER=yes - KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093,EXTERNAL://:9094 - KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092,EXTERNAL://kafka:9094 - KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,EXTERNAL:PLAINTEXT,PLAINTEXT:PLAINTEXT - KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE=true - KAFKA_CFG_NODE_ID=0 - KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=0@kafka:9093 - KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER - KAFKA_CFG_PROCESS_ROLES=controller,broker volumes: kafka_data: driver: local
核心原因
首次测试时,Kafka会自动创建主题(因为KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE=true),但主题创建需要完成元数据同步、分区分配、副本初始化等操作,这个过程存在延迟。虽然producer.flush()确保消息被发送到Broker,且acks='all'要求所有副本确认,但消费者启动时可能主题还未完全就绪,或者消费者本地的元数据还没同步完成,导致在consumer_timeout_ms=1000的超时时间内无法拉取到消息。后续测试时主题已经存在,没有了创建延迟,所以能正常消费。
解决方法
方法1:等待Producer所有发送任务完成
仅用producer.flush()不足以确保所有消息的发送确认,应该等待producer.send()返回的Future对象完成,确保Broker已经持久化消息。修改发送逻辑:
# 替换原来的发送代码 futures = [] for message in hits: future = producer.send(topic, message) futures.append(future) # 等待所有消息发送完成 for future in futures: future.get(timeout=5) # 设置合理超时时间 producer.flush()
方法2:增加消费者拉取的重试逻辑
消费者启动后,可能需要一点时间同步元数据,增加重试拉取的逻辑,避免单次超时没拿到消息:
import time consumer = KafkaConsumer(topic, bootstrap_servers=[kafka_host], group_id=f"{topic}_grp", auto_offset_reset='earliest', consumer_timeout_ms=500) # 缩短单次超时时间 received_messages = [] # 重试拉取3次 for _ in range(3): received_messages.extend([msg.value.decode() for msg in consumer]) if len(received_messages) >= len(hits): break time.sleep(0.5) # 等待0.5秒后重试
方法3:提前创建测试主题
在测试开始前,先创建好所有需要用到的测试主题,避免动态创建的延迟。可以用kafka-topics.sh命令,或者在测试代码中通过AdminClient创建:
from kafka.admin import KafkaAdminClient, NewTopic def setup_class(self): kafka_host = 'localhost:9094' admin_client = KafkaAdminClient(bootstrap_servers=[kafka_host]) # 定义需要创建的主题 topics = [ NewTopic(name="topic1", num_partitions=1, replication_factor=1), NewTopic(name="topic2", num_partitions=1, replication_factor=1), NewTopic(name="topic3", num_partitions=1, replication_factor=1), NewTopic(name="topic4", num_partitions=1, replication_factor=1), ] # 创建主题(忽略已存在的主题) try: admin_client.create_topics(new_topics=topics, validate_only=False) except Exception as e: if "Topic already exists" not in str(e): raise e
方法4:延长消费者超时时间
简单粗暴的方式,把consumer_timeout_ms从1000调整到更大的值(比如3000),给足够的时间让消费者同步元数据并拉取消息:
consumer = KafkaConsumer(topic, bootstrap_servers=[kafka_host], group_id=f"{topic}_grp", auto_offset_reset='earliest', consumer_timeout_ms=3000)
内容的提问来源于stack exchange,提问作者shameoff

