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

如何将异步流式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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 12:03:30