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

Airflow中同消费组的两个KafkaConsumer仅一个消费消息问题

Kafka消费者组分区分配异常排查

搭建了一套包含两个独立进程的系统:一个向Kafka的TopicA生产消息,另一个运行两个同属group-id为my-group的KafkaConsumer实例消费消息。TopicA通过Kafdrop创建,包含4个分区,预期两个消费者应分配分区共同消费,但实际仅一个消费者处理消息,另一个无任何消费动作。以下是相关代码及日志,请求排查原因。


消息生产DAG脚本

import airflow
from airflow.models import DAG
from airflow.operators.python_operator import PythonOperator
from preprocessing.testers import create_msg, print_current_time
import random

args = {
    'owner': 'airflow',
    'start_date': airflow.utils.dates.days_ago(1),
    'provide_context': True,                     
}

dag = DAG(
    dag_id='creating_msgs',
    default_args=args,
    schedule_interval= '@once',
    catchup=True,
)

create_msgs = [
    PythonOperator(
        task_id='create_msg_'+str(i+1),
        python_callable=create_msg,
        op_kwargs={'task_id': i+1, 'nbr': random.randrange(50,100)},
        dag=dag) for i in range(2)
]

print_time_task = PythonOperator(
        task_id='print_time',
        python_callable=print_current_time,
        dag=dag
    )

create_msgs >> print_time_task

消息生产逻辑代码

def create_msg(**kwargs):
    producer = KafkaProducer(bootstrap_servers=['kafka:9092'],                              # set up Producer
    value_serializer=lambda x: json.dumps(x).encode('utf-8'))

    for i in range(kwargs['nbr']):
        task_name = 'Task' + str(kwargs['task_id']) + '_' + str(i+1)
        logging.warn('Producing message of task: ' + task_name)
        producer.send('TopicA', {'name': task_name})
        sleep_time = random.randrange(1,5)
        logging.warn('Sleeping for ' + str(sleep_time))
        sleep(sleep_time)
    producer.close()

消息消费DAG脚本

import airflow
from airflow.models import DAG
from airflow.operators.python_operator import PythonOperator
from preprocessing.testers import print_current_time
from crawling.crawler import start_task
import random

args = { 
    'owner': 'airflow',
    'start_date': airflow.utils.dates.days_ago(1), 
    'provide_context': True,
}

dag = DAG(
    dag_id='crawl_listings',
    default_args=args,
    schedule_interval= '@once',
    catchup=True,
)

start_tasks = [
    PythonOperator(
        task_id='start_task'+str(i+1),
        python_callable=create_task,
        op_kwargs={'task_id': i+1},
        dag=dag) for i in range(2)
]

print_time_task = PythonOperator(
        task_id='print_time',
        python_callable=print_current_time,
        dag=dag
    )

start_tasks >> print_time_task

消息消费逻辑代码(KafkaConsumer实例)

def start_task(**kwargs):
    consumer = KafkaConsumer(
        'TopicA',
        bootstrap_servers=['kafka:9092'],
        consumer_timeout_ms=3000,
        auto_offset_reset='earliest',
        enable_auto_commit=True,
        group_id='my-group',
        value_deserializer=lambda x: json.loads(x.decode('utf-8')))

    try:
        for message in consumer:
            logging.info('Processing: ' + message.value['name'])
            sleep_time = random.randrange(1,5)
            logging.warn('Sleeping for ' + str(sleep_time))
            sleep(sleep_time)
            
        
        logging.info('Done')
        consumer.close()
    except Exception as e:
        print(e)
        logging.error('Error: ' + e)

无消息消费的KafkaConsumer任务Airflow日志

