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
相关产品推荐
相关产品推荐

