咨询基于Docker的Cloudera与Spark Jobs环境搭建方案及示例
项目结构
repo-project-job-spark-2 ├── docker-compose.yml ├── docker │ └── conf │ ├── spark-defaults.conf │ └── hdfs-site.xml ├── logs │ ├── cloudera │ └── spark └── src ├── application.py └── jobs ├── job_daily_1.py └── job_daily_2.py
Docker Compose 配置(docker-compose.yml)
version: '3.8' services: cloudera: image: cloudera/quickstart:latest container_name: cloudera-hadoop hostname: quickstart.cloudera privileged: true ports: - "8088:8088" # YARN ResourceManager - "50070:50070" # HDFS NameNode - "8888:8888" # Hue volumes: - cloudera_data:/var/lib/hadoop-hdfs - cloudera_logs:/var/log/hadoop - ./logs/cloudera:/opt/spark-job-logs restart: always networks: - spark-cloudera-net spark-job-runner: image: bitnami/spark:3.3.0 container_name: spark-job-runner hostname: spark-job-runner volumes: - ./src:/opt/spark-app/src - ./docker/conf:/opt/spark/conf - ./logs/spark:/opt/spark/logs depends_on: - cloudera networks: - spark-cloudera-net entrypoint: | bash -c " echo '0 0 * * * /opt/spark/bin/spark-submit --master local[*] /opt/spark-app/src/jobs/job_daily_1.py' >> /var/spool/cron/crontabs/root echo '0 1 * * * /opt/spark/bin/spark-submit --master local[*] /opt/spark-app/src/jobs/job_daily_2.py' >> /var/spool/cron/crontabs/root cron -f " volumes: cloudera_data: cloudera_logs: networks: spark-cloudera-net: driver: bridge
核心文件内容
1. Spark 配置文件(docker/conf/spark-defaults.conf)
spark.hadoop.fs.defaultFS hdfs://quickstart.cloudera:8020 spark.hadoop.dfs.replication 1 spark.driver.extraJavaOptions -Dlog4j.configuration=file:/opt/spark/conf/log4j.properties spark.executor.extraJavaOptions -Dlog4j.configuration=file:/opt/spark/conf/log4j.properties
2. 每日作业示例1(src/jobs/job_daily_1.py)
from pyspark.sql import SparkSession import logging import os # 配置日志输出到Cloudera指定目录 log_dir = "/opt/spark-job-logs" os.makedirs(log_dir, exist_ok=True) logging.basicConfig( filename=f"{log_dir}/job_daily_1.log", level=logging.INFO, format="%(asctime)s - %(levelname)s - %(message)s" ) def main(): spark = SparkSession.builder \ .appName("DailyJob1") \ .getOrCreate() # 生成示例业务数据 data = [("Alice", 34), ("Bob", 45), ("Charlie", 29)] df = spark.createDataFrame(data, ["name", "age"]) # 写入Cloudera HDFS hdfs_path = "hdfs://quickstart.cloudera:8020/user/spark/jobs/output/daily_job1.csv" df.write.mode("overwrite").csv(hdfs_path, header=True) logging.info(f"数据已成功写入HDFS路径: {hdfs_path}") spark.stop() if __name__ == "__main__": main()
3. 每日作业示例2(src/jobs/job_daily_2.py)
from pyspark.sql import SparkSession import logging import os log_dir = "/opt/spark-job-logs" os.makedirs(log_dir, exist_ok=True) logging.basicConfig( filename=f"{log_dir}/job_daily_2.log", level=logging.INFO, format="%(asctime)s - %(levelname)s - %(message)s" ) def main(): spark = SparkSession.builder \ .appName("DailyJob2") \ .getOrCreate() # 生成不同维度的示例数据 data = [("NewYork", 8419000), ("LosAngeles", 3971000), ("Chicago", 2716000)] df = spark.createDataFrame(data, ["city", "population"]) # 写入Cloudera HDFS hdfs_path = "hdfs://quickstart.cloudera:8020/user/spark/jobs/output/daily_job2.csv" df.write.mode("overwrite").csv(hdfs_path, header=True) logging.info(f"数据已成功写入HDFS路径: {hdfs_path}") spark.stop() if __name__ == "__main__": main()
关键配置说明
- Cloudera 容器:通过
restart: always确保全天候运行,挂载持久化卷保存HDFS数据和系统日志,避免容器重启后数据丢失;开放常用端口方便外部管理。 - Spark 作业容器:通过
entrypoint配置cron定时任务,实现每日自动运行不同作业;挂载本地代码和配置目录,方便修改作业逻辑和Spark参数;通过自定义Docker网络与Cloudera容器互通,直接使用容器hostname访问HDFS。 - 日志与数据交互:Spark作业将日志写入Cloudera容器挂载的共享目录,同时直接通过HDFS地址写入业务数据,无需额外配置存储共享。
内容的提问来源于stack exchange,提问作者Felipe
相关产品推荐
相关产品推荐

