本文目录导读:

我来为您提供一个自动配置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
配置优化建议
-
内存配置
- Executor内存:系统总内存的40-60%
- Driver内存:系统总内存的10-20%
- 预留20-30%给操作系统
-
并行度配置
- Shuffle分区数:CPU核心数×2
- Executor核心数:2-4个
- Executor数量:CPU核心数/Executor核心数
-
性能优化
- 启用自适应查询执行
- 启用CBO优化
- 启用动态资源分配
- 使用Kryo序列化
这个自动配置脚本能够根据系统资源自动优化SparkSQL引擎配置,提高查询性能。