使用PySpark将Kafka数据写入HDFS时无法实现目标目录结构
PySpark Kafka到HDFS数据存储问题修复
问题概述
当前代码运行后存在以下问题:
- CSV目录下生成了
file_type=csv/和file_type=image/分区子目录,不符合直接在csv/下存储Parquet文件的预期结构 - 视频文件未写入HDFS
- CSV文件未按「原文件名+.parquet」的格式存储
修复后代码
import logging from pyspark.sql import SparkSession from pyspark.sql.functions import * from pyspark.sql.types import * import pyarrow.hdfs as hdfs import os # 配置日志 logging.getLogger("py4j").setLevel(logging.ERROR) logging.getLogger("org.apache.kafka").setLevel(logging.WARN) # 初始化SparkSession spark = SparkSession.builder \ .appName("KafkaToHDFS") \ .config("spark.jars.packages", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.4.0") \ .getOrCreate() # Kafka broker配置 kafka_server = "sandbox-hdf.hortonworks.com:6667" kafka_topic = "airbnbdata-final" # 从Kafka读取数据 df = spark.read \ .format("kafka") \ .option("kafka.bootstrap.servers", kafka_server) \ .option("subscribe", kafka_topic) \ .load() # 创建输出目录(带时间维度) def create_directory(directory): current_time = spark.sql("SELECT current_timestamp()").collect()[0][0] year_str = current_time.strftime("%Y") month_str = current_time.strftime("%m") day_str = current_time.strftime("%d") hour_str = current_time.strftime("%H") fulldirectory = f"{directory}/{year_str}/{month_str}/{day_str}/{hour_str}/" # 确保根目录及子目录存在 fs = hdfs.connect(host="sandbox-hdp.hortonworks.com", port=8020) for sub_dir in ["csv", "images", "videos"]: full_sub_dir = f"{fulldirectory}{sub_dir}" if not fs.exists(full_sub_dir): fs.mkdir(full_sub_dir, recursive=True) return fulldirectory # 处理Kafka消息,解析JSON def process_batch(df): schema = StructType([ StructField("file_type", StringType()), StructField("file_name", StringType()), StructField("data", StringType()) ]) extracted_df = df.selectExpr( "CAST(key AS STRING)", "CAST(value AS STRING)", "topic", "partition", "offset", "timestamp", "timestampType" ) extracted_df = extracted_df.withColumn("value_json", from_json(col("value"), schema)) \ .select( col("value_json.file_type").alias("file_type"), col("value_json.file_name").alias("file_name"), col("value_json.data").alias("data") ) extracted_df.printSchema() return extracted_df # 自定义写入文件到HDFS def write_file(row, base_dir): fs = hdfs.connect(host="sandbox-hdp.hortonworks.com", port=8020) file_type = row["file_type"] original_name = row["file_name"] data = row["data"] # 确定目标路径和文件名 if file_type == "csv": # CSV文件追加.parquet后缀 target_path = f"{base_dir}csv/{original_name}.parquet" # 将CSV字符串转为DataFrame写入Parquet temp_df = spark.createDataFrame([(data,)], ["csv_content"]) temp_df.write.mode("append").parquet(target_path) elif file_type in ["jpg", "png"]: target_path = f"{base_dir}images/{original_name}" with fs.open(target_path, "wb") as f: f.write(data.encode('utf-8') if isinstance(data, str) else data) elif file_type in ["mp4", "mov"]: target_path = f"{base_dir}videos/{original_name}" with fs.open(target_path, "wb") as f: f.write(data.encode('utf-8') if isinstance(data, str) else data) # 处理所有文件类型 def process_all_files(df, base_dir): # 遍历所有行写入对应目录 df.foreach(lambda row: write_file(row, base_dir)) # 主流程 directory_path = "hdfs://sandbox-hdp.hortonworks.com:8020/user/hadoop/OUTPUT_FINAL1" fulldirectory = create_directory(directory_path) processed_df = process_batch(df) process_all_files(processed_df, fulldirectory) spark.stop()
关键修改说明
- 移除分区写入,自定义文件命名
- 去掉原代码中
partitionBy("file_type", "file_name")的写法,改用foreach逐行处理,确保CSV文件以「原文件名+.parquet」格式存放在csv/目录下,不再生成分区子目录
- 去掉原代码中
- 统一文件分类处理逻辑
- 在
write_file函数中按file_type分别处理CSV、图片、视频,确保图片和视频写入对应images/、videos/目录,且保留原后缀
- 在
- 提前创建所有子目录
- 在
create_directory函数中提前创建csv/、images/、videos/子目录,避免写入时因目录不存在导致失败
- 在
- 修正CSV数据写入逻辑
- 针对CSV数据,将字符串内容转为临时DataFrame后写入Parquet文件,保证格式符合要求
- 合并处理流程
- 将原代码中分开的CSV写入和多媒体文件写入合并为一个
process_all_files函数,简化流程并确保所有文件类型都被处理
- 将原代码中分开的CSV写入和多媒体文件写入合并为一个
运行命令
保持原运行命令不变:
spark-submit --jars /tmp/Airbnb_Data/jsr305-3.0.2.jar,/tmp/Airbnb_Data/snappy-java-1.1.8.4.jar,/tmp/Airbnb_Data/kafka-clients-3.3.2.jar,/tmp/Airbnb_Data/spark-sql-kafka-0-10_2.11-2.4.0.jar,/tmp/Airbnb_Data/commons-logging-1.1.3.jar,/tmp/Airbnb_Data/spark-streaming-kafka-0-10_2.12-3.4.0.jar,/tmp/Airbnb_Data/hadoop-client-api-3.3.4.jar,/tmp/Airbnb_Data/lz4-java-1.8.0.jar,/tmp/Airbnb_Data/hadoop-client-runtime-3.3.4.jar,/tmp/Airbnb_Data/scala-library-2.12.17.jar,/tmp/Airbnb_Data/spark-tags_2.12-3.4.0.jar,/tmp/Airbnb_Data/spark-token-provider-kafka-0-10_2.12-3.4.0.jar,/tmp/Airbnb_Data/slf4j-api-1.7.36.jar KafkaToHDFS.py
内容的提问来源于stack exchange,提问作者user31081998
相关产品推荐
相关产品推荐

