在Foundry中如何解析JIRA AVRO数据内的JSON格式response列
核心问题排查
你的原有代码存在4个影响运行的核心错误:
- 调用
collect()把全量response列数据拉到驱动节点内存,数据量稍大就会触发OOM,完全不适合生产场景使用 - 定义的JSON Schema和实际JIRA返回结构不匹配:
clauseNames是字符串数组类型、schema是嵌套结构体类型,你都错误声明为了字符串类型 parse_json函数中提取的a、b字段在JIRA返回的实际结构中不存在,会直接抛出字段不存在的异常- 手动初始化
SQLContext、SparkContext属于冗余操作,Foundry转换环境已经内置了可用的Spark上下文
修正后代码
推荐使用PySpark内置函数实现,比自定义Python UDF性能高10倍以上,完全避免JVM和Python之间的序列化开销:
from transforms.api import transform_df, Input, Output from pyspark.sql.functions import from_json, explode, col import pyspark.sql.types as T @transform_df( Output("json output"), json_raw=Input("json input"), ) def my_compute_function(json_raw): # 定义和实际JIRA返回数据完全匹配的schema jira_field_schema = T.StructType([ T.StructField("id", T.StringType(), nullable=True), T.StructField("name", T.StringType(), nullable=True), T.StructField("custom", T.BooleanType(), nullable=True), T.StructField("orderable", T.BooleanType(), nullable=True), T.StructField("navigable", T.BooleanType(), nullable=True), T.StructField("searchable", T.BooleanType(), nullable=True), T.StructField("clauseNames", T.ArrayType(T.StringType()), nullable=True), T.StructField("schema", T.StructType([ T.StructField("type", T.StringType(), nullable=True), T.StructField("custom", T.StringType(), nullable=True), T.StructField("customId", T.IntegerType(), nullable=True) ]), nullable=True) ]) # response整体是数组结构,外层套ArrayType response_schema = T.ArrayType(jira_field_schema) # 1. 解析字符串格式的response为结构化数组 df_parsed = json_raw.withColumn("parsed_response", from_json(col("response"), response_schema)) # 2. 炸开数组,每个自定义字段对应独立行 df_exploded = df_parsed.withColumn("field_item", explode(col("parsed_response"))) # 3. 展开结构体字段为独立列,可根据业务需求增减字段 df_final = df_exploded.select( col("field_item.id").alias("field_id"), col("field_item.name").alias("field_name"), col("field_item.custom").alias("is_custom_field"), col("field_item.orderable").alias("is_orderable"), col("field_item.navigable").alias("is_navigable"), col("field_item.searchable").alias("is_searchable"), col("field_item.clauseNames").alias("clause_names"), col("field_item.schema.type").alias("field_data_type"), col("field_item.schema.custom").alias("custom_field_type"), col("field_item.schema.customId").alias("custom_field_id") # 可补充保留原始表的其他字段,比如导入时间、请求ID等 ) return df_final
补充说明
如果你的response列经Magritte导入后已经是AVRO结构体类型而非字符串类型,可以直接跳过from_json解析步骤,直接调用explode函数炸开数组即可。
内容的提问来源于stack exchange,提问作者Robert F
相关产品推荐
相关产品推荐

