《Spark-shell 交互式使用技巧》

《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
阅读 — · 全站 —
🎸 我的歌单 0 首