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

如何用Python定时持续获取PostgreSQL表数据(每分钟一行)

如何用Python定时逐行获取PostgreSQL表数据

你的现有代码是一次性查询并获取全量/单行数据,执行后立即关闭连接,无法实现每隔1分钟获取下一行的持续操作需求。下面从连接复用、定时触发、逐行定位三个方面给出调整方案:

方案一:基于游标持续逐行获取(适合静态数据)

通过一次建立连接和游标,执行查询后循环调用fetchone()逐行获取,每次获取后等待1分钟。适合数据不会新增/修改的只读场景。

import psycopg2
import time
import psycopg2.errors

def continuous_fetch_from_cursor():
    conn = None
    cursor = None
    try:
        # 建立数据库连接(仅执行一次)
        conn = psycopg2.connect(
            database="mydb", user='postgres', password='password', host='127.0.0.1', port='5432'
        )
        conn.autocommit = True
        cursor = conn.cursor()
        
        # 执行查询,按id排序保证顺序
        cursor.execute("SELECT * FROM EMPLOYEE ORDER BY id")
        
        while True:
            row = cursor.fetchone()
            if row:
                print("获取到的行:", row)
            else:
                print("已遍历完所有数据,等待1分钟后重试...")
            
            # 等待1分钟
            time.sleep(60)
            
    except psycopg2.errors.Error as e:
        print(f"数据库异常: {str(e)}")
    finally:
        # 关闭资源
        if cursor:
            cursor.close()
        if conn:
            conn.close()

if __name__ == "__main__":
    continuous_fetch_from_cursor()

方案一注意点

  • 若数据库连接因闲置超时断开,程序会抛出异常终止,需额外添加重连逻辑
  • 若表中有新数据插入,已执行的游标无法自动捕获新数据,需重新执行查询

方案二:基于主键定位逐行获取(适合动态新增数据)

每次查询时以上一次获取的记录主键为基准,获取下一行数据。即使程序重启或连接断开,也能从正确位置继续,适合数据持续新增的场景。

import psycopg2
import time
import psycopg2.errors

def get_next_row(last_id=0):
    conn = None
    try:
        conn = psycopg2.connect(
            database="mydb", user='postgres', password='password', host='127.0.0.1', port='5432'
        )
        conn.autocommit = True
        cursor = conn.cursor()
        
        # 按id大于上一次的记录,获取最小id的行(即下一行)
        cursor.execute(
            "SELECT * FROM EMPLOYEE WHERE id > %s ORDER BY id LIMIT 1",
            (last_id,)
        )
        row = cursor.fetchone()
        
        # 返回获取到的行和最新的id,若无数据则返回原last_id
        return row, row[0] if row else last_id
        
    except psycopg2.errors.Error as e:
        print(f"数据库异常: {str(e)}")
        return None, last_id
    finally:
        if conn:
            conn.close()

def main():
    # 初始last_id设为0,若需重启续传可从文件/数据库读取上次的id
    last_id = 0
    while True:
        row, last_id = get_next_row(last_id)
        if row:
            print("获取到的行:", row)
        else:
            print("暂无新数据,等待1分钟后重试...")
        
        time.sleep(60)

if __name__ == "__main__":
    main()

方案二注意点

  • 需保证表存在可排序的唯一键(如自增主键id、时间戳create_time),否则无法准确定位下一行
  • 若需程序重启后从上次停止位置继续,可将last_id写入本地文件或专用配置表中

通用注意事项

  • 异常处理:必须捕获psycopg2相关异常,避免因数据库波动导致程序直接崩溃
  • 连接复用:频繁开关连接会增加数据库负载,方案一的连接复用更高效,但需处理连接断开重连
  • 定时精度:time.sleep(60)会让程序暂停60秒,若查询耗时较长,实际间隔会超过1分钟。如需精确控制,可通过datetime计算下次执行时间,调整sleep时长

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 13:55:20