本文目录导读:

我来详细说明如何批处理分布式Spark任务:
Spark任务批处理基础架构
1 任务提交脚本模板
#!/bin/bash
# submit_spark_jobs.sh
# Spark配置
SPARK_HOME=/opt/spark
MASTER=spark://master-node:7077
DEPLOY_MODE=cluster
# 任务队列文件
JOB_QUEUE="job_queue.txt"
# 处理单个任务
submit_job() {
local job_name=$1
local main_class=$2
local jar_path=$3
local config_file=$4
echo "[$(date '+%Y-%m-%d %H:%M:%S')] 提交任务: $job_name"
$SPARK_HOME/bin/spark-submit \
--master $MASTER \
--deploy-mode $DEPLOY_MODE \
--name "$job_name" \
--class "$main_class" \
--conf spark.default.parallelism=200 \
--conf spark.executor.memory=4g \
--conf spark.executor.cores=4 \
--conf spark.driver.memory=2g \
--conf spark.speculation=true \
--conf spark.yarn.maxAppAttempts=2 \
$jar_path \
--config $config_file
local status=$?
if [ $status -eq 0 ]; then
echo "[$(date '+%Y-%m-%d %H:%M:%S')] 任务成功: $job_name"
else
echo "[$(date '+%Y-%m-%d %H:%M:%S')] 任务失败: $job_name"
return 1
fi
}
# 批量提交任务
batch_submit_jobs() {
local job_queue=$1
local max_concurrent=${2:-3} # 最大并发任务数
# 读取任务队列
while IFS=',' read -r job_name main_class jar_path config_file; do
# 跳过注释和空行
[[ "$job_name" =~ ^#.*$ ]] && continue
[[ -z "$job_name" ]] && continue
# 等待资源
while [ $(jobs -r | wc -l) -ge $max_concurrent ]; do
sleep 5
done
# 后台提交任务
submit_job "$job_name" "$main_class" "$jar_path" "$config_file" &
done < "$job_queue"
# 等待所有后台任务完成
wait
echo "所有任务已完成"
}
2 任务队列文件格式
# job_queue.txt # 格式: job_name,main_class,jar_path,config_file data_ingestion,com.example.DataIngestionJob,/app/jobs/ingestion.jar,/app/config/ingestion.conf data_cleaning,com.example.DataCleaningJob,/app/jobs/cleaning.jar,/app/config/cleaning.conf feature_engineering,com.example.FeatureEngineeringJob,/app/jobs/features.jar,/app/config/features.conf model_training,com.example.ModelTrainingJob,/app/jobs/training.jar,/app/config/training.conf
Python批处理框架
1 任务调度管理器
# spark_batch_manager.py
import os
import sys
import time
import json
import logging
import subprocess
from datetime import datetime
from typing import List, Dict, Any, Optional
from concurrent.futures import ThreadPoolExecutor, as_completed
class SparkJob:
"""Spark任务类"""
def __init__(self, job_config: Dict):
self.name = job_config['name']
self.main_class = job_config['main_class']
self.jar_path = job_config['jar_path']
self.config = job_config.get('config', {})
self.dependencies = job_config.get('dependencies', [])
self.retry_count = job_config.get('retry_count', 3)
self.timeout = job_config.get('timeout', 3600)
self.status = 'pending'
self.start_time = None
self.end_time = None
self.log_file = f"logs/{self.name}_{datetime.now().strftime('%Y%m%d_%H%M%S')}.log"
class BatchSparkManager:
"""批处理Spark任务管理器"""
def __init__(self, config_file: str):
self.logger = self._setup_logging()
self.config = self._load_config(config_file)
self.spark_home = self.config.get('spark_home', '/opt/spark')
self.master = self.config.get('master', 'yarn')
self.deploy_mode = self.config.get('deploy_mode', 'cluster')
self.max_concurrent = self.config.get('max_concurrent_jobs', 5)
def _setup_logging(self):
"""设置日志"""
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
handlers=[
logging.FileHandler('spark_batch_manager.log'),
logging.StreamHandler()
]
)
return logging.getLogger(__name__)
def _load_config(self, config_file: str) -> Dict:
"""加载配置文件"""
with open(config_file, 'r') as f:
return json.load(f)
def submit_job(self, job: SparkJob) -> bool:
"""提交单个Spark任务"""
cmd = [
f"{self.spark_home}/bin/spark-submit",
"--master", self.master,
"--deploy-mode", self.deploy_mode,
"--name", job.name,
"--class", job.main_class,
f"--conf", "spark.executor.memory=4g",
f"--conf", "spark.executor.cores=4",
f"--conf", "spark.driver.memory=2g",
job.jar_path
]
# 添加额外的配置
for key, value in job.config.items():
cmd.extend([f"--conf", f"{key}={value}"])
job.start_time = datetime.now()
job.status = 'running'
self.logger.info(f"提交任务: {job.name}")
self.logger.info(f"命令: {' '.join(cmd)}")
try:
with open(job.log_file, 'w') as f:
process = subprocess.Popen(
cmd,
stdout=f,
stderr=subprocess.STDOUT,
universal_newlines=True
)
# 等待任务完成或超时
process.wait(timeout=job.timeout)
if process.returncode == 0:
job.status = 'completed'
self.logger.info(f"任务完成: {job.name}")
return True
else:
job.status = 'failed'
self.logger.error(f"任务失败: {job.name}, 返回码: {process.returncode}")
return False
except subprocess.TimeoutExpired:
process.kill()
job.status = 'timeout'
self.logger.error(f"任务超时: {job.name}")
return False
except Exception as e:
job.status = 'failed'
self.logger.error(f"任务异常: {job.name}, 错误: {str(e)}")
return False
finally:
job.end_time = datetime.now()
def batch_submit(self, jobs: List[SparkJob]):
"""批量提交任务"""
# 按依赖关系排序
sorted_jobs = self._topological_sort(jobs)
completed_jobs = set()
with ThreadPoolExecutor(max_workers=self.max_concurrent) as executor:
future_to_job = {}
while sorted_jobs:
# 查找没有依赖或依赖已完成的作业
ready_jobs = []
for job in sorted_jobs:
if all(dep in completed_jobs for dep in job.dependencies):
if job.status == 'pending':
ready_jobs.append(job)
if not ready_jobs:
if not future_to_job:
self.logger.error("死锁检测: 所有任务都在等待依赖完成")
break
else:
# 等待正在运行的任务完成
pass
else:
# 提交就绪的任务
for job in ready_jobs:
future = executor.submit(self.submit_job_with_retry, job)
future_to_job[future] = job
sorted_jobs.remove(job)
# 等待完成的任务
for future in as_completed(future_to_job):
job = future_to_job[future]
if future.result():
completed_jobs.add(job.name)
self.logger.info(f"添加完成的任务: {job.name}")
del future_to_job[future]
def submit_job_with_retry(self, job: SparkJob) -> bool:
"""带重试的任务提交"""
for attempt in range(1, job.retry_count + 1):
self.logger.info(f"尝试 {attempt}/{job.retry_count}: {job.name}")
if self.submit_job(job):
return True
if attempt < job.retry_count:
wait_time = 2 ** attempt # 指数退避
self.logger.info(f"等待 {wait_time} 秒后重试")
time.sleep(wait_time)
return False
def _topological_sort(self, jobs: List[SparkJob]) -> List[SparkJob]:
"""拓扑排序处理依赖关系"""
# 简化实现,实际应考虑循环依赖检测
return jobs
2 配置文件示例
{
"spark_home": "/opt/spark",
"master": "yarn",
"deploy_mode": "cluster",
"max_concurrent_jobs": 5,
"jobs": [
{
"name": "data_ingestion",
"main_class": "com.example.DataIngestionJob",
"jar_path": "/app/jobs/ingestion.jar",
"dependencies": [],
"retry_count": 3,
"timeout": 3600,
"config": {
"spark.executor.memory": "8g",
"spark.executor.cores": "8"
}
},
{
"name": "data_cleaning",
"main_class": "com.example.DataCleaningJob",
"jar_path": "/app/jobs/cleaning.jar",
"dependencies": ["data_ingestion"],
"retry_count": 2,
"timeout": 7200
},
{
"name": "feature_engineering",
"main_class": "com.example.FeatureEngineeringJob",
"jar_path": "/app/jobs/features.jar",
"dependencies": ["data_cleaning"],
"retry_count": 2,
"timeout": 5400
}
]
}
3 使用示例
# main_batch.py
from spark_batch_manager import BatchSparkManager, SparkJob
def main():
# 初始化批处理管理器
manager = BatchSparkManager('batch_config.json')
# 从配置文件创建任务
jobs = []
for job_config in manager.config['jobs']:
job = SparkJob(job_config)
jobs.append(job)
# 执行批处理
manager.batch_submit(jobs)
if __name__ == "__main__":
main()
任务编排与监控
1 Airflow集成
# spark_dag.py
from airflow import DAG
from airflow.contrib.operators.spark_submit_operator import SparkSubmitOperator
from datetime import datetime, timedelta
default_args = {
'owner': 'data_team',
'depends_on_past': False,
'start_date': datetime(2024, 1, 1),
'email_on_failure': True,
'email_on_retry': False,
'retries': 3,
'retry_delay': timedelta(minutes=5)
}
dag = DAG(
'spark_batch_pipeline',
default_args=default_args,
description='Spark批处理流水线',
schedule_interval='0 2 * * *', # 每天凌晨2点运行
catchup=False
)
# 数据导入任务
data_ingestion = SparkSubmitOperator(
task_id='data_ingestion',
application='/app/jobs/ingestion.jar',
conn_id='spark_default',
java_class='com.example.DataIngestionJob',
conf={
'spark.executor.memory': '8g',
'spark.executor.cores': '4',
'spark.sql.shuffle.partitions': '200'
},
driver_memory='4g',
executor_memory='8g',
executor_cores=4,
num_executors=10,
dag=dag
)
# 数据清洗任务
data_cleaning = SparkSubmitOperator(
task_id='data_cleaning',
application='/app/jobs/cleaning.jar',
conn_id='spark_default',
java_class='com.example.DataCleaningJob',
conf={
'spark.executor.memory': '8g',
'spark.sql.shuffle.partitions': '200'
},
executor_memory='8g',
executor_cores=4,
num_executors=10,
dag=dag
)
# 特征工程任务
feature_engineering = SparkSubmitOperator(
task_id='feature_engineering',
application='/app/jobs/features.jar',
conn_id='spark_default',
java_class='com.example.FeatureEngineeringJob',
dag=dag
)
# 模型训练任务
model_training = SparkSubmitOperator(
task_id='model_training',
application='/app/jobs/training.jar',
conn_id='spark_default',
java_class='com.example.ModelTrainingJob',
executor_memory='16g',
executor_cores=8,
num_executors=20,
dag=dag
)
# 设置依赖关系
data_ingestion >> data_cleaning >> feature_engineering >> model_training
2 监控仪表板
# spark_monitor.py
import yaml
import requests
from datetime import datetime
from tabulate import tabulate
class SparkMonitor:
"""Spark任务监控器"""
def __init__(self, history_server_url: str = "http://localhost:18080"):
self.history_server_url = history_server_url
def get_active_applications(self) -> List[Dict]:
"""获取活动应用"""
response = requests.get(f"{self.history_server_url}/api/v1/applications")
return response.json() if response.status_code == 200 else []
def get_application_details(self, app_id: str) -> Dict:
"""获取应用详情"""
response = requests.get(
f"{self.history_server_url}/api/v1/applications/{app_id}"
)
return response.json() if response.status_code == 200 else {}
def display_jobs_status(self):
"""显示任务状态"""
apps = self.get_active_applications()
table_data = []
for app in apps:
app_id = app['id']
app_name = app['name']
start_time = datetime.fromtimestamp(app['attempts'][0]['startTimeEpoch']/1000)
end_time = datetime.fromtimestamp(app['attempts'][0]['endTimeEpoch']/1000)
duration = (end_time - start_time).total_seconds() if app['attempts'][0]['completed'] else "运行中"
table_data.append([
app_id,
app_name,
start_time.strftime('%Y-%m-%d %H:%M:%S'),
str(duration) if isinstance(duration, float) else duration,
app['attempts'][0]['sparkUser']
])
headers = ['App ID', 'Name', 'Start Time', 'Duration', 'User']
print(tabulate(table_data, headers=headers, tablefmt='grid'))
def check_failed_applications(self):
"""检查失败的应用"""
apps = self.get_active_applications()
failed_apps = [app for app in apps if any(
attempt['completed'] and attempt['endTime'] == 'failed'
for attempt in app['attempts']
)]
if failed_apps:
print("失败的应用:")
for app in failed_apps:
print(f" - {app['name']} ({app['id']})")
# 使用示例
monitor = SparkMonitor()
monitor.display_jobs_status()
monitor.check_failed_applications()
最佳实践
1 资源优化策略
# resource_optimizer.py
class SparkResourceOptimizer:
"""Spark资源优化器"""
@staticmethod
def calculate_optimal_resources(data_size_gb: float,
total_cores: int,
total_memory_gb: int) -> Dict:
"""计算最优资源分配"""
# 估算所需分区数
estimated_partitions = max(200, int(data_size_gb * 10))
# 每个执行器的核心数(建议4-5)
executor_cores = min(5, total_cores // (total_memory_gb // 4))
# 每个执行器的内存(建议4-8GB)
executor_memory_gb = min(8, max(4, total_memory_gb // (total_cores // executor_cores)))
# 执行器数量
num_executors = min(
total_cores // executor_cores,
total_memory_gb // executor_memory_gb
)
return {
'num_executors': num_executors,
'executor_cores': executor_cores,
'executor_memory': f"{executor_memory_gb}g",
'driver_memory': f"{max(2, executor_memory_gb // 2)}g",
'shuffle_partitions': estimated_partitions
}
2 错误处理与重试策略
# retry_handler.py
class RetryHandler:
"""重试处理器"""
def __init__(self, max_retries: int = 3, backoff_factor: float = 2.0):
self.max_retries = max_retries
self.backoff_factor = backoff_factor
def execute_with_retry(self, func, *args, **kwargs):
"""带重试的执行"""
last_exception = None
for attempt in range(self.max_retries):
try:
return func(*args, **kwargs)
except Exception as e:
last_exception = e
wait_time = self.backoff_factor ** attempt
print(f"尝试 {attempt + 1} 失败: {str(e)}")
print(f"等待 {wait_time} 秒后重试...")
time.sleep(wait_time)
raise last_exception
# 使用示例
retry_handler = RetryHandler(max_retries=3)
try:
result = retry_handler.execute_with_retry(
spark_manager.submit_job,
job_object
)
except Exception as e:
print(f"最终失败: {str(e)}")
# 发送告警
alert_system.send_alert(f"Spark任务失败: {str(e)}")
这些方案可以帮助你构建高效、可靠的Spark批处理系统,根据实际需求选择合适的工具和策略。