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

PySpark写入S3时将RECORDS_*字段的JSON字符串转为JSON数组

问题描述

我有一个PySpark DataFrame,其中以RECORDS_开头的字段存储的是JSON格式的字符串。使用foreach()调用自定义write_final_json函数将数据写入S3时,这些字段仍以转义的JSON字符串形式保存。现需修改该函数,将RECORDS_*字段的JSON字符串转换为JSON数组后再写入S3。

DataFrame示例

+--------------------+--------------------+--------------------+-----------------+-----------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+
|     input_file_name|              ROW_ID|  FAILED_VALIDATIONS|VALIDATION_STATUS|  MEMBER_ID|      RECORDS_HEADER|         RECORDS_001|         RECORDS_002|         RECORDS_003|         RECORDS_004|         RECORDS_005|         RECORDS_006|         RECORDS_007|            batch_id|
+--------------------+--------------------+--------------------+-----------------+-----------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+
|s3://gov-solution...|a64fbbea-7cf3-43e...|                null|           PASSED|U7742139901|[{"RECORD_TYPE":"...|[{"BATCH_BEGIN_DA...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|d0f87fd5-fca1-4e2...|
|s3://gov-solution...|a64fbbea-7cf3-43e...|MEMBER_ID is missing|           FAILED|           |[{"RECORD_TYPE":"...|[{"BATCH_BEGIN_DA...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|d0f87fd5-fca1-4e2...|
|s3://gov-solution...|a64fbbea-7cf3-43e...|MEMBER_ID is miss...|           FAILED|           |[{"RECORD_TYPE":"...|[{"BATCH_BEGIN_DA...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|d0f87fd5-fca1-4e2...|
|s3://gov-solution...|a64fbbea-7cf3-43e...|                null|           PASSED|U6881487301|[{"RECORD_TYPE":"...|[{"BATCH_BEGIN_DA...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|[{"RECORD_TYPE":"...|                null|d0f87fd5-fca1-4e2...|
+--------------------+--------------------+--------------------+-----------------+-----------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+--------------------+

当前写入后的格式示例

{
    "input_file_name": "s3://gov-solutions-dev-fulfillment2-files/Input-Data/EOBLetter/20240327/MONTHLY_EOB_U6881487301_20231218180516.json",
    "ROW_ID": "1f6df8f3-7cb7-4a29-999c-ea8f5923168a",
    "FAILED_VALIDATIONS": null,
    "VALIDATION_STATUS": "PASSED",
    "MEMBER_ID": "U6881487301",
    "RECORDS_HEADER": "[{\"RECORD_TYPE\":\"HEADER\",\"COUNT_OF_OBJECTS_RECORD_TYPE_001\":1,\"COUNT_OF_OBJECTS_RECORD_TYPE_002\":3,\"COUNT_OF_OBJECTS_RECORD_TYPE_003\":6,\"COUNT_OF_OBJECTS_RECORD_TYPE_004\":3,\"COUNT_OF_OBJECTS_RECORD_TYPE_005\":1,\"COUNT_OF_OBJECTS_RECORD_TYPE_006\":1,\"COUNT_OF_OBJECTS_RECORD_TYPE_007\":0}]",
    "RECORDS_001": "[{\"BATCH_BEGIN_DATE\":\"01/01/2023\",\"BATCH_END_DATE\":\"11/30/2023\",\"BATCH_ID\":1000000078,\"BATCH_RUN_DATE\":\"12/18/2023\",\"BATCH_STATUS\":\"COMPLETE\",\"COLATERAL_TYPE\":\"EOB\",\"INDIV_ID\":\"3003790\",\"MBR_PREF_LANG\":\"English\",\"MBR_PREF_LARGE_PRINT\":\"N\",\"MEMBER_ID\":\"U6881487301\",\"SOURCE_SYSTEM_ID\":\"GBSF\"}]",
    "RECORDS_007": null,
    "batch_id": "302bdcc8-dec1-445a-86a6-0409d3959b75"
}

当前函数代码

