如何配置Spark应用实现多节点并行拉取单个Kafka主题数据
问题处理方案
核心配置修正
必须将R53单主机名配置替换为逗号分隔的Kafka broker列表,这是导致单节点消费的核心配置错误。
- R53单域名无论配置的是轮询解析还是单IP转发,Spark Streaming Kafka连接器初始化时,仅能通过首次连接拿到单个broker端点,后续分区分配阶段极易触发分配倾斜,最终只有首个成功建立稳定连接的executor会拉取所有分区数据。
- 正确配置格式参考:
broker1:9092,broker2:9092,broker3:9092,不需要列出集群所有broker,至少配置2~3个存活的broker地址即可保证元数据拉取正常。
其余排查方向(20分区仍单节点消费的其他诱因)
Spark Streaming Direct API的消费并行度和Kafka主题分区数一一对应,每个Kafka分区会对应生成一个消费Task,出现单节点消费基本都是调度、资源、残留锁三类问题:
- 资源配置不足
检查Spark任务提交时的总executor核数配置:总可用核数必须≥你期望的并行消费数。比如20个分区要全并行消费,总executor核数至少要配置到20,核数不足时所有消费Task会排队调度到少量可用核上,表现为单节点高负载、其余节点空闲。同时不要随意调大spark.task.cpus参数,默认1即可满足消费任务需求。 - 消费者组残留锁
停掉Spark任务后,执行Kafka自带命令检查消费者组状态:
如果之前的任务异常退出,会残留分区占用标记,新启动的任务无法正常获取所有分区的消费权,极端情况下只有1个executor能抢到分区。这类问题等待kafka-consumer-groups.sh --bootstrap-server <你的broker地址> --describe --group <你的消费者组ID>session.timeout.ms(默认10s)超时自动释放,或者手动删除旧消费者组的偏移量记录即可解决。 - 数据本地性调度倾斜
Spark默认会优先将任务调度到和数据存储节点同位置的executor上,如果你的Kafka所有分区副本都落在同一个节点,且spark.locality.wait配置的等待时间过长,所有消费任务都会被调度到该节点。可以提交任务时增加配置spark.locality.wait=0关闭本地性等待,验证是否为该问题。
补充说明:Spark 3.2.1版本的Spark Streaming Kafka连接器不需要依赖DFS存储偏移量,你当前未使用DFS的架构不会影响消费并行度。
内容的提问来源于stack exchange,提问作者Kevin
相关产品推荐
相关产品推荐

