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

使用psycopg2批量导入文件时数据重复且日期不匹配问题求助

问题:导入文本文件到PostgreSQL时数据重复且日期匹配错误

我编程经验尚浅,目前尝试用psycopg2将目录下格式统一的文本文件导入PostgreSQL,希望为每条数据添加对应文件的日期列。但导入后出现异常:始终重复最后一个文件的数据,却搭配其他文件的日期。

示例文件内容

file_01/01/2024

VDS_ID,FWY,co
1001,1,LA
1002,1,LA
1003,2,LA
1004,1,LA
1005,1,LA

file_03/01/2024

VDS_ID,FWY,co
1001,1,LA
1002,1,LA
1003,2,LA
1005,1,LA

file_05/01/2024

VDS_ID,FWY,co
1001,1,LA
1002,1,LA
1003,2,LA

导入后异常结果

VDS_ID,FWY,config_date
1001,1,LA,01/01/2024
1002,1,LA,01/01/2024
1003,2,LA,01/01/2024
1001,1,LA,03/01/2024
1002,1,LA,03/01/2024
1003,2,LA,03/01/2024
1001,1,LA,05/01/2024
1002,1,LA,05/01/2024
1003,2,LA,05/01/2024

我的代码

import psycopg2
from psycopg2.extras import execute_values
import time
import os
import subprocess
import csv

def main():
    mypath = r'C:\Downloads\PEMS\Meta'
    tableName = 'pems_config'
    temptableName = 'temp_pems_config'

    # Connect to an existing database (password blanked out question
    conn = psycopg2.connect(database='Test', user='postgres', password ='****')

    # Open a cursor to perform database operations
    cur = conn.cursor()

    start_time = time.time()
    print ("Start time: {0}", time.asctime(time.localtime(start_time)))

    textfiles = []

    for file in os.listdir(mypath):
        if file.endswith(".txt"):
##           print(os.path.join(mypath, file))
           textfiles.append(file)

    create_table_sql = """CREATE TABLE pems_config
(
    "VDS_ID" integer,
    "Fwy" smallint,
    "Dir" character(1),
    "District" smallint,
    "County" integer,
    "City" integer,
    "State_PM" varchar(9),
    "Abs_PM" real,
    "Latitude" real,
    "Longitude" real,
    "Length" real,
    "Type" character(2),
    "Lanes" smallint,
    "Name" text,
    "User_ID_1" text,
    "User_ID_2" text,
    "User_ID_3" text,
    "User_ID_4" text,
    "pems_config_date" date
);"""

    create_temp_table_sql = """CREATE TEMPORARY TABLE temp_pems_config
    (
        "VDS_ID" integer,
        "Fwy" smallint,
        "Dir" character(1),
        "District" smallint,
        "County" integer,
        "City" integer,
        "State_PM" varchar(9),
        "Abs_PM" real,
        "Latitude" real,
        "Longitude" real,
        "Length" real,
        "Type" character(2),
        "Lanes" smallint,
        "Name" text,
        "User_ID_1" text,
        "User_ID_2" text,
        "User_ID_3" text,
        "User_ID_4" text
    );"""

    drop_table_sql = """DROP TABLE IF EXISTS temp_pems_config"""

    copy_sql = """COPY temp_pems_config FROM stdin WITH DELIMITER '\t' NULL AS '' csv HEADER"""

    add_column_sql = """ALTER TABLE temp_pems_config ADD COLUMN pems_config_date date;"""

    update_date_sql = """UPDATE temp_pems_config SET pems_config_date = %s ;"""

    append_table_sql = """INSERT INTO pems_config select * FROM temp_pems_config;"""

    final_table_sql = """CREATE TABLE final_pems_config AS Select "VDS_ID","Dir","District","County","City","State_PM","Abs_PM","Latitude","Longitude","Length","Type",
    "Lanes","Name","User_ID_1","User_ID_2","User_ID_3","User_ID_4",max("pems_config_date") AS "pems_config_date" from pems_config
    GROUP BY "VDS_ID","Dir","District","County","City","State_PM","Abs_PM","Latitude","Longitude","Length","Type",
    "Lanes","Name","User_ID_1","User_ID_2","User_ID_3","User_ID_4";"""

    cur.execute("SELECT EXISTS (SELECT 1 AS result FROM pg_tables WHERE schemaname = 'public' AND tablename = 'pems_config');")
    tableExists = cur.fetchone()[0]
    print ("{0} Exists: {1}".format(tableName, str(tableExists)))

    if tableExists == False:
        print ("Creating {0} Table:".format(tableName))
        #Execute a command: this creates a new table
        cur.execute(create_table_sql)
        print (cur.statusmessage)
        print ("Time to create {0} Table: {1} s".format(tableName,time.time()-start_time))


    for tfile in textfiles:

        tfile_name = os.path.splitext(tfile)[0]
        config_date = tfile_name[14:18]+"-"+tfile_name[19:21]+"-"+tfile_name[22:25]