def write_final_json(row):
    print("in write_final_json")
    if row.VALIDATION_STATUS == 'PASSED':
        # Initialize S3 client
        s3 = boto3.client('s3')
    
        # S3 bucket and directory
        bucket_name = 'gov-solutions-dev-fulfillment2-files'
        output_dir = 'Extracted-Data/EOBLetter/Ready'
        
        # Create JSON string from row
        json_data = json.dumps(row.asDict())
        
        # Define S3 key
        s3_key = f"{output_dir}/{row.batch_id}/{row.MEMBER_ID}.json"
    
        # Write JSON data to S3
        s3.put_object(Body=json_data, Bucket=bucket_name, Key=s3_key)
解决方案

核心思路是将RECORDS_开头字段的JSON字符串解析为Python原生对象(列表/字典),再整体序列化为JSON,避免转义字符出现。修改后的函数如下:

import json
import boto3

def write_final_json(row):
    print("in write_final_json")
    if row.VALIDATION_STATUS == 'PASSED':
        # Initialize S3 client
        s3 = boto3.client('s3')
    
        # S3 bucket and directory
        bucket_name = 'gov-solutions-dev-fulfillment2-files'
        output_dir = 'Extracted-Data/EOBLetter/Ready'
        
        # 将Row对象转为字典
        row_dict = row.asDict()
        
        # 遍历处理所有RECORDS_开头的字段
        for key in list(row_dict.keys()):
            if key.startswith('RECORDS_'):
                json_str = row_dict[key]
                if json_str is not None:
                    try:
                        # 解析JSON字符串为Python原生对象
                        row_dict[key] = json.loads(json_str)
                    except json.JSONDecodeError:
                        # 解析失败时保留原字符串(可根据需求调整异常处理逻辑)
                        pass
        
        # 序列化为JSON,此时RECORDS_字段已转为原生JSON结构
        json_data = json.dumps(row_dict)
        
        # 定义S3存储路径
        s3_key = f"{output_dir}/{row.batch_id}/{row.MEMBER_ID}.json"
    
        # 写入S3
        s3.put_object(Body=json_data, Bucket=bucket_name, Key=s3_key)

修改说明

  • 先将Row对象转为字典,遍历所有键筛选出RECORDS_开头的字段
  • 对每个目标字段,判断非空后用json.loads()解析JSON字符串为Python列表/字典
  • 最后将处理后的字典序列化为JSON,此时RECORDS_字段会以原生JSON数组/对象的形式输出,不会出现转义字符
  • 增加异常捕获逻辑,避免因无效JSON导致整个写入流程失败(可根据业务需求调整异常处理方式)

处理后,RECORDS_HEADER和RECORDS_001等字段会变成原生JSON数组,示例如下:

{
    "input_file_name": "s3://gov-solutions-dev-fulfillment2-files/Input-Data/EOBLetter/20240327/MONTHLY_EOB_U6881487301_20231218180516.json",
    "ROW_ID": "1f6df8f3-7cb7-4a29-999c-ea8f5923168a",
    "FAILED_VALIDATIONS": null,
    "VALIDATION_STATUS": "PASSED",
    "MEMBER_ID": "U6881487301",
    "RECORDS_HEADER": [{"RECORD_TYPE":"HEADER","COUNT_OF_OBJECTS_RECORD_TYPE_001":1,"COUNT_OF_OBJECTS_RECORD_TYPE_002":3,"COUNT_OF_OBJECTS_RECORD_TYPE_003":6,"COUNT_OF_OBJECTS_RECORD_TYPE_004":3,"COUNT_OF_OBJECTS_RECORD_TYPE_005":1,"COUNT_OF_OBJECTS_RECORD_TYPE_006":1,"COUNT_OF_OBJECTS_RECORD_TYPE_007":0}],
    "RECORDS_001": [{"BATCH_BEGIN_DATE":"01/01/2023","BATCH_END_DATE":"11/30/2023","BATCH_ID":1000000078,"BATCH_RUN_DATE":"12/18/2023","BATCH_STATUS":"COMPLETE","COLATERAL_TYPE":"EOB","INDIV_ID":"3003790","MBR_PREF_LANG":"English","MBR_PREF_LARGE_PRINT":"N","MEMBER_ID":"U6881487301","SOURCE_SYSTEM_ID":"GBSF"}],
    "RECORDS_007": null,
    "batch_id": "302bdcc8-dec1-445a-86a6-0409d3959b75"
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 13:27:02