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

如何使用Python实时读取TXT文件并更新PostgreSQL数据库

实现步骤说明

1. 现有代码优化点

你当前的基础代码有两个明显问题:

  • 每次循环都重读整个文件,会重复处理已经入库的历史行
  • 没有等待间隔,会占满CPU资源,性能极低
    我们先把读取逻辑改成记录文件偏移量,只读取新增行的模式,再对接PostgreSQL操作即可。

2. 依赖准备

先安装PostgreSQL Python驱动:

pip install psycopg2-binary

3. 完整实现代码

import time
import psycopg2
from psycopg2 import OperationalError

# 数据库连接配置,改成你自己的实际信息
DB_CONFIG = {
    "dbname": "你的库名",
    "user": "你的用户名",
    "password": "你的密码",
    "host": "你的数据库地址",
    "port": "5432"
}
# 要监控的TXT文件路径
FILE_PATH = "你的txt文件路径"
# 记录上次读取的文件偏移量
last_offset = 0

def get_db_connection():
    """获取数据库连接,自动处理重连"""
    try:
        conn = psycopg2.connect(**DB_CONFIG)
        return conn
    except OperationalError as e:
        print(f"数据库连接失败: {e},5秒后重试")
        time.sleep(5)
        return get_db_connection()

if __name__ == "__main__":
    conn = get_db_connection()
    cursor = conn.cursor()
    while True:
        try:
            with open(FILE_PATH, 'r', encoding='utf-8') as f:
                # 移动到上次读取的末尾位置
                f.seek(last_offset)
                new_lines = f.readlines()
                # 更新偏移量到当前文件末尾
                last_offset = f.tell()

            for line in new_lines:
                line = line.strip()
                if not line:
                    continue
                # 按分号拆分字段
                fields = line.split(';')
                if len(fields) < 7:
                    print(f"跳过格式错误行: {line}")
                    continue
                # 取第4个字段(索引从0开始,所以是3)
                field4 = fields[3]
                # 这里默认取第6个字段(员工号,索引5)作为更新的匹配条件,可根据实际业务修改
                staff_no = fields[5]

                # 编写你的更新SQL,表名、字段名改成你实际的
                update_sql = """
                UPDATE 你的表名 
                SET 要更新的字段名 = %s 
                WHERE 员工号字段名 = %s
                """
                try:
                    cursor.execute(update_sql, (field4, staff_no))
                    conn.commit()
                    print(f"已更新员工{staff_no}的状态为{field4}")
                except Exception as e:
                    conn.rollback()
                    print(f"更新数据库失败,行内容: {line},错误: {e}")
            
            # 无新增数据时休眠1秒,减少CPU占用,可根据实时性要求调整
            time.sleep(1)
        except Exception as e:
            print(f"程序运行出错: {e},5秒后重启循环")
            # 连接异常时重置连接
            if conn.closed:
                conn = get_db_connection()
                cursor = conn.cursor()
            time.sleep(5)

4. 调整说明

  • 上述代码里的SQL更新逻辑是示例,你需要根据自己的业务需求修改匹配条件、更新的字段和表名
  • 如果TXT的编码不是utf-8,改open函数里的encoding参数即可
  • 生产环境建议把print改成日志记录,方便后续排查问题
  • 可以根据你的实时性要求调整sleep的时间,比如要更低延迟可以改成0.5秒

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 07:27:04