##        #check if table exists
        cur.execute("SELECT EXISTS (SELECT 1 AS result FROM pg_tables WHERE schemaname = 'public' AND tablename = 'temp_pems_config');")
        tableExists = cur.fetchone()[0]
        print ("{0} Exists: {1}".format(temptableName, str(tableExists)))
    ##
    ##        #if table does not exist then create it
        if tableExists == False:
           print ("Creating {0} Table:".format(temptableName))
           #Execute a command: this creates a new table
           cur.execute(drop_table_sql)
           cur.execute(create_temp_table_sql)
           print (cur.statusmessage)
           print ("Time to create {0} Table: {1} s".format(tableName,time.time()-start_time))
        else:
            print ("Deleting {0} Table:".format(tableName))
            cur.execute(drop_table_sql)
            print ("Creating {0} Table:".format(tableName))
            cur.execute(create_table_sql)
            print (cur.statusmessage)
            print ("Time to create {0} Table: {1} s".format(tableName,time.time()-start_time))

        with open(os.path.join(mypath,file), 'r') as f:
            copyDataFile_time = time.time()
             #copy zfile into CHTV
            execute_values 
            cur.copy_expert(sql=copy_sql, file=f)
            cur.execute(add_column_sql)
            cur.execute(update_date_sql, (config_date,))
            cur.execute(append_table_sql)
            print ("Time to copy {0} Table: {1} mins".format(tableName,(time.time()-copyDataFile_time)/60))
        print (cur.statusmessage)
    cur.execute(final_table_sql)
    conn.commit()
    cur.close()
    Endtime = time.time()
    print ("Start time: {0}".format(time.asctime(time.localtime(start_time))))
    print ("End time: {0} ".format(time.asctime(time.localtime(Endtime))))
    print ("Total processing time: {0} Minutes".format((Endtime-start_time)/60))

if __name__ == '__main__':
    main()

异常原因

  1. 文件读取变量错误:循环遍历textfiles时用的是tfile变量,但打开文件时却用了外层遍历目录的file变量,导致每次循环都读取最后一个遍历到的文件,而日期是当前tfile对应的日期,所以所有日期都绑定同一个文件的数据。
  2. 临时表逻辑混乱:当临时表存在时,错误地删除临时表后重建了主表pems_config,直接清空之前导入的数据,每次循环重复这个错误操作,进一步导致数据异常。
  3. 无效代码残留:execute_values这行代码没有实际作用,属于冗余代码。

解决方法

  1. 修正文件读取变量:将open(os.path.join(mypath,file), 'r')改为open(os.path.join(mypath,tfile), 'r'),确保每次循环读取当前迭代的目标文件。
  2. 简化临时表处理:删除复杂的存在性判断逻辑,每次循环直接删除并重建临时表,避免逻辑冲突:
    # 替换原有的临时表判断代码
    cur.execute(drop_table_sql)
    cur.execute(create_temp_table_sql)
    print(f"重建临时表 {temptableName}")
    
  3. 移除错误的主表重建代码:删除else分支中重建pems_config主表的代码,防止已导入数据被反复清空。
  4. 清理无效代码:删除execute_values这一行无用代码。
  5. 优化日期提取逻辑:PostgreSQL的date类型要求格式为YYYY-MM-DD,建议改用更可靠的方式提取日期,避免索引截取错误:
    # 假设文件名格式为file_dd/mm/yyyy.txt
    date_part = tfile_name.split('_')[1]  # 提取dd/mm/yyyy部分
    day, month, year = date_part.split('/')
    config_date = f"{year}-{month}-{day}"  # 转换为标准日期格式
    

修正后的核心循环代码

for tfile in textfiles:
    tfile_name = os.path.splitext(tfile)[0]
    # 优化日期提取逻辑,适配文件名格式
    date_part = tfile_name.split('_')[1]
    day, month, year = date_part.split('/')
    config_date = f"{year}-{month}-{day}"

    # 每次循环重建临时表
    cur.execute(drop_table_sql)
    cur.execute(create_temp_table_sql)
    print(f"重建临时表 {temptableName}")

    with open(os.path.join(mypath, tfile), 'r') as f:
        copyDataFile_time = time.time()
        cur.copy_expert(sql=copy_sql, file=f)
        cur.execute(add_column_sql)
        cur.execute(update_date_sql, (config_date,))
        cur.execute(append_table_sql)
        print(f"导入文件 {tfile} 耗时: {(time.time()-copyDataFile_time)/60} 分钟")
    print(cur.statusmessage)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 00:28:14