如何用Pandas合并MySQL表JSON列所有行成单个JSON以调用REST API
实现方案与代码示例
1. 依赖安装
先安装所需的Python库:
pip install pandas sqlalchemy requests schedule python-dotenv
2. 完整代码实现
import pandas as pd from sqlalchemy import create_engine import requests import schedule import time import json from dotenv import load_dotenv import os # 加载环境变量(避免硬编码敏感信息) load_dotenv() # 配置参数 DB_CONFIG = { 'user': os.getenv('DB_USER'), 'password': os.getenv('DB_PASSWORD'), 'host': os.getenv('DB_HOST'), 'database': os.getenv('DB_NAME'), 'port': int(os.getenv('DB_PORT', 3306)) } API_URL = os.getenv('API_URL') INTERVAL_MINUTES = 10 def process_and_upload(): try: # 连接MySQL读取数据 engine = create_engine(f"mysql+pymysql://{DB_CONFIG['user']}:{DB_CONFIG['password']}@{DB_CONFIG['host']}:{DB_CONFIG['port']}/{DB_CONFIG['database']}") # 替换为你的表名和查询语句 df = pd.read_sql("SELECT json_data FROM your_target_table", engine) engine.dispose() if df.empty: print("当前无待上传数据") return # 合并JSON数据为有效数组 merged_data = [] for row_json in df['json_data']: try: # 如果MySQL中json_data是字符串类型,需要解析;若是JSON类型可直接append(row_json) parsed = json.loads(row_json) merged_data.append(parsed) except json.JSONDecodeError as e: print(f"跳过无效JSON行:{row_json},错误:{str(e)}") continue if not merged_data: print("无有效JSON数据可上传") return # 转换为JSON字符串 payload = json.dumps(merged_data) # 发送POST请求 headers = {'Content-Type': 'application/json'} resp = requests.post(API_URL, data=payload, headers=headers) resp.raise_for_status() print(f"上传成功,API响应:{resp.status_code} | {resp.text}") except Exception as e: print(f"任务执行失败:{str(e)}") # 设置定时任务 schedule.every(INTERVAL_MINUTES).minutes.do(process_and_upload) if __name__ == "__main__": print(f"定时任务启动,每{INTERVAL_MINUTES}分钟执行一次") # 首次启动立即执行一次 process_and_upload() # 维持任务循环 while True: schedule.run_pending() time.sleep(60)
3. 关键说明
- 配置管理:创建
.env文件存储敏感信息,示例内容:DB_USER=your_db_username DB_PASSWORD=your_db_password DB_HOST=127.0.0.1 DB_NAME=your_database API_URL=https://your-api-domain.com/upload-endpoint - JSON解析适配:如果你的MySQL表中
json_data列是原生JSON类型(而非字符串),可以去掉json.loads()步骤,直接用merged_data = df['json_data'].tolist()合并。 - 异常处理:代码包含了数据库连接、JSON解析、API请求的异常捕获,可根据需求添加重试逻辑(比如API请求失败时重试2-3次)。
- 定时任务:使用
schedule库实现简单的周期调度,若需要更复杂的调度(如指定时间段执行),可替换为APScheduler库。
内容的提问来源于stack exchange,提问作者Nicolas
相关产品推荐
相关产品推荐

