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

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()

YARN界面


核心原因分析

出现单个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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 03:20:53