spark.read.format 是 Spark SQL 中最通用的数据读取接口,通过指定数据源格式(如 csvjsonparquetjdbc 等),并配合 option 配置参数,可灵活读取几乎所有 Spark 支持的数据源。以下是其核心语法和典型场景示例:

基本语法(PySpark)

python

运行

df = spark.read.format(source)  # 指定数据源格式(字符串类型,如 "csv"、"parquet" 等)
    .option(key, value)        # 配置该数据源的特定参数(键值对形式,按需添加)
    .schema(custom_schema)     # 可选,手动指定数据结构(StructType对象)
    .load(path)                # 加载数据的路径(部分数据源如jdbc无需路径)

典型场景示例

1. 读取 CSV 文件(等效于 spark.read.csv

python

运行

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

spark = SparkSession.builder.appName("FormatExample").getOrCreate()

# 定义自定义Schema(可选,推荐用于生产环境)
csv_schema = StructType([
    StructField("id", IntegerType(), nullable=True),
    StructField("name", StringType(), nullable=True),
    StructField("age", IntegerType(), nullable=True)
])

# 用format读取CSV
df_csv = spark.read.format("csv")
    .option("header", "true")       # 第一行为表头
    .option("sep", ",")             # 字段分隔符(默认逗号,可改为"\t"等)
    .option("nullValue", "NA")      # 将"NA"视为null值
    .schema(csv_schema)             # 应用自定义Schema
    .load("hdfs:///user/data/people.csv")  # 本地路径用"file:///..."

df_csv.show()
2. 读取 JSON 文件(等效于 spark.read.json

python

运行

# 读取多行JSON(一个JSON对象跨多行)
df_json = spark.read.format("json")
    .option("multiLine", "true")    # 开启多行JSON解析
    .option("inferSchema", "false") # 关闭自动推断(如需手动指定)
    .load("file:///local/data/multi_line_data.json")

df_json.printSchema()  # 查看解析后的结构
3. 读取 Parquet 文件(列式存储格式)

python

运行

# Parquet是Spark默认列式格式,支持嵌套结构,无需额外配置
df_parquet = spark.read.format("parquet")
    .load("hdfs:///user/data/logs_parquet/")  # 可读取目录下所有parquet文件

df_parquet.show(truncate=False)
4. 读取 JDBC 数据库(如 MySQL)

python

运行

df_jdbc = spark.read.format("jdbc")
    .option("url", "jdbc:mysql://localhost:3306/testdb?useSSL=false")  # 数据库连接地址
    .option("dbtable", "users")  # 表名,或子查询如"(select * from users where age>18) t"
    .option("user", "root")      # 数据库用户名
    .option("password", "123456")# 数据库密码
    .option("fetchsize", "1000") # 每次拉取的行数(优化读取性能)
    .load()  # JDBC无需路径,通过url和dbtable定位数据

df_jdbc.createOrReplaceTempView("jdbc_users")  # 注册为临时表,可用于SQL查询
spark.sql("select name, age from jdbc_users where age < 30").show()
5. 读取 Hive 表

python

运行

# 需初始化SparkSession时启用Hive支持(enableHiveSupport())
spark = SparkSession.builder.appName("ReadHive") \
    .enableHiveSupport() \
    .getOrCreate()

# 读取Hive表
df_hive = spark.read.format("hive")
    .table("default.employee")  # 格式:数据库名.表名

df_hive.show()

核心特点

  • 通用性:支持所有 Spark 兼容的数据源(内置或第三方),如 csvjsonparquetorcjdbchivekafka 等。
  • 灵活性:通过 option 可配置数据源的所有参数(不同格式参数不同,需参考对应数据源文档)。
  • 扩展性:可通过自定义 DataSourceV2 接口扩展支持新数据源,只需在 format 中指定自定义数据源名称。

使用时需注意:不同数据源的 option 参数差异较大(如 JDBC 需要数据库连接信息,CSV 需要分隔符配置),需根据具体数据源查阅官方文档。

更多推荐