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

如何使用PySpark按Schema读取JSON文件并转换为CSV文件

How to Read JSON with a Specified Schema and Convert to CSV in PySpark

Hey there! As someone who’s fumbled through PySpark’s learning curve just like you, let’s break this down into simple, actionable steps that you can follow right away. No jargon overload—just straight-to-the-point guidance.

Step 1: Set Up Your SparkSession

First, you need a SparkSession to work with PySpark DataFrames. This is your gateway to all data operations.

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, IntegerType, StringType, DoubleType

# Initialize a SparkSession
spark = SparkSession.builder \
    .appName("JSONtoCSVWithCustomSchema") \
    .getOrCreate()

Step 2: Define Your Exact Schema

Instead of letting PySpark guess the schema (which can be slow for large datasets or lead to unexpected type mismatches), we’ll explicitly define the structure that matches your JSON data.

Let’s say your JSON entries look like this:

{"order_id": 101, "customer_name": "Alice", "total_amount": 45.99, "order_date": "2024-05-20"}

Here’s how you’d map that to a PySpark schema:

# Define the custom schema using StructType and StructField
custom_schema = StructType([
    StructField("order_id", IntegerType(), nullable=False),  # Non-null integer field
    StructField("customer_name", StringType(), nullable=True),
    StructField("total_amount", DoubleType(), nullable=True),
    StructField("order_date", StringType(), nullable=True)  # Use TimestampType if you want to parse dates later
])
  • Tweak field names, data types, and nullable flags to match your actual JSON structure. For other types like TimestampType or BooleanType, just import them from pyspark.sql.types.

Step 3: Read the JSON File with Your Schema

Now use the schema we defined to read your JSON data. This ensures PySpark parses everything exactly as you expect.

# Read JSON file with the custom schema
df = spark.read.schema(custom_schema).json("/path/to/your/json/files")

# If your JSON uses multi-line records (one full JSON object per line is standard, but adjust if needed):
# df = spark.read.schema(custom_schema).option("multiLine", "true").json("/path/to/your/json/files")
  • Replace /path/to/your/json/files with your actual file/directory path (works for local storage or distributed systems like S3/HDFS).

Step 4: Convert and Write to CSV

Finally, export the DataFrame to CSV format. We’ll include headers (so your CSV has column names) and set a write mode to handle existing files.

# Write DataFrame to CSV
df.write \
    .mode("overwrite")  # Options: overwrite, append, ignore, errorifexists
    .option("header", "true")  # Include column headers in the output
    .option("sep", ",")  # Set delimiter (default is comma; use "\t" for tab-separated)
    .csv("/path/to/output/csv")
  • Mode breakdown:
    • overwrite: Replaces any existing data at the output path
    • append: Adds new data to existing files
    • ignore: Does nothing if data already exists
    • errorifexists: Throws an error if data is present (default behavior)

Quick Beginner Tips

  • Always validate your DataFrame after reading: Use df.printSchema() to confirm the schema matches your definition, and df.show(5) to preview the first 5 rows.
  • If you get type mismatch errors, double-check that your schema’s data types align with the actual values in your JSON (e.g., don’t use IntegerType for a field that sometimes has string values).
  • Explicit schemas are way faster than inferred ones for large datasets—they save PySpark from doing an extra pass over the data to guess types.

内容的提问来源于stack exchange,提问作者shahin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 08:17:38