如何在Jupyter Notebook(Python3.5内核)中集群模式运行PySpark应用?
Hey there! Let's get your PySpark application scaling across all 6 nodes in your cluster. You're already on the right track with setting master("yarn")—here's the full breakdown of code tweaks and Jupyter configurations you need:
1. Update Your PySpark Initialization Code
The key is to explicitly configure Spark to use YARN with sufficient resources, and ensure consistency between your Jupyter Python environment and the cluster's.
Example Spark 2.x+ (Recommended: Use SparkSession)
SparkSession simplifies context management and is the standard for newer versions:
from pyspark.sql import SparkSession # Configure Spark for YARN with cluster resource settings conf = SparkConf() \ .setMaster("yarn") \ .setAppName("test") \ .set("spark.executor.instances", "5") # Adjust based on your cluster (e.g., 5 executors for 6 nodes, leaving one for driver) .set("spark.executor.cores", "2") # Number of CPU cores per executor (match your node's capacity) .set("spark.executor.memory", "4g") # Memory allocated to each executor .set("spark.driver.memory", "2g") # Memory for the Jupyter driver process .set("spark.pyspark.python", "/usr/bin/python3.5") # Path to Python 3.5 on cluster nodes (critical for version consistency) # Initialize SparkSession, SparkContext, and SQLContext spark = SparkSession.builder.config(conf=conf).getOrCreate() sc = spark.sparkContext sqlContext = spark.sqlContext
For Older Spark Versions (Spark 1.x)
If you're stuck with Spark 1.x, use SparkContext directly:
from pyspark import SparkConf, SparkContext from pyspark.sql import SQLContext conf = SparkConf() \ .setMaster("yarn") \ .setAppName("test") \ .set("spark.executor.instances", "5") \ .set("spark.executor.cores", "2") \ .set("spark.executor.memory", "4g") \ .set("spark.driver.memory", "2g") \ .set("spark.pyspark.python", "/usr/bin/python3.5") sc = SparkContext(conf=conf) sqlContext = SQLContext(sc)
Important Notes:
- Adjust
spark.executor.instances,cores, andmemoryvalues based on your cluster's hardware (e.g., if each node has 8 cores and 16GB RAM, you could run 2 executors per node with 4 cores and 6GB RAM each). - If your cluster nodes don't have Python 3.5 installed, you'll need to package your local Python environment and distribute it via
spark.yarn.dist.filesor use a cluster-wide Python environment manager like Conda.
2. Configure Jupyter Notebook Environment
Ensure Jupyter knows where to find Spark and uses the correct Python version:
Option 1: Set Environment Variables Before Starting Jupyter
Run these commands in your terminal before launching Jupyter:
export SPARK_HOME="/path/to/your/spark/installation" export PYSPARK_PYTHON="/usr/bin/python3.5" export PYSPARK_DRIVER_PYTHON="jupyter" export PYSPARK_DRIVER_PYTHON_OPTS="notebook" # Launch Jupyter jupyter notebook
Option 2: Set Variables Inside Jupyter
Add this to the first cell of your notebook to configure the environment on-the-fly:
import os # Point to your Spark installation os.environ['SPARK_HOME'] = '/path/to/your/spark/installation' # Ensure Jupyter uses the same Python version as the cluster os.environ['PYSPARK_PYTHON'] = '/usr/bin/python3.5'
3. Verify Cluster Connectivity
- Make sure the YARN Resource Manager (default port 8088) is accessible from the node running Jupyter.
- Ensure HDFS is running and properly configured—Spark on YARN relies on HDFS to store intermediate data and application jars.
- Check the YARN UI after launching your app to confirm executors are spawned across multiple nodes.
Once you've made these changes, your PySpark application should distribute tasks across all available nodes in your cluster!
内容的提问来源于stack exchange,提问作者Ravi Kiran G

