ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

PySpark 调试与性能剖析完全指南:从远程断点到 UDF 性能剖析

PySpark 调试与性能剖析完全指南:从远程断点到 UDF 性能剖析 大数据数据分析批处理流处理机器学习图计算【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址https://gitcode.com/gh_mirrors/sp/spark点击查看免费下载本指南以 Apache Spark 仓库中 python/docs/source/development/debugging.rst 为骨架系统讲解 PySpark Python 侧的调试与剖析方法包括借助 PyCharm Professional 对 Driver 与 Executor 两端进行远程断点调试、用top/ps检查进程资源、用memory_profiler与 Python 内置cProfile定位内存与 CPU 热点以及 PySpark 常见异常与错误码的排查修复。读完本文你将掌握一套从进程定位 → 断点调试 → 性能剖析 → 异常排除的完整 PySpark 排查链路。理解 PySpark 的进程模型Driver 与 Executor 两侧的调试差异PySpark 以 Spark 为计算引擎并通过Py4J在 Python 与 JVM 之间通信Driver 侧PySpark 通过Py4J与 JVM 中的 Driver 通信。当pyspark.sql.SparkSession或pyspark.SparkContext创建并初始化时PySpark 会启动一个 JVM 用于通信。因此 Driver 侧的 Python 代码本质上是一个普通 Python 进程除非以 YARN 集群模式把 Driver 运行在远端机器上调试方式与常规 Python 程序无异。Executor 侧Python worker 进程负责执行 Python 原生函数或处理 Python 数据。它们是惰性启动的——只有应用确实需要 Python worker 与 JVM 交互时才被拉起例如执行 pandas UDF 或 PySpark RDD API 时。这与 Driver 侧始终存在的进程不同因此两侧的调试方式必须分开讨论。本文聚焦于 PySpark Python 侧的调试不涉及 JVM 侧剖析JVM 的 profiling/debugging 工具在 Spark 官方 Developer Tools 文档中另有说明。另外请注意两个前提在本地运行时Driver 侧可以直接用 IDE 调试无需远程调试特性IDE 配置方式参见 python/docs/source/development/setting_ide.rst其中包含 PyCharm 相关章节。除本文方案外也可以使用开源pydevd的远程调试器Remote Debugger替代 PyCharm Professional 来完成类似工作。远程调试PyCharm Professional 双端断点实战本节以单机演示为例分别展示 Driver 与 Executor 两侧的远程调试若要调试其他机器上的 PySpark 应用请参照 PyCharm 官方remote debugging with product完整说明。第一步创建 Python Debug Server 配置从Run菜单选择Edit Configuration...打开Run/Debug Configurations dialog点击工具栏上的新建配置在可用配置列表中选择Python Debug Server输入配置名称例如MyRemoteDebugger并指定端口号例如12345。之后需要在所有将要连接该调试器的机器上安装与本地 PyCharm 版本匹配的pydevd-pycharm包上一对话框会直接给出安装命令pip install pydevd-pycharm~version of PyCharm on the local machineDriver 侧远程调试Driver 侧调试时应用需要能连接到调试服务器。把对话框生成的pydevd_pycharm.settrace代码复制到 PySpark 脚本最顶部。假设脚本名为app.pyecho #Copy and paste from the previous dialog import pydevd_pycharm pydevd_pycharm.settrace(localhost, port12345, stdoutToServerTrue, stderrToServerTrue) # # Your PySpark application codes: from pyspark.sql import SparkSession spark SparkSession.builder.getOrCreate() spark.range(10).show() app.py然后点击运行MyRemoteDebugger调试配置再提交应用spark-submit app.py应用提交后即会连接 PyCharm 调试服务器从而在 Driver 侧实现远程断点调试。Executor 侧远程调试Executor 侧无法像 Driver 那样直接在脚本顶部插入代码因为 Python worker 的入口是 Spark 内部的daemon/worker模块。解法是自定义一个 Python daemon 模块用pydevd_pycharm.settrace包装 worker 主函数再通过spark.python.daemon.module配置让 Spark 使用它。先在当前工作目录准备remote_debug.pyecho from pyspark import daemon, worker def remote_debug_wrapped(*args, **kwargs): #Copy and paste from the previous dialog import pydevd_pycharm pydevd_pycharm.settrace(localhost, port12345, stdoutToServerTrue, stderrToServerTrue) # worker.main(*args, **kwargs) daemon.worker_main remote_debug_wrapped if __name__ __main__: daemon.manager() remote_debug.py该文件将 Python worker 替换为remote_debug_wrapped每次 worker 被调用前先settrace连接到调试服务器再转发给真正的worker.main。随后用该配置启动pysparkshellpyspark --conf spark.python.daemon.moduleremote_debug启动MyRemoteDebugger调试服务器后运行一个能创建 Python worker 的作业即可触发断点例如spark.range(10).repartition(1).rdd.map(lambda x: x).collect()源码佐证daemon 模块如何被加载从源码看spark.python.daemon.module是 Spark 官方定义的核心配置项core/src/main/scala/org/apache/spark/internal/config/Python.scala 中声明了PYTHON_DAEMON_MODULE ConfigBuilder(spark.python.daemon.module)与PYTHON_WORKER_MODULEspark.python.worker.module。在 core/src/main/scala/org/apache/spark/api/python/PythonRunner.scala 中daemonModule读取该配置若设置了自定义模块则用它启动 daemon 进程未设置时使用默认的pyspark.daemon。同时pyspark/daemon.py见 python/pyspark/daemon.py中manager()正是 daemon 进程的事件循环入口daemon.worker_main则被替换为包装函数以注入调试逻辑。检查资源占用top与ps定位 Python 进程Driver 与 Executor 两侧的 Python 进程都可以用常规的top、ps命令查看。Driver 侧直接取 PID在 PySpark shell 中直接获取当前进程 ID import os; os.getpid() 18482再查看该进程的详细信息与资源占用ps -fe 18482UID PID PPID C STIME TTY TIME CMD 000 18482 12345 0 0:00PM ttys001 0:00.00 /.../pythonExecutor 侧grep 定位 daemon 派生的 workerPython worker 是从pyspark.daemonfork 出来的因此直接 grep 即可得到进程 ID 与资源占用情况ps -fe | grep pyspark.daemon000 12345 1 0 0:00PM ttys000 0:00.00 /.../python -m pyspark.daemon 000 12345 1 0 0:00PM ttys000 0:00.00 /.../python -m pyspark.daemon 000 12345 1 0 0:00PM ttys000 0:00.00 /.../python -m pyspark.daemon 000 12345 1 0 0:00PM ttys000 0:00.00 /.../python -m pyspark.daemon ...当 Executor 数量较大时可以结合top按 CPU/内存排序快速锁定异常 worker 进程。剖析内存使用memory_profiler 逐行定位memory_profiler是逐行检查内存占用的剖析器支持 Driver 侧普通脚本也支持 PySpark 为 UDF 提供的内置远程内存剖析。Driver 侧装饰profile逐行统计只要 Driver 不在远端机器如 YARN 集群模式就能直接对 Driver 进程做内存剖析。假设脚本名为profile_memory.pyecho from pyspark.sql import SparkSession #Your function should be decorated with profile from memory_profiler import profile profile # def my_func(): session SparkSession.builder.getOrCreate() df session.range(10000) return df.collect() if __name__ __main__: my_func() profile_memory.py用 memory_profiler 模块方式运行python -m memory_profiler profile_memory.py输出逐行内存报表Filename: profile_memory.py Line # Mem usage Increment Line Contents ... 6 def my_func(): 7 51.5 MiB 0.6 MiB session SparkSession.builder.getOrCreate() 8 51.5 MiB 0.0 MiB df session.range(10000) 9 54.4 MiB 2.8 MiB return df.collect()可以看到创建 SparkSession 后基线内存约 51.5 MiBcollect()拉取 10000 行结果到 Driver 增加了约 2.8 MiB。对大数据量的collect场景这类报表能直观暴露内存峰值来源。Python/Pandas/Arrow UDF 内存剖析内置远程 ProfilerPySpark 为 Python/Pandas/Arrow UDF 提供了内置的远程 memory_profiler适用于带行号的编辑器如 Jupyter Notebook但不支持生成器函数generator functions。通过设置运行时 SQL 配置spark.sql.pyspark.udf.profiler为memory即可开启from pyspark.sql.functions import pandas_udf df spark.range(10) pandas_udf(long) def add1(x): return x 1 spark.conf.set(spark.sql.pyspark.udf.profiler, memory) added df.select(add1(id)) added.show() spark.profile.show(typememory)输出结果 Profile of UDFid2 Filename: ... Line # Mem usage Increment Occurrences Line Contents 4 974.0 MiB 974.0 MiB 10 pandas_udf(long) 5 def add1(x): 6 974.4 MiB 0.4 MiB 10 return x 1UDF 的 ID 可以从查询计划中看到——例如物理计划里ArrowEvalPython节点中的add1(...)#2Ladded.explain() Physical Plan *(2) Project [pythonUDF0#11L AS add1(id)#3L] - ArrowEvalPython [add1(id#0L)#2L], [pythonUDF0#11L], 200 - *(1) Range (0, 10, step1, splits16)剖析结果还支持自定义渲染与清理def do_render(codemap): # Your custom rendering logic ... spark.profile.render(id2, typememory, rendererdo_render) spark.profile.clear(id2, typememory)源码佐证SparkSession 级 Profiler APIspark.profile是 PySpark 4.0 引入的会话级剖析入口实现在 python/pyspark/sql/profiler.pyWorkerMemoryProfilerprofiler.py基于UDFLineProfilerV2按调用次数统计逐行内存并将结果经_ProfileResultsParamV2累加器SQL_UDF_PROFIER_V2通道聚合回 DriverProfile类profiler.py提供show/dump/render/clear四个方法type参数取perf或memory未指定时默认同时展示/清理两类结果剖析数据通过 Spark Accumulator 机制从各 executor 汇总因此show展示的是所有 Python 执行任务的聚合结果例如 8 个 task 各处理 1000 行则展示 8000 行的累计统计。定位热点Python 内置 ProfilerscProfilePython 标准库的Python Profilers提供确定性剖析能给出丰富的统计信息用于定位 Driver 侧与 UDF 侧的昂贵热点代码路径。Driver 侧当作普通 Python 程序剖析Driver 侧本质是普通 Python 进程除非运行在 YARN 集群模式等远端因此按常规方式使用即可echo from pyspark.sql import SparkSession spark SparkSession.builder.getOrCreate() spark.range(10).show() app.pypython -m cProfile app.py输出示例... 129215 function calls (125446 primitive calls) in 5.926 seconds Ordered by: standard name ncalls tottime percall cumtime percall filename:lineno(function) 1198/405 0.001 0.000 0.083 0.000 frozen importlib._bootstrap:1009(_handle_fromlist) 561 0.001 0.000 0.001 0.000 frozen importlib._bootstrap:103(release) 276 0.000 0.000 0.000 0.000 frozen importlib._bootstrap:143(__init__) 276 0.000 0.000 0.002 0.000 frozen importlib._bootstrap:147(__enter__) ...Python/Pandas/Arrow UDF 性能剖析内置远程 Profiler与内存剖析类似PySpark 为 UDF 提供内置的远程 Python Profilers同样不支持生成器函数。将spark.sql.pyspark.udf.profiler设为perf即可 from pyspark.sql.functions import pandas_udf df spark.range(10) pandas_udf(long) ... def add1(x): ... return x 1 ... added df.select(add1(id)) spark.conf.set(spark.sql.pyspark.udf.profiler, perf) added.show() -------- |add1(id)| -------- ... -------- spark.profile.show(typeperf) Profile of UDFid2 2300 function calls (2270 primitive calls) in 0.006 seconds Ordered by: internal time, cumulative time ncalls tottime percall cumtime percall filename:lineno(function) 10 0.001 0.000 0.005 0.001 series.py:5515(_arith_method) 10 0.001 0.000 0.001 0.000 _ufunc_config.py:425(__init__) 10 0.000 0.000 0.000 0.000 {built-in method _operator.add} 10 0.000 0.000 0.002 0.000 series.py:315(__init__) ...UDF ID 同样可以从查询计划ArrowEvalPython节点中读取 added.explain() Physical Plan *(2) Project [pythonUDF0#11L AS add1(id)#3L] - ArrowEvalPython [add1(id#0L)#2L], [pythonUDF0#11L], 200 - *(1) Range (0, 10, step1, splits16)render方法默认使用flameprof渲染器输出火焰图在 IPython 环境返回IPython.display.HTML否则返回 SVG 源字符串也可传入自定义渲染函数 spark.profile.render(id2, typeperf) # rendererflameprof by default def do_render(stats): ... # Your custom rendering logic ... ... ... spark.profile.render(id2, typeperf, rendererdo_render)清理结果 spark.profile.clear(id2, typeperf)源码佐证UDF 性能剖析的实现WorkerPerfProfilerpython/pyspark/sql/profiler.py在 worker 侧基于cProfile.Profile启停采集save()时通过pstats.Stats生成可 picklable 的统计并写入累加器_render_flameprofprofiler.py将pstats渲染为 SVG且_renderers注册表中默认注册了(perf, None)与(perf, flameprof)两条记录——这正是rendererflameprof by default的由来。使用火焰图前需安装flameprof版本不低于 0.4否则会抛出PACKAGE_NOT_INSTALLED错误。常见异常与错误排查手册PySpark SQL 层异常AnalysisException——分析 SQL 查询计划失败时抛出典型场景是引用了不存在的列 df spark.range(1) df[bad_key] Traceback (most recent call last): ... pyspark.errors.exceptions.AnalysisException: Cannot resolve column name bad_key among (id)解决改用存在的列名。 df[id] ColumnidParseException——解析 SQL 命令失败时抛出 spark.sql(select * 1) Traceback (most recent call last): ... pyspark.errors.exceptions.ParseException: [PARSE_SYNTAX_ERROR] Syntax error at or near 1: extra input 1.(line 1, pos 9) SQL select * 1 ---------^^^解决修正 SQL 语法。 spark.sql(select *) DataFrame[]IllegalArgumentException——传入非法或不合适的参数时抛出 spark.range(1).sample(-1.0) Traceback (most recent call last): ... pyspark.errors.exceptions.IllegalArgumentException: requirement failed: Sampling fraction (-1.0) must be on interval [0, 1] without replacement解决传参需满足取值范围约束例如采样比例必须在[0, 1]区间。 spark.range(1).sample(1.0) DataFrame[id: bigint]PythonException——来自 Python worker 的异常外层能看到 worker 中抛出的异常类型与堆栈如下例中的TypeError import pyspark.sql.functions as sf from pyspark.sql.functions import udf def f(x): ... return sf.abs(x) ... spark.range(-1, 1).withColumn(abs, udf(f)(id)).collect() 22/04/12 14:52:31 ERROR Executor: Exception in task 7.0 in stage 37.0 (TID 232) org.apache.spark.api.python.PythonException: Traceback (most recent call last): ... TypeError: Invalid argument, not a string or column: -1 of type class int. For column literals, use lit, array, struct or create_map function.解决在 UDF 内部使用纯 Python 逻辑不要混用 DataFrame 列函数。 def f(x): ... return abs(x) ... spark.range(-1, 1).withColumn(abs, udf(f)(id)).collect() [Row(id-1, abs1), Row(id0, abs0)]StreamingQueryException——StreamingQuery 失败时抛出多数情况下它由 Python worker 抛出并包装为PythonException sdf spark.readStream.format(text).load(python/test_support/sql/streaming) from pyspark.sql.functions import col, udf bad_udf udf(lambda x: 1 / 0) (sdf.select(bad_udf(col(value))).writeStream.format(memory).queryName(q1).start()).processAllAvailable() Traceback (most recent call last): ... org.apache.spark.api.python.PythonException: Traceback (most recent call last): File stdin, line 1, in lambda ZeroDivisionError: division by zero ... pyspark.errors.exceptions.StreamingQueryException: [STREAM_FAILED] Query [id 74eb53a8-89bd-49b0-9313-14d29eed03aa, runId 9f2d5cf6-a373-478d-b718-2c2b6d8a0f24] terminated with exception: Job aborted解决修复 StreamingQuery 中的错误并重新执行工作流。SparkUpgradeException——因 Spark 升级而抛出典型如 Spark 3.0 起对日期时间格式解析行为的变化 from pyspark.sql.functions import to_date, unix_timestamp, from_unixtime df spark.createDataFrame([(2014-31-12,)], [date_str]) df2 df.select(date_str, to_date(from_unixtime(unix_timestamp(date_str, yyyy-dd-aa)))) df2.collect() Traceback (most recent call last): ... pyspark.sql.utils.SparkUpgradeException: You may get a different result due to the upgrading to Spark 3.0: Fail to recognize yyyy-dd-aa pattern in the DateTimeFormatter. 1) You can set spark.sql.legacy.timeParserPolicy to LEGACY to restore the behavior before Spark 3.0. 2) You can form a valid datetime pattern with the guide from https://spark.apache.org/docs/latest/sql-ref-datetime-pattern.html解决将spark.sql.legacy.timeParserPolicy设为LEGACY恢复旧行为或改用符合新规范的合法时间格式。 spark.conf.set(spark.sql.legacy.timeParserPolicy, LEGACY) df2 df.select(date_str, to_date(from_unixtime(unix_timestamp(date_str, yyyy-dd-aa)))) df2.collect() [Row(date_str2014-31-12, to_date(from_unixtime(unix_timestamp(date_str, yyyy-dd-aa), yyyy-MM-dd HH:mm:ss))None)]pandas API on Spark 常见异常ValueError: Cannot combine the series or dataframe because it comes from a different dataframe当一次操作涉及多个不同 DataFrame 的 Series/DataFrame且compute.ops_on_diff_frames处于禁用状态默认禁用时会抛出该错误。此类操作因需要 join 底层 Spark 帧而代价高昂仅在必要时才应开启该选项 ps.Series([1, 2]) ps.Series([3, 4]) Traceback (most recent call last): ... ValueError: Cannot combine the series or dataframe because it comes from a different dataframe. In order to allow this operation, enable compute.ops_on_diff_frames option.解决在option_context中显式开启后执行 with ps.option_context(compute.ops_on_diff_frames, True): ... ps.Series([1, 2]) ps.Series([3, 4]) ... 0 4 1 6 dtype: int64PySparkRuntimeError: [RESULT_ROWS_MISMATCH] The number of output rows must match the number of input rowstransform等要求输入输出行数一致的操作返回行数不一致时抛出 def f(x) - ps.Series[np.int32]: ... return x[:-1] ... ps.DataFrame({x:[1, 2], y:[3, 4]}).transform(f) 22/04/12 13:46:39 ERROR Executor: Exception in task 2.0 in stage 16.0 (TID 88) org.apache.spark.api.python.PythonException: Traceback (most recent call last): ... pyspark.errors.exceptions.base.PySparkRuntimeError: [RESULT_ROWS_MISMATCH] The number of output rows (0) must match the number of input rows (1).解决确保变换函数保持行数一致 def f(x) - ps.Series[np.int32]: ... return x ... ps.DataFrame({x:[1, 2], y:[3, 4]}).transform(f) x y 0 1 3 1 2 4Py4J 层异常Py4JJavaError——Java 客户端代码中发生异常时抛出可以看到 Java 侧抛出的异常类型及堆栈如java.lang.NullPointerException spark.sparkContext._jvm.java.lang.String(None) Traceback (most recent call last): ... py4j.protocol.Py4JJavaError: An error occurred while calling None.java.lang.String. : java.lang.NullPointerException ..解决传入合法的参数。 spark.sparkContext._jvm.java.lang.String(x) xPy4JError——其他类错误时抛出例如 Python 客户端访问了一个在 Java 侧已不存在的对象 from pyspark.ml.linalg import Vectors from pyspark.ml.regression import LinearRegression df spark.createDataFrame( ... [(1.0, 2.0, Vectors.dense(1.0)), (0.0, 2.0, Vectors.sparse(1, [], []))], ... [label, weight, features], ... ) lr LinearRegression( ... maxIter1, regParam0.0, solvernormal, weightColweight, fitInterceptFalse ... ) model lr.fit(df) model LinearRegressionModel: uidLinearRegression_eb7bc1d4bf25, numFeatures1 model.__del__() model Traceback (most recent call last): ... py4j.protocol.Py4JError: An error occurred while calling o531.toString. Trace: py4j.Py4JException: Target Object ID does not exist for this gateway :o531 ...解决访问仍存在于 Java 侧的对象。Py4JNetworkError——网络传输出现问题时抛出如连接断开。此时应排查网络链路并重建连接。控制堆栈追踪的两个关键配置Spark 提供两个配置项控制异常堆栈的展示方式spark.sql.execution.pyspark.udf.simplifiedTraceback.enabled默认 true简化 Python UDF 与数据源Data Sources的 traceback避免输出冗长的内部调用链spark.sql.pyspark.jvmStacktrace.enabled默认 false隐藏 JVM 堆栈只展示 Python 友好的异常信息方便直接定位到 Python 侧的错误位置。需要注意上述配置与日志级别设置相互独立。日志级别请通过pyspark.SparkContext.setLogLevel控制。总结一条完整的 PySpark 排查链路把以上方法串联起来即可形成一套可落地的排查流程进程定位用os.getpid()ps定位 Driver 进程用ps -fe | grep pyspark.daemon定位 executor worker 进程配合top观察资源占用断点调试Driver 侧在脚本顶部插入pydevd_pycharm.settrace并配合spark-submitExecutor 侧通过自定义模块 spark.python.daemon.module配置包装 worker 入口spark.python.worker.module也可类似地替换 worker 模块实现两端断点性能剖析Driver 侧用memory_profiler逐行统计内存、cProfile统计 CPU 热点UDF 侧设置运行时配置spark.sql.pyspark.udf.profiler为memory或perf通过spark.profile的show/dump/render/clear查看、导出、可视化火焰图并清理结果异常排除按异常类型AnalysisException、ParseException、IllegalArgumentException、PythonException、StreamingQueryException、SparkUpgradeException、Py4J 系列等对照错误信息定位根因必要时调整spark.sql.legacy.timeParserPolicy、compute.ops_on_diff_frames、spark.sql.execution.pyspark.udf.simplifiedTraceback.enabled与spark.sql.pyspark.jvmStacktrace.enabled等配置。上述所有配置与 API 均可在当前仓库中找到实现依据配置定义、daemon 加载逻辑、Profile API 实现并结合 测试用例 与 内存剖析测试 进一步验证使用方式。赞分享大数据数据分析批处理流处理机器学习图计算【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址https://gitcode.com/gh_mirrors/sp/spark点击查看免费下载相关推荐5分钟快速上手TDengine物联网时序数据库的完整部署指南 5分钟快速上手TDengine物联网时序数据库的完整部署指南 你是否正在为海量物联网设备数据存储而烦恼面对每秒数万条传感器数据传统数据库早已不堪重负数据库时序数据库大数据物联网云原生gVisor 沙箱调试完全指南从日志、strace 到栈转储与 pprof 性能剖析gVisor 沙箱调试完全指南从日志、strace 到栈转储与 pprof 性能剖析 本文是面向 gVisorApplication Kernel for云原生容器运行时操作系统应用安全Cilium 调试完全指南从 Delve 附加调试到 toFQDNs 排障、锁竞争分析与 pprof 性能剖析Cilium 调试完全指南从 Delve 附加调试到 toFQDNs 排障、锁竞争分析与 pprof 性能剖析 本文以 Cilium 官方调试文档为核心系统云原生网络服务网格可观测性网络安全eBPF上一篇MyTinySTL中的函数绑定参数重排与占位符下一篇Bilibili-EvolvedTypeScript接口与类型别名区别与使用创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表