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

Python读取EventStream API后的数据清洗、存储及代码优化咨询

处理EventStream API数据:清洗、存储与代码优化

一、数据清洗为JSON格式

你当前手动遍历响应行的方式没有利用SSE规范的解析逻辑,导致需要自行处理杂乱格式。既然已经导入了sseclient,直接用它解析EventStream流会更高效,自动拆分事件、过滤无用的心跳信息:

import requests
import urllib3
import json
import sseclient

# 禁用SSL警告(仅测试环境使用,生产环境建议配置合法证书)
urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning)
requests.packages.urllib3.util.ssl_.DEFAULT_CIPHERS += ':HIGH:!DH:!aNULL'

# 配置信息
HEADERS = {"Cookie": "你的Cookie内容"}
API_URL = "你的EventStream API地址"

# 建立流式请求并检查状态
with requests.get(API_URL, headers=HEADERS, verify=False, stream=True) as response:
    response.raise_for_status()
    # 用sseclient解析流
    client = sseclient.SSEClient(response)
    for event in client.events():
        # 过滤心跳事件和空数据
        if event.event == "heartbeat" or not event.data.strip():
            continue
        # 解析为JSON对象
        try:
            clean_data = json.loads(event.data)
            print(clean_data)  # 此处可替换为存储逻辑
        except json.JSONDecodeError as e:
            print(f"JSON解析失败: {e}, 原始数据: {event.data}")

运行后就能直接得到结构化的JSON对象,跳过无效的心跳事件。

二、合适的数据存储方法

根据使用场景,推荐以下几种存储方案:

1. JSON Lines格式文件(轻量易读)

适合小规模数据存储,每行一个JSON对象,后续可直接用Pandas读取分析:

# 在循环内添加存储逻辑
with open("event_data.jsonl", "a+", encoding="utf-8") as f:
    json.dump(clean_data, f, ensure_ascii=False)
    f.write("\n")

2. SQLite数据库(轻量持久化)

无需额外服务,适合需要结构化查询的场景:

import sqlite3

# 初始化数据库连接(不存在则自动创建)
with sqlite3.connect("event_stream.db") as conn:
    cursor = conn.cursor()
    # 创建数据表(首次运行执行)
    cursor.execute("""
    CREATE TABLE IF NOT EXISTS events (
        id TEXT PRIMARY KEY,
        timestamp INTEGER,
        heure INTEGER,
        sens INTEGER,
        action TEXT
    )
    """)
    conn.commit()

    # 在循环内插入数据
    try:
        cursor.execute("""
        INSERT OR REPLACE INTO events (id, timestamp, heure, sens, action)
        VALUES (?, ?, ?, ?, ?)
        """, (clean_data["id"], clean_data["timestamp"], clean_data["heure"], clean_data["sens"], clean_data["action"]))
        conn.commit()
    except sqlite3.Error as e:
        print(f"数据库操作失败: {e}")

3. Redis(实时流数据缓存)

适合高并发实时场景,支持快速读写:

import redis

# 连接Redis(需提前安装redis库:pip install redis)
r = redis.Redis(host="localhost", port=6379, db=0)

# 存储为字符串,或添加到列表实现队列
r.set(f"event:{clean_data['id']}", json.dumps(clean_data))
r.rpush("event_stream", json.dumps(clean_data))

三、代码优化建议

  1. 用sseclient替代手动解析:符合SSE规范,减少重复造轮子的出错概率。
  2. 完善异常处理:捕获网络请求错误、JSON解析错误、数据库操作错误,避免程序意外崩溃。
  3. 配置分离:把URL、Cookie等敏感/可变配置放到单独变量或.env文件(用python-dotenv库),不要硬编码。
  4. 使用上下文管理器:处理文件、数据库连接、请求响应,确保资源自动释放,避免内存泄漏。
  5. 替换print为logging:生产环境更易调试和记录日志,便于问题排查:
    import logging
    logging.basicConfig(level=logging.INFO, format="%(asctime)s - %(levelname)s - %(message)s")
    logging.info(f"成功解析事件: {clean_data}")
    logging.error(f"JSON解析失败: {e}")
    
  6. 谨慎禁用SSL警告:测试环境可临时使用,生产环境建议配置合法的SSL证书,不要全局禁用警告。

内容的提问来源于stack exchange,提问作者flopab

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 16:36:22