You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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()

关键修改说明

  1. 移除分区写入,自定义文件命名
    • 去掉原代码中partitionBy("file_type", "file_name")的写法,改用foreach逐行处理,确保CSV文件以「原文件名+.parquet」格式存放在csv/目录下,不再生成分区子目录
  2. 统一文件分类处理逻辑
    • 在write_file函数中按file_type分别处理CSV、图片、视频,确保图片和视频写入对应images/、videos/目录,且保留原后缀
  3. 提前创建所有子目录
    • 在create_directory函数中提前创建csv/、images/、videos/子目录,避免写入时因目录不存在导致失败
  4. 修正CSV数据写入逻辑
    • 针对CSV数据,将字符串内容转为临时DataFrame后写入Parquet文件,保证格式符合要求
  5. 合并处理流程
    • 将原代码中分开的CSV写入和多媒体文件写入合并为一个process_all_files函数,简化流程并确保所有文件类型都被处理

运行命令

保持原运行命令不变:

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.19 11:52:37