[2022-09-15 13:23:11,737] {{taskinstance.py:887}} INFO - Executing <Task(PythonOperator): start_crawling_task2> on 2022-09-14T00:00:00+00:00
[2022-09-15 13:23:11,739] {{standard_task_runner.py:53}} INFO - Started process 3319 to run task
[2022-09-15 13:23:11,785] {{logging_mixin.py:112}} INFO - Running %s on host %s <TaskInstance: crawl_listings.start_crawling_task2 2022-09-14T00:00:00+00:00 [running]> f68c0ccde6c0
[2022-09-15 13:23:11,805] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:11,804] {{conn.py:372}} INFO - <BrokerConnection node_id=bootstrap-0 host=kafka:9092 <connecting> [IPv4 ('172.19.0.6', 9092)]>: connecting to kafka:9092 [('172.19.0.6', 9092) IPv4]
[2022-09-15 13:23:11,805] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:11,805] {{conn.py:1049}} INFO - Probing node bootstrap-0 broker version
[2022-09-15 13:23:11,807] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:11,806] {{conn.py:401}} INFO - <BrokerConnection node_id=bootstrap-0 host=kafka:9092 <connecting> [IPv4 ('172.19.0.6', 9092)]>: Connection complete.
[2022-09-15 13:23:11,913] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:11,913] {{conn.py:1106}} INFO - Broker version identifed as 1.0.0
[2022-09-15 13:23:11,914] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:11,913] {{conn.py:1108}} INFO - Set configuration api_version=(1, 0, 0) to skip auto check_version requests on startup
[2022-09-15 13:23:11,914] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:11,914] {{subscription_state.py:171}} INFO - Updating subscribed topics to: ('TopicA',)
[2022-09-15 13:23:12,121] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:12,121] {{conn.py:372}} INFO - <BrokerConnection node_id=1 host=kafka:9092 <connecting> [IPv4 ('172.19.0.6', 9092)]>: connecting to kafka:9092 [('172.19.0.6', 9092) IPv4]
[2022-09-15 13:23:12,222] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:12,222] {{conn.py:401}} INFO - <BrokerConnection node_id=1 host=kafka:9092 <connecting> [IPv4 ('172.19.0.6', 9092)]>: Connection complete.
[2022-09-15 13:23:12,222] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:12,222] {{conn.py:811}} INFO - <BrokerConnection node_id=bootstrap-0 host=kafka:9092 <connected> [IPv4 ('172.19.0.6', 9092)]>: Closing connection. 
[2022-09-15 13:23:12,427] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:12,427] {{cluster.py:376}} INFO - Group coordinator for my-group is BrokerMetadata(nodeId=1, host='kafka', port=9092, rack=None)
[2022-09-15 13:23:12,428] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:12,427] {{base.py:688}} INFO - Discovered coordinator 1 for group my-group
[2022-09-15 13:23:12,428] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:12,428] {{base.py:735}} INFO - Starting new heartbeat thread
[2022-09-15 13:23:12,428] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:12,428] {{consumer.py:342}} INFO - Revoking previously assigned partitions set() for group my-group
[2022-09-15 13:23:12,429] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:12,429] {{base.py:447}} INFO - (Re-)joining group my-group
[2022-09-15 13:23:16,647] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:16,647] {{base.py:333}} INFO - Successfully joined group my-group with generation 2
[2022-09-15 13:23:16,647] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:16,647] {{subscription_state.py:257}} INFO - Updated partition assignment: [TopicPartition(topic='TopicA', partition=2), TopicPartition(topic='TopicA', partition=3)]
[2022-09-15 13:23:16,648] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:16,648] {{consumer.py:239}} INFO - Setting newly assigned partitions {TopicPartition(topic='TopicA', partition=2), TopicPartition(topic='TopicA', partition=3)} for group my-group
[2022-09-15 13:23:16,651] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:16,650] {{crawler.py:25}} INFO - Done
[2022-09-15 13:23:16,654] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:16,654] {{base.py:742}} INFO - Stopping heartbeat thread
[2022-09-15 13:23:16,654] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:16,654] {{base.py:767}} INFO - Leaving consumer group (my-group).
[2022-09-15 13:23:16,661] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:16,661] {{conn.py:811}} INFO - <BrokerConnection node_id=1 host=kafka:9092 <connected> [IPv4 ('172.19.0.6', 9092)]>: Closing connection. 
[2022-09-15 13:23:16,662] {{python_operator.py:114}} INFO - Done. Returned value was: None
[2022-09-15 13:23:16,670] {{taskinstance.py:1048}} INFO - Marking task as SUCCESS.dag_id=crawl_listings, task_id=start_crawling_task2, execution_date=20220914T000000, start_date=20220915T132311, end_date=20220915T132316
[2022-09-15 13:23:21,700] {{logging_mixin.py:112}} INFO - [2022-09-15 13:23:21,700] {{local_task_job.py:103}} INFO - Task exited with return code 0

问题排查分析

从日志可以看到,无消费的消费者成功加入了消费者组并被分配了分区2和3,但紧接着就打印了Done并退出,根本没有进入消息循环处理消息。核心原因在于:

  1. consumer_timeout_ms设置过短:该参数定义了消费者在没有新消息时等待的超时时间,当前设置为3000ms(3秒)。而消费者加入组、完成分区分配的过程已经花费了约4秒(从13:23:12到13:23:16),此时消费者刚完成分配就触发了超时,直接退出循环,没有机会处理消息。

  2. 生产和消费的执行时序问题:生产DAG和消费DAG都是@once调度,可能消费任务启动时,生产任务还没完成消息写入,或者消息还未同步到分区,消费者在超时时间内没读到消息就退出了。

解决方案

  • 调整consumer_timeout_ms参数:如果是要持续消费消息,建议移除该参数(让消费者一直等待新消息);如果是测试一次性消费,可将超时时间延长至足够覆盖分区分配和消息读取的时间,比如设置为30000(30秒)。
  • 确保消费任务在生产任务完成后启动:通过Airflow的DAG依赖,让消费DAG依赖生产DAG的成功执行,或者在消费任务中增加等待逻辑,确保消息已经生产完成。
  • 检查消息的分区分布:生产消息时如果没有指定分区键,Kafka会轮询分配分区,但如果生产任务在消费任务启动前已经完成,消费者分配到分区后应该能读到历史消息(因为设置了auto_offset_reset='earliest'),所以主要问题还是超时时间过短。

修改后的消费代码示例:

def start_task(**kwargs):
    consumer = KafkaConsumer(
        'TopicA',
        bootstrap_servers=['kafka:9092'],
        # 移除超时参数,持续消费;或延长超时时间
        # consumer_timeout_ms=30000,
        auto_offset_reset='earliest',
        enable_auto_commit=True,
        group_id='my-group',
        value_deserializer=lambda x: json.loads(x.decode('utf-8')))

    try:
        for message in consumer:
            logging.info('Processing: ' + message.value['name'])
            sleep_time = random.randrange(1,5)
            logging.warn('Sleeping for ' + str(sleep_time))
            sleep(sleep_time)
            
        logging.info('Done')
        consumer.close()
    except Exception
相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 13:15:31