自动配置SparkSQL引擎的脚本

wen 实用脚本 26

本文目录导读:

自动配置SparkSQL引擎的脚本

  1. 基础自动配置脚本 (auto_config_sparksql.sh)
  2. Docker环境配置脚本
  3. Python配置管理脚本
  4. 使用说明
  5. 配置优化建议

我来为您提供一个自动配置SparkSQL引擎的脚本,包含多种配置方案和环境适配:

基础自动配置脚本 (auto_config_sparksql.sh)

#!/bin/bash
# ============================================
# SparkSQL引擎自动配置脚本
# 版本: 1.0
# 适用环境: Linux/Unix (Spark 2.x/3.x)
# ============================================
# 颜色输出
RED='\033[0;31m'
GREEN='\033[0;32m'
YELLOW='\033[1;33m'
NC='\033[0m' # No Color
# 日志函数
log_info() { echo -e "${GREEN}[INFO]${NC} $1"; }
log_warn() { echo -e "${YELLOW}[WARN]${NC} $1"; }
log_error() { echo -e "${RED}[ERROR]${NC} $1"; }
# 环境检测
detect_spark_env() {
    log_info "检测Spark环境..."
    # 检查SPARK_HOME
    if [ -z "$SPARK_HOME" ]; then
        log_warn "SPARK_HOME未设置,尝试自动检测..."
        SPARK_HOME=$(find / -name "spark-submit" -type f 2>/dev/null | head -1 | xargs dirname | xargs dirname)
        if [ -z "$SPARK_HOME" ]; then
            log_error "无法找到Spark安装路径"
            exit 1
        fi
    fi
    log_info "Spark路径: $SPARK_HOME"
    export SPARK_HOME
}
# 检测系统资源
detect_system_resources() {
    log_info "检测系统资源..."
    TOTAL_MEM=$(free -g | awk '/^Mem:/{print $2}')
    CPU_CORES=$(nproc)
    DISK_SPACE=$(df -h /tmp | awk 'NR==2{print $4}' | sed 's/G//')
    log_info "总内存: ${TOTAL_MEM}GB"
    log_info "CPU核心数: ${CPU_CORES}"
    log_info "可用磁盘空间: ${DISK_SPACE}GB"
    # 自动计算Spark配置参数
    EXECUTOR_MEM=$((TOTAL_MEM * 60 / 100))  # 60% 给executor
    EXECUTOR_CORES=$((CPU_CORES / 2))
    DRIVER_MEM=$((TOTAL_MEM * 20 / 100))    # 20% 给driver
    if [ $EXECUTOR_MEM -lt 1 ]; then
        EXECUTOR_MEM=1
    fi
    if [ $DRIVER_MEM -lt 1 ]; then
        DRIVER_MEM=1
    fi
}
# 生成SparkSQL配置
generate_sparksql_config() {
    local config_file="$SPARK_HOME/conf/spark-defaults.conf"
    log_info "生成SparkSQL配置文件..."
    # 备份原配置文件
    if [ -f "$config_file" ]; then
        cp "$config_file" "${config_file}.backup.$(date +%Y%m%d_%H%M%S)"
        log_info "已备份原配置文件"
    fi
    cat > "$config_file" << EOF
# ============================================
# SparkSQL自动配置 - 生成时间: $(date)
# 根据当前系统资源自动优化
# ============================================
# 基本配置
spark.master                     yarn
spark.submit.deployMode          client
spark.app.name                   AutoConfiguredSparkSQL
# 资源分配
spark.executor.memory            ${EXECUTOR_MEM}g
spark.executor.cores             ${EXECUTOR_CORES}
spark.executor.instances         $((CPU_CORES / EXECUTOR_CORES))
spark.driver.memory              ${DRIVER_MEM}g
spark.driver.cores               2
# SparkSQL配置
spark.sql.adaptive.enabled              true
spark.sql.adaptive.coalescePartitions.enabled  true
spark.sql.adaptive.join.enabled         true
spark.sql.adaptive.skewJoin.enabled     true
# 性能优化
spark.sql.shuffle.partitions            $((CPU_CORES * 2))
spark.sql.autoBroadcastJoinThreshold    10485760
spark.sql.broadcastTimeout              600
spark.sql.files.maxPartitionBytes       536870912  # 512MB
# 动态分配
spark.dynamicAllocation.enabled         true
spark.dynamicAllocation.minExecutors    1
spark.dynamicAllocation.maxExecutors    $((CPU_CORES * 2))
spark.shuffle.service.enabled           true
# 序列化配置
spark.serializer                        org.apache.spark.serializer.KryoSerializer
spark.kryoserializer.buffer.max         256m
spark.kryoserializer.buffer             64m
# 内存管理
spark.memory.offHeap.enabled            true
spark.memory.offHeap.size               1g
spark.memory.fraction                   0.75
spark.memory.storageFraction            0.5
# CBO优化
spark.sql.cbo.enabled                   true
spark.sql.cbo.joinReorder.enabled       true
spark.sql.cbo.joinReorder.dp.threshold  5
spark.sql.cbo.starSchemaDetection       true
# 执行配置
spark.sql.execution.arrow.enabled       true
spark.sql.codegen.wholeStage            true
spark.sql.adaptive.advisoryPartitionSizeInBytes 64m
spark.sql.adaptive.coalescePartitions.minPartitionNum $((CPU_CORES * 2))
EOF
    log_info "SparkSQL配置已生成: $config_file"
}
# 配置Spark环境变量
configure_spark_env() {
    local env_file="$SPARK_HOME/conf/spark-env.sh"
    log_info "配置Spark环境变量..."
    cat > "$env_file" << EOF
#!/usr/bin/env bash
# Spark环境变量配置
export JAVA_HOME=${JAVA_HOME:-/usr/lib/jvm/java-8-openjdk-amd64}
export HADOOP_HOME=${HADOOP_HOME:-/usr/local/hadoop}
export HADOOP_CONF_DIR=${HADOOP_CONF_DIR:-\$HADOOP_HOME/etc/hadoop}
export SPARK_HOME=$SPARK_HOME
export SPARK_CONF_DIR=\$SPARK_HOME/conf
# Classpath
export SPARK_DIST_CLASSPATH=\$(\$HADOOP_HOME/bin/hadoop classpath)
# 日志配置
export SPARK_LOG_DIR=\$SPARK_HOME/logs
export SPARK_PID_DIR=\$SPARK_HOME/pids
# YARN配置
export SPARK_YARN_AM_JAVA_OPTS="-Dspark.yarn.app.container.log.dir=/var/log/spark"
EOF
    chmod +x "$env_file"
    log_info "Spark环境变量配置完成"
}
# 配置SparkSQL Hive集成
configure_hive_integration() {
    log_info "配置Hive集成..."
    # 检查Hive配置
    if [ -n "$HIVE_HOME" ]; then
        # 复制Hive配置文件
        cp "$HIVE_HOME/conf/hive-site.xml" "$SPARK_HOME/conf/" 2>/dev/null || {
            log_warn "未找到hive-site.xml,创建默认配置"
            create_hive_site_config
        }
        # 添加Hive依赖
        echo "spark.sql.catalogImplementation   hive" >> "$SPARK_HOME/conf/spark-defaults.conf"
        echo "spark.sql.hive.metastore.version   2.3.7" >> "$SPARK_HOME/conf/spark-defaults.conf"
        echo "spark.sql.hive.metastore.jars      /usr/local/hive/lib/*" >> "$SPARK_HOME/conf/spark-defaults.conf"
        log_info "Hive集成配置完成"
    else
        log_warn "HIVE_HOME未设置,跳过Hive集成配置"
    fi
}
# 创建默认Hive配置
create_hive_site_config() {
    cat > "$SPARK_HOME/conf/hive-site.xml" << EOF
<?xml version="1.0"?>
<?xml-stylesheet type="text/xsl" href="configuration.xsl"?>
<configuration>
    <property>
        <name>hive.metastore.warehouse.dir</name>
        <value>/user/hive/warehouse</value>
    </property>
    <property>
        <name>hive.metastore.local</name>
        <value>true</value>
    </property>
</configuration>
EOF
}
# 测试SparkSQL配置
test_sparksql_config() {
    log_info "测试SparkSQL配置..."
    # 创建测试SQL文件
    cat > /tmp/test_sparksql.sql << EOF
SELECT 'SparkSQL Configuration Test' as message;
EOF
    # 运行测试
    $SPARK_HOME/bin/spark-sql -f /tmp/test_sparksql.sql 2>&1 || {
        log_error "SparkSQL配置测试失败"
        return 1
    }
    log_info "SparkSQL配置测试成功"
    return 0
}
# 生成配置报告
generate_config_report() {
    local report_file="sparksql_config_report_$(date +%Y%m%d_%H%M%S).txt"
    cat > "$report_file" << EOF
=========================================
   SparkSQL自动配置报告
   生成时间: $(date)
=========================================
系统信息:
  - CPU核心数: ${CPU_CORES}
  - 总内存: ${TOTAL_MEM}GB
  - 可用磁盘: ${DISK_SPACE}GB
Spark配置:
  - SPARK_HOME: ${SPARK_HOME}
  - Executor内存: ${EXECUTOR_MEM}GB
  - Executor核心: ${EXECUTOR_CORES}
  - Driver内存: ${DRIVER_MEM}GB
  - 并行度: $((CPU_CORES * 2))
功能状态:
  - 自适应执行: 已启用
  - CBO优化: 已启用
  - 动态分配: 已启用
  - Hive集成: $([ -n "$HIVE_HOME" ] && echo "已配置" || echo "未配置")
配置文件位置:
  - spark-defaults.conf: ${SPARK_HOME}/conf/spark-defaults.conf
  - spark-env.sh: ${SPARK_HOME}/conf/spark-env.sh
  - hive-site.xml: ${SPARK_HOME}/conf/hive-site.xml
=========================================
EOF
    log_info "配置报告已生成: $report_file"
}
# 主函数
main() {
    log_info "开始自动配置SparkSQL引擎..."
    echo "==========================================="
    # 检测环境
    detect_spark_env
    detect_system_resources
    # 生成配置
    generate_sparksql_config
    configure_spark_env
    configure_hive_integration
    # 测试配置
    if test_sparksql_config; then
        generate_config_report
        log_info "SparkSQL自动配置完成!"
    else
        log_error "SparkSQL配置失败,请检查日志"
        exit 1
    fi
}
# 执行主函数
main "$@"

Docker环境配置脚本

# Dockerfile for SparkSQL
FROM openjdk:8-jre-slim
# 设置环境变量
ENV SPARK_VERSION=3.3.0
ENV HADOOP_VERSION=3.3.4
# 安装依赖
RUN apt-get update && apt-get install -y \
    wget \
    tar \
    python3 \
    python3-pip \
    && rm -rf /var/lib/apt/lists/*
# 下载并安装Spark
RUN wget https://archive.apache.org/dist/spark/spark-${SPARK_VERSION}/spark-${SPARK_VERSION}-bin-hadoop3.tgz \
    && tar -xzf spark-${SPARK_VERSION}-bin-hadoop3.tgz \
    && mv spark-${SPARK_VERSION}-bin-hadoop3 /opt/spark \
    && rm spark-${SPARK_VERSION}-bin-hadoop3.tgz
# 设置环境变量
ENV SPARK_HOME=/opt/spark
ENV PATH=$PATH:$SPARK_HOME/bin
# 创建用户
RUN useradd -m -s /bin/bash spark
USER spark
# 暴露端口
EXPOSE 4040 8080 7077
# 启动脚本
COPY entrypoint.sh /entrypoint.sh
ENTRYPOINT ["/entrypoint.sh"]

Python配置管理脚本

#!/usr/bin/env python3
"""
SparkSQL Configuration Manager
自动优化SparkSQL配置的Python工具
"""
import os
import psutil
import json
import logging
from datetime import datetime
class SparkSQLConfigurator:
    def __init__(self, spark_home=None):
        self.spark_home = spark_home or os.environ.get('SPARK_HOME', '/opt/spark')
        self.config = {}
        self.logger = self._setup_logger()
    def _setup_logger(self):
        logging.basicConfig(
            level=logging.INFO,
            format='%(asctime)s - %(levelname)s - %(message)s'
        )
        return logging.getLogger(__name__)
    def detect_resources(self):
        """检测系统资源"""
        cpu_count = psutil.cpu_count()
        memory_gb = psutil.virtual_memory().total / (1024**3)
        disk_gb = psutil.disk_usage('/').free / (1024**3)
        return {
            'cpu_cores': cpu_count,
            'memory_gb': round(memory_gb, 2),
            'disk_gb': round(disk_gb, 2)
        }
    def calculate_spark_config(self, resources):
        """自动计算Spark配置参数"""
        cpu_cores = resources['cpu_cores']
        memory_gb = resources['memory_gb']
        executor_memory = int(memory_gb * 0.6)
        executor_cores = max(1, cpu_cores // 2)
        driver_memory = int(memory_gb * 0.2)
        config = {
            'spark.executor.memory': f'{executor_memory}g',
            'spark.executor.cores': str(executor_cores),
            'spark.executor.instances': str(max(1, cpu_cores // executor_cores)),
            'spark.driver.memory': f'{driver_memory}g',
            'spark.driver.cores': '2',
            'spark.sql.shuffle.partitions': str(cpu_cores * 2),
        }
        return config
    def generate_optimized_config(self):
        """生成优化的配置"""
        resources = self.detect_resources()
        self.logger.info(f"Detected resources: {resources}")
        base_config = {
            # 自适应查询执行
            'spark.sql.adaptive.enabled': 'true',
            'spark.sql.adaptive.coalescePartitions.enabled': 'true',
            'spark.sql.adaptive.join.enabled': 'true',
            'spark.sql.adaptive.skewJoin.enabled': 'true',
            # 性能优化
            'spark.sql.autoBroadcastJoinThreshold': '10485760',
            'spark.sql.broadcastTimeout': '600',
            'spark.sql.files.maxPartitionBytes': '536870912',
            # 动态分配
            'spark.dynamicAllocation.enabled': 'true',
            'spark.dynamicAllocation.minExecutors': '1',
            'spark.dynamicAllocation.maxExecutors': str(resources['cpu_cores'] * 2),
            'spark.shuffle.service.enabled': 'true',
            # 序列化
            'spark.serializer': 'org.apache.spark.serializer.KryoSerializer',
            'spark.kryoserializer.buffer.max': '256m',
            # CBO优化
            'spark.sql.cbo.enabled': 'true',
            'spark.sql.cbo.joinReorder.enabled': 'true',
            'spark.sql.cbo.starSchemaDetection': 'true',
        }
        # 合并计算出的配置
        calculated_config = self.calculate_spark_config(resources)
        self.config = {**base_config, **calculated_config}
        return self.config
    def save_config(self, filepath=None):
        """保存配置到文件"""
        if filepath is None:
            filepath = os.path.join(self.spark_home, 'conf', 'spark-defaults.conf')
        with open(filepath, 'w') as f:
            f.write(f"# Auto-generated SparkSQL configuration\n")
            f.write(f"# Generated at: {datetime.now()}\n")
            f.write(f"# Resources: {self.detect_resources()}\n\n")
            for key, value in self.config.items():
                f.write(f"{key}\t{value}\n")
        self.logger.info(f"Configuration saved to {filepath}")
        return filepath
    def validate_config(self):
        """验证配置完整性"""
        required_keys = [
            'spark.executor.memory',
            'spark.executor.cores',
            'spark.driver.memory',
            'spark.sql.adaptive.enabled'
        ]
        missing_keys = [k for k in required_keys if k not in self.config]
        if missing_keys:
            self.logger.warning(f"Missing configuration keys: {missing_keys}")
            return False
        return True
    def apply_config(self):
        """应用配置到Spark"""
        if not self.validate_config():
            self.logger.error("Invalid configuration")
            return False
        # 设置环境变量
        for key, value in self.config.items():
            spark_key = key.replace('.', '_').upper()
            os.environ[spark_key] = value
        self.logger.info("Configuration applied successfully")
        return True
    def generate_report(self):
        """生成配置报告"""
        resources = self.detect_resources()
        report = {
            'timestamp': datetime.now().isoformat(),
            'system_resources': resources,
            'spark_home': self.spark_home,
            'configuration': self.config,
            'validation': self.validate_config()
        }
        report_file = f"sparksql_report_{datetime.now().strftime('%Y%m%d_%H%M%S')}.json"
        with open(report_file, 'w') as f:
            json.dump(report, f, indent=2)
        self.logger.info(f"Report generated: {report_file}")
        return report
def main():
    """主函数"""
    configurator = SparkSQLConfigurator()
    print("="*50)
    print("SparkSQL Configuration Manager")
    print("="*50)
    # 检测资源
    resources = configurator.detect_resources()
    print(f"\nSystem Resources:")
    print(f"  CPU Cores: {resources['cpu_cores']}")
    print(f"  Memory: {resources['memory_gb']}GB")
    print(f"  Disk: {resources['disk_gb']}GB")
    # 生成配置
    print("\nGenerating optimized configuration...")
    config = configurator.generate_optimized_config()
    print("\nOptimized Configuration:")
    for key, value in config.items():
        print(f"  {key}: {value}")
    # 保存配置
    config_file = configurator.save_config()
    print(f"\nConfiguration saved to: {config_file}")
    # 验证配置
    if configurator.validate_config():
        print("\n✓ Configuration validation passed")
    else:
        print("\n✗ Configuration validation failed")
    # 生成报告
    report = configurator.generate_report()
    print(f"\nConfiguration report generated successfully")
if __name__ == "__main__":
    main()

使用说明

基础脚本使用

# 给脚本执行权限
chmod +x auto_config_sparksql.sh
# 运行脚本
./auto_config_sparksql.sh
# 或指定Spark路径
SPARK_HOME=/usr/local/spark ./auto_config_sparksql.sh

Python配置工具使用

# 安装依赖
pip install psutil
# 运行配置工具
python3 spark_sql_configurator.py

Docker部署

# 构建镜像
docker build -t sparksql-auto-config .
# 运行容器
docker run -d --name sparksql \
  -p 4040:4040 \
  -p 8080:8080 \
  sparksql-auto-config

配置优化建议

  1. 内存配置

    • Executor内存:系统总内存的40-60%
    • Driver内存:系统总内存的10-20%
    • 预留20-30%给操作系统
  2. 并行度配置

    • Shuffle分区数:CPU核心数×2
    • Executor核心数:2-4个
    • Executor数量:CPU核心数/Executor核心数
  3. 性能优化

    • 启用自适应查询执行
    • 启用CBO优化
    • 启用动态资源分配
    • 使用Kryo序列化

这个自动配置脚本能够根据系统资源自动优化SparkSQL引擎配置,提高查询性能。

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