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并退出,根本没有进入消息循环处理消息。核心原因在于:
consumer_timeout_ms设置过短:该参数定义了消费者在没有新消息时等待的超时时间,当前设置为3000ms(3秒)。而消费者加入组、完成分区分配的过程已经花费了约4秒(从13:23:12到13:23:16),此时消费者刚完成分配就触发了超时,直接退出循环,没有机会处理消息。生产和消费的执行时序问题:生产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
相关产品推荐
相关产品推荐

