批处理分布式Spark任务

wen java案例 3

本文目录导读:

批处理分布式Spark任务

  1. Spark任务批处理基础架构
  2. Python批处理框架
  3. 任务编排与监控
  4. 最佳实践

我来详细说明如何批处理分布式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批处理系统,根据实际需求选择合适的工具和策略。

抱歉,评论功能暂时关闭!