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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 17:45:25