如何将异步流式JSON数据按货币分组合并存储至对象存储
解决方案:异步流式JSON按币种分组并追加到对象存储
先修正一个关键格式问题
你给出的最终期望结果存在JSON语法错误:同一个prices对象里不能有多个同名的price键,解析时会直接覆盖掉前面的条目。正确的结构应该把价格条目放在数组里,比如:
{ "prices": { "header": { "currency": "EUR" }, "priceList": [ {"productId":"0000A", "amount":"60.00"}, {"productId":"0000B", "amount":"120.00"} ] } }
实现流程
1. 单条数据解析
每次收到异步JSON后,做这几步:
- 提取
header.currency作为分组标识 - 把
price里的value字段重命名为amount,整理成独立的价格条目 - 以币种为键,确定对象存储中对应的文件路径(比如
prices/EUR.json)
2. 对象存储的读写逻辑
因为是异步追加,必须遵循先读、修改、再写回的流程:
- 检查对象存储中是否存在对应币种的文件:
- 不存在:直接创建包含当前价格条目的新JSON结构,写入存储
- 已存在:读取已有JSON,把新价格条目追加到数组里,再覆盖写回原文件
3. 并发冲突处理
如果有多个异步请求同时处理同币种数据,必须加分布式锁(比如利用对象存储的ETag校验、或者Redis锁),避免并发写入导致数据丢失。
代码示例(Python + 对象存储)
import json # 这里以AWS S3为例,其他对象存储API逻辑类似 import boto3 from botocore.exceptions import ClientError s3_client = boto3.client('s3') BUCKET_NAME = "your-bucket-name" def process_single_price(raw_json_str): # 解析输入的单条数据(注意确保输入是合法JSON,原始示例缺少外层大括号) raw_data = json.loads(f"{{{raw_json_str.strip()}}}") currency = raw_data['prices']['header']['currency'] # 整理价格条目 price_item = { "productId": raw_data['prices']['price']['productId'], "amount": raw_data['prices']['price']['value'] } object_key = f"price_groups/{currency}.json" try: # 读取已有分组文件 response = s3_client.get_object(Bucket=BUCKET_NAME, Key=object_key) existing_group = json.loads(response['Body'].read().decode('utf-8')) # 追加新条目 existing_group['prices']['priceList'].append(price_item) except ClientError as e: if e.response['Error']['Code'] == "NoSuchKey": # 不存在则创建新分组结构 existing_group = { "prices": { "header": {"currency": currency}, "priceList": [price_item] } } else: # 其他错误直接抛出 raise # 写回对象存储 s3_client.put_object( Bucket=BUCKET_NAME, Key=object_key, Body=json.dumps(existing_group, indent=2), ContentType="application/json" ) # 测试调用示例 sample_data = ''' "prices": { "header": { "currency": "EUR" }, "price":{ "productId":"0000A", "value":"60.00" } } ''' process_single_price(sample_data)
优化建议
- 批量处理:如果数据流频率较高,可以攒一批同币种数据再写入,减少存储读写次数,降低成本
- 版本控制:开启对象存储的版本控制,防止误写导致数据丢失
- 重试机制:给存储读写操作添加重试逻辑,应对网络波动等临时问题
内容的提问来源于stack exchange,提问作者DanBo
相关产品推荐
相关产品推荐

