PySpark读取Kafka执行count仅用单Executor,如何优化并行度?
问题:Spark读取Kafka执行count()仅用单个Executor,如何提升并行效率?
我尝试从Kafka读取数据并统计总量,执行df.count()时发现仅使用单个Executor,处理耗时过长。当前Spark配置和代码如下:
spark = SparkSession.builder.appName('oracle_read_test') \ .config("spark.driver.memory", "30g") \ .config("spark.driver.maxResultSize", "64g") \ .config("spark.executor.cores", "10") \ .config("spark.executor.instances", "15") \ .config('spark.executor.memory', '30g') \ .config('num-executors', '20') \ .config('spark.yarn.executor.memoryOverhead', '32g') \ .config("hive.exec.dynamic.partition", "true") \ .config("orc.compress", "ZLIB") \ .config("hive.merge.smallfiles.avgsize", "40000000") \ .config("hive.merge.size.per.task", "209715200") \ .config("dfs.blocksize", "268435456") \ .config("hive.metastore.try.direct.sql", "true") \ .config("spark.sql.orc.enabled", "true") \ .config("spark.dynamicAllocation.enabled", "false") \ .config("spark.sql.sources.partitionOverwriteMode","dynamic") \ .getOrCreate() df = spark.read.format("kafka") \ .option("kafka.bootstrap.servers","localhost:9092") \ .option("includeHeaders","true") \ .option("subscribe","test") \ .load() df.count()

核心原因分析
出现单个Executor工作的本质问题是Spark任务并行度不足:
- Spark读取Kafka时,默认会创建与Kafka主题分区数一致的RDD分区。如果
test主题只有1个分区,Spark只会启动1个任务,其余Executor会处于空闲状态。 - 配置存在参数冲突:同时设置了
spark.executor.instances(15)和旧版参数num-executors(20),Spark会优先使用spark.executor.instances,但冗余配置可能引发资源调度异常。
具体调整方案
1. 扩容Kafka主题分区数(最关键)
Kafka主题分区数是Spark并行度的基础,先检查并扩容主题分区:
# 查看test主题的分区数 kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic test # 扩容到N个分区(N建议等于或略大于Executor总核数) kafka-topics.sh --bootstrap-server localhost:9092 --alter --topic test --partitions 150
2. 修正Spark配置,消除参数冲突
移除冗余的num-executors参数,同时优化资源配置(避免单Executor核数过高):
spark = SparkSession.builder.appName('oracle_read_test') \ .config("spark.driver.memory", "30g") \ .config("spark.driver.maxResultSize", "64g") \ .config("spark.executor.cores", "5") # 降低单Executor核数,提升任务并行性 .config("spark.executor.instances", "15") \ .config('spark.executor.memory', '30g') \ .config('spark.yarn.executor.memoryOverhead', '8g') # 内存Overhead无需过大,一般为executor内存的10%-20% .config("hive.exec.dynamic.partition", "true") \ .config("orc.compress", "ZLIB") \ .config("hive.merge.smallfiles.avgsize", "40000000") \ .config("hive.merge.size.per.task", "209715200") \ .config("dfs.blocksize", "268435456") \ .config("hive.metastore.try.direct.sql", "true") \ .config("spark.sql.orc.enabled", "true") \ .config("spark.dynamicAllocation.enabled", "false") \ .config("spark.sql.sources.partitionOverwriteMode","dynamic") \ .getOrCreate()
3. 强制提升Spark读取Kafka的并行度(可选)
如果无法修改Kafka主题分区,可通过Spark参数强制增加RDD分区数:
df = spark.read.format("kafka") \ .option("kafka.bootstrap.servers","localhost:9092") \ .option("includeHeaders","true") \ .option("subscribe","test") \ .option("minPartitions", "15") # 强制设置最小分区数,需与Executor数量匹配 .load() # 或者读取后重分区 df = df.repartition(15) df.count()
4. 验证并行度
执行以下代码确认RDD分区数,确保与预期并行度一致:
print(df.rdd.getNumPartitions())
内容的提问来源于stack exchange,提问作者veganzombie
相关产品推荐
相关产品推荐

