JSON Lines文件解析写入及与指定Schema DataFrame互转方法咨询
实现方法
以下分别提供PySpark和Pandas两种常用场景的实现代码:
1. JSON Lines文件转指定Schema的DataFrame
PySpark实现
因为原文件每行是JSON数组字符串,无法直接用spark.read.json解析,需要先按文本读取再做格式转换:
from pyspark.sql import SparkSession from pyspark.sql.functions import from_json, col from pyspark.sql.types import StructType, StructField, IntegerType, StringType, ArrayType spark = SparkSession.builder.appName("JsonlConvert").getOrCreate() # 定义JSON数组对应的结构 array_schema = ArrayType( StructType([ StructField("_c0", IntegerType(), nullable=True), StructField("_c1", StringType(), nullable=True), StructField("_c2", IntegerType(), nullable=True), StructField("_c3", StringType(), nullable=True), StructField("_c4", StringType(), nullable=True) ]) ) # 读取原始文件为文本格式 raw_df = spark.read.text("your_input_path.jsonl") # 解析JSON数组并提取对应字段 result_df = raw_df.select( from_json(col("value"), array_schema).alias("data_arr") ).select( col("data_arr")[0].cast(IntegerType()).alias("id"), col("data_arr")[1].cast(StringType()).alias("name"), col("data_arr")[2].cast(IntegerType()).alias("age"), col("data_arr")[3].cast(StringType()).alias("gender"), col("data_arr")[4].cast(StringType()).alias("time") ) # 输出Schema验证 result_df.printSchema()
Pandas实现
import pandas as pd import json data_list = [] with open("your_input_path.jsonl", "r", encoding="utf-8") as f: for line in f: line = line.strip() if not line: continue arr = json.loads(line) data_list.append({ "id": arr[0], "name": arr[1], "age": arr[2], "gender": arr[3], "time": arr[4] }) df = pd.DataFrame(data_list)
2. DataFrame转目标格式JSON Lines文件
PySpark实现
将字段按顺序拼接为数组后转JSON字符串,再按文本格式写出即可:
from pyspark.sql.functions import array, to_json output_df = result_df.select( to_json(array( col("id"), col("name"), col("age"), col("gender"), col("time") )).alias("value") ) # 写出文件,如需合并为单个文件可添加.coalesce(1) output_df.write.mode("overwrite").text("your_output_path.jsonl")
Pandas实现
with open("your_output_path.jsonl", "w", encoding="utf-8") as f: for _, row in df.iterrows(): arr = [row["id"], row["name"], row["age"], row["gender"], row["time"]] f.write(json.dumps(arr, ensure_ascii=False) + "\n")
内容的提问来源于stack exchange,提问作者will.wang
相关产品推荐
相关产品推荐

