《Spark-shell 交互式使用技巧》
交互式环境 spark-shell 是 Spark 开发与调试最常用的入口。本文整理启动参数、依赖 jar 的引入方式以及 :load 脚本执行法。
1 启动时指定依赖的 jar 包
spark-shell \
--jars ./espark-connector-2.11-5.5.3.jar,./my-spark-hello-1.0-SNAPSHOT-jar-with-dependencies.jar \
--files client.keystore.jks,client.truststore.jks \
--conf spark.serializer=org.apache.spark.serializer.KryoSerializer \
--conf spark.streaming.concurrentJobs=8 \
--conf spark.scheduler.mode=FIFO \
--driver-memory 16G \
--conf spark.driver.maxResultSize=15g \
--conf spark.hadoop.hadoop.log.dir=./logs \
master yarn --num-executors=20 --executor-cores=4
注意:多个 jar 使用逗号 , 分隔,中间不能有空格。
2 查看配置,确认 jar 被引入
scala> :settings
// 当前的配置
-Yrepl-class-based = true
-Yrepl-outdir = /tmp/spark-.../repl-...
-classpath = /home/disk3/work/espark-connector-2.11-5.5.3.jar:/home/disk3/work/my-spark-hello-1.0-SNAPSHOT-jar-with-dependencies.jar
-d = .
-encoding = UTF-8
-nowarn = false
3 通过 :load 执行本目录下的 scala 脚本
scala> :load cmd.scala
cmd.scala 内容示例(批量写入 Elasticsearch):
import org.apache.spark._
import org.apache.spark.sql
import org.apache.spark.sql._
import org.apache.spark.sql.SaveMode.Append
import org.elasticsearch.spark._
val resourceName = "my-index/doc"
val esOptions = Map(
"es.nodes" -> "http://es-server.example.com:8200",
"es.net.http.auth.user" -> "user",
"es.net.http.auth.pass" -> "pass",
"es.resource" -> resourceName,
"es.index.auto.create" -> "true",
"es.http.timeout" -> "10m",
"es.http.retries" -> "5",
"es.batch.size.bytes" -> "20mb",
"es.batch.size.entries" -> "100000",
"es.batch.write.refresh" -> "false",
"es.batch.write.retry.count" -> "10",
"es.batch.write.retry.wait" -> "15s"
)
// 需要执行的 sql 语句
val ipCntSqlQuery = """
SELECT
cookies['uid'] AS uid,
COUNT(DISTINCT cli_ip) AS count_ip
FROM
bdm.bdm_sample_log
WHERE
day = '20200101'
GROUP BY cookies['uid']
HAVING count_ip > 1
"""
// ipCntDf.show
ipCntDf.write.format("org.elasticsearch.spark.sql").options(esOptions).mode(SaveMode.Append).save
4 常见错误:Java heap space
报错示例:
org.apache.spark.SparkException: Job aborted due to stage failure: Task 14 in stage 3.0 failed 1 times,
most recent failure: Lost task 14.0 in stage 3.0 (TID 4014, localhost, executor driver):
java.lang.OutOfMemoryError: Java heap space
解决方案:在启动 spark-shell 时适当调大 driver 的内存。
spark-shell \
--driver-memory 16G \
--conf spark.driver.maxResultSize=15g
5 flow-shell 支持的启动参数
Usage: ./bin/spark-shell [options]
Options:
--master MASTER_URL spark://host:port, mesos://host:port, yarn,
k8s://https://host:port, or local (Default: local[*]).
--deploy-mode DEPLOY_MODE Whether to launch the driver program locally ("client") or
on one of the worker machines inside the cluster ("cluster")
(Default: client).
--class CLASS_NAME Your application's main class (for Java / Scala apps).
--name NAME A name of your application.
--jars JARS Comma-separated list of jars to include on the driver
and executor classpaths.
--packages Comma-separated list of maven coordinates of jars to include
on the driver and executor classpaths.
--repositories Comma-separated list of additional remote repositories to
search for the maven coordinates given with --packages.
--py-files PY_FILES Comma-separated list of .zip, .egg, or .py files to place
on the PYTHONPATH for Python apps.
--files FILES Comma-separated list of files to be placed in the working
directory of each executor.
--conf PROP=VALUE Arbitrary Spark configuration property.
--properties-file FILE Path to a file from which to load extra properties.
--driver-memory MEM Memory for driver (e.g. 1000M, 2G) (Default: 1024M).
--driver-java-options Extra Java options to pass to the driver.
--driver-library-path Extra library path entries to pass to the driver.
--driver-class-path Extra class path entries to pass to the driver.
--executor-memory MEM Memory per executor (e.g. 1000M, 2G) (Default: 1G).
--proxy-user NAME User to impersonate when submitting the application.
--verbose, -v Print additional debug output.
Cluster deploy mode only:
--driver-cores NUM Number of cores used by the driver, only in cluster mode.
Spark standalone or Mesos with cluster deploy mode only:
--supervise If given, restarts the driver on failure.
Spark standalone and YARN only:
--executor-cores NUM Number of cores per executor.
YARN-only:
--queue QUEUE_NAME The YARN queue to submit to (Default: "default").
--num-executors NUM Number of executors to launch (Default: 2).
--principal PRINCIPAL Principal to be used to login to KDC, while running on secure HDFS.
--keytab KEYTAB The full path to the file that contains the keytab.
6 交互式环境中支持的命令
scala> :help
:edit <id>|<line> edit history
:help [command] print this summary or command-specific help
:history [num] show the history
:imports show import history, identifying sources of names
:javap <path|class> disassemble a file or class name
:load <path> interpret lines in a file
:paste [-raw] [path] enter paste mode or paste a file
:quit exit the interpreter
:replay [options] reset the repl and replay all previous commands
:require <path> add a jar to the classpath
:reset [options] reset the repl to its initial state
:save <path> save replayable session to a file
:sh <command line> run a shell command
:settings <options> update compiler options, if possible; see reset
:silent disable/enable automatic printing of results
:type [-v] <expr> display the type of an expression without evaluating it
:warnings show the suppressed warnings from the most recent line
阅读 —
·
全站 —