如何从AWS Athena结果的Python DataFrame生成指定格式JSON请求体
解决方案:将AWS Athena结果转换为REST API请求体
核心思路
利用Pandas的groupby分组功能按order_id聚合订单数据,把同一订单的多产品合并到contents数组,提取订单公共信息后构造符合要求的JSON请求体,最终调用REST API。
完整代码示例
import pandas as pd import json import requests import logging # 初始化日志(保持你原有的日志配置) logging.basicConfig(level=logging.DEBUG) logger = logging.getLogger(__name__) def get_order_info(company_name): # 替换为你实际的Athena查询逻辑,以下是模拟数据示例 data = { "order_id": ["123", "123", "123", "234"], "order_value": [22.0, 22.0, 22.0, 45.0], "product_currency": ["INR", "INR", "INR", "EUR"], "order_email": ["abc@gmail.com", "abc@gmail.com", "abc@gmail.com", "xyz@gmail.com"], "product_id": ["1A", "2A", "3E", "5T"], "product_quantity": [1, 2, 2, 1], "product_price": [12.0, 4.0, 1.0, 45.0], "product_name": ["Piano-set", "Ludo", "Toffee", "Dumbbell-set"] } return pd.DataFrame(data) def build_invoice_request(order_group, company_name): """构造单个订单的API请求体""" # 提取订单公共信息(取分组内第一条数据的公共字段) order_base = order_group.iloc[0] # 组装contents产品数组 contents = [ { "quantity": row["product_quantity"], "product_price": row["product_price"], "title": row["product_name"] } for _, row in order_group.iterrows() ] # 构造完整请求结构 request_body = { "data": { "value": f"{order_base['order_value']:.2f}", # 格式化为两位小数的字符串 "currency": order_base["product_currency"], "orderId": str(order_base["order_id"]), # 转为字符串类型 "contents": contents, "company": company_name } } return request_body def send_invoice_request(request_body, api_url="https://your-api-domain/invoice"): """发送POST请求到REST API""" headers = {"content-type": "application/json"} try: response = requests.post(api_url, headers=headers, json=request_body) response.raise_for_status() # 捕获HTTP错误 logger.debug(f"订单 {request_body['data']['orderId']} 请求成功,响应: {response.text}") return response except requests.exceptions.RequestException as e: logger.error(f"订单 {request_body['data']['orderId']} 请求失败: {str(e)}") return None # 主执行逻辑 if __name__ == "__main__": company_name = "ABC" # 替换为实际公司名称 order_query_results_df = get_order_info(company_name) # 按order_id分组处理每个订单 for order_id, group in order_query_results_df.groupby("order_id"): logger.debug(f"开始处理订单: {order_id}") # 生成请求体 invoice_request = build_invoice_request(group, company_name) logger.debug(f"生成的请求体:\n{json.dumps(invoice_request, indent=2)}") # 发送API请求 send_invoice_request(invoice_request)
关键代码说明
- 分组聚合:用
groupby("order_id")将同一订单的行聚合,避免逐行遍历的低效操作,同时确保同一订单的产品被归类到一起。 - 公共字段提取:通过
order_group.iloc[0]获取分组内第一条数据的公共字段(如order_value、product_currency),因为同一订单的这些字段是重复的。 - Contents数组构造:用列表推导式遍历分组内的每一行产品数据,转换为API要求的字典格式。
- 数据格式处理:
value字段格式化为保留两位小数的字符串,匹配API示例要求;orderId转为字符串类型,避免数字类型的潜在兼容性问题。
- API调用:使用
requests.post发送JSON请求,携带Content-Type头,同时捕获并处理请求异常。
注意事项
- 替换
get_order_info内的模拟数据为你实际的Athena查询逻辑; - 修改
send_invoice_request中的api_url为你的REST API实际地址; - 若需批量处理大量订单,可引入线程池/进程池提升处理效率。
内容的提问来源于stack exchange,提问作者azaveri7
相关产品推荐
相关产品推荐

