PySpark中如何让from_json自动将字符串数字转为整数?
问题描述
我有一个PySpark DataFrame,其中quantity字段被指定为IntegerType。但当JSON数据中包含字符串形式的数字(例如"30")时,该记录会被移至corrupt_records中。
复现代码
from pyspark.sql.functions import from_json, col, when from pyspark.sql.types import StringType, StructType, StructField, IntegerType data = [ ("{'fruit':'Apple', 'quantity':10}",), ("{'fruit':'Banana', 'quantity':20}",), ("{'fruit':'Cherry', 'quantity':'30'}",), ("{'fruit':'Date', 'quantity':'40'}",), ("{'fruit':'Elderberry', 'quantity':'50'}",) ] # 定义输入Schema schema = StructType([ StructField("json_column", StringType(), True) ]) # 创建DataFrame input_df = spark.createDataFrame(data, schema) input_df.show(truncate=False)
输入DataFrame输出
+---------------------------------------+ |json_column | +---------------------------------------+ |{'fruit':'Apple', 'quantity':10} | |{'fruit':'Banana', 'quantity':20} | |{'fruit':'Cherry', 'quantity':'30'} | |{'fruit':'Date', 'quantity':'40'} | |{'fruit':'Elderberry', 'quantity':'50'}| +---------------------------------------+
解析JSON的代码
json_options = {"columnNameOfCorruptRecord":"corrupt_json"} json_schema = StructType([ StructField('fruit', StringType(), True), StructField('quantity', IntegerType(), True), StructField('corrupt_json', StringType(), True) ]) json_df = input_df.select(from_json("json_column", schema=json_schema, options=json_options).alias("data")) df = json_df.select("data.*") df.show(truncate=False)
解析后输出
+----------+--------+---------------------------------------+ |fruit |quantity|corrupt_json | +----------+--------+---------------------------------------+ |Apple |10 |NULL | |Banana |20 |NULL | |Cherry |NULL |{'fruit':'Cherry', 'quantity':'30'} | |Date |NULL |{'fruit':'Date', 'quantity':'40'} | |Elderberry|NULL |{'fruit':'Elderberry', 'quantity':'50'}| +----------+--------+---------------------------------------+
需求
有没有办法让from_json函数自动将字符串数字转换为整数,使包含字符串数字的记录不会被移至corrupt_records中?
解决方案
方法1:先解析为字符串,再安全转换为整数
核心思路是先把quantity按字符串类型解析,再用try_cast函数安全转换为整数——该函数会自动处理可转换的字符串数字,无法转换的返回NULL但不会触发脏数据记录。
from pyspark.sql.functions import from_json, col, try_cast from pyspark.sql.types import StringType, StructType, StructField, IntegerType # 修改Schema,将quantity设为StringType json_schema = StructType([ StructField('fruit', StringType(), True), StructField('quantity', StringType(), True), StructField('corrupt_json', StringType(), True) ]) json_options = {"columnNameOfCorruptRecord":"corrupt_json"} json_df = input_df.select(from_json("json_column", schema=json_schema, options=json_options).alias("data")) # 用try_cast安全转换为整数 df = json_df.select( "data.fruit", try_cast(col("data.quantity").cast(StringType())).alias("quantity"), "data.corrupt_json" ) df.show(truncate=False)
输出结果
+----------+--------+-----------+ |fruit |quantity|corrupt_json| +----------+--------+-----------+ |Apple |10 |NULL | |Banana |20 |NULL | |Cherry |30 |NULL | |Date |40 |NULL | |Elderberry|50 |NULL | +----------+--------+-----------+
方法2:自定义UDF处理JSON解析
如果需要更灵活的逻辑,可以写一个自定义UDF,手动解析JSON字符串并处理类型转换:
import json from pyspark.sql.functions import udf from pyspark.sql.types import StructType, StructField, StringType, IntegerType # 定义输出Schema output_schema = StructType([ StructField('fruit', StringType(), True), StructField('quantity', IntegerType(), True), StructField('corrupt_json', StringType(), True) ]) @udf(output_schema) def parse_json_udf(json_str): try: data = json.loads(json_str.replace("'", "\"")) # 处理单引号格式的JSON # 尝试将quantity转为整数 quantity = int(data.get("quantity")) if data.get("quantity") is not None else None return (data.get("fruit"), quantity, None) except Exception as e: # 解析失败时返回原字符串到corrupt_json return (None, None, json_str) # 应用UDF df = input_df.select(parse_json_udf(col("json_column")).alias("data")).select("data.*") df.show(truncate=False)
输出结果
+----------+--------+-----------+ |fruit |quantity|corrupt_json| +----------+--------+-----------+ |Apple |10 |NULL | |Banana |20 |NULL | |Cherry |30 |NULL | |Date |40 |NULL | |Elderberry|50 |NULL | +----------+--------+-----------+
内容的提问来源于stack exchange,提问作者DumbCoder
相关产品推荐
相关产品推荐

