You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 09:42